Files
synapbus/internal/webhooks/ratelimit.go
T
Algis DumbrisandClaude Opus 4.6 42775df65f feat: webhooks & Kubernetes Job runner for event-driven agents
- 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>
2026-03-14 18:14:04 +02:00

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)
}