Some checks failed
check / check (push) Has been cancelled
A delivery could reach a bad end without the engine recording why, and a retrying delivery could fail to reach an end at all. An unknown target type marked the delivery failed and wrote no DeliveryResult, so the event log showed "failed" with no attempts and the only account of why was one line in the server log. It now records a result naming the type before failing the delivery. A deleted target left its retrying deliveries stranded. Both recovery and the sweep began with a scoped loadTarget, which cannot see a soft deleted row, so both logged and returned: the delivery stayed retrying for the life of the database while the sweep repeated the same error every minute. Both now terminalise it with a recorded reason. Deleting a target also did not stop deliveries to it. A scheduled retry is a time.AfterFunc holding the target's configuration from when the chain began, and nothing on that path read the target row, so the timer kept firing and kept sending to the destination the operator had removed for the rest of the backoff chain; terminalising in recovery and the sweep alone would only have caught it after a restart. processRetryTask now confirms the target still exists before it attempts, and abandons the chain when it does not. Only a target confirmed gone stops anything. A lookup that fails for any other reason is the main database being unreadable, which is transient, and every path leaves the delivery exactly as it was rather than failing it. The reason text comes from one Unscoped lookup confined to these terminal paths, because a soft deleted row is what distinguishes a target the operator deleted from an id that never named one. The engine's normal target loading stays scoped, or deleting a target would stop nothing. Terminal writes reached from recovery keep going through the existing retainIdle ownership gate; the retry path writes directly, as a target's own Deliver does, because the worker already holds that delivery. Nothing was added to either sweep dispatch arm. Retry fixtures that drove processRetryTask for a target id with no row in the main database now create one. That state is not reachable in service: the handler reads the target to build the task.
1443 lines
28 KiB
Go
1443 lines
28 KiB
Go
package delivery_test
|
|
|
|
import (
|
|
"context"
|
|
"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",
|
|
)
|
|
|
|
// Opened the way the service opens the main database, so these
|
|
// tests cannot pass against journal and locking settings
|
|
// production does not use.
|
|
sqlDB, err := database.OpenSQLite(
|
|
dbPath, database.SQLiteModeCreate,
|
|
)
|
|
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)
|
|
|
|
// The target row exists because the engine confirms a scheduled
|
|
// retry's target has not been deleted before it runs it. A retry
|
|
// task whose target id names no row at all is a state the service
|
|
// does not produce: the handler read that target to build the
|
|
// task. See https://git.eeqj.de/sneak/webhooker/issues/107.
|
|
iCreateTarget(
|
|
t, s.MainDB, targetID, s.WebhookID, "retry-target",
|
|
database.TargetTypeHTTP, cfg, 5,
|
|
)
|
|
|
|
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)
|
|
|
|
iCreateTarget(
|
|
t, s.MainDB, targetID, s.WebhookID, "retry-large",
|
|
database.TargetTypeHTTP, cfg, 5,
|
|
)
|
|
|
|
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)
|
|
|
|
iCreateTarget(
|
|
t, s.MainDB, targetID, s.WebhookID, "retry-chan-test",
|
|
database.TargetTypeHTTP, cfg, 5,
|
|
)
|
|
|
|
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)
|
|
}
|
|
}
|
|
}
|