diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 674f33b..6011558 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -3,15 +3,10 @@ package delivery import ( - "bytes" "context" - "encoding/json" - "errors" "fmt" - "io" "log/slog" "net/http" - "strings" "sync" "time" @@ -71,25 +66,13 @@ const ( httpSuccessMax = 300 ) -// Sentinel errors returned by config parsers. -var ( - errEmptyTargetConfig = errors.New( - "empty target config", - ) - errMissingWebhookURL = errors.New( - "webhook_url is required", - ) - errMissingTargetURL = errors.New( - "target URL is required", - ) -) - // Task contains everything needed to deliver an event to a // single target. type Task struct { - DeliveryID string - EventID string - WebhookID string + DeliveryID string + EventID string + WebhookID string + EntrypointID string TargetID string TargetName string @@ -111,20 +94,6 @@ type Notifier interface { Notify(tasks []Task) } -// HTTPTargetConfig holds configuration for http target -// types. -type HTTPTargetConfig struct { - URL string `json:"url"` - Headers map[string]string `json:"headers,omitempty"` - Timeout int `json:"timeout,omitempty"` -} - -// SlackTargetConfig holds configuration for slack target -// types. -type SlackTargetConfig struct { - WebhookURL string `json:"webhookUrl"` -} - // EngineParams are the fx dependencies for the delivery // engine. type EngineParams struct { @@ -136,21 +105,28 @@ type EngineParams struct { } // Engine processes queued deliveries in the background -// using a bounded worker pool architecture. +// using a bounded worker pool architecture. It owns only +// the cross-target machinery: the worker pool, the queue +// and retry channels, restart recovery, and the persistence +// helpers and Scheduler that individual targets rely on. +// Each target type owns its own delivery, including retries, +// backoff, and circuit breaking. type Engine struct { database *database.Database dbManager *database.WebhookDBManager log *slog.Logger - client *http.Client cancel context.CancelFunc wg sync.WaitGroup deliveryCh chan Task retryCh chan Task workers int - // circuitBreakers stores a *CircuitBreaker per target - // ID. - circuitBreakers sync.Map + // targets maps each target type to its implementation. + targets map[database.TargetType]Target + + // httpTarget is retained so tests can reach the HTTP + // target's shared client and circuit breakers. + httpTarget *httpTarget } // New creates and registers the delivery engine with the @@ -160,18 +136,19 @@ func New( params EngineParams, ) *Engine { e := &Engine{ - database: params.DB, - dbManager: params.DBManager, - log: params.Logger.Get(), - client: &http.Client{ - Timeout: httpClientTimeout, - Transport: NewSSRFSafeTransport(), - }, + database: params.DB, + dbManager: params.DBManager, + log: params.Logger.Get(), deliveryCh: make(chan Task, deliveryChannelSize), retryCh: make(chan Task, retryChannelSize), workers: defaultWorkers, } + e.initTargets(&http.Client{ + Timeout: httpClientTimeout, + Transport: NewSSRFSafeTransport(), + }) + lc.Append(fx.Hook{ OnStart: func(ctx context.Context) error { e.start(ctx) @@ -205,52 +182,32 @@ func (e *Engine) Notify(tasks []Task) { } } -// FormatSlackMessage builds a Slack-compatible message -// string from a webhook event. -func FormatSlackMessage( - event *database.Event, -) string { - var b strings.Builder - - b.WriteString("*Webhook Event Received*\n") - - fmt.Fprintf( - &b, "*Method:* `%s`\n", event.Method, +// ScheduleRetry schedules a task to be re-enqueued onto the +// retry channel after delay. It implements the Scheduler +// interface the targets use to own their durable retries. +func (e *Engine) ScheduleRetry( + task Task, delay time.Duration, +) { + e.log.Debug( + "scheduling delivery retry", + "webhook_id", task.WebhookID, + "delivery_id", task.DeliveryID, + "delay", delay, + "next_attempt", task.AttemptNum, ) - fmt.Fprintf( - &b, - "*Content-Type:* `%s`\n", - event.ContentType, - ) - - fmt.Fprintf( - &b, - "*Timestamp:* `%s`\n", - event.CreatedAt.UTC().Format(time.RFC3339), - ) - - fmt.Fprintf( - &b, - "*Body Size:* %d bytes\n", - len(event.Body), - ) - - if event.Body == "" { - b.WriteString("\n_(empty body)_\n") - - return b.String() - } - - if formatted := formatJSONBody(event.Body); formatted != "" { - b.WriteString(formatted) - - return b.String() - } - - formatRawBody(&b, event.Body) - - return b.String() + time.AfterFunc(delay, func() { + select { + case e.retryCh <- task: + default: + e.log.Warn( + "retry channel full, delivery "+ + "will be recovered by periodic sweep", + "delivery_id", task.DeliveryID, + "webhook_id", task.WebhookID, + ) + } + }) } func (e *Engine) start(ctx context.Context) { @@ -493,13 +450,36 @@ func (e *Engine) recoverRetryingDeliveries( } } +// recoverSingleRetry hands an orphaned retrying delivery back +// to its target to recompute the remaining backoff, then +// reschedules it. Targets that do not own durable retries +// (fire-and-forget) never produce retrying deliveries, so +// they are skipped. func (e *Engine) recoverSingleRetry( webhookDB *gorm.DB, webhookID string, d *database.Delivery, ) { + target, err := e.loadTarget(d.TargetID) + if err != nil { + e.log.Error( + "failed to load target for retrying "+ + "delivery recovery", + "delivery_id", d.ID, + "target_id", d.TargetID, + "error", err, + ) + + return + } + + rs, ok := e.targets[target.Type].(rescheduler) + if !ok { + return + } + attemptNum := e.countAttempts(webhookDB, d.ID) - remaining := e.calcRemainingBackoff( + remaining := rs.remainingBackoff( webhookDB, d.ID, attemptNum, ) @@ -516,19 +496,6 @@ func (e *Engine) recoverSingleRetry( return } - target, err := e.loadTarget(d.TargetID) - if err != nil { - e.log.Error( - "failed to load target for retrying "+ - "delivery recovery", - "delivery_id", d.ID, - "target_id", d.TargetID, - "error", err, - ) - - return - } - task := buildRecoveryTask( d, webhookID, &event, &target, attemptNum+1, ) @@ -541,7 +508,7 @@ func (e *Engine) recoverSingleRetry( "remaining_backoff", remaining, ) - e.scheduleRetry(task, remaining) + e.ScheduleRetry(task, remaining) } func (e *Engine) recoverPendingDeliveries( @@ -586,31 +553,6 @@ func (e *Engine) recoverPendingDeliveries( ) } -func (e *Engine) scheduleRetry( - task Task, delay time.Duration, -) { - e.log.Debug( - "scheduling delivery retry", - "webhook_id", task.WebhookID, - "delivery_id", task.DeliveryID, - "delay", delay, - "next_attempt", task.AttemptNum, - ) - - time.AfterFunc(delay, func() { - select { - case e.retryCh <- task: - default: - e.log.Warn( - "retry channel full, delivery "+ - "will be recovered by periodic sweep", - "delivery_id", task.DeliveryID, - "webhook_id", task.WebhookID, - ) - } - }) -} - func (e *Engine) retrySweep(ctx context.Context) { defer e.wg.Done() @@ -705,14 +647,35 @@ func (e *Engine) sweepWebhookRetries( } } +// sweepSingleRetry re-enqueues an orphaned retrying delivery +// whose backoff window has elapsed, delegating the backoff +// decision to the delivery's target. Targets that do not own +// durable retries are skipped. func (e *Engine) sweepSingleRetry( webhookDB *gorm.DB, webhookID string, d *database.Delivery, ) { + target, err := e.loadTarget(d.TargetID) + if err != nil { + e.log.Error( + "retry sweep: failed to load target", + "delivery_id", d.ID, + "target_id", d.TargetID, + "error", err, + ) + + return + } + + rs, ok := e.targets[target.Type].(rescheduler) + if !ok { + return + } + attemptNum := e.countAttempts(webhookDB, d.ID) - if !e.backoffElapsed( + if !rs.backoffElapsed( webhookDB, d.ID, attemptNum, ) { return @@ -730,18 +693,6 @@ func (e *Engine) sweepSingleRetry( return } - target, err := e.loadTarget(d.TargetID) - if err != nil { - e.log.Error( - "retry sweep: failed to load target", - "delivery_id", d.ID, - "target_id", d.TargetID, - "error", err, - ) - - return - } - task := buildRecoveryTask( d, webhookID, &event, &target, attemptNum+1, ) @@ -759,22 +710,16 @@ func (e *Engine) sweepSingleRetry( } } +// processDelivery dispatches a delivery to the target that +// owns its type. Unknown target types fail the delivery. func (e *Engine) processDelivery( ctx context.Context, webhookDB *gorm.DB, d *database.Delivery, task *Task, ) { - switch d.Target.Type { - case database.TargetTypeHTTP: - e.deliverHTTP(ctx, webhookDB, d, task) - case database.TargetTypeDatabase: - e.deliverDatabase(webhookDB, d) - case database.TargetTypeLog: - e.deliverLog(webhookDB, d) - case database.TargetTypeSlack: - e.deliverSlack(ctx, webhookDB, d) - default: + target, ok := e.targets[d.Target.Type] + if !ok { e.log.Error( "unknown target type", "target_id", d.TargetID, @@ -784,481 +729,16 @@ func (e *Engine) processDelivery( e.updateDeliveryStatus( webhookDB, d, database.DeliveryStatusFailed, ) - } -} - -func (e *Engine) deliverHTTP( - ctx context.Context, - webhookDB *gorm.DB, - d *database.Delivery, - task *Task, -) { - cfg, err := e.parseHTTPConfig(d.Target.Config) - if err != nil { - e.log.Error( - "invalid HTTP target config", - "target_id", d.TargetID, - "error", err, - ) - - e.recordResult( - webhookDB, d, task.AttemptNum, - false, 0, "", err.Error(), 0, - ) - - e.updateDeliveryStatus( - webhookDB, d, database.DeliveryStatusFailed, - ) return } - if d.Target.MaxRetries == 0 { - e.deliverHTTPFireAndForget( - ctx, webhookDB, d, cfg, - ) - - return - } - - e.deliverHTTPWithRetry( - ctx, webhookDB, d, task, cfg, - ) -} - -func (e *Engine) deliverHTTPFireAndForget( - ctx context.Context, - webhookDB *gorm.DB, - d *database.Delivery, - cfg *HTTPTargetConfig, -) { - statusCode, respBody, duration, reqErr := - e.doHTTPRequest(ctx, cfg, &d.Event) - - success := reqErr == nil && - statusCode >= httpSuccessMin && - statusCode < httpSuccessMax - - errMsg := "" - if reqErr != nil { - errMsg = reqErr.Error() - } - - e.recordResult( - webhookDB, d, 1, success, - statusCode, respBody, errMsg, duration, - ) - - if success { - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusDelivered, - ) - } else { - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) - } -} - -func (e *Engine) deliverHTTPWithRetry( - ctx context.Context, - webhookDB *gorm.DB, - d *database.Delivery, - task *Task, - cfg *HTTPTargetConfig, -) { - cb := e.getCircuitBreaker(task.TargetID) - if e.circuitBreakerBlock( - webhookDB, d, task, cb, - ) { - return - } - - attemptNum := task.AttemptNum - - statusCode, respBody, duration, reqErr := - e.doHTTPRequest(ctx, cfg, &d.Event) - - success := reqErr == nil && - statusCode >= httpSuccessMin && - statusCode < httpSuccessMax - - errMsg := "" - if reqErr != nil { - errMsg = reqErr.Error() - } - - e.recordResult( - webhookDB, d, attemptNum, success, - statusCode, respBody, errMsg, duration, - ) - - if success { - cb.RecordSuccess() - - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusDelivered, - ) - - return - } - - cb.RecordFailure() - e.handleHTTPRetry(webhookDB, d, task, attemptNum) -} - -func (e *Engine) circuitBreakerBlock( - webhookDB *gorm.DB, - d *database.Delivery, - task *Task, - cb *CircuitBreaker, -) bool { - if cb.Allow() { - return false - } - - remaining := cb.CooldownRemaining() - - e.log.Info( - "circuit breaker open, skipping delivery", - "target_id", task.TargetID, - "target_name", task.TargetName, - "delivery_id", d.ID, - "cooldown_remaining", remaining, - ) - - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusRetrying, - ) - - retryTask := *task - e.scheduleRetry(retryTask, remaining) - - return true -} - -func (e *Engine) handleHTTPRetry( - webhookDB *gorm.DB, - d *database.Delivery, - task *Task, - attemptNum int, -) { - if attemptNum >= d.Target.MaxRetries { - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) - - return - } - - e.updateDeliveryStatus( - webhookDB, d, database.DeliveryStatusRetrying, - ) - - backoff := calcBackoff(attemptNum) - - retryTask := *task - retryTask.AttemptNum = attemptNum + 1 - e.scheduleRetry(retryTask, backoff) -} - -func (e *Engine) getCircuitBreaker( - targetID string, -) *CircuitBreaker { - if val, ok := e.circuitBreakers.Load(targetID); ok { - cb, _ := val.(*CircuitBreaker) - - return cb - } - - fresh := NewCircuitBreaker() - - actual, _ := e.circuitBreakers.LoadOrStore( - targetID, fresh, - ) - - cb, _ := actual.(*CircuitBreaker) - - return cb -} - -func (e *Engine) deliverDatabase( - webhookDB *gorm.DB, d *database.Delivery, -) { - e.recordResult( - webhookDB, d, 1, true, 0, "", "", 0, - ) - - e.updateDeliveryStatus( - webhookDB, d, database.DeliveryStatusDelivered, - ) -} - -func (e *Engine) deliverLog( - webhookDB *gorm.DB, d *database.Delivery, -) { - e.log.Info( - "webhook event delivered to log target", - "delivery_id", d.ID, - "event_id", d.EventID, - "target_id", d.TargetID, - "target_name", d.Target.Name, - "method", d.Event.Method, - "content_type", d.Event.ContentType, - "body_length", len(d.Event.Body), - ) - - e.recordResult( - webhookDB, d, 1, true, 0, "", "", 0, - ) - - e.updateDeliveryStatus( - webhookDB, d, database.DeliveryStatusDelivered, - ) -} - -func (e *Engine) deliverSlack( - ctx context.Context, - webhookDB *gorm.DB, - d *database.Delivery, -) { - cfg, err := e.parseSlackConfig(d.Target.Config) - if err != nil { - e.log.Error( - "invalid Slack target config", - "target_id", d.TargetID, - "error", err, - ) - - e.recordResult( - webhookDB, d, 1, - false, 0, "", err.Error(), 0, - ) - - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) - - return - } - - msg := FormatSlackMessage(&d.Event) - - payload, err := json.Marshal( - map[string]string{"text": msg}, - ) - if err != nil { - e.log.Error( - "failed to marshal Slack payload", - "target_id", d.TargetID, - "error", err, - ) - - e.recordResult( - webhookDB, d, 1, - false, 0, "", err.Error(), 0, - ) - - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) - - return - } - - e.sendSlackRequest( - ctx, webhookDB, d, cfg, payload, - ) -} - -func (e *Engine) sendSlackRequest( - ctx context.Context, - webhookDB *gorm.DB, - d *database.Delivery, - cfg *SlackTargetConfig, - payload []byte, -) { - start := time.Now() - - req, err := http.NewRequestWithContext( - ctx, - http.MethodPost, - cfg.WebhookURL, - bytes.NewReader(payload), - ) - if err != nil { - e.failSlackDelivery( - webhookDB, d, err.Error(), 0, - ) - - return - } - - req.Header.Set("Content-Type", "application/json") - req.Header.Set("User-Agent", "webhooker/1.0") - - resp, doErr := e.executeRequest(req) - durationMs := time.Since(start).Milliseconds() - - if doErr != nil { - errStr := fmt.Errorf( - "sending request: %w", doErr, - ).Error() - - e.failSlackDelivery( - webhookDB, d, errStr, durationMs, - ) - - return - } - - defer func() { _ = resp.Body.Close() }() - - e.handleSlackResponse( - webhookDB, d, resp, durationMs, - ) -} - -func (e *Engine) failSlackDelivery( - webhookDB *gorm.DB, - d *database.Delivery, - errMsg string, - durationMs int64, -) { - e.recordResult( - webhookDB, d, 1, - false, 0, "", errMsg, durationMs, - ) - - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) -} - -func (e *Engine) handleSlackResponse( - webhookDB *gorm.DB, - d *database.Delivery, - resp *http.Response, - durationMs int64, -) { - body, readErr := io.ReadAll( - io.LimitReader(resp.Body, maxBodyLog), - ) - if readErr != nil { - e.log.Error( - "failed to read Slack response body", - "error", readErr, - ) - } - - respBody := string(body) - - success := resp.StatusCode >= httpSuccessMin && - resp.StatusCode < httpSuccessMax - - errMsg := "" - if !success { - errMsg = fmt.Sprintf("HTTP %d", resp.StatusCode) - } - - e.recordResult( - webhookDB, d, 1, success, - resp.StatusCode, respBody, errMsg, durationMs, - ) - - if success { - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusDelivered, - ) - } else { - e.updateDeliveryStatus( - webhookDB, d, - database.DeliveryStatusFailed, - ) - } -} - -func (e *Engine) parseSlackConfig( - configJSON string, -) (*SlackTargetConfig, error) { - if configJSON == "" { - return nil, errEmptyTargetConfig - } - - var cfg SlackTargetConfig - - err := json.Unmarshal( - []byte(configJSON), &cfg, - ) - if err != nil { - return nil, fmt.Errorf( - "parsing config JSON: %w", err, - ) - } - - if cfg.WebhookURL == "" { - return nil, errMissingWebhookURL - } - - return &cfg, nil -} - -func (e *Engine) doHTTPRequest( - ctx context.Context, - cfg *HTTPTargetConfig, - event *database.Event, -) (int, string, int64, error) { - start := time.Now() - - req, reqErr := http.NewRequestWithContext( - ctx, - http.MethodPost, - cfg.URL, - bytes.NewReader([]byte(event.Body)), - ) - if reqErr != nil { - return 0, "", 0, fmt.Errorf( - "creating request: %w", reqErr, - ) - } - - applyRequestHeaders(req, event, cfg) - - client := e.clientForConfig(cfg) - - resp, doErr := executeHTTPRequest(client, req) - - dur := time.Since(start).Milliseconds() - if doErr != nil { - return 0, "", dur, fmt.Errorf( - "sending request: %w", doErr, - ) - } - - defer func() { _ = resp.Body.Close() }() - - body, readErr := io.ReadAll( - io.LimitReader(resp.Body, maxBodyLog), - ) - if readErr != nil { - return resp.StatusCode, "", dur, - fmt.Errorf( - "reading response body: %w", readErr, - ) - } - - return resp.StatusCode, string(body), dur, nil + target.Deliver(ctx, webhookDB, d, task, e) } +// recordResult persists a DeliveryResult row describing a +// single attempt. It is a cross-target helper the targets +// call. func (e *Engine) recordResult( webhookDB *gorm.DB, d *database.Delivery, @@ -1288,6 +768,8 @@ func (e *Engine) recordResult( } } +// updateDeliveryStatus persists a new status for a delivery. +// It is a cross-target helper the targets call. func (e *Engine) updateDeliveryStatus( webhookDB *gorm.DB, d *database.Delivery, @@ -1305,45 +787,6 @@ func (e *Engine) updateDeliveryStatus( } } -func (e *Engine) parseHTTPConfig( - configJSON string, -) (*HTTPTargetConfig, error) { - if configJSON == "" { - return nil, errEmptyTargetConfig - } - - var cfg HTTPTargetConfig - - err := json.Unmarshal( - []byte(configJSON), &cfg, - ) - if err != nil { - return nil, fmt.Errorf( - "parsing config JSON: %w", err, - ) - } - - if cfg.URL == "" { - return nil, errMissingTargetURL - } - - return &cfg, nil -} - -// isForwardableHeader returns true if the header should -// be forwarded to targets. -func isForwardableHeader(name string) bool { - switch http.CanonicalHeaderKey(name) { - case "Host", "Connection", "Keep-Alive", - "Transfer-Encoding", "Te", "Trailer", - "Upgrade", "Proxy-Authorization", - "Proxy-Connection", "Content-Length": - return false - default: - return true - } -} - func truncate(s string, maxLen int) string { if len(s) <= maxLen { return s @@ -1356,9 +799,10 @@ func truncate(s string, maxLen int) string { func buildEventFromTask(task *Task) database.Event { event := database.Event{ - Method: task.Method, - Headers: task.Headers, - ContentType: task.ContentType, + EntrypointID: task.EntrypointID, + Method: task.Method, + Headers: task.Headers, + ContentType: task.ContentType, } event.ID = task.EventID @@ -1466,55 +910,6 @@ func (e *Engine) loadTarget( return target, nil } -func calcBackoff(attemptNum int) time.Duration { - shift := max(attemptNum-1, 0) - shift = min(shift, maxBackoffShift) - - return time.Duration(1<= backoff -} - func buildRecoveryTask( d *database.Delivery, webhookID string, @@ -1533,6 +928,7 @@ func buildRecoveryTask( DeliveryID: d.ID, EventID: d.EventID, WebhookID: webhookID, + EntrypointID: event.EntrypointID, TargetID: target.ID, TargetName: target.Name, TargetType: target.Type, @@ -1628,120 +1024,3 @@ func (e *Engine) sendRecoveredDeliveries( } } } - -func formatJSONBody(body string) string { - var parsed json.RawMessage - if json.Unmarshal([]byte(body), &parsed) != nil { - return "" - } - - var pretty bytes.Buffer - if json.Indent(&pretty, parsed, "", " ") != nil { - return "" - } - - var b strings.Builder - - b.WriteString("\n```\n") - - prettyStr := pretty.String() - - const maxPayloadDisplay = 3500 - if len(prettyStr) > maxPayloadDisplay { - b.WriteString(prettyStr[:maxPayloadDisplay]) - b.WriteString("\n... (truncated)") - } else { - b.WriteString(prettyStr) - } - - b.WriteString("\n```\n") - - return b.String() -} - -func formatRawBody(b *strings.Builder, body string) { - b.WriteString("\n```\n") - - const maxRawDisplay = 3500 - if len(body) > maxRawDisplay { - b.WriteString(body[:maxRawDisplay]) - b.WriteString("\n... (truncated)") - } else { - b.WriteString(body) - } - - b.WriteString("\n```\n") -} - -func applyRequestHeaders( - req *http.Request, - event *database.Event, - cfg *HTTPTargetConfig, -) { - if event.ContentType != "" { - req.Header.Set( - "Content-Type", event.ContentType, - ) - } - - var originalHeaders map[string][]string - - if event.Headers != "" { - jsonErr := json.Unmarshal( - []byte(event.Headers), - &originalHeaders, - ) - if jsonErr == nil { - for k, vals := range originalHeaders { - if isForwardableHeader(k) { - for _, v := range vals { - req.Header.Add(k, v) - } - } - } - } - } - - for k, v := range cfg.Headers { - req.Header.Set(k, v) - } - - req.Header.Set("User-Agent", "webhooker/1.0") -} - -func (e *Engine) clientForConfig( - cfg *HTTPTargetConfig, -) *http.Client { - if cfg.Timeout > 0 { - // Reuse the shared client's SSRF-safe transport so - // a per-target timeout does not drop the - // request-time private-IP guard. Only the timeout - // is overridden. - return &http.Client{ - Timeout: time.Duration( - cfg.Timeout, - ) * time.Second, - Transport: e.client.Transport, - } - } - - return e.client -} - -// executeRequest sends an HTTP request using the engine's -// default client. URLs are validated by SSRF-safe -// transport and config parsers before reaching here. -func (e *Engine) executeRequest( - req *http.Request, -) (*http.Response, error) { - return e.client.Do(req) //#nosec G704 -- URL validated by parseSlackConfig and SSRF-safe transport -} - -// executeHTTPRequest sends an HTTP request using the -// provided client. URLs are validated by config parsers -// and SSRF-safe transport before reaching here. -func executeHTTPRequest( - client *http.Client, req *http.Request, -) (*http.Response, error) { - return client.Do(req) //#nosec G704 -- URL validated by parseHTTPConfig and SSRF-safe transport -} diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index 00ce06b..25b2c8b 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -1,6 +1,7 @@ package delivery_test import ( + "bytes" "context" "database/sql" "encoding/json" @@ -1652,6 +1653,179 @@ func TestProcessDelivery_RoutesToSlack(t *testing.T) { ) } +// newLogCaptureEngine builds a test engine whose logger +// writes to the returned buffer, for inspecting log output. +func newLogCaptureEngine( + t *testing.T, +) (*delivery.Engine, *bytes.Buffer) { + t.Helper() + + var buf bytes.Buffer + + log := slog.New(slog.NewTextHandler( + &buf, + &slog.HandlerOptions{Level: slog.LevelDebug}, + )) + + e := delivery.NewTestEngine( + log, &http.Client{Timeout: 5 * time.Second}, 1, + ) + + return e, &buf +} + +// assertLogLineComplete asserts the captured log output +// carries the full inbound webhook content and ids. +func assertLogLineComplete( + t *testing.T, out string, event database.Event, +) { + t.Helper() + + assert.Contains(t, out, "log-body-marker", + "log line must contain the full request body", + ) + + assert.Contains(t, out, "Content-Type", + "log line must contain the full request headers", + ) + + assert.Contains(t, out, event.EntrypointID, + "log line must contain the entrypoint id", + ) + + assert.Contains(t, out, event.WebhookID, + "log line must contain the webhook id", + ) + + assert.Contains(t, out, "application/json", + "log line must contain the content type", + ) +} + +func TestDeliverLog_LogsFullContent(t *testing.T) { + t.Parallel() + + db := testWebhookDB(t) + e, buf := newLogCaptureEngine(t) + + event := seedEvent( + t, db, `{"log-body-marker":"abc123"}`, + ) + + dlv := seedDelivery( + t, db, event.ID, uuid.New().String(), + database.DeliveryStatusPending, + ) + + d := &database.Delivery{ + EventID: event.ID, + TargetID: dlv.TargetID, + Status: database.DeliveryStatusPending, + Event: event, + Target: database.Target{ + Name: "test-log-full", + Type: database.TargetTypeLog, + }, + } + d.ID = dlv.ID + + e.ExportDeliverLog(db, d) + + assertLogLineComplete(t, buf.String(), event) + + assertDeliveryStatus(t, db, dlv.ID, + database.DeliveryStatusDelivered, + ) +} + +// buildSlackRetryDelivery builds a Slack delivery whose +// target is configured with retries enabled. +func buildSlackRetryDelivery( + dlv database.Delivery, + event database.Event, + targetID, cfg string, +) *database.Delivery { + d := &database.Delivery{ + EventID: event.ID, + TargetID: targetID, + Status: database.DeliveryStatusPending, + Event: event, + Target: database.Target{ + Name: "test-slack-retry", + Type: database.TargetTypeSlack, + Config: cfg, + MaxRetries: 5, + }, + } + d.ID = dlv.ID + + return d +} + +func TestDeliverSlack_WithRetries_SchedulesRetry( + t *testing.T, +) { + t.Parallel() + + db := testWebhookDB(t) + ts := newStatusServer(t, http.StatusServiceUnavailable) + e := testEngine(t, 1) + targetID := uuid.New().String() + + slackCfg, err := json.Marshal( + delivery.SlackTargetConfig{WebhookURL: ts.URL}, + ) + require.NoError(t, err) + + event := seedEvent(t, db, `{"slack":"retry"}`) + + dlv := seedDelivery( + t, db, event.ID, targetID, + database.DeliveryStatusPending, + ) + + d := buildSlackRetryDelivery( + dlv, event, targetID, string(slackCfg), + ) + + task := &delivery.Task{ + DeliveryID: dlv.ID, + TargetID: targetID, + TargetType: database.TargetTypeSlack, + MaxRetries: 5, + AttemptNum: 1, + } + + e.ExportProcessDelivery(context.TODO(), db, d, task) + + assertDeliveryStatus(t, db, dlv.ID, + database.DeliveryStatusRetrying, + ) + + assertDeliveryResult( + t, db, dlv.ID, false, + http.StatusServiceUnavailable, + ) +} + +// newStatusServer starts a test server that always responds +// with the given status code. +func newStatusServer( + t *testing.T, code int, +) *httptest.Server { + t.Helper() + + ts := httptest.NewServer(http.HandlerFunc( + func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(code) + }, + )) + + t.Cleanup(ts.Close) + + return ts +} + // readAll is a small helper to avoid importing io in // a test handler inline. func readAll(r interface { diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index f9b987d..ec5de60 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -39,37 +39,50 @@ func ExportTruncate(s string, maxLen int) string { return truncate(s, maxLen) } -// ExportDeliverHTTP exposes deliverHTTP for testing. +// ExportDeliverHTTP delivers via the http target for testing. func (e *Engine) ExportDeliverHTTP( ctx context.Context, webhookDB *gorm.DB, d *database.Delivery, task *Task, ) { - e.deliverHTTP(ctx, webhookDB, d, task) + e.httpTarget.Deliver(ctx, webhookDB, d, task, e) } -// ExportDeliverDatabase exposes deliverDatabase. +// ExportDeliverDatabase delivers via the database target. func (e *Engine) ExportDeliverDatabase( webhookDB *gorm.DB, d *database.Delivery, ) { - e.deliverDatabase(webhookDB, d) + e.targets[database.TargetTypeDatabase].Deliver( + context.Background(), webhookDB, d, &Task{}, e, + ) } -// ExportDeliverLog exposes deliverLog for testing. +// ExportDeliverLog delivers via the log target for testing. func (e *Engine) ExportDeliverLog( webhookDB *gorm.DB, d *database.Delivery, ) { - e.deliverLog(webhookDB, d) + e.targets[database.TargetTypeLog].Deliver( + context.Background(), webhookDB, d, &Task{}, e, + ) } -// ExportDeliverSlack exposes deliverSlack for testing. +// ExportDeliverSlack delivers via the slack target for +// testing. func (e *Engine) ExportDeliverSlack( ctx context.Context, webhookDB *gorm.DB, d *database.Delivery, ) { - e.deliverSlack(ctx, webhookDB, d) + task := &Task{ + DeliveryID: d.ID, + TargetID: d.TargetID, + AttemptNum: 1, + } + + e.targets[database.TargetTypeSlack].Deliver( + ctx, webhookDB, d, task, e, + ) } // ExportProcessNewTask exposes processNewTask. @@ -96,53 +109,56 @@ func (e *Engine) ExportProcessDelivery( e.processDelivery(ctx, webhookDB, d, task) } -// ExportGetCircuitBreaker exposes getCircuitBreaker. +// ExportGetCircuitBreaker exposes the http target's +// getCircuitBreaker. func (e *Engine) ExportGetCircuitBreaker( targetID string, ) *CircuitBreaker { - return e.getCircuitBreaker(targetID) + return e.httpTarget.getCircuitBreaker(targetID) } // ExportParseHTTPConfig exposes parseHTTPConfig. func (e *Engine) ExportParseHTTPConfig( configJSON string, ) (*HTTPTargetConfig, error) { - return e.parseHTTPConfig(configJSON) + return parseHTTPConfig(configJSON) } // ExportParseSlackConfig exposes parseSlackConfig. func (e *Engine) ExportParseSlackConfig( configJSON string, ) (*SlackTargetConfig, error) { - return e.parseSlackConfig(configJSON) + return parseSlackConfig(configJSON) } -// ExportDoHTTPRequest exposes doHTTPRequest. +// ExportDoHTTPRequest exposes the http target's +// doHTTPRequest. func (e *Engine) ExportDoHTTPRequest( ctx context.Context, cfg *HTTPTargetConfig, event *database.Event, ) (int, string, int64, error) { - return e.doHTTPRequest(ctx, cfg, event) + return e.httpTarget.doHTTPRequest(ctx, cfg, event) } -// ExportClientForConfig exposes clientForConfig. +// ExportClientForConfig exposes the http target's +// clientForConfig. func (e *Engine) ExportClientForConfig( cfg *HTTPTargetConfig, ) *http.Client { - return e.clientForConfig(cfg) + return e.httpTarget.clientForConfig(cfg) } -// ExportClient returns the engine's shared HTTP client. +// ExportClient returns the http target's shared HTTP client. func (e *Engine) ExportClient() *http.Client { - return e.client + return e.httpTarget.client } -// ExportScheduleRetry exposes scheduleRetry. +// ExportScheduleRetry exposes ScheduleRetry. func (e *Engine) ExportScheduleRetry( task Task, delay time.Duration, ) { - e.scheduleRetry(task, delay) + e.ScheduleRetry(task, delay) } // ExportRecoverPendingDeliveries exposes @@ -199,13 +215,15 @@ func NewTestEngine( client *http.Client, workers int, ) *Engine { - return &Engine{ + e := &Engine{ log: log, - client: client, deliveryCh: make(chan Task, deliveryChannelSize), retryCh: make(chan Task, retryChannelSize), workers: workers, } + e.initTargets(client) + + return e } // NewTestEngineSmallRetry creates an Engine with a tiny @@ -213,10 +231,13 @@ func NewTestEngine( func NewTestEngineSmallRetry( log *slog.Logger, ) *Engine { - return &Engine{ + e := &Engine{ log: log, retryCh: make(chan Task, 1), } + e.initTargets(nil) + + return e } // NewTestEngineWithDB creates an Engine with a real @@ -228,15 +249,17 @@ func NewTestEngineWithDB( client *http.Client, workers int, ) *Engine { - return &Engine{ + e := &Engine{ database: db, dbManager: dbMgr, log: log, - client: client, deliveryCh: make(chan Task, deliveryChannelSize), retryCh: make(chan Task, retryChannelSize), workers: workers, } + e.initTargets(client) + + return e } // NewTestCircuitBreaker creates a CircuitBreaker with diff --git a/internal/delivery/target.go b/internal/delivery/target.go new file mode 100644 index 0000000..264ce2f --- /dev/null +++ b/internal/delivery/target.go @@ -0,0 +1,101 @@ +package delivery + +import ( + "context" + "net/http" + "time" + + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +// Scheduler re-enqueues a task for a future delivery attempt. +// The engine provides one to each target so a target can own +// its retries durably: it records the attempt, marks the +// delivery retrying, and asks the Scheduler to deliver the +// next attempt after delay — exactly what the engine does for +// its own restart recovery. +type Scheduler interface { + ScheduleRetry(task Task, delay time.Duration) +} + +// Target delivers an event to one target type. Each type is +// an implementation. A Target owns its whole delivery: it +// makes the attempt, records the DeliveryResult and updates +// the DeliveryStatus, and — for targets that retry — decides +// whether to retry, computes its own backoff, gates with its +// own circuit breaker, and reschedules via the injected +// Scheduler. Fire-and-forget targets simply record a single +// attempt. +type Target interface { + Deliver( + ctx context.Context, + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, + ) +} + +// rescheduler is implemented by targets that own durable +// retries. The engine's restart recovery and periodic sweep +// use it to let the target recompute the schedule for an +// orphaned retrying delivery, keeping the retry schedule +// target-owned. Fire-and-forget targets do not implement it +// and their (never-occurring) retrying deliveries are +// skipped. +type rescheduler interface { + // remainingBackoff returns how long to wait before the + // next attempt of a recovered retrying delivery. + remainingBackoff( + webhookDB *gorm.DB, + deliveryID string, + attemptNum int, + ) time.Duration + + // backoffElapsed reports whether the backoff window for + // the last attempt has already passed, so the periodic + // sweep can re-enqueue the delivery now. + backoffElapsed( + webhookDB *gorm.DB, + deliveryID string, + attemptNum int, + ) bool +} + +// attemptResult is the outcome of a single delivery attempt, +// as reported by a target's per-attempt function to the +// shared retry core. +type attemptResult struct { + statusCode int + respBody string + duration int64 + success bool + errMsg string +} + +// initTargets builds the target registry, wiring each target +// to the engine's persistence helpers and giving the HTTP and +// Slack targets the shared SSRF-safe client. It is called by +// both New and the test constructors so the registry is +// always populated. +func (e *Engine) initTargets(client *http.Client) { + httpT := &httpTarget{ + httpCore: &httpCore{eng: e}, + client: client, + } + + slackT := &slackTarget{ + httpCore: &httpCore{eng: e}, + client: client, + } + + e.httpTarget = httpT + + e.targets = map[database.TargetType]Target{ + database.TargetTypeHTTP: httpT, + database.TargetTypeSlack: slackT, + database.TargetTypeDatabase: &databaseTarget{eng: e}, + database.TargetTypeLog: &logTarget{eng: e}, + } +} diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go new file mode 100644 index 0000000..a8b1d6a --- /dev/null +++ b/internal/delivery/target_database.go @@ -0,0 +1,34 @@ +package delivery + +import ( + "context" + + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +// databaseTarget is a fire-and-forget target: the event is +// already persisted in the per-webhook database by the time +// delivery runs, so the target records a single successful +// attempt. (Durable archiving to a separate store is tracked +// as its own work.) +type databaseTarget struct { + eng *Engine +} + +// Deliver implements Target. +func (t *databaseTarget) Deliver( + _ context.Context, + webhookDB *gorm.DB, + d *database.Delivery, + _ *Task, + _ Scheduler, +) { + t.eng.recordResult( + webhookDB, d, 1, true, 0, "", "", 0, + ) + + t.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusDelivered, + ) +} diff --git a/internal/delivery/target_http.go b/internal/delivery/target_http.go new file mode 100644 index 0000000..9a58a8f --- /dev/null +++ b/internal/delivery/target_http.go @@ -0,0 +1,499 @@ +package delivery + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "sync" + "time" + + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +// Sentinel errors returned by the config parsers. +var ( + errEmptyTargetConfig = errors.New( + "empty target config", + ) + errMissingTargetURL = errors.New( + "target URL is required", + ) +) + +// HTTPTargetConfig holds configuration for http target +// types. +type HTTPTargetConfig struct { + URL string `json:"url"` + Headers map[string]string `json:"headers,omitempty"` + Timeout int `json:"timeout,omitempty"` +} + +// httpCore holds the retry, backoff, and circuit-breaker +// machinery shared by the HTTP and Slack targets. Each of +// those targets owns its own httpCore instance (and thus its +// own circuit breakers); the per-attempt request differs +// between them and is supplied as a closure. +type httpCore struct { + eng *Engine + + // circuitBreakers stores a *CircuitBreaker per target ID. + circuitBreakers sync.Map +} + +// deliver runs one delivery attempt through the retry core. +// A maxRetries of 0 is fire-and-forget: a single attempt is +// recorded and no circuit breaker is consulted. A positive +// maxRetries gates the attempt on the circuit breaker and +// schedules a backed-off retry on failure. +func (c *httpCore) deliver( + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, + maxRetries int, + attempt func() attemptResult, +) { + if maxRetries == 0 { + c.fireAndForget(webhookDB, d, attempt()) + + return + } + + c.withRetry( + webhookDB, d, task, sched, maxRetries, attempt, + ) +} + +func (c *httpCore) fireAndForget( + webhookDB *gorm.DB, + d *database.Delivery, + res attemptResult, +) { + c.eng.recordResult( + webhookDB, d, 1, res.success, + res.statusCode, res.respBody, res.errMsg, + res.duration, + ) + + if res.success { + c.eng.updateDeliveryStatus( + webhookDB, d, + database.DeliveryStatusDelivered, + ) + + return + } + + c.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusFailed, + ) +} + +func (c *httpCore) withRetry( + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, + maxRetries int, + attempt func() attemptResult, +) { + cb := c.getCircuitBreaker(task.TargetID) + if c.circuitBreakerBlock(webhookDB, d, task, sched, cb) { + return + } + + attemptNum := task.AttemptNum + + res := attempt() + + c.eng.recordResult( + webhookDB, d, attemptNum, res.success, + res.statusCode, res.respBody, res.errMsg, + res.duration, + ) + + if res.success { + cb.RecordSuccess() + + c.eng.updateDeliveryStatus( + webhookDB, d, + database.DeliveryStatusDelivered, + ) + + return + } + + cb.RecordFailure() + + c.handleRetry( + webhookDB, d, task, sched, maxRetries, attemptNum, + ) +} + +func (c *httpCore) circuitBreakerBlock( + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, + cb *CircuitBreaker, +) bool { + if cb.Allow() { + return false + } + + remaining := cb.CooldownRemaining() + + c.eng.log.Info( + "circuit breaker open, skipping delivery", + "target_id", task.TargetID, + "target_name", task.TargetName, + "delivery_id", d.ID, + "cooldown_remaining", remaining, + ) + + c.eng.updateDeliveryStatus( + webhookDB, d, + database.DeliveryStatusRetrying, + ) + + retryTask := *task + sched.ScheduleRetry(retryTask, remaining) + + return true +} + +func (c *httpCore) handleRetry( + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, + maxRetries int, + attemptNum int, +) { + if attemptNum >= maxRetries { + c.eng.updateDeliveryStatus( + webhookDB, d, + database.DeliveryStatusFailed, + ) + + return + } + + c.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusRetrying, + ) + + backoff := calcBackoff(attemptNum) + + retryTask := *task + retryTask.AttemptNum = attemptNum + 1 + sched.ScheduleRetry(retryTask, backoff) +} + +func (c *httpCore) getCircuitBreaker( + targetID string, +) *CircuitBreaker { + if val, ok := c.circuitBreakers.Load(targetID); ok { + cb, _ := val.(*CircuitBreaker) + + return cb + } + + fresh := NewCircuitBreaker() + + actual, _ := c.circuitBreakers.LoadOrStore( + targetID, fresh, + ) + + cb, _ := actual.(*CircuitBreaker) + + return cb +} + +// remainingBackoff returns how long remains of the backoff +// window for the last attempt of a recovered retrying +// delivery. It implements rescheduler. +func (c *httpCore) remainingBackoff( + webhookDB *gorm.DB, + deliveryID string, + attemptNum int, +) time.Duration { + var lastResult database.DeliveryResult + + err := webhookDB. + Where("delivery_id = ?", deliveryID). + Order("created_at DESC"). + First(&lastResult).Error + if err != nil { + return 0 + } + + backoff := calcBackoff(attemptNum) + elapsed := time.Since(lastResult.CreatedAt) + remaining := backoff - elapsed + + return max(remaining, 0) +} + +// backoffElapsed reports whether the backoff window for the +// last attempt of a retrying delivery has passed. It +// implements rescheduler. +func (c *httpCore) backoffElapsed( + webhookDB *gorm.DB, + deliveryID string, + attemptNum int, +) bool { + var lastResult database.DeliveryResult + + err := webhookDB. + Where("delivery_id = ?", deliveryID). + Order("created_at DESC"). + First(&lastResult).Error + if err != nil { + return true + } + + backoff := calcBackoff(attemptNum) + + return time.Since(lastResult.CreatedAt) >= backoff +} + +func calcBackoff(attemptNum int) time.Duration { + shift := max(attemptNum-1, 0) + shift = min(shift, maxBackoffShift) + + return time.Duration(1<= httpSuccessMin && + statusCode < httpSuccessMax + + errMsg := "" + if reqErr != nil { + errMsg = reqErr.Error() + } + + return attemptResult{ + statusCode: statusCode, + respBody: respBody, + duration: duration, + success: success, + errMsg: errMsg, + } +} + +func (t *httpTarget) doHTTPRequest( + ctx context.Context, + cfg *HTTPTargetConfig, + event *database.Event, +) (int, string, int64, error) { + start := time.Now() + + req, reqErr := http.NewRequestWithContext( + ctx, + http.MethodPost, + cfg.URL, + bytes.NewReader([]byte(event.Body)), + ) + if reqErr != nil { + return 0, "", 0, fmt.Errorf( + "creating request: %w", reqErr, + ) + } + + applyRequestHeaders(req, event, cfg) + + client := t.clientForConfig(cfg) + + resp, doErr := executeHTTPRequest(client, req) + + dur := time.Since(start).Milliseconds() + if doErr != nil { + return 0, "", dur, fmt.Errorf( + "sending request: %w", doErr, + ) + } + + defer func() { _ = resp.Body.Close() }() + + body, readErr := io.ReadAll( + io.LimitReader(resp.Body, maxBodyLog), + ) + if readErr != nil { + return resp.StatusCode, "", dur, + fmt.Errorf( + "reading response body: %w", readErr, + ) + } + + return resp.StatusCode, string(body), dur, nil +} + +func (t *httpTarget) clientForConfig( + cfg *HTTPTargetConfig, +) *http.Client { + if cfg.Timeout > 0 { + // Reuse the shared client's SSRF-safe transport so + // a per-target timeout does not drop the + // request-time private-IP guard. Only the timeout + // is overridden. + return &http.Client{ + Timeout: time.Duration( + cfg.Timeout, + ) * time.Second, + Transport: t.client.Transport, + } + } + + return t.client +} + +func parseHTTPConfig( + configJSON string, +) (*HTTPTargetConfig, error) { + if configJSON == "" { + return nil, errEmptyTargetConfig + } + + var cfg HTTPTargetConfig + + err := json.Unmarshal( + []byte(configJSON), &cfg, + ) + if err != nil { + return nil, fmt.Errorf( + "parsing config JSON: %w", err, + ) + } + + if cfg.URL == "" { + return nil, errMissingTargetURL + } + + return &cfg, nil +} + +// isForwardableHeader returns true if the header should +// be forwarded to targets. +func isForwardableHeader(name string) bool { + switch http.CanonicalHeaderKey(name) { + case "Host", "Connection", "Keep-Alive", + "Transfer-Encoding", "Te", "Trailer", + "Upgrade", "Proxy-Authorization", + "Proxy-Connection", "Content-Length": + return false + default: + return true + } +} + +func applyRequestHeaders( + req *http.Request, + event *database.Event, + cfg *HTTPTargetConfig, +) { + if event.ContentType != "" { + req.Header.Set( + "Content-Type", event.ContentType, + ) + } + + var originalHeaders map[string][]string + + if event.Headers != "" { + jsonErr := json.Unmarshal( + []byte(event.Headers), + &originalHeaders, + ) + if jsonErr == nil { + for k, vals := range originalHeaders { + if isForwardableHeader(k) { + for _, v := range vals { + req.Header.Add(k, v) + } + } + } + } + } + + for k, v := range cfg.Headers { + req.Header.Set(k, v) + } + + req.Header.Set("User-Agent", "webhooker/1.0") +} + +// executeHTTPRequest sends an HTTP request using the provided +// client. URLs are validated by the config parsers and the +// SSRF-safe transport before reaching here. +func executeHTTPRequest( + client *http.Client, req *http.Request, +) (*http.Response, error) { + return client.Do(req) //#nosec G704 -- URL validated by parseHTTPConfig/parseSlackConfig and SSRF-safe transport +} diff --git a/internal/delivery/target_log.go b/internal/delivery/target_log.go new file mode 100644 index 0000000..2484bec --- /dev/null +++ b/internal/delivery/target_log.go @@ -0,0 +1,47 @@ +package delivery + +import ( + "context" + + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +// logTarget is a fire-and-forget target that logs the entire +// inbound webhook — the full request body and headers, plus +// the method, content type, and the webhook and entrypoint +// ids — then records a single successful attempt. +type logTarget struct { + eng *Engine +} + +// Deliver implements Target. +func (t *logTarget) Deliver( + _ context.Context, + webhookDB *gorm.DB, + d *database.Delivery, + _ *Task, + _ Scheduler, +) { + t.eng.log.Info( + "webhook event delivered to log target", + "delivery_id", d.ID, + "event_id", d.EventID, + "target_id", d.TargetID, + "target_name", d.Target.Name, + "webhook_id", d.Event.WebhookID, + "entrypoint_id", d.Event.EntrypointID, + "method", d.Event.Method, + "content_type", d.Event.ContentType, + "headers", d.Event.Headers, + "body", d.Event.Body, + ) + + t.eng.recordResult( + webhookDB, d, 1, true, 0, "", "", 0, + ) + + t.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusDelivered, + ) +} diff --git a/internal/delivery/target_slack.go b/internal/delivery/target_slack.go new file mode 100644 index 0000000..fdb95f6 --- /dev/null +++ b/internal/delivery/target_slack.go @@ -0,0 +1,299 @@ +package delivery + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "strings" + "time" + + "gorm.io/gorm" + "sneak.berlin/go/webhooker/internal/database" +) + +// errMissingWebhookURL is returned when a Slack target config +// omits its webhook URL. +var errMissingWebhookURL = errors.New( + "webhook_url is required", +) + +// SlackTargetConfig holds configuration for slack target +// types. +type SlackTargetConfig struct { + WebhookURL string `json:"webhookUrl"` +} + +// slackTarget delivers events to Slack incoming webhooks. It +// formats the event into a Slack message and posts it as +// JSON. It shares the retry core with the HTTP target: a +// MaxRetries of 0 stays single-attempt fire-and-forget +// (preserving existing Slack targets), while a positive +// MaxRetries adds backoff and circuit breaking. +type slackTarget struct { + *httpCore + + client *http.Client +} + +// Deliver implements Target. +func (t *slackTarget) Deliver( + ctx context.Context, + webhookDB *gorm.DB, + d *database.Delivery, + task *Task, + sched Scheduler, +) { + cfg, err := parseSlackConfig(d.Target.Config) + if err != nil { + t.eng.log.Error( + "invalid Slack target config", + "target_id", d.TargetID, + "error", err, + ) + + t.failConfig(webhookDB, d, err) + + return + } + + msg := FormatSlackMessage(&d.Event) + + payload, err := json.Marshal( + map[string]string{"text": msg}, + ) + if err != nil { + t.eng.log.Error( + "failed to marshal Slack payload", + "target_id", d.TargetID, + "error", err, + ) + + t.failConfig(webhookDB, d, err) + + return + } + + attempt := func() attemptResult { + return t.attempt(ctx, cfg, payload) + } + + t.deliver( + webhookDB, d, task, sched, + d.Target.MaxRetries, attempt, + ) +} + +// failConfig records a first-attempt failure for a delivery +// that could not be prepared (bad config or unmarshalable +// payload) and marks it failed. +func (t *slackTarget) failConfig( + webhookDB *gorm.DB, + d *database.Delivery, + err error, +) { + t.eng.recordResult( + webhookDB, d, 1, + false, 0, "", err.Error(), 0, + ) + + t.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusFailed, + ) +} + +// attempt performs a single Slack POST and derives its +// outcome, preserving the engine's original semantics: a +// non-2xx response records an "HTTP " error string and +// a transport error records a "sending request" error. +func (t *slackTarget) attempt( + ctx context.Context, + cfg *SlackTargetConfig, + payload []byte, +) attemptResult { + start := time.Now() + + req, err := http.NewRequestWithContext( + ctx, + http.MethodPost, + cfg.WebhookURL, + bytes.NewReader(payload), + ) + if err != nil { + return attemptResult{ + success: false, + errMsg: err.Error(), + } + } + + req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", "webhooker/1.0") + + resp, doErr := executeHTTPRequest(t.client, req) + durationMs := time.Since(start).Milliseconds() + + if doErr != nil { + return attemptResult{ + success: false, + duration: durationMs, + errMsg: fmt.Errorf( + "sending request: %w", doErr, + ).Error(), + } + } + + defer func() { _ = resp.Body.Close() }() + + return t.readSlackResponse(resp, durationMs) +} + +func (t *slackTarget) readSlackResponse( + resp *http.Response, + durationMs int64, +) attemptResult { + body, readErr := io.ReadAll( + io.LimitReader(resp.Body, maxBodyLog), + ) + if readErr != nil { + t.eng.log.Error( + "failed to read Slack response body", + "error", readErr, + ) + } + + success := resp.StatusCode >= httpSuccessMin && + resp.StatusCode < httpSuccessMax + + errMsg := "" + if !success { + errMsg = fmt.Sprintf("HTTP %d", resp.StatusCode) + } + + return attemptResult{ + statusCode: resp.StatusCode, + respBody: string(body), + duration: durationMs, + success: success, + errMsg: errMsg, + } +} + +func parseSlackConfig( + configJSON string, +) (*SlackTargetConfig, error) { + if configJSON == "" { + return nil, errEmptyTargetConfig + } + + var cfg SlackTargetConfig + + err := json.Unmarshal( + []byte(configJSON), &cfg, + ) + if err != nil { + return nil, fmt.Errorf( + "parsing config JSON: %w", err, + ) + } + + if cfg.WebhookURL == "" { + return nil, errMissingWebhookURL + } + + return &cfg, nil +} + +// FormatSlackMessage builds a Slack-compatible message +// string from a webhook event. +func FormatSlackMessage( + event *database.Event, +) string { + var b strings.Builder + + b.WriteString("*Webhook Event Received*\n") + + fmt.Fprintf( + &b, "*Method:* `%s`\n", event.Method, + ) + + fmt.Fprintf( + &b, + "*Content-Type:* `%s`\n", + event.ContentType, + ) + + fmt.Fprintf( + &b, + "*Timestamp:* `%s`\n", + event.CreatedAt.UTC().Format(time.RFC3339), + ) + + fmt.Fprintf( + &b, + "*Body Size:* %d bytes\n", + len(event.Body), + ) + + if event.Body == "" { + b.WriteString("\n_(empty body)_\n") + + return b.String() + } + + if formatted := formatJSONBody(event.Body); formatted != "" { + b.WriteString(formatted) + + return b.String() + } + + formatRawBody(&b, event.Body) + + return b.String() +} + +func formatJSONBody(body string) string { + var parsed json.RawMessage + if json.Unmarshal([]byte(body), &parsed) != nil { + return "" + } + + var pretty bytes.Buffer + if json.Indent(&pretty, parsed, "", " ") != nil { + return "" + } + + var b strings.Builder + + b.WriteString("\n```\n") + + prettyStr := pretty.String() + + const maxPayloadDisplay = 3500 + if len(prettyStr) > maxPayloadDisplay { + b.WriteString(prettyStr[:maxPayloadDisplay]) + b.WriteString("\n... (truncated)") + } else { + b.WriteString(prettyStr) + } + + b.WriteString("\n```\n") + + return b.String() +} + +func formatRawBody(b *strings.Builder, body string) { + b.WriteString("\n```\n") + + const maxRawDisplay = 3500 + if len(body) > maxRawDisplay { + b.WriteString(body[:maxRawDisplay]) + b.WriteString("\n... (truncated)") + } else { + b.WriteString(body) + } + + b.WriteString("\n```\n") +} diff --git a/internal/handlers/webhook.go b/internal/handlers/webhook.go index e7c59cd..d84ce90 100644 --- a/internal/handlers/webhook.go +++ b/internal/handlers/webhook.go @@ -329,6 +329,7 @@ func (h *Handlers) buildDeliveryTasks( DeliveryID: dlv.ID, EventID: event.ID, WebhookID: entrypoint.WebhookID, + EntrypointID: entrypoint.ID, TargetID: targets[i].ID, TargetName: targets[i].Name, TargetType: targets[i].Type,