Files
webhooker/internal/delivery/engine_integration_test.go
sneak 3296b166b1
All checks were successful
check / check (push) Successful in 3m18s
Render delivery attempt detail in the event log (closes #202)
Expanding a delivery on the event log page now shows each recorded
attempt: attempt number, outcome, status code, duration, error and
response body. Previously a failure rendered as "target: failed" and
diagnosing it meant opening the per-webhook SQLite file by hand.

The response body is cut by SQLite via substr over a blob cast, the
same projection the event body uses, so an oversized stored response
never becomes a Go string. The page reports the cut with a marker.

Response bodies and errors are remote content, so both go through a
new delivery.Redactor that strips the target's own destination URL,
path, query and userinfo, plus the values of credential-shaped
request headers, before rendering. Target configuration keeps
reaching the template only as a TargetView.

A body that reaches the cap is treated as cut whether or not SQLite
is what cut it. The delivery engine stops reading a response at its
own cap, which is the same number of bytes this page renders, and
the row it writes records that cut length as the whole length, so
nothing in the row separates a response that ended at the cap from
one severed there. Such a body goes through RedactCut, which drops
any tail that is a proper prefix of a secret: the remote chooses the
padding in front of a credential it echoes, so it chooses where the
cut falls inside that credential. Its marker says the response
reached the recording limit rather than quoting a total the row does
not know.

Redactors are built from an unscoped target load. Deleting a target
only soft deletes the row while its deliveries survive, and a scoped
load would leave exactly those deliveries rendering unredacted. The
views the page lists stay scoped.

Attempt loading is chunked so the IN clause cannot exceed SQLite's
bound-parameter limit, and a chunk that fails fails the page rather
than rendering the deliveries it covered as never having run. The
page renders at most 20 attempts per delivery, counting what it
leaves out.

static/css/tailwind.css is regenerated with tailwindcss for the
utility classes the new markup uses.
2026-08-20 06:06:44 +00:00

1420 lines
27 KiB
Go

package delivery_test
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
_ "modernc.org/sqlite"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
// iSetup holds common integration test dependencies.
type iSetup struct {
MainDB *gorm.DB
DBMgr *database.WebhookDBManager
WebhookID string
WebhookDB *gorm.DB
Engine *delivery.Engine
}
func newISetup(t *testing.T) iSetup {
t.Helper()
mainDB := iMainDB(t)
dbMgr := iDBManager(t)
wID := uuid.New().String()
wDB := iSeedWebhookDB(t, dbMgr, wID)
return iSetup{
MainDB: mainDB,
DBMgr: dbMgr,
WebhookID: wID,
WebhookDB: wDB,
Engine: delivery.NewTestEngineWithDB(
database.NewTestDatabase(mainDB),
dbMgr,
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{
Level: slog.LevelDebug,
},
)),
&http.Client{Timeout: 5 * time.Second},
2,
),
}
}
func iMainDB(t *testing.T) *gorm.DB {
t.Helper()
dbPath := filepath.Join(
t.TempDir(), "main-test.db",
)
dsn := fmt.Sprintf(
"file:%s?cache=shared&mode=rwc", dbPath,
)
sqlDB, err := sql.Open("sqlite", dsn)
require.NoError(t, err)
t.Cleanup(func() { _ = sqlDB.Close() })
db, err := gorm.Open(
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
)
require.NoError(t, err)
require.NoError(t, db.AutoMigrate(
&database.Webhook{},
&database.Target{},
&database.User{},
&database.Setting{},
))
return db
}
func iDBManager(
t *testing.T,
) *database.WebhookDBManager {
t.Helper()
return database.NewTestWebhookDBManager(t.TempDir())
}
func iSeedWebhookDB(
t *testing.T,
mgr *database.WebhookDBManager,
webhookID string,
) *gorm.DB {
t.Helper()
db, err := mgr.GetDB(webhookID)
require.NoError(t, err)
return db
}
func iHTTPConfig(url string) string {
cfg := delivery.HTTPTargetConfig{URL: url}
data, err := json.Marshal(cfg)
if err != nil {
panic("failed to marshal HTTPTargetConfig")
}
return string(data)
}
func iEngine(
t *testing.T, workers int,
) *delivery.Engine {
t.Helper()
return delivery.NewTestEngine(
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
)),
&http.Client{Timeout: 5 * time.Second},
workers,
)
}
// iSeedEvent creates a test event in the database.
func iSeedEvent(
t *testing.T,
db *gorm.DB,
webhookID, body string,
) database.Event {
t.Helper()
event := database.Event{
WebhookID: webhookID,
EntrypointID: uuid.New().String(),
Method: http.MethodPost,
Headers: `{}`,
Body: body,
ContentType: testContentType,
}
require.NoError(t, db.Create(&event).Error)
return event
}
// iSeedDelivery creates a test delivery record.
func iSeedDelivery(
t *testing.T,
db *gorm.DB,
eventID, targetID string,
status database.DeliveryStatus,
) database.Delivery {
t.Helper()
d := database.Delivery{
EventID: eventID,
TargetID: targetID,
Status: status,
}
require.NoError(t, db.Create(&d).Error)
return d
}
// iTask builds a delivery.Task for integration tests.
func iTask(
d database.Delivery,
event database.Event,
webhookID, targetID, name, config string,
maxRetries, attemptNum int,
body *string,
) delivery.Task {
return delivery.Task{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: webhookID,
TargetID: targetID,
TargetName: name,
TargetType: database.TargetTypeHTTP,
TargetConfig: config,
MaxRetries: maxRetries,
Method: event.Method,
Headers: event.Headers,
ContentType: event.ContentType,
Body: body,
AttemptNum: attemptNum,
}
}
// iAssertStatus checks the delivery status.
func iAssertStatus(
t *testing.T,
db *gorm.DB,
deliveryID string,
expected database.DeliveryStatus,
) {
t.Helper()
var updated database.Delivery
require.NoError(t, db.First(
&updated, "id = ?", deliveryID,
).Error)
assert.Equal(t, expected, updated.Status)
}
// --- processNewTask Tests ---
func TestProcessNewTask_InlineBody(t *testing.T) {
t.Parallel()
s := newISetup(t)
var received atomic.Bool
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
received.Store(true)
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"hello":"world"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
cfg := iHTTPConfig(ts.URL)
task := iTask(
d, event, s.WebhookID, targetID,
"test-target", cfg, 0, 1, &bodyStr,
)
s.Engine.ExportProcessNewTask(
context.TODO(), &task,
)
assert.True(t, received.Load())
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
func TestProcessNewTask_LargeBody_FetchFromDB(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
largeBody := strings.Repeat(
"x", delivery.MaxInlineBodySize+100,
)
var receivedBody string
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
receivedBody = string(body)
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, largeBody,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
cfg := iHTTPConfig(ts.URL)
task := iTask(
d, event, s.WebhookID, targetID,
"test-large", cfg, 0, 1, nil,
)
s.Engine.ExportProcessNewTask(
context.TODO(), &task,
)
assert.Equal(t, largeBody, receivedBody)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
func TestProcessNewTask_InvalidWebhookID(t *testing.T) {
t.Parallel()
s := newISetup(t)
task := delivery.Task{
DeliveryID: uuid.New().String(),
EventID: uuid.New().String(),
WebhookID: uuid.New().String(),
TargetID: uuid.New().String(),
TargetName: "test",
TargetType: database.TargetTypeHTTP,
TargetConfig: iHTTPConfig("http://localhost:9999"),
MaxRetries: 0,
Body: nil,
AttemptNum: 1,
}
s.Engine.ExportProcessNewTask(
context.TODO(), &task,
)
}
// --- processRetryTask Tests ---
func TestProcessRetryTask_SuccessfulRetry(t *testing.T) {
t.Parallel()
s := newISetup(t)
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"retry":"test"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
bodyStr := event.Body
cfg := iHTTPConfig(ts.URL)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-target", cfg, 5, 2, &bodyStr,
)
s.Engine.ExportProcessRetryTask(
context.TODO(), &task,
)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
func TestProcessRetryTask_SkipsNonRetryingDelivery(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"skip":"test"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusDelivered,
)
bodyStr := event.Body
cfg := iHTTPConfig("http://localhost:1")
task := iTask(
d, event, s.WebhookID, targetID,
"skip-target", cfg, 5, 2, &bodyStr,
)
s.Engine.ExportProcessRetryTask(
context.TODO(), &task,
)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
func TestProcessRetryTask_LargeBody_FetchFromDB(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
largeBody := strings.Repeat(
"z", delivery.MaxInlineBodySize+50,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, largeBody,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
cfg := iHTTPConfig(ts.URL)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-large", cfg, 5, 2, nil,
)
s.Engine.ExportProcessRetryTask(
context.TODO(), &task,
)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusDelivered,
)
}
// --- Worker Lifecycle Tests ---
func TestWorkerLifecycle_StartStop(t *testing.T) {
t.Parallel()
s := newISetup(t)
s.Engine.ExportStart()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"lifecycle":"test"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
task := iTask(
d, event, s.WebhookID, targetID,
"lifecycle-test", "",
0, 1, &bodyStr,
)
task.TargetType = database.TargetTypeLog
s.Engine.Notify([]delivery.Task{task})
iWaitForDelivered(t, s.WebhookDB, d.ID)
require.NoError(t, s.Engine.ExportStop(context.Background()))
}
// iWaitForDelivered polls until the delivery reaches the
// delivered status.
func iWaitForDelivered(
t *testing.T,
db *gorm.DB,
deliveryID string,
) {
t.Helper()
require.Eventually(t, func() bool {
var d database.Delivery
err := db.First(
&d, "id = ?", deliveryID,
).Error
if err != nil {
return false
}
return d.Status == database.DeliveryStatusDelivered
}, 5*time.Second, 50*time.Millisecond)
}
func TestWorkerLifecycle_ProcessesRetryChannel(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"retry-chan":"test"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
s.Engine.ExportStart()
bodyStr := event.Body
cfg := iHTTPConfig(ts.URL)
task := iTask(
d, event, s.WebhookID, targetID,
"retry-chan-test", cfg, 5, 2, &bodyStr,
)
s.Engine.ExportRetryCh() <- task
iWaitForDelivered(t, s.WebhookDB, d.ID)
require.NoError(t, s.Engine.ExportStop(context.Background()))
}
// --- processDelivery: unknown target type ---
func TestProcessDelivery_UnknownTargetType(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"unknown":"type"}`,
)
del := iSeedDelivery(
t, s.WebhookDB, event.ID,
uuid.New().String(),
database.DeliveryStatusPending,
)
d := &database.Delivery{
EventID: event.ID,
TargetID: del.TargetID,
Status: database.DeliveryStatusPending,
Event: event,
Target: database.Target{
Name: "unknown",
Type: database.TargetType("unknown"),
},
}
d.ID = del.ID
task := &delivery.Task{
DeliveryID: del.ID,
TargetType: database.TargetType("unknown"),
}
s.Engine.ExportProcessDelivery(
context.TODO(), s.WebhookDB, d, task,
)
iAssertStatus(t, s.WebhookDB, del.ID,
database.DeliveryStatusFailed,
)
}
// --- Recovery Tests ---
func TestRecoverPendingDeliveries(t *testing.T) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "recovery-target",
database.TargetTypeLog, "", 0,
)
iSeedPendingDeliveries(
t, s.WebhookDB, s.WebhookID, targetID, 3,
)
s.Engine.ExportRecoverPendingDeliveries(
context.Background(), s.WebhookDB,
s.WebhookID,
)
for i := range 3 {
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(t, targetID, task.TargetID)
assert.Equal(t,
database.TargetTypeLog,
task.TargetType,
)
case <-time.After(2 * time.Second):
t.Fatalf("expected task %d", i)
}
}
}
// iCreateTarget creates a target in the main DB.
func iCreateTarget(
t *testing.T,
mainDB *gorm.DB,
targetID, webhookID, name string,
targetType database.TargetType,
config string,
maxRetries int,
) {
t.Helper()
target := database.Target{
WebhookID: webhookID,
Name: name,
Type: targetType,
Active: true,
Config: config,
MaxRetries: maxRetries,
}
target.ID = targetID
require.NoError(t, mainDB.Create(&target).Error)
}
// iSeedPendingDeliveries creates n pending deliveries.
func iSeedPendingDeliveries(
t *testing.T,
webhookDB *gorm.DB,
webhookID, targetID string,
n int,
) {
t.Helper()
for i := range n {
event := iSeedEvent(
t, webhookDB, webhookID,
fmt.Sprintf(`{"recovery":%d}`, i),
)
iSeedDelivery(
t, webhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
}
}
func TestRecoverWebhookDeliveries_RetryingDeliveries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "retry-recovery",
database.TargetTypeHTTP,
iHTTPConfig("http://example.com/hook"), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"retry-recovery":"test"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "test-webhook",
)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
select {
case task := <-s.Engine.ExportRetryCh():
assert.Equal(t, d.ID, task.DeliveryID)
assert.Equal(t, targetID, task.TargetID)
assert.Equal(t, 2, task.AttemptNum)
case <-time.After(5 * time.Second):
t.Fatal("expected retry task from recovery")
}
// Regression guard: a target that still supports retries
// must be rescheduled, never terminally failed, and must
// not gain a synthetic result row.
iAssertStatus(
t, s.WebhookDB, d.ID,
database.DeliveryStatusRetrying,
)
assert.Len(t, iResults(t, s.WebhookDB, d.ID), 1)
}
// --- Retrying deliveries whose target type changed ---
// iSeedRetryingWithType seeds a retrying delivery with one
// recorded failed attempt against a target of the given type,
// standing in for a target whose type was edited in the main
// database while the delivery was still retrying.
func iSeedRetryingWithType(
t *testing.T,
s iSetup,
targetType database.TargetType,
) string {
t.Helper()
targetID := uuid.New().String()
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "mutated-target", targetType,
iHTTPConfig("http://example.com/hook"), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"orphaned":"retry"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
return d.ID
}
// iResults loads a delivery's results in attempt order.
func iResults(
t *testing.T, db *gorm.DB, deliveryID string,
) []database.DeliveryResult {
t.Helper()
var results []database.DeliveryResult
require.NoError(t, db.
Where("delivery_id = ?", deliveryID).
Order("attempt_num").
Find(&results).Error)
return results
}
// iAssertTerminallyFailed asserts the delivery ended failed
// with a result row recording why, and was not rescheduled.
func iAssertTerminallyFailed(
t *testing.T,
s iSetup,
deliveryID string,
targetType database.TargetType,
) {
t.Helper()
iAssertStatus(
t, s.WebhookDB, deliveryID,
database.DeliveryStatusFailed,
)
results := iResults(t, s.WebhookDB, deliveryID)
require.Len(t, results, 2)
last := results[1]
assert.False(t, last.Success)
assert.Equal(t, 2, last.AttemptNum)
assert.Contains(
t, last.Error, string(targetType),
)
assert.Contains(
t, last.Error, "does not support retries",
)
assert.Empty(t, s.Engine.ExportRetryCh())
}
func TestRecoverSingleRetry_TypeNoLongerRetries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "mutated-type",
)
deliveryID := iSeedRetryingWithType(
t, s, database.TargetTypeLog,
)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(
t, s, deliveryID, database.TargetTypeLog,
)
}
func TestSweepSingleRetry_TypeNoLongerRetries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "mutated-type-sweep",
)
deliveryID := iSeedRetryingWithType(
t, s, database.TargetTypeDatabase,
)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(
t, s, deliveryID, database.TargetTypeDatabase,
)
}
// TestFailUnretryableRetry_WritesNoTargetRow proves the
// orphaned-retry terminal path leaves no target row — and so no
// plaintext target config — in the per-webhook event database.
//
// That path loads the delivery without its Target relation on
// purpose. Populating d.Target makes GORM's SaveBeforeAssociations
// upsert the whole target row on the status UPDATE, which for a slack
// target writes the incoming-webhook credential into events-*.db.
// See https://git.eeqj.de/sneak/webhooker/issues/206.
func TestFailUnretryableRetry_WritesNoTargetRow(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "no-target-row",
)
targetID := uuid.New().String()
// A Slack incoming-webhook URL: the target config IS the
// credential, which is what makes a leaked target row a
// disclosure rather than a curiosity.
hookURL := "https://hooks.slack.com/services/T00/B00/x"
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "credential-bearing",
database.TargetTypeLog, iHTTPConfig(hookURL), 5,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"orphaned":"retry"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusRetrying,
)
iSeedFailedResult(t, s.WebhookDB, d.ID)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertStatus(t, s.WebhookDB, d.ID,
database.DeliveryStatusFailed,
)
// The table exists in the per-webhook database because GORM
// migrates the Delivery relation's model alongside it. It must
// stay empty.
var targetRows int64
require.NoError(t, s.WebhookDB.
Table("targets").
Count(&targetRows).Error)
assert.Zero(t, targetRows,
"orphaned-retry terminal failure wrote a target row "+
"into the per-webhook event database",
)
var configs []string
require.NoError(t, s.WebhookDB.
Table("targets").
Pluck("config", &configs).Error)
assert.NotContains(
t, strings.Join(configs, " "), hookURL,
)
}
func TestRecoverSingleRetry_UnknownTargetType(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "unknown-type",
)
unknown := database.TargetType("not-a-target-type")
deliveryID := iSeedRetryingWithType(t, s, unknown)
s.Engine.ExportRecoverWebhookDeliveries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(t, s, deliveryID, unknown)
}
func TestSweepSingleRetry_UnknownTargetType(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
iCreateWebhook(
t, s.MainDB, s.WebhookID, "unknown-type-sweep",
)
unknown := database.TargetType("not-a-target-type")
deliveryID := iSeedRetryingWithType(t, s, unknown)
s.Engine.ExportSweepWebhookRetries(
context.Background(), s.WebhookID,
)
iAssertTerminallyFailed(t, s, deliveryID, unknown)
}
// iSeedFailedResult creates a failed delivery result.
func iSeedFailedResult(
t *testing.T,
db *gorm.DB,
deliveryID string,
) {
t.Helper()
result := database.DeliveryResult{
DeliveryID: deliveryID,
AttemptNum: 1,
Success: false,
StatusCode: 500,
Error: "server error",
}
require.NoError(t, db.Create(&result).Error)
}
// iCreateWebhook creates a webhook record in main DB.
func iCreateWebhook(
t *testing.T,
mainDB *gorm.DB,
webhookID, name string,
) {
t.Helper()
webhook := database.Webhook{
UserID: uuid.New().String(),
Name: name,
}
webhook.ID = webhookID
require.NoError(t, mainDB.Create(&webhook).Error)
}
// --- recoverInFlight Tests ---
func TestRecoverInFlight_NoWebhooks(t *testing.T) {
t.Parallel()
s := newISetup(t)
s.Engine.ExportRecoverInFlight(
context.Background(),
)
}
func TestRecoverInFlight_WithPendingDeliveries(
t *testing.T,
) {
t.Parallel()
s := newISetup(t)
targetID := uuid.New().String()
iCreateWebhook(
t, s.MainDB, s.WebhookID, "recover-test",
)
iCreateTarget(t, s.MainDB, targetID,
s.WebhookID, "recover-target",
database.TargetTypeLog, "", 0,
)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"recover":"inflight"}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportRecoverInFlight(
context.Background(),
)
select {
case task := <-s.Engine.ExportDeliveryCh():
assert.Equal(t, d.ID, task.DeliveryID)
assert.Equal(t,
database.TargetTypeLog, task.TargetType,
)
case <-time.After(2 * time.Second):
t.Fatal("expected task from recoverInFlight")
}
}
// --- HTTP Config with custom headers ---
func TestDeliverHTTP_CustomTargetHeaders(t *testing.T) {
t.Parallel()
s := newISetup(t)
var receivedAuth string
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
receivedAuth = r.Header.Get(
"Authorization",
)
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
cfg := delivery.HTTPTargetConfig{
URL: ts.URL,
Headers: map[string]string{
"Authorization": "Bearer secret-token",
},
}
cfgJSON, err := json.Marshal(cfg)
require.NoError(t, err)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID,
`{"auth":"test"}`,
)
targetID := uuid.New().String()
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
task := iTask(
d, event, s.WebhookID, targetID,
"auth-target", string(cfgJSON),
0, 1, &bodyStr,
)
s.Engine.ExportProcessNewTask(
context.TODO(), &task,
)
assert.Equal(t,
"Bearer secret-token", receivedAuth,
)
}
// --- HTTP delivery with custom timeout ---
func TestDeliverHTTP_TargetTimeout(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := iEngine(t, 1)
ts := httptest.NewServer(
http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
time.Sleep(2 * time.Second)
w.WriteHeader(http.StatusOK)
},
),
)
defer ts.Close()
cfg := delivery.HTTPTargetConfig{
URL: ts.URL,
Timeout: 1,
}
cfgJSON, err := json.Marshal(cfg)
require.NoError(t, err)
event, del := iSeedEventAndDelivery(
t, db, `{"timeout":"test"}`,
string(cfgJSON),
)
task, d := iHTTPTaskAndDelivery(
event, del, "timeout-target",
string(cfgJSON), 0, 1,
)
e.ExportDeliverHTTP(context.TODO(), db, d, task)
iAssertStatus(t, db, del.ID,
database.DeliveryStatusFailed,
)
iAssertResultFailed(t, db, del.ID)
}
// TestDeliverHTTP_CutsStoredResponseAtMaxBodyLog pins the size
// this engine stores for an oversized response, because the
// event log's redaction is written against it: the row holds
// exactly maxBodyLog bytes and records nothing about how much
// more the remote sent, so a credential echoed across that
// boundary reaches the database already severed and no reader
// of the row can tell the cut happened.
func TestDeliverHTTP_CutsStoredResponseAtMaxBodyLog(
t *testing.T,
) {
t.Parallel()
// Padded so the cut falls five bytes before the end of the
// echoed webhook URL.
const (
severedTail = 5
overshoot = 100000
)
sent := strings.Repeat(
"A",
delivery.ExportMaxBodyLog-len(slackWebhookURL)+
severedTail,
) + slackWebhookURL + strings.Repeat("Z", overshoot)
s := newISetup(t)
ts := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusBadGateway)
_, _ = io.WriteString(w, sent)
},
))
defer ts.Close()
cfgJSON := iHTTPConfig(ts.URL)
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"cut":"test"}`,
)
targetID := uuid.New().String()
del := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
bodyStr := event.Body
task := iTask(
del, event, s.WebhookID, targetID,
"cut-target", cfgJSON, 0, 1, &bodyStr,
)
s.Engine.ExportProcessNewTask(context.TODO(), &task)
results := iResults(t, s.WebhookDB, del.ID)
require.Len(t, results, 1)
stored := results[0].ResponseBody
assert.Len(
t, stored, delivery.ExportMaxBodyLog,
"an oversized response is stored at exactly the cap",
)
assert.Equal(
t, sent[:delivery.ExportMaxBodyLog], stored,
)
assert.NotContains(
t, stored, slackWebhookURL,
"the echoed URL is severed by the cut",
)
assert.Contains(
t, stored, "T00000000",
"the severed prefix still carries the credential",
)
}
// iSeedEventAndDelivery creates event + delivery
// for standalone tests.
func iSeedEventAndDelivery(
t *testing.T,
db *gorm.DB,
body, _ string,
) (database.Event, database.Delivery) {
t.Helper()
event := database.Event{
WebhookID: uuid.New().String(),
EntrypointID: uuid.New().String(),
Method: http.MethodPost,
Headers: `{"Content-Type":["application/json"]}`,
Body: body,
ContentType: testContentType,
}
require.NoError(t, db.Create(&event).Error)
d := database.Delivery{
EventID: event.ID,
TargetID: uuid.New().String(),
Status: database.DeliveryStatusPending,
}
require.NoError(t, db.Create(&d).Error)
return event, d
}
// iHTTPTaskAndDelivery builds a task/delivery pair for
// standalone HTTP tests.
func iHTTPTaskAndDelivery(
event database.Event,
del database.Delivery,
name, config string,
maxRetries, attemptNum int,
) (*delivery.Task, *database.Delivery) {
task := &delivery.Task{
DeliveryID: del.ID,
EventID: event.ID,
WebhookID: event.WebhookID,
TargetID: del.TargetID,
TargetName: name,
TargetType: database.TargetTypeHTTP,
TargetConfig: config,
MaxRetries: maxRetries,
AttemptNum: attemptNum,
}
d := &database.Delivery{
EventID: event.ID,
TargetID: del.TargetID,
Status: database.DeliveryStatusPending,
Event: event,
Target: database.Target{
Name: name,
Type: database.TargetTypeHTTP,
Config: config,
},
}
d.ID = del.ID
return task, d
}
// iAssertResultFailed checks that a failed delivery
// result exists.
func iAssertResultFailed(
t *testing.T,
db *gorm.DB,
deliveryID string,
) {
t.Helper()
var result database.DeliveryResult
require.NoError(t, db.Where(
"delivery_id = ?", deliveryID,
).First(&result).Error)
assert.False(t, result.Success)
assert.NotEmpty(t, result.Error)
}
// --- HTTP request with invalid config ---
func TestDeliverHTTP_InvalidConfig(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := iEngine(t, 1)
event, del := iSeedEventAndDelivery(
t, db, `{"config":"invalid"}`, "",
)
task, d := iHTTPTaskAndDelivery(
event, del, "bad-config", `not-json`, 0, 1,
)
e.ExportDeliverHTTP(context.TODO(), db, d, task)
iAssertStatus(t, db, del.ID,
database.DeliveryStatusFailed,
)
}
// --- Notify batching ---
func TestNotify_MultipleTasks(t *testing.T) {
t.Parallel()
e := iEngine(t, 1)
tasks := make([]delivery.Task, 5)
for i := range tasks {
tasks[i] = delivery.Task{
DeliveryID: fmt.Sprintf("task-%d", i),
}
}
e.Notify(tasks)
for i := range 5 {
select {
case task := <-e.ExportDeliveryCh():
assert.Equal(t,
fmt.Sprintf("task-%d", i),
task.DeliveryID,
)
case <-time.After(time.Second):
t.Fatalf("expected task %d", i)
}
}
}