Initial commit
Test / test (push) Successful in 7m5s
Release / gates (push) Successful in 7m28s
Release / build (amd64, freebsd) (push) Successful in 2m52s
Release / build (amd64, linux) (push) Successful in 2m46s
Release / build (arm64, freebsd) (push) Successful in 2m22s
Release / build (arm64, linux) (push) Successful in 2m38s
Release / build (loong64, linux) (push) Successful in 2m7s
Release / build (riscv64, linux) (push) Successful in 2m17s
Release / release (push) Successful in 1m0s

Assisted-by: GLM 5.3
This commit is contained in:
2026-09-29 10:03:32 +02:00
commit f8ed33df83
206 changed files with 44165 additions and 0 deletions
+86
View File
@@ -0,0 +1,86 @@
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
// SPDX-License-Identifier: PolyForm-Noncommercial-1.0.0
package webhooks
import (
"fmt"
"os"
"sourcedock.dev/petrbalvin/interpres/v2"
"sourcedock.dev/petrbalvin/volumen/internal/tomlfile"
)
// File is the admin-managed hook store: webhooks.toml beside the users
// file. The admin edits it through the settings forms, so a change
// applies without a restart; the file is written atomically with
// owner-only permissions, like the other state files.
type File struct {
Path string
}
// LoadFile reads the hook store. A missing file means no hooks, not an
// error; a malformed one is, because silently dropping the operator's
// endpoints would hide a real fault.
func LoadFile(path string) ([]Webhook, error) {
raw, err := os.ReadFile(path)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("read %s: %w", path, err)
}
data, err := interpres.ParseMap(raw)
if err != nil {
return nil, fmt.Errorf("parse %s: %w", path, err)
}
var out []Webhook
for _, entry := range tomlfile.Tables(data["webhooks"]) {
url := tomlfile.String(entry["url"])
if url == "" {
continue
}
out = append(out, Webhook{
URL: url,
Secret: tomlfile.String(entry["secret"]),
Events: stringList(entry["events"]),
Enabled: tomlfile.Bool(entry["enabled"], true),
})
}
return out, nil
}
// SaveFile replaces the hook store with the given hooks.
func SaveFile(path string, hooks []Webhook) error {
rows := make([]map[string]any, 0, len(hooks))
for _, hook := range hooks {
rows = append(rows, map[string]any{
"url": hook.URL,
"secret": hook.Secret,
"events": hook.Events,
"enabled": hook.Enabled,
})
}
if err := tomlfile.Write(path, "webhooks", rows); err != nil {
return fmt.Errorf("write %s: %w", path, err)
}
return nil
}
// stringList reads a string array that may arrive as either the []any
// the parser produces or a typed slice.
func stringList(value any) []string {
switch list := value.(type) {
case []string:
return list
case []any:
out := make([]string, 0, len(list))
for _, item := range list {
if s, ok := item.(string); ok && s != "" {
out = append(out, s)
}
}
return out
}
return nil
}
+268
View File
@@ -0,0 +1,268 @@
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (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=<hex digest of the raw body>
//
// 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)
}
+274
View File
@@ -0,0 +1,274 @@
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
// SPDX-License-Identifier: PolyForm-Noncommercial-1.0.0
package webhooks
import (
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"path/filepath"
"slices"
"strings"
"sync/atomic"
"testing"
"time"
)
// A hook with an event filter receives only those events, and one with no
// filter receives all of them. A disabled hook receives none.
func TestWebhookAccepts(t *testing.T) {
hook := Webhook{
URL: "https://example.com/hook", Secret: "s3cret",
Events: []string{"post.created"}, Enabled: true,
}
if !hook.Accepts("post.created") || hook.Accepts("post.deleted") {
t.Fatal("event filter broken")
}
all := Webhook{URL: "https://x.example", Enabled: true}
if !all.Accepts("anything") {
t.Fatalf("hook = %+v", all)
}
disabled := Webhook{URL: "https://x.example"}
if disabled.Accepts("post.created") {
t.Fatal("disabled hook accepts events")
}
}
func TestFireDeliversSignedPayload(t *testing.T) {
var gotBody []byte
var gotEvent, gotSig, gotDelivery string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotBody, _ = io.ReadAll(r.Body)
gotEvent = r.Header.Get("X-Volumen-Event")
gotSig = r.Header.Get("X-Volumen-Signature")
gotDelivery = r.Header.Get("X-Volumen-Delivery")
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
m := NewManager([]Webhook{{URL: srv.URL, Secret: "k", Enabled: true}}, "0.0.0-test")
if n := m.Fire("post.created", map[string]any{"post": map[string]any{"slug": "x"}}, true); n != 1 {
t.Fatalf("targets = %d", n)
}
if gotEvent != "post.created" || gotDelivery == "" {
t.Fatalf("headers: event=%q delivery=%q", gotEvent, gotDelivery)
}
if !strings.HasPrefix(gotSig, "sha256=") {
t.Fatalf("signature = %q", gotSig)
}
var payload map[string]any
if err := json.Unmarshal(gotBody, &payload); err != nil {
t.Fatalf("body: %v", err)
}
if payload["event"] != "post.created" || payload["version"] != "0.0.0-test" {
t.Fatalf("payload = %v", payload)
}
deliveries := m.Deliveries("")
if len(deliveries) != 1 || deliveries[0].Status != "ok" || deliveries[0].Attempts != 1 {
t.Fatalf("deliveries = %+v", deliveries)
}
}
func TestFireSkipsNonMatchingHooks(t *testing.T) {
var hits atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
hits.Add(1)
}))
defer srv.Close()
m := NewManager([]Webhook{{URL: srv.URL, Events: []string{"post.deleted"}, Enabled: true}}, "t")
if n := m.Fire("post.created", nil, true); n != 0 {
t.Fatalf("targets = %d", n)
}
if hits.Load() != 0 {
t.Fatal("non-matching hook was contacted")
}
}
// withoutBackoff removes the retry sleeps for the duration of one test,
// so the suite does not spend the configured backoff in real time.
func withoutBackoff(t *testing.T) {
t.Helper()
original := retryBackoff
retryBackoff = make([]time.Duration, len(original))
t.Cleanup(func() { retryBackoff = original })
}
func TestFireRetriesAndRecordsFailure(t *testing.T) {
withoutBackoff(t)
var attempts int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
atomic.AddInt32(&attempts, 1)
w.WriteHeader(http.StatusInternalServerError)
}))
defer srv.Close()
m := NewManager([]Webhook{{URL: srv.URL, Enabled: true}}, "t")
m.Fire("post.updated", nil, true)
if atomic.LoadInt32(&attempts) != maxAttempts {
t.Fatalf("attempts = %d, want %d", attempts, maxAttempts)
}
d := m.Deliveries("")
if len(d) != 1 || d[0].Status != "failed" || d[0].Attempts != maxAttempts {
t.Fatalf("deliveries = %+v", d)
}
if !strings.Contains(d[0].Error, "HTTP 500") {
t.Fatalf("error = %q", d[0].Error)
}
}
func TestFireRecordsConnectionError(t *testing.T) {
withoutBackoff(t)
m := NewManager([]Webhook{{URL: "http://127.0.0.1:1/none", Enabled: true}}, "t")
m.Fire("post.deleted", nil, true)
d := m.Deliveries("")
if len(d) != 1 || d[0].Status != "failed" || d[0].Error == "" {
t.Fatalf("deliveries = %+v", d)
}
}
func TestTestHookPingIgnoresFilters(t *testing.T) {
var gotEvent string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotEvent = r.Header.Get("X-Volumen-Event")
}))
defer srv.Close()
hook := Webhook{URL: srv.URL, Events: []string{"post.created"}, Enabled: true}
m := NewManager([]Webhook{hook}, "t")
m.TestHook(hook)
if gotEvent != "ping" {
t.Fatalf("event = %q", gotEvent)
}
}
func TestDeliveriesFilterByHook(t *testing.T) {
ok := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {}))
defer ok.Close()
m := NewManager([]Webhook{
{URL: ok.URL + "/a", Enabled: true},
{URL: ok.URL + "/b", Enabled: true},
}, "t")
m.Fire("post.created", nil, true)
if got := len(m.Deliveries(ok.URL + "/a")); got != 1 {
t.Fatalf("filtered deliveries = %d", got)
}
if got := len(m.Deliveries("")); got != 2 {
t.Fatalf("all deliveries = %d", got)
}
}
func TestHistoryIsCapped(t *testing.T) {
ok := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {}))
defer ok.Close()
m := NewManager([]Webhook{{URL: ok.URL, Enabled: true}}, "t")
for range historyLimit + 10 {
m.Fire("ping", nil, true)
}
if got := len(m.Deliveries("")); got != historyLimit {
t.Fatalf("history = %d, want %d", got, historyLimit)
}
}
// Production fires asynchronously; the delivery must still be recorded,
// and Wait must not return before it is.
func TestFireAsyncIsRecordedAndWaitedFor(t *testing.T) {
delivered := make(chan struct{}, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
select {
case delivered <- struct{}{}:
default:
}
}))
defer srv.Close()
m := NewManager([]Webhook{{URL: srv.URL, Enabled: true}}, "t")
if n := m.Fire("post.created", nil, false); n != 1 {
t.Fatalf("targets = %d", n)
}
m.Wait()
select {
case <-delivered:
default:
t.Fatal("Wait returned before the delivery reached the hook")
}
if d := m.Deliveries(""); len(d) != 1 || d[0].Status != "ok" {
t.Fatalf("deliveries = %+v", d)
}
}
// The hooks are handed out as a copy, so a caller cannot change what the
// manager delivers.
func TestHooksIsACopy(t *testing.T) {
m := NewManager([]Webhook{{URL: "https://example.com/h", Enabled: true}}, "t")
hooks := m.Hooks()
hooks[0].URL = "https://evil.example"
if got := m.Hooks()[0].URL; got != "https://example.com/h" {
t.Fatalf("hook mutated through the returned slice: %q", got)
}
}
// SetHooks swaps the delivery targets without building a new manager,
// which is how an admin settings change applies without a restart.
func TestSetHooksAppliesImmediately(t *testing.T) {
delivered := make(chan string, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
delivered <- r.Header.Get("X-Volumen-Event") + ":" + string(body)
}))
defer srv.Close()
m := NewManager(nil, "t")
if n := m.Fire("post.created", nil, false); n != 0 {
t.Fatalf("targets with no hooks = %d", n)
}
m.SetHooks([]Webhook{{URL: srv.URL, Enabled: true}})
if n := m.Fire("post.created", nil, false); n != 1 {
t.Fatalf("targets after SetHooks = %d", n)
}
m.Wait()
select {
case <-delivered:
default:
t.Fatal("the new hook did not receive the event")
}
}
// The admin-managed store round-trips through webhooks.toml: the fields
// survive a save and a missing file means no hooks, not an error.
func TestFileRoundTrip(t *testing.T) {
path := filepath.Join(t.TempDir(), "webhooks.toml")
hooks, err := LoadFile(path)
if err != nil || hooks != nil {
t.Fatalf("missing file: hooks = %+v, err = %v", hooks, err)
}
want := []Webhook{
{URL: "https://example.com/one", Secret: "k", Events: []string{"post.created"}, Enabled: true},
{URL: "https://example.com/two", Enabled: false},
}
if err := SaveFile(path, want); err != nil {
t.Fatalf("SaveFile: %v", err)
}
got, err := LoadFile(path)
if err != nil {
t.Fatalf("LoadFile: %v", err)
}
if len(got) != len(want) {
t.Fatalf("hooks = %+v, want %+v", got, want)
}
for i := range want {
if got[i].URL != want[i].URL || got[i].Secret != want[i].Secret ||
!slices.Equal(got[i].Events, want[i].Events) || got[i].Enabled != want[i].Enabled {
t.Fatalf("hook %d = %+v, want %+v", i, got[i], want[i])
}
}
}