// Copyright (c) 2026 Petr BalvĂ­n (https://petrbalvin.org) // SPDX-License-Identifier: PolyForm-Noncommercial-1.0.0 // Package webhooks notifies external services about post changes. // // Hooks come from two places with the same shape: [[webhooks]] tables in // config.toml (read-only, applied at start) and webhooks.toml beside the // users file, which the admin edits and which applies without a // restart. Each event payload is POSTed as JSON, signed with HMAC-SHA256 // when the hook carries a secret: // // X-Volumen-Event: post.created // X-Volumen-Delivery: 0f3b1c9e2a4d // X-Volumen-Signature: sha256= // // Delivery runs in background goroutines so admin requests never block // on a third-party endpoint; failures are retried with a short backoff // and recorded in an in-memory delivery history (lost on restart). package webhooks import ( "bytes" "crypto/hmac" "crypto/rand" "crypto/sha256" "encoding/hex" json "encoding/json/v2" "fmt" "io" "log/slog" "maps" "net/http" "slices" "strings" "sync" "time" ) const ( requestTimeout = 10 * time.Second maxAttempts = 3 historyLimit = 50 // maxConcurrentDeliveries bounds the deliveries in flight at once. maxConcurrentDeliveries = 8 ) var retryBackoff = []time.Duration{time.Second, 4 * time.Second} // Webhook is one configured endpoint. type Webhook struct { URL string Secret string Events []string // empty = every event Enabled bool } // Accepts reports whether the hook wants event (ignoring ping, which // is always delivered on demand). func (h Webhook) Accepts(event string) bool { if !h.Enabled { return false } if len(h.Events) == 0 { return true } return slices.Contains(h.Events, event) } // Delivery is the outcome of one delivery attempt series. type Delivery struct { ID string HookURL string Event string Timestamp string Status string // "ok" | "failed" StatusCode int Attempts int Error string } // Manager delivers signed JSON payloads to the configured webhooks. type Manager struct { mu sync.Mutex hooks []Webhook version string client *http.Client // deliveries bounds how many deliveries run at once, so a burst of // post writes cannot spawn an unbounded number of goroutines that // each retry with a backoff. deliveries chan struct{} inflight sync.WaitGroup history []Delivery } // NewManager builds the manager from the configured hooks. func NewManager(hooks []Webhook, version string) *Manager { return &Manager{ hooks: hooks, version: version, client: &http.Client{Timeout: requestTimeout}, deliveries: make(chan struct{}, maxConcurrentDeliveries), } } // SetHooks replaces the hook set: the admin edits the persisted list and // the change applies without a restart. Deliveries already in flight // finish against the hook they hold. func (m *Manager) SetHooks(hooks []Webhook) { m.mu.Lock() defer m.mu.Unlock() m.hooks = hooks } // Hooks returns a copy of the configured hooks. func (m *Manager) Hooks() []Webhook { m.mu.Lock() defer m.mu.Unlock() return slices.Clone(m.hooks) } // Wait blocks until the asynchronous deliveries started by Fire have // finished, so a short-lived command can flush before it exits. func (m *Manager) Wait() { m.inflight.Wait() } // Deliveries returns recorded deliveries, newest first, optionally // filtered by hook URL. func (m *Manager) Deliveries(hookURL string) []Delivery { m.mu.Lock() defer m.mu.Unlock() out := make([]Delivery, 0, len(m.history)) for _, v := range slices.Backward(m.history) { if hookURL == "" || v.HookURL == hookURL { out = append(out, v) } } return out } // Fire schedules event delivery to every matching hook and returns the // number of hooks contacted. With wait the deliveries run inline, so // short-lived CLI commands finish before the process exits. func (m *Manager) Fire(event string, payload map[string]any, wait bool) int { body := map[string]any{ "event": event, "timestamp": nowISO(), "version": m.version, } maps.Copy(body, payload) encoded, err := json.Marshal(body) if err != nil { slog.Warn("webhook: cannot encode payload", "error", err) return 0 } targets := 0 for _, hook := range m.Hooks() { if event != "ping" && !hook.Accepts(event) { continue } targets++ if wait { m.deliver(hook, event, encoded) continue } m.inflight.Go(func() { m.deliveries <- struct{}{} defer func() { <-m.deliveries }() m.deliver(hook, event, encoded) }) } return targets } // TestHook delivers a ping event to one hook, inline, regardless of // the hook's filters, and returns the delivery outcome. func (m *Manager) TestHook(hook Webhook) Delivery { encoded, err := json.Marshal(map[string]any{ "event": "ping", "timestamp": nowISO(), "version": m.version, }) if err != nil { return Delivery{HookURL: hook.URL, Event: "ping", Status: "failed", Error: err.Error()} } return m.deliver(hook, "ping", encoded) } func (m *Manager) deliver(hook Webhook, event string, body []byte) Delivery { deliveryID := randomID() startedAt := nowISO() lastError := "" statusCode := 0 for attempt := 1; attempt <= maxAttempts; attempt++ { code, err := m.post(hook, event, body, deliveryID) statusCode = code if err == nil && code >= 200 && code < 300 { return m.record(Delivery{ ID: deliveryID, HookURL: hook.URL, Event: event, Timestamp: startedAt, Status: "ok", StatusCode: code, Attempts: attempt, }) } if err != nil { lastError = err.Error() } else { lastError = fmt.Sprintf("HTTP %d", code) } if attempt < maxAttempts { time.Sleep(retryBackoff[attempt-1]) } } return m.record(Delivery{ ID: deliveryID, HookURL: hook.URL, Event: event, Timestamp: startedAt, Status: "failed", StatusCode: statusCode, Attempts: maxAttempts, Error: lastError, }) } func (m *Manager) post(hook Webhook, event string, body []byte, deliveryID string) (int, error) { req, err := http.NewRequest(http.MethodPost, hook.URL, bytes.NewReader(body)) if err != nil { return 0, err } req.Header.Set("Content-Type", "application/json") req.Header.Set("User-Agent", "volumen/"+m.version) req.Header.Set("X-Volumen-Event", event) req.Header.Set("X-Volumen-Delivery", deliveryID) if hook.Secret != "" { mac := hmac.New(sha256.New, []byte(hook.Secret)) mac.Write(body) req.Header.Set("X-Volumen-Signature", "sha256="+hex.EncodeToString(mac.Sum(nil))) } resp, err := m.client.Do(req) if err != nil { return 0, err } defer resp.Body.Close() _, _ = io.Copy(io.Discard, resp.Body) return resp.StatusCode, nil } func (m *Manager) record(delivery Delivery) Delivery { m.mu.Lock() defer m.mu.Unlock() m.history = append(m.history, delivery) if len(m.history) > historyLimit { m.history = m.history[len(m.history)-historyLimit:] } if delivery.Status == "ok" { slog.Info("webhook delivered", "url", delivery.HookURL, "event", delivery.Event, "status", delivery.StatusCode, "attempt", delivery.Attempts) } else { slog.Warn("webhook failed", "url", delivery.HookURL, "event", delivery.Event, "error", delivery.Error) } return delivery } func nowISO() string { return time.Now().UTC().Format("2006-01-02T15:04:05-07:00") } func randomID() string { buf := make([]byte, 6) if _, err := rand.Read(buf); err != nil { return strings.Repeat("0", 12) } return hex.EncodeToString(buf) }