- Webhook registration via MCP (register_webhook, list_webhooks, delete_webhook) - HMAC-SHA256 payload signing, SSRF-safe HTTP client, loop detection (depth 5) - 8-worker goroutine delivery pool with exponential backoff retry (1s/5s/30s) - Dead letter queue with auto-purge, auto-disable after 50 consecutive failures - Per-agent rate limiting (60 deliveries/min) - K8s Job runner (register_k8s_handler, list_k8s_handlers, delete_k8s_handler) - Auto-detect in-cluster via InClusterConfig, NoopRunner fallback - REST API for webhook deliveries, dead letters, K8s job runs and logs - Web UI: webhook management, K8s handler pages, dead letters view - MultiDispatcher fan-out pattern for webhook + K8s event dispatch - SQLite migration 009: webhooks, webhook_deliveries, k8s_handlers, k8s_job_runs - 51 tests across 9 test packages, all passing Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
54 lines
1.5 KiB
Go
54 lines
1.5 KiB
Go
package webhooks
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"time"
|
|
|
|
"golang.org/x/time/rate"
|
|
)
|
|
|
|
// AgentRateLimiter manages per-agent rate limiters for webhook delivery.
|
|
// Each agent gets a token bucket allowing 60 deliveries per minute (1/second sustained)
|
|
// with burst capacity of 60.
|
|
type AgentRateLimiter struct {
|
|
limiters sync.Map // map[string]*rate.Limiter
|
|
rate rate.Limit
|
|
burst int
|
|
}
|
|
|
|
// NewAgentRateLimiter creates a new rate limiter with the given rate and burst.
|
|
// Default: 1 per second sustained, burst of 60.
|
|
func NewAgentRateLimiter(perMinute int) *AgentRateLimiter {
|
|
return &AgentRateLimiter{
|
|
rate: rate.Every(time.Minute / time.Duration(perMinute)),
|
|
burst: perMinute,
|
|
}
|
|
}
|
|
|
|
// Wait blocks until the agent is allowed to make a delivery, or the context is cancelled.
|
|
func (r *AgentRateLimiter) Wait(ctx context.Context, agentName string) error {
|
|
limiter := r.getLimiter(agentName)
|
|
return limiter.Wait(ctx)
|
|
}
|
|
|
|
// Allow checks if the agent can make a delivery without blocking.
|
|
func (r *AgentRateLimiter) Allow(agentName string) bool {
|
|
limiter := r.getLimiter(agentName)
|
|
return limiter.Allow()
|
|
}
|
|
|
|
// Remove removes the rate limiter for an agent (cleanup when agent has no webhooks).
|
|
func (r *AgentRateLimiter) Remove(agentName string) {
|
|
r.limiters.Delete(agentName)
|
|
}
|
|
|
|
func (r *AgentRateLimiter) getLimiter(agentName string) *rate.Limiter {
|
|
if v, ok := r.limiters.Load(agentName); ok {
|
|
return v.(*rate.Limiter)
|
|
}
|
|
limiter := rate.NewLimiter(r.rate, r.burst)
|
|
actual, _ := r.limiters.LoadOrStore(agentName, limiter)
|
|
return actual.(*rate.Limiter)
|
|
}
|