Compare commits
2 Commits
feat/recei
...
36a1bacf11
| Author | SHA1 | Date | |
|---|---|---|---|
| 36a1bacf11 | |||
| ee7c626071 |
@@ -1,5 +1,9 @@
|
||||
version: "2"
|
||||
|
||||
# Config schema uses the golangci-lint v2 layout (settings live under
|
||||
# linters.settings, not top-level linters-settings) so that the
|
||||
# thresholds below are actually applied by golangci-lint >= v2.
|
||||
|
||||
run:
|
||||
timeout: 5m
|
||||
modules-download-mode: readonly
|
||||
@@ -14,8 +18,7 @@ linters:
|
||||
- wsl # Deprecated, replaced by wsl_v5
|
||||
- wrapcheck # Too verbose for internal packages
|
||||
- varnamelen # Short names like db, id are idiomatic Go
|
||||
|
||||
linters-settings:
|
||||
settings:
|
||||
lll:
|
||||
line-length: 88
|
||||
funlen:
|
||||
@@ -27,6 +30,5 @@ linters-settings:
|
||||
threshold: 100
|
||||
|
||||
issues:
|
||||
exclude-use-default: false
|
||||
max-issues-per-linter: 0
|
||||
max-same-issues: 0
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
# Lint stage
|
||||
# golangci/golangci-lint:v2.11.3 (Debian-based), 2026-03-17
|
||||
# golangci/golangci-lint:v2.12.2 (Debian-based), 2026-08-07
|
||||
# Using Debian-based image because mattn/go-sqlite3 (CGO) does not
|
||||
# compile on Alpine musl (off64_t is a glibc type).
|
||||
FROM golangci/golangci-lint:v2.11.3@sha256:e838e8ab68aaefe83e2408691510867ade9329c0e0b895a3fb35eb93d1c2a4ba AS lint
|
||||
FROM golangci/golangci-lint:v2.12.2@sha256:5cceeef04e53efe1470638d4b4b4f5ceefd574955ab3941b2d9a68a8c9ad5240 AS lint
|
||||
|
||||
RUN apt-get update && apt-get install -y --no-install-recommends make && rm -rf /var/lib/apt/lists/*
|
||||
|
||||
|
||||
31
README.md
31
README.md
@@ -363,10 +363,12 @@ events should be forwarded.
|
||||
greater than 0, failed deliveries are retried with exponential backoff
|
||||
up to `max_retries` attempts, protected by a per-target circuit
|
||||
breaker.
|
||||
- **`database`** — Confirm the event is stored in the webhook's
|
||||
per-webhook database (no external delivery). Since events are always
|
||||
written to the per-webhook DB on ingestion, this target marks delivery
|
||||
as immediately successful. Useful for ensuring durable event archival.
|
||||
- **`database`** — Archive the full event as a row into a separate
|
||||
per-webhook archive database (`archive-{webhookID}.db`) for long-term
|
||||
retention, with an optional creation-validated expiry (default: keep
|
||||
forever). No external delivery and no retries; an archive write
|
||||
failure fails the delivery. See the database target section under
|
||||
"Per-Webhook Event Databases" for the full semantics.
|
||||
- **`log`** — Write the event to the application log (stdout). Useful
|
||||
for debugging.
|
||||
|
||||
@@ -512,11 +514,22 @@ This separation provides:
|
||||
page cache, and its own lock, so concurrent event ingestion across
|
||||
webhooks won't contend.
|
||||
|
||||
The **database target type** leverages this architecture: since events
|
||||
are already stored in the per-webhook database by design, the database
|
||||
target simply marks the delivery as immediately successful. The
|
||||
per-webhook DB IS the dedicated event database — that's the whole point
|
||||
of the database target type.
|
||||
The **database target type** builds on this architecture to provide
|
||||
long-term archiving, separate from the per-webhook event database (which
|
||||
may prune events under its own retention). Delivering to a database
|
||||
target writes the full event — body, headers, method, content type, and
|
||||
webhook/entrypoint/event identifiers — as a row into a dedicated archive
|
||||
database, `archive-{webhookID}.db`, stored under the data directory
|
||||
beside the event database. After each write the archive handle is closed
|
||||
and reopened, debounced to at most once per second, so an operator can
|
||||
move the archive file away for offline archiving without stopping the
|
||||
service; a moved or removed archive file is recreated automatically on
|
||||
the next write. An optional `expiry` in the target's config JSON (e.g.
|
||||
`{"expiry":"720h"}`) is validated when the target is created — the
|
||||
default (unset or the literal `never`) keeps rows forever — and rows
|
||||
older than the expiry are pruned each time the archive is (re)opened. An
|
||||
archive write failure is never silent success: the delivery records a
|
||||
failed attempt with the error and is marked failed.
|
||||
|
||||
The **Slack target type** sends webhook events as formatted messages to
|
||||
any Slack-compatible incoming webhook URL (works with Slack, Mattermost,
|
||||
|
||||
5
TODO.md
5
TODO.md
@@ -28,6 +28,11 @@ databases currently grow without bound.
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-08-07 Update golangci-lint to v2.12.2 (Docker image digest in
|
||||
`Dockerfile`, release-archive sha256 pins in `script/bootstrap`),
|
||||
adopt the canonical `.golangci.yml` (v2 `linters.settings` layout so
|
||||
`lll`/`funlen`/`cyclop`/`dupl` thresholds actually apply), and fix
|
||||
all newly surfaced lint findings
|
||||
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
|
||||
Makefile shims, README Entrypoints section
|
||||
- 2026-03-25 pin golangci-lint Docker image for linting (#55)
|
||||
|
||||
@@ -11,6 +11,15 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
)
|
||||
|
||||
const (
|
||||
// testAppname is the Globals.Appname used in tests.
|
||||
testAppname = "webhooker-test"
|
||||
// testVersion is the Globals.Version used in tests.
|
||||
testVersion = "test"
|
||||
// testContentType is the event content type used in tests.
|
||||
testContentType = "application/json"
|
||||
)
|
||||
|
||||
func setupTestDB(
|
||||
t *testing.T,
|
||||
) (*database.Database, *fxtest.Lifecycle) {
|
||||
@@ -19,8 +28,8 @@ func setupTestDB(
|
||||
lc := fxtest.NewLifecycle(t)
|
||||
|
||||
g := &globals.Globals{
|
||||
Appname: "webhooker-test",
|
||||
Version: "test",
|
||||
Appname: testAppname,
|
||||
Version: testVersion,
|
||||
}
|
||||
|
||||
l, err := logger.New(
|
||||
|
||||
@@ -5,7 +5,10 @@ type Entrypoint struct {
|
||||
BaseModel
|
||||
|
||||
WebhookID string `gorm:"type:uuid;not null" json:"webhookId"`
|
||||
Path string `gorm:"uniqueIndex;not null" json:"path"` // URL path for this entrypoint
|
||||
|
||||
// Path is the URL path for this entrypoint.
|
||||
Path string `gorm:"uniqueIndex;not null" json:"path"`
|
||||
|
||||
Description string `json:"description"`
|
||||
Active bool `gorm:"default:true" json:"active"`
|
||||
|
||||
|
||||
@@ -23,7 +23,8 @@ type Target struct {
|
||||
// Configuration fields (JSON stored based on type)
|
||||
Config string `gorm:"type:text" json:"config"` // JSON configuration
|
||||
|
||||
// For HTTP targets (max_retries=0 means fire-and-forget, >0 enables retries with backoff)
|
||||
// For HTTP targets (max_retries=0 means fire-and-forget,
|
||||
// >0 enables retries with backoff)
|
||||
MaxRetries int `json:"maxRetries,omitempty"`
|
||||
MaxQueueSize int `json:"maxQueueSize,omitempty"`
|
||||
|
||||
|
||||
@@ -7,7 +7,9 @@ type Webhook struct {
|
||||
UserID string `gorm:"type:uuid;not null" json:"userId"`
|
||||
Name string `gorm:"not null" json:"name"`
|
||||
Description string `json:"description"`
|
||||
RetentionDays int `gorm:"default:30" json:"retentionDays"` // Days to retain events
|
||||
|
||||
// RetentionDays is the number of days to retain events.
|
||||
RetentionDays int `gorm:"default:30" json:"retentionDays"`
|
||||
|
||||
// Relations
|
||||
User User `json:"user,omitzero"`
|
||||
|
||||
@@ -2,6 +2,7 @@ package database_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -30,8 +31,8 @@ func setupRetentionTest(t *testing.T) *retentionTestEnv {
|
||||
lc := fxtest.NewLifecycle(t)
|
||||
|
||||
g := &globals.Globals{
|
||||
Appname: "webhooker-test",
|
||||
Version: "test",
|
||||
Appname: testAppname,
|
||||
Version: testVersion,
|
||||
}
|
||||
|
||||
l, err := logger.New(lc, logger.LoggerParams{Globals: g})
|
||||
@@ -117,9 +118,9 @@ func seedEventChain(
|
||||
event := &database.Event{
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Body: `{"seed": true}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
event.CreatedAt = createdAt
|
||||
require.NoError(t, db.Create(event).Error)
|
||||
|
||||
@@ -14,7 +14,10 @@ import (
|
||||
func NewTestDatabase(db *gorm.DB) *Database {
|
||||
return &Database{
|
||||
db: db,
|
||||
log: slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelDebug})),
|
||||
log: slog.New(slog.NewTextHandler(
|
||||
os.Stderr,
|
||||
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -23,6 +26,9 @@ func NewTestDatabase(db *gorm.DB) *Database {
|
||||
func NewTestWebhookDBManager(dataDir string) *WebhookDBManager {
|
||||
return &WebhookDBManager{
|
||||
dataDir: dataDir,
|
||||
log: slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelDebug})),
|
||||
log: slog.New(slog.NewTextHandler(
|
||||
os.Stderr,
|
||||
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package database_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
@@ -25,8 +26,8 @@ func setupTestWebhookDBManager(
|
||||
lc := fxtest.NewLifecycle(t)
|
||||
|
||||
g := &globals.Globals{
|
||||
Appname: "webhooker-test",
|
||||
Version: "test",
|
||||
Appname: testAppname,
|
||||
Version: testVersion,
|
||||
}
|
||||
|
||||
l, err := logger.New(
|
||||
@@ -83,10 +84,10 @@ func TestWebhookDBManager_CreateAndGetDB(t *testing.T) {
|
||||
event := &database.Event{
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{"Content-Type":["application/json"]}`,
|
||||
Body: `{"test": true}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
require.NoError(t, db.Create(event).Error)
|
||||
assert.NotEmpty(t, event.ID)
|
||||
@@ -99,7 +100,7 @@ func TestWebhookDBManager_CreateAndGetDB(t *testing.T) {
|
||||
db.First(&readEvent, "id = ?", event.ID).Error,
|
||||
)
|
||||
assert.Equal(t, webhookID, readEvent.WebhookID)
|
||||
assert.Equal(t, "POST", readEvent.Method)
|
||||
assert.Equal(t, http.MethodPost, readEvent.Method)
|
||||
assert.Equal(t, `{"test": true}`, readEvent.Body)
|
||||
}
|
||||
|
||||
@@ -123,9 +124,9 @@ func TestWebhookDBManager_DeleteDB(t *testing.T) {
|
||||
event := &database.Event{
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Body: `{"test": true}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
require.NoError(t, db.Create(event).Error)
|
||||
|
||||
@@ -196,10 +197,10 @@ func seedDeliveryWorkflow(
|
||||
event := &database.Event{
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{"Content-Type":["application/json"]}`,
|
||||
Body: `{"payload": "test"}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
require.NoError(t, db.Create(event).Error)
|
||||
|
||||
@@ -231,7 +232,7 @@ func verifyPendingDeliveries(
|
||||
)
|
||||
require.Len(t, pending, 1)
|
||||
assert.Equal(t, event.ID, pending[0].EventID)
|
||||
assert.Equal(t, "POST", pending[0].Event.Method)
|
||||
assert.Equal(t, http.MethodPost, pending[0].Event.Method)
|
||||
}
|
||||
|
||||
func completeDelivery(
|
||||
@@ -303,16 +304,16 @@ func TestWebhookDBManager_MultipleWebhooks(t *testing.T) {
|
||||
event1 := &database.Event{
|
||||
WebhookID: webhook1,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Body: `{"webhook": 1}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
event2 := &database.Event{
|
||||
WebhookID: webhook2,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "PUT",
|
||||
Method: http.MethodPut,
|
||||
Body: `{"webhook": 2}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
|
||||
require.NoError(t, db1.Create(event1).Error)
|
||||
|
||||
@@ -126,36 +126,6 @@ func iHTTPConfig(url string) string {
|
||||
return string(data)
|
||||
}
|
||||
|
||||
func iWebhookDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
|
||||
dbPath := filepath.Join(
|
||||
t.TempDir(), "events-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.Event{},
|
||||
&database.Delivery{},
|
||||
&database.DeliveryResult{},
|
||||
))
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
func iEngine(
|
||||
t *testing.T, workers int,
|
||||
) *delivery.Engine {
|
||||
@@ -182,10 +152,10 @@ func iSeedEvent(
|
||||
event := database.Event{
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{}`,
|
||||
Body: body,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
|
||||
require.NoError(t, db.Create(&event).Error)
|
||||
@@ -935,7 +905,7 @@ func TestDeliverHTTP_CustomTargetHeaders(t *testing.T) {
|
||||
func TestDeliverHTTP_TargetTimeout(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := iWebhookDB(t)
|
||||
db := testWebhookDB(t)
|
||||
e := iEngine(t, 1)
|
||||
|
||||
ts := httptest.NewServer(
|
||||
@@ -987,10 +957,10 @@ func iSeedEventAndDelivery(
|
||||
event := database.Event{
|
||||
WebhookID: uuid.New().String(),
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{"Content-Type":["application/json"]}`,
|
||||
Body: body,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
|
||||
require.NoError(t, db.Create(&event).Error)
|
||||
@@ -1067,7 +1037,7 @@ func iAssertResultFailed(
|
||||
func TestDeliverHTTP_InvalidConfig(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := iWebhookDB(t)
|
||||
db := testWebhookDB(t)
|
||||
e := iEngine(t, 1)
|
||||
|
||||
event, del := iSeedEventAndDelivery(
|
||||
|
||||
@@ -27,6 +27,9 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
)
|
||||
|
||||
// testContentType is the event content type used in tests.
|
||||
const testContentType = "application/json"
|
||||
|
||||
func testWebhookDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
|
||||
@@ -94,10 +97,10 @@ func seedEvent(
|
||||
event := database.Event{
|
||||
WebhookID: uuid.New().String(),
|
||||
EntrypointID: uuid.New().String(),
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{"Content-Type":["application/json"]}`,
|
||||
Body: body,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
|
||||
require.NoError(t, db.Create(&event).Error)
|
||||
@@ -342,33 +345,29 @@ func TestDeliverDatabase_ImmediateSuccess(
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
|
||||
event := seedEvent(t, db, `{"db":"target"}`)
|
||||
|
||||
dlv := seedDelivery(
|
||||
t, db, event.ID, uuid.New().String(),
|
||||
database.DeliveryStatusPending,
|
||||
// The database target archives for real now, so the engine
|
||||
// needs a webhook DB manager to locate the data directory.
|
||||
e := delivery.NewTestEngineWithDB(
|
||||
nil,
|
||||
database.NewTestWebhookDBManager(t.TempDir()),
|
||||
slog.New(slog.NewTextHandler(
|
||||
os.Stderr,
|
||||
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||
)),
|
||||
&http.Client{Timeout: 5 * time.Second},
|
||||
1,
|
||||
)
|
||||
|
||||
d := &database.Delivery{
|
||||
EventID: event.ID,
|
||||
TargetID: dlv.TargetID,
|
||||
Status: database.DeliveryStatusPending,
|
||||
Event: event,
|
||||
Target: database.Target{
|
||||
Name: "test-db",
|
||||
Type: database.TargetTypeDatabase,
|
||||
},
|
||||
}
|
||||
d.ID = dlv.ID
|
||||
event := seedEvent(t, db, `{"db":"target"}`)
|
||||
d := seedDatabaseTargetDelivery(t, db, event, "")
|
||||
|
||||
e.ExportDeliverDatabase(db, d)
|
||||
|
||||
var updated database.Delivery
|
||||
|
||||
require.NoError(t, db.First(
|
||||
&updated, "id = ?", dlv.ID,
|
||||
&updated, "id = ?", d.ID,
|
||||
).Error)
|
||||
|
||||
assert.Equal(t,
|
||||
@@ -379,7 +378,7 @@ func TestDeliverDatabase_ImmediateSuccess(
|
||||
var result database.DeliveryResult
|
||||
|
||||
require.NoError(t, db.Where(
|
||||
"delivery_id = ?", dlv.ID,
|
||||
"delivery_id = ?", d.ID,
|
||||
).First(&result).Error)
|
||||
|
||||
assert.True(t, result.Success)
|
||||
@@ -1117,10 +1116,10 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
|
||||
}
|
||||
|
||||
event := &database.Event{
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
Headers: `{"X-Custom":["value1"],"Content-Type":["application/json"]}`,
|
||||
Body: `{"test":true}`,
|
||||
ContentType: "application/json",
|
||||
ContentType: testContentType,
|
||||
}
|
||||
|
||||
statusCode, _, _, err := e.ExportDoHTTPRequest(
|
||||
@@ -1142,7 +1141,7 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
|
||||
)
|
||||
|
||||
assert.Equal(t,
|
||||
"application/json",
|
||||
testContentType,
|
||||
receivedHeaders.Get("Content-Type"),
|
||||
)
|
||||
|
||||
@@ -1158,7 +1157,19 @@ func TestProcessDelivery_RoutesToCorrectHandler(
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
|
||||
// The database target archives for real now, so the engine
|
||||
// needs a webhook DB manager to locate the data directory.
|
||||
e := delivery.NewTestEngineWithDB(
|
||||
nil,
|
||||
database.NewTestWebhookDBManager(t.TempDir()),
|
||||
slog.New(slog.NewTextHandler(
|
||||
os.Stderr,
|
||||
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||
)),
|
||||
&http.Client{Timeout: 5 * time.Second},
|
||||
1,
|
||||
)
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
@@ -1289,8 +1300,8 @@ func TestFormatSlackMessage_JSONBody(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
event := &database.Event{
|
||||
Method: "POST",
|
||||
ContentType: "application/json",
|
||||
Method: http.MethodPost,
|
||||
ContentType: testContentType,
|
||||
Body: `{"action":"push",` +
|
||||
`"repo":"test/repo",` +
|
||||
`"ref":"refs/heads/main"}`,
|
||||
@@ -1315,7 +1326,7 @@ func TestFormatSlackMessage_NonJSONBody(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
event := &database.Event{
|
||||
Method: "POST",
|
||||
Method: http.MethodPost,
|
||||
ContentType: "text/plain",
|
||||
Body: "hello world plain text",
|
||||
}
|
||||
@@ -1338,8 +1349,8 @@ func TestFormatSlackMessage_EmptyBody(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
event := &database.Event{
|
||||
Method: "POST",
|
||||
ContentType: "application/json",
|
||||
Method: http.MethodPost,
|
||||
ContentType: testContentType,
|
||||
Body: "",
|
||||
}
|
||||
event.CreatedAt = time.Date(
|
||||
@@ -1367,8 +1378,8 @@ func TestFormatSlackMessage_LargeJSONTruncated(
|
||||
require.NoError(t, err)
|
||||
|
||||
event := &database.Event{
|
||||
Method: "POST",
|
||||
ContentType: "application/json",
|
||||
Method: http.MethodPost,
|
||||
ContentType: testContentType,
|
||||
Body: string(largeJSON),
|
||||
}
|
||||
event.CreatedAt = time.Date(
|
||||
@@ -1697,7 +1708,7 @@ func assertLogLineComplete(
|
||||
"log line must contain the webhook id",
|
||||
)
|
||||
|
||||
assert.Contains(t, out, "application/json",
|
||||
assert.Contains(t, out, testContentType,
|
||||
"log line must contain the content type",
|
||||
)
|
||||
}
|
||||
|
||||
@@ -273,3 +273,64 @@ func NewTestCircuitBreaker(
|
||||
cooldown: cooldown,
|
||||
}
|
||||
}
|
||||
|
||||
// ExportArchivedEvent aliases the archive row type so black-box
|
||||
// tests can construct and read archive rows.
|
||||
type ExportArchivedEvent = archivedEvent
|
||||
|
||||
// ExportArchiveWriter wraps an archiveWriter so black-box tests
|
||||
// can exercise the per-webhook archive file mechanics.
|
||||
type ExportArchiveWriter struct {
|
||||
w *archiveWriter
|
||||
}
|
||||
|
||||
// NewExportArchiveWriter builds an archive writer for tests,
|
||||
// optionally overriding the reopen debounce (a non-positive
|
||||
// debounce keeps the production default).
|
||||
func NewExportArchiveWriter(
|
||||
path string, log *slog.Logger, debounce time.Duration,
|
||||
) *ExportArchiveWriter {
|
||||
w := newArchiveWriter(path, log)
|
||||
if debounce > 0 {
|
||||
w.debounce = debounce
|
||||
}
|
||||
|
||||
return &ExportArchiveWriter{w: w}
|
||||
}
|
||||
|
||||
// Write archives a row through the writer.
|
||||
func (e *ExportArchiveWriter) Write(
|
||||
row ExportArchivedEvent, expiry time.Duration,
|
||||
) error {
|
||||
return e.w.write(row, expiry)
|
||||
}
|
||||
|
||||
// Open opens the archive file, pruning when expiry is positive.
|
||||
func (e *ExportArchiveWriter) Open(expiry time.Duration) error {
|
||||
return e.w.open(expiry)
|
||||
}
|
||||
|
||||
// Reopen closes and reopens the archive file.
|
||||
func (e *ExportArchiveWriter) Reopen(
|
||||
expiry time.Duration,
|
||||
) error {
|
||||
return e.w.reopen(expiry)
|
||||
}
|
||||
|
||||
// Reopens reports how many times the file has been opened.
|
||||
func (e *ExportArchiveWriter) Reopens() int {
|
||||
return e.w.reopens
|
||||
}
|
||||
|
||||
// DB returns the writer's current open handle for row
|
||||
// inspection in tests.
|
||||
func (e *ExportArchiveWriter) DB() *gorm.DB {
|
||||
return e.w.db
|
||||
}
|
||||
|
||||
// ExportParseArchiveExpiry exposes parseArchiveExpiry.
|
||||
func ExportParseArchiveExpiry(
|
||||
configJSON string,
|
||||
) (time.Duration, error) {
|
||||
return parseArchiveExpiry(configJSON)
|
||||
}
|
||||
|
||||
@@ -2,21 +2,38 @@ package delivery
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
|
||||
"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.)
|
||||
// databaseTarget is a no-retry target that archives the
|
||||
// full inbound event into a per-webhook archive SQLite file,
|
||||
// separate from the per-webhook event database. The event is
|
||||
// already persisted in the per-webhook event DB by the time
|
||||
// delivery runs; the database target additionally writes a
|
||||
// durable long-term copy into archive-{webhookID}.db and then
|
||||
// records a single attempt whose outcome reflects whether the
|
||||
// archive write succeeded. See archiveWriter for the
|
||||
// close/reopen, auto-recreate, and expiry semantics.
|
||||
type databaseTarget struct {
|
||||
eng *Engine
|
||||
|
||||
mu sync.Mutex
|
||||
writers map[string]*archiveWriter
|
||||
}
|
||||
|
||||
// Deliver implements Target.
|
||||
// Deliver implements Target. It archives the event, then
|
||||
// records one successful attempt and marks the delivery
|
||||
// delivered. An archiving error fails the delivery: the
|
||||
// attempt is recorded as failed with the error and the
|
||||
// delivery is marked failed, so a target that could not do
|
||||
// its one job (archiving) never reports success. The target
|
||||
// does not retry; the event remains durably stored in the
|
||||
// per-webhook event database.
|
||||
func (t *databaseTarget) Deliver(
|
||||
_ context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
@@ -24,6 +41,27 @@ func (t *databaseTarget) Deliver(
|
||||
_ *Task,
|
||||
_ Scheduler,
|
||||
) {
|
||||
err := t.archive(d)
|
||||
if err != nil {
|
||||
t.eng.log.Error(
|
||||
"failed to archive event to database target",
|
||||
"delivery_id", d.ID,
|
||||
"event_id", d.EventID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
t.eng.recordResult(
|
||||
webhookDB, d, 1, false, 0, "",
|
||||
err.Error(), 0,
|
||||
)
|
||||
|
||||
t.eng.updateDeliveryStatus(
|
||||
webhookDB, d, database.DeliveryStatusFailed,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
t.eng.recordResult(
|
||||
webhookDB, d, 1, true, 0, "", "", 0,
|
||||
)
|
||||
@@ -32,3 +70,68 @@ func (t *databaseTarget) Deliver(
|
||||
webhookDB, d, database.DeliveryStatusDelivered,
|
||||
)
|
||||
}
|
||||
|
||||
// archive writes the full event as a row into the webhook's
|
||||
// archive database, honouring the optional per-target expiry
|
||||
// parsed from the target config JSON.
|
||||
func (t *databaseTarget) archive(d *database.Delivery) error {
|
||||
webhookID := d.Event.WebhookID
|
||||
if webhookID == "" {
|
||||
return errArchiveMissingWebhookID
|
||||
}
|
||||
|
||||
expiry, err := parseArchiveExpiry(d.Target.Config)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
w, err := t.writerFor(webhookID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
row := archivedEvent{
|
||||
EventID: d.Event.ID,
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: d.Event.EntrypointID,
|
||||
Method: d.Event.Method,
|
||||
Headers: d.Event.Headers,
|
||||
Body: d.Event.Body,
|
||||
ContentType: d.Event.ContentType,
|
||||
}
|
||||
|
||||
return w.write(row, expiry)
|
||||
}
|
||||
|
||||
// writerFor returns the archiveWriter for a webhook, creating
|
||||
// and caching it on first use. Each webhook has one writer so
|
||||
// its close/reopen debounce state is shared across concurrent
|
||||
// deliveries. The archive file lives beside the per-webhook
|
||||
// event database in the data directory.
|
||||
func (t *databaseTarget) writerFor(
|
||||
webhookID string,
|
||||
) (*archiveWriter, error) {
|
||||
if t.eng.dbManager == nil {
|
||||
return nil, errArchiveNoDataDir
|
||||
}
|
||||
|
||||
dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID))
|
||||
path := filepath.Join(
|
||||
dir, fmt.Sprintf("archive-%s.db", webhookID),
|
||||
)
|
||||
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
|
||||
if t.writers == nil {
|
||||
t.writers = make(map[string]*archiveWriter)
|
||||
}
|
||||
|
||||
w, ok := t.writers[webhookID]
|
||||
if !ok {
|
||||
w = newArchiveWriter(path, t.eng.log)
|
||||
t.writers[webhookID] = w
|
||||
}
|
||||
|
||||
return w, nil
|
||||
}
|
||||
|
||||
312
internal/delivery/target_database_archive.go
Normal file
312
internal/delivery/target_database_archive.go
Normal file
@@ -0,0 +1,312 @@
|
||||
package delivery
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// archiveExpiryNever is the expiry sentinel (and default) that
|
||||
// disables pruning so archived rows are kept forever.
|
||||
const archiveExpiryNever = "never"
|
||||
|
||||
// archiveReopenDebounce bounds how often an archive file is
|
||||
// closed and reopened. After each write the handle is closed
|
||||
// and reopened so an operator can move the file away for
|
||||
// offline archiving, but never more than once per this window.
|
||||
const archiveReopenDebounce = time.Second
|
||||
|
||||
var (
|
||||
// errArchiveMissingWebhookID is returned when an event to
|
||||
// archive has no webhook id to key its archive file on.
|
||||
errArchiveMissingWebhookID = errors.New(
|
||||
"cannot archive event without a webhook id",
|
||||
)
|
||||
|
||||
// errArchiveNoDataDir is returned when the database target
|
||||
// has no webhook database manager and so cannot locate the
|
||||
// data directory for archive files.
|
||||
errArchiveNoDataDir = errors.New(
|
||||
"database target has no data directory",
|
||||
)
|
||||
|
||||
// errArchiveExpiryNotPositive is returned when a
|
||||
// user-supplied archive expiry parses as a duration but is
|
||||
// zero or negative; "never" is the way to disable pruning.
|
||||
errArchiveExpiryNotPositive = errors.New(
|
||||
"expiry must be a positive duration or \"never\"",
|
||||
)
|
||||
)
|
||||
|
||||
// databaseTargetConfig is the optional per-target JSON config
|
||||
// for a database (archive) target.
|
||||
type databaseTargetConfig struct {
|
||||
// Expiry is a Go duration (e.g. "720h") after which
|
||||
// archived rows are pruned, or "never" (the default) to
|
||||
// keep them forever.
|
||||
Expiry string `json:"expiry"`
|
||||
}
|
||||
|
||||
// archivedEvent is one fully captured webhook event stored in a
|
||||
// per-webhook archive database for long-term retention. It is a
|
||||
// self-contained copy — independent of the per-webhook event
|
||||
// database, which may prune events under its own retention.
|
||||
type archivedEvent struct {
|
||||
ID uint `gorm:"primaryKey;autoIncrement"`
|
||||
EventID string `gorm:"index"`
|
||||
WebhookID string
|
||||
EntrypointID string
|
||||
Method string
|
||||
Headers string
|
||||
Body string
|
||||
ContentType string
|
||||
|
||||
// ArchivedAt is when the row was archived and is the age
|
||||
// basis for expiry pruning.
|
||||
ArchivedAt time.Time `gorm:"index"`
|
||||
}
|
||||
|
||||
// parseArchiveExpiry reads the optional expiry from a database
|
||||
// target's config JSON. An empty config, an empty expiry, or
|
||||
// the literal "never" all mean keep forever, returned as a zero
|
||||
// duration. Any other value must parse as a positive Go
|
||||
// duration; a set-but-invalid value (unparseable, zero, or
|
||||
// negative) is an error rather than a silent default, matching
|
||||
// ValidateArchiveExpiry at target creation.
|
||||
func parseArchiveExpiry(
|
||||
configJSON string,
|
||||
) (time.Duration, error) {
|
||||
if configJSON == "" {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
var cfg databaseTargetConfig
|
||||
|
||||
err := json.Unmarshal([]byte(configJSON), &cfg)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"parsing database target config: %w", err,
|
||||
)
|
||||
}
|
||||
|
||||
if cfg.Expiry == "" || cfg.Expiry == archiveExpiryNever {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
dur, err := time.ParseDuration(cfg.Expiry)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"parsing archive expiry %q: %w", cfg.Expiry, err,
|
||||
)
|
||||
}
|
||||
|
||||
if dur <= 0 {
|
||||
return 0, fmt.Errorf(
|
||||
"%w: %q", errArchiveExpiryNotPositive, cfg.Expiry,
|
||||
)
|
||||
}
|
||||
|
||||
return dur, nil
|
||||
}
|
||||
|
||||
// ValidateArchiveExpiry checks a user-supplied archive expiry
|
||||
// for a database target at configuration time. Valid values are
|
||||
// empty, "never" (both meaning keep forever), or a positive Go
|
||||
// duration such as "720h". Anything else is an error, so a bad
|
||||
// expiry is rejected when the target is created rather than
|
||||
// failing every subsequent delivery.
|
||||
func ValidateArchiveExpiry(expiry string) error {
|
||||
if expiry == "" || expiry == archiveExpiryNever {
|
||||
return nil
|
||||
}
|
||||
|
||||
dur, err := time.ParseDuration(expiry)
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"expiry must be %q or a Go duration "+
|
||||
"such as \"720h\": %w",
|
||||
archiveExpiryNever, err,
|
||||
)
|
||||
}
|
||||
|
||||
if dur <= 0 {
|
||||
return fmt.Errorf(
|
||||
"%w: %q", errArchiveExpiryNotPositive, expiry,
|
||||
)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// archiveWriter owns one per-webhook archive SQLite file. It
|
||||
// serialises writes, and after each write closes and reopens
|
||||
// the file (debounced to at most once per debounce window) so
|
||||
// an operator can move the file away for offline archiving. The
|
||||
// next write recreates a moved or removed file, because the
|
||||
// file is opened create-if-missing and its schema is migrated
|
||||
// on every open.
|
||||
type archiveWriter struct {
|
||||
mu sync.Mutex
|
||||
path string
|
||||
log *slog.Logger
|
||||
debounce time.Duration
|
||||
db *gorm.DB
|
||||
lastReopen time.Time
|
||||
reopens int
|
||||
}
|
||||
|
||||
// newArchiveWriter builds an archiveWriter for a file path with
|
||||
// the default reopen debounce.
|
||||
func newArchiveWriter(
|
||||
path string, log *slog.Logger,
|
||||
) *archiveWriter {
|
||||
return &archiveWriter{
|
||||
path: path,
|
||||
log: log,
|
||||
debounce: archiveReopenDebounce,
|
||||
}
|
||||
}
|
||||
|
||||
// write appends the event as a row, then applies the debounced
|
||||
// close/reopen. It recreates the archive file if it was moved
|
||||
// or removed since the last open. A positive expiry prunes rows
|
||||
// older than it on each (re)open.
|
||||
func (w *archiveWriter) write(
|
||||
row archivedEvent, expiry time.Duration,
|
||||
) error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
if w.db == nil || !fileExists(w.path) {
|
||||
err := w.reopen(expiry)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
row.ArchivedAt = time.Now()
|
||||
|
||||
err := w.db.Create(&row).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"archiving event to %s: %w", w.path, err,
|
||||
)
|
||||
}
|
||||
|
||||
if time.Since(w.lastReopen) >= w.debounce {
|
||||
return w.reopen(expiry)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// open opens (creating if missing) the archive file, migrates
|
||||
// its schema, records the reopen time, and prunes expired rows
|
||||
// when expiry is positive.
|
||||
func (w *archiveWriter) open(expiry time.Duration) error {
|
||||
dbURL := fmt.Sprintf("file:%s?mode=rwc", w.path)
|
||||
|
||||
sqlDB, err := sql.Open("sqlite", dbURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"opening archive database %s: %w", w.path, err,
|
||||
)
|
||||
}
|
||||
|
||||
gdb, err := gorm.Open(
|
||||
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||
)
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
|
||||
return fmt.Errorf(
|
||||
"connecting to archive database %s: %w",
|
||||
w.path, err,
|
||||
)
|
||||
}
|
||||
|
||||
err = gdb.AutoMigrate(&archivedEvent{})
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
|
||||
return fmt.Errorf(
|
||||
"migrating archive database %s: %w", w.path, err,
|
||||
)
|
||||
}
|
||||
|
||||
w.db = gdb
|
||||
w.lastReopen = time.Now()
|
||||
w.reopens++
|
||||
|
||||
if expiry > 0 {
|
||||
w.prune(expiry)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// reopen closes any open handle and opens the file afresh. The
|
||||
// fresh open recreates the file if it was moved away.
|
||||
func (w *archiveWriter) reopen(expiry time.Duration) error {
|
||||
w.close()
|
||||
|
||||
return w.open(expiry)
|
||||
}
|
||||
|
||||
// close closes the underlying handle, if any.
|
||||
func (w *archiveWriter) close() {
|
||||
if w.db == nil {
|
||||
return
|
||||
}
|
||||
|
||||
sqlDB, err := w.db.DB()
|
||||
if err == nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
|
||||
w.db = nil
|
||||
}
|
||||
|
||||
// prune deletes archived rows older than expiry, measured from
|
||||
// each row's archived time. It runs on every (re)open, and
|
||||
// because the file is reopened after writes this keeps the
|
||||
// archive swept without a separate background sweeper. Failures
|
||||
// are logged, not fatal: a prune error must not stop archiving.
|
||||
func (w *archiveWriter) prune(expiry time.Duration) {
|
||||
cutoff := time.Now().Add(-expiry)
|
||||
|
||||
res := w.db.Where("archived_at < ?", cutoff).
|
||||
Delete(&archivedEvent{})
|
||||
if res.Error != nil {
|
||||
w.log.Error(
|
||||
"failed to prune expired archive rows",
|
||||
"path", w.path,
|
||||
"error", res.Error,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if res.RowsAffected > 0 {
|
||||
w.log.Info(
|
||||
"pruned expired archive rows",
|
||||
"path", w.path,
|
||||
"rows_deleted", res.RowsAffected,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// fileExists reports whether a path currently exists.
|
||||
func fileExists(path string) bool {
|
||||
_, err := os.Stat(path)
|
||||
|
||||
return err == nil
|
||||
}
|
||||
395
internal/delivery/target_database_test.go
Normal file
395
internal/delivery/target_database_test.go
Normal file
@@ -0,0 +1,395 @@
|
||||
package delivery_test
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"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" // Pure Go SQLite driver.
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
)
|
||||
|
||||
func archiveTestLogger() *slog.Logger {
|
||||
return slog.New(slog.NewTextHandler(
|
||||
os.Stderr,
|
||||
&slog.HandlerOptions{Level: slog.LevelDebug},
|
||||
))
|
||||
}
|
||||
|
||||
// openArchiveDBForRead opens an archive file read-only so a
|
||||
// test can inspect the rows the writer persisted.
|
||||
func openArchiveDBForRead(
|
||||
t *testing.T, path string,
|
||||
) *gorm.DB {
|
||||
t.Helper()
|
||||
|
||||
sqlDB, err := sql.Open(
|
||||
"sqlite",
|
||||
fmt.Sprintf("file:%s?mode=ro", path),
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||
|
||||
gdb, err := gorm.Open(
|
||||
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{},
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
return gdb
|
||||
}
|
||||
|
||||
// removeArchiveFiles simulates an operator moving the archive
|
||||
// away by deleting the SQLite file and its sidecar files.
|
||||
func removeArchiveFiles(t *testing.T, path string) {
|
||||
t.Helper()
|
||||
|
||||
for _, suffix := range []string{
|
||||
"", "-wal", "-shm", "-journal",
|
||||
} {
|
||||
err := os.Remove(path + suffix)
|
||||
if err != nil && !os.IsNotExist(err) {
|
||||
t.Fatalf("removing %s%s: %v", path, suffix, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeliverDatabase_ArchivesEvent verifies that delivering to
|
||||
// a database target marks the delivery delivered and archives
|
||||
// the full event into a separate per-webhook archive file.
|
||||
func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dataDir := t.TempDir()
|
||||
dbMgr := database.NewTestWebhookDBManager(dataDir)
|
||||
|
||||
e := delivery.NewTestEngineWithDB(
|
||||
nil, dbMgr,
|
||||
archiveTestLogger(),
|
||||
&http.Client{Timeout: 5 * time.Second},
|
||||
1,
|
||||
)
|
||||
|
||||
webhookDB := testWebhookDB(t)
|
||||
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
||||
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
|
||||
|
||||
e.ExportDeliverDatabase(webhookDB, d)
|
||||
|
||||
var updated database.Delivery
|
||||
|
||||
require.NoError(t, webhookDB.First(
|
||||
&updated, "id = ?", d.ID,
|
||||
).Error)
|
||||
assert.Equal(t,
|
||||
database.DeliveryStatusDelivered, updated.Status,
|
||||
"database target should mark the delivery delivered",
|
||||
)
|
||||
|
||||
archivePath := filepath.Join(
|
||||
dataDir,
|
||||
fmt.Sprintf("archive-%s.db", event.WebhookID),
|
||||
)
|
||||
assert.FileExists(t, archivePath)
|
||||
|
||||
rdb := openArchiveDBForRead(t, archivePath)
|
||||
|
||||
var rows []delivery.ExportArchivedEvent
|
||||
|
||||
require.NoError(t, rdb.Find(&rows).Error)
|
||||
require.Len(t, rows, 1)
|
||||
assert.Equal(t, event.ID, rows[0].EventID)
|
||||
assert.Equal(t, event.WebhookID, rows[0].WebhookID)
|
||||
assert.Equal(t, event.Method, rows[0].Method)
|
||||
assert.JSONEq(t, `{"archived":true}`, rows[0].Body)
|
||||
}
|
||||
|
||||
func TestArchiveWriter_WritesRow(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "archive-wh.db")
|
||||
w := delivery.NewExportArchiveWriter(
|
||||
path, archiveTestLogger(), 0,
|
||||
)
|
||||
|
||||
row := delivery.ExportArchivedEvent{
|
||||
EventID: "ev-1",
|
||||
WebhookID: "wh-1",
|
||||
EntrypointID: "ep-1",
|
||||
Method: "POST",
|
||||
Headers: `{"X":"Y"}`,
|
||||
Body: `{"hello":"world"}`,
|
||||
ContentType: "application/json",
|
||||
}
|
||||
|
||||
require.NoError(t, w.Write(row, 0))
|
||||
assert.FileExists(t, path)
|
||||
|
||||
var got []delivery.ExportArchivedEvent
|
||||
|
||||
require.NoError(t, w.DB().Find(&got).Error)
|
||||
require.Len(t, got, 1)
|
||||
assert.Equal(t, "ev-1", got[0].EventID)
|
||||
assert.Equal(t, "wh-1", got[0].WebhookID)
|
||||
assert.Equal(t, "ep-1", got[0].EntrypointID)
|
||||
assert.Equal(t, row.Method, got[0].Method)
|
||||
assert.Equal(t, row.ContentType, got[0].ContentType)
|
||||
assert.JSONEq(t, `{"hello":"world"}`, got[0].Body)
|
||||
assert.False(t, got[0].ArchivedAt.IsZero())
|
||||
}
|
||||
|
||||
func TestArchiveWriter_RecreatesAfterRemoval(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "archive-wh.db")
|
||||
w := delivery.NewExportArchiveWriter(
|
||||
path, archiveTestLogger(), 0,
|
||||
)
|
||||
|
||||
require.NoError(t, w.Write(
|
||||
delivery.ExportArchivedEvent{EventID: "a"}, 0,
|
||||
))
|
||||
assert.FileExists(t, path)
|
||||
|
||||
// The operator moves the archive away while the handle is
|
||||
// still open.
|
||||
removeArchiveFiles(t, path)
|
||||
require.NoFileExists(t, path)
|
||||
|
||||
// The next write recreates the file with a fresh schema and
|
||||
// only the new row.
|
||||
require.NoError(t, w.Write(
|
||||
delivery.ExportArchivedEvent{EventID: "b"}, 0,
|
||||
))
|
||||
assert.FileExists(t, path)
|
||||
|
||||
var got []delivery.ExportArchivedEvent
|
||||
|
||||
require.NoError(t, w.DB().Find(&got).Error)
|
||||
require.Len(t, got, 1)
|
||||
assert.Equal(t, "b", got[0].EventID)
|
||||
}
|
||||
|
||||
func TestArchiveWriter_ReopenDebounce(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// A generous debounce keeps the two rapid writes inside
|
||||
// the window even on a heavily loaded test machine.
|
||||
path := filepath.Join(t.TempDir(), "archive-wh.db")
|
||||
w := delivery.NewExportArchiveWriter(
|
||||
path, archiveTestLogger(), 2*time.Second,
|
||||
)
|
||||
|
||||
require.NoError(t, w.Write(
|
||||
delivery.ExportArchivedEvent{EventID: "a"}, 0,
|
||||
))
|
||||
require.NoError(t, w.Write(
|
||||
delivery.ExportArchivedEvent{EventID: "b"}, 0,
|
||||
))
|
||||
|
||||
// Two writes inside the debounce window trigger only the
|
||||
// initial open — no extra close/reopen.
|
||||
assert.Equal(t, 1, w.Reopens())
|
||||
|
||||
time.Sleep(2100 * time.Millisecond)
|
||||
|
||||
require.NoError(t, w.Write(
|
||||
delivery.ExportArchivedEvent{EventID: "c"}, 0,
|
||||
))
|
||||
|
||||
// A write after the window elapses closes and reopens once.
|
||||
assert.Equal(t, 2, w.Reopens())
|
||||
}
|
||||
|
||||
func TestArchiveWriter_ExpiryPrune(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "archive-wh.db")
|
||||
w := delivery.NewExportArchiveWriter(
|
||||
path, archiveTestLogger(), 0,
|
||||
)
|
||||
|
||||
require.NoError(t, w.Open(0))
|
||||
|
||||
old := delivery.ExportArchivedEvent{
|
||||
EventID: "old",
|
||||
ArchivedAt: time.Now().Add(-2 * time.Hour),
|
||||
}
|
||||
fresh := delivery.ExportArchivedEvent{
|
||||
EventID: "fresh",
|
||||
ArchivedAt: time.Now(),
|
||||
}
|
||||
|
||||
require.NoError(t, w.DB().Create(&old).Error)
|
||||
require.NoError(t, w.DB().Create(&fresh).Error)
|
||||
|
||||
// Reopening with a one-hour expiry prunes the old row.
|
||||
require.NoError(t, w.Reopen(time.Hour))
|
||||
|
||||
var got []delivery.ExportArchivedEvent
|
||||
|
||||
require.NoError(t, w.DB().Find(&got).Error)
|
||||
require.Len(t, got, 1)
|
||||
assert.Equal(t, "fresh", got[0].EventID)
|
||||
}
|
||||
|
||||
func TestParseArchiveExpiry(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
in string
|
||||
want time.Duration
|
||||
wantErr bool
|
||||
}{
|
||||
{"empty config", "", 0, false},
|
||||
{"explicit never", `{"expiry":"never"}`, 0, false},
|
||||
{"empty expiry", `{"expiry":""}`, 0, false},
|
||||
{"duration", `{"expiry":"1h"}`, time.Hour, false},
|
||||
{"unparseable", `{"expiry":"nonsense"}`, 0, true},
|
||||
{"zero duration", `{"expiry":"0s"}`, 0, true},
|
||||
{"negative duration", `{"expiry":"-5h"}`, 0, true},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
got, err := delivery.ExportParseArchiveExpiry(tc.in)
|
||||
if tc.wantErr {
|
||||
require.Error(t, err)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, tc.want, got)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// seedDatabaseTargetDelivery seeds a pending delivery for a
|
||||
// database target with the given config JSON and returns the
|
||||
// in-memory delivery the target handler is invoked with.
|
||||
func seedDatabaseTargetDelivery(
|
||||
t *testing.T,
|
||||
webhookDB *gorm.DB,
|
||||
event database.Event,
|
||||
config string,
|
||||
) *database.Delivery {
|
||||
t.Helper()
|
||||
|
||||
dlv := seedDelivery(
|
||||
t, webhookDB, 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-db",
|
||||
Type: database.TargetTypeDatabase,
|
||||
Config: config,
|
||||
},
|
||||
}
|
||||
d.ID = dlv.ID
|
||||
|
||||
return d
|
||||
}
|
||||
|
||||
// TestDeliverDatabase_ArchiveFailureFailsDelivery verifies that
|
||||
// an archive error (here: an unparseable expiry in the target
|
||||
// config) fails the delivery loudly: the attempt is recorded as
|
||||
// failed with the error and the delivery is marked failed, not
|
||||
// delivered.
|
||||
func TestDeliverDatabase_ArchiveFailureFailsDelivery(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
dataDir := t.TempDir()
|
||||
|
||||
e := delivery.NewTestEngineWithDB(
|
||||
nil, database.NewTestWebhookDBManager(dataDir),
|
||||
archiveTestLogger(),
|
||||
&http.Client{Timeout: 5 * time.Second},
|
||||
1,
|
||||
)
|
||||
|
||||
webhookDB := testWebhookDB(t)
|
||||
event := seedEvent(t, webhookDB, `{"archived":false}`)
|
||||
d := seedDatabaseTargetDelivery(
|
||||
t, webhookDB, event, `{"expiry":"nonsense"}`,
|
||||
)
|
||||
|
||||
e.ExportDeliverDatabase(webhookDB, d)
|
||||
|
||||
var updated database.Delivery
|
||||
|
||||
require.NoError(t, webhookDB.First(
|
||||
&updated, "id = ?", d.ID,
|
||||
).Error)
|
||||
assert.Equal(t,
|
||||
database.DeliveryStatusFailed, updated.Status,
|
||||
"archive failure must mark the delivery failed",
|
||||
)
|
||||
|
||||
var results []database.DeliveryResult
|
||||
|
||||
require.NoError(t, webhookDB.Where(
|
||||
"delivery_id = ?", d.ID,
|
||||
).Find(&results).Error)
|
||||
require.Len(t, results, 1)
|
||||
assert.False(t,
|
||||
results[0].Success,
|
||||
"the attempt must be recorded as failed",
|
||||
)
|
||||
assert.Contains(t,
|
||||
results[0].Error, "nonsense",
|
||||
"the archive error must be recorded on the attempt",
|
||||
)
|
||||
|
||||
assert.NoFileExists(t,
|
||||
filepath.Join(
|
||||
dataDir,
|
||||
fmt.Sprintf("archive-%s.db", event.WebhookID),
|
||||
),
|
||||
"no archive file should exist for a failed config",
|
||||
)
|
||||
}
|
||||
|
||||
func TestValidateArchiveExpiry(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
valid := []string{"", "never", "1h", "720h", "30m"}
|
||||
for _, in := range valid {
|
||||
require.NoError(t,
|
||||
delivery.ValidateArchiveExpiry(in),
|
||||
"expiry %q should be accepted", in,
|
||||
)
|
||||
}
|
||||
|
||||
invalid := []string{"nonsense", "7d", "-5h", "0s", "0"}
|
||||
for _, in := range invalid {
|
||||
require.Error(t,
|
||||
delivery.ValidateArchiveExpiry(in),
|
||||
"expiry %q should be rejected", in,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -495,5 +495,5 @@ func applyRequestHeaders(
|
||||
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
|
||||
return client.Do(req) //#nosec G704 -- validated URL, SSRF-safe transport
|
||||
}
|
||||
|
||||
@@ -19,7 +19,7 @@ func (h *Handlers) HandleLoginPage() http.HandlerFunc {
|
||||
|
||||
// Render login page
|
||||
data := map[string]any{
|
||||
"Error": "",
|
||||
tmplKeyError: "",
|
||||
}
|
||||
|
||||
h.renderTemplate(w, r, "login.html", data)
|
||||
@@ -86,7 +86,7 @@ func (h *Handlers) renderLoginError(
|
||||
status int,
|
||||
) {
|
||||
data := map[string]any{
|
||||
"Error": msg,
|
||||
tmplKeyError: msg,
|
||||
}
|
||||
|
||||
w.WriteHeader(status)
|
||||
|
||||
@@ -13,12 +13,26 @@ func (s *Handlers) RenderTemplateForTest(
|
||||
s.renderTemplate(w, r, pageTemplate, data)
|
||||
}
|
||||
|
||||
// BuildSlackTargetConfigForTest exposes buildSlackTargetConfig
|
||||
// for use in the handlers_test package.
|
||||
// BuildSlackTargetConfigForTest exposes buildURLTargetConfig
|
||||
// with the Slack target parameters for use in the
|
||||
// handlers_test package.
|
||||
func (s *Handlers) BuildSlackTargetConfigForTest(
|
||||
w http.ResponseWriter,
|
||||
r *http.Request,
|
||||
targetURL string,
|
||||
) (string, error) {
|
||||
return s.buildSlackTargetConfig(w, r, targetURL)
|
||||
return s.buildURLTargetConfig(
|
||||
w, r, targetURL, "webhookUrl",
|
||||
"Webhook URL is required for Slack targets",
|
||||
)
|
||||
}
|
||||
|
||||
// BuildDatabaseTargetConfigForTest exposes
|
||||
// buildDatabaseTargetConfig for use in the handlers_test
|
||||
// package.
|
||||
func (s *Handlers) BuildDatabaseTargetConfigForTest(
|
||||
w http.ResponseWriter,
|
||||
expiry string,
|
||||
) (string, error) {
|
||||
return s.buildDatabaseTargetConfig(w, expiry)
|
||||
}
|
||||
|
||||
@@ -30,6 +30,11 @@ const (
|
||||
defaultRetentionDays = 30
|
||||
// paginationPerPage is the number of items per page.
|
||||
paginationPerPage = 25
|
||||
|
||||
// tmplKeyError is the template data key for an error message.
|
||||
tmplKeyError = "Error"
|
||||
// tmplKeyWebhook is the template data key for a webhook.
|
||||
tmplKeyWebhook = "Webhook"
|
||||
)
|
||||
|
||||
// errInvalidPassword is returned when a password does not match.
|
||||
|
||||
@@ -186,3 +186,57 @@ func TestRenderTemplate(t *testing.T) {
|
||||
t, http.StatusInternalServerError, w.Code,
|
||||
)
|
||||
}
|
||||
|
||||
func TestBuildDatabaseTargetConfig_Valid(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var h *handlers.Handlers
|
||||
|
||||
app := newTestApp(t, &h)
|
||||
app.RequireStart()
|
||||
|
||||
t.Cleanup(app.RequireStop)
|
||||
|
||||
// Empty expiry: the keep-forever default, empty config.
|
||||
w := httptest.NewRecorder()
|
||||
cfg, err := h.BuildDatabaseTargetConfigForTest(w, "")
|
||||
require.NoError(t, err)
|
||||
assert.Empty(t, cfg)
|
||||
|
||||
// Explicit never is stored as config.
|
||||
w = httptest.NewRecorder()
|
||||
cfg, err = h.BuildDatabaseTargetConfigForTest(w, "never")
|
||||
require.NoError(t, err)
|
||||
assert.JSONEq(t, `{"expiry":"never"}`, cfg)
|
||||
|
||||
// A positive duration is stored as config.
|
||||
w = httptest.NewRecorder()
|
||||
cfg, err = h.BuildDatabaseTargetConfigForTest(w, "720h")
|
||||
require.NoError(t, err)
|
||||
assert.JSONEq(t, `{"expiry":"720h"}`, cfg)
|
||||
}
|
||||
|
||||
func TestBuildDatabaseTargetConfig_RejectsBadExpiry(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
var h *handlers.Handlers
|
||||
|
||||
app := newTestApp(t, &h)
|
||||
app.RequireStart()
|
||||
|
||||
t.Cleanup(app.RequireStop)
|
||||
|
||||
for _, bad := range []string{"nonsense", "7d", "-5h"} {
|
||||
w := httptest.NewRecorder()
|
||||
cfg, err := h.BuildDatabaseTargetConfigForTest(w, bad)
|
||||
|
||||
require.Error(t, err, "expiry %q", bad)
|
||||
assert.Empty(t, cfg)
|
||||
assert.Equal(
|
||||
t, http.StatusBadRequest, w.Code,
|
||||
"expiry %q should be rejected with 400", bad,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/go-chi/chi"
|
||||
"github.com/google/uuid"
|
||||
@@ -106,7 +107,7 @@ func (h *Handlers) buildWebhookListItems(
|
||||
func (h *Handlers) HandleSourceCreate() http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
data := map[string]any{
|
||||
"Error": "",
|
||||
tmplKeyError: "",
|
||||
}
|
||||
|
||||
h.renderTemplate(w, r, "sources_new.html", data)
|
||||
@@ -145,7 +146,7 @@ func (h *Handlers) HandleSourceCreateSubmit() http.HandlerFunc {
|
||||
|
||||
if name == "" {
|
||||
data := map[string]any{
|
||||
"Error": "Name is required",
|
||||
tmplKeyError: "Name is required",
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
@@ -315,7 +316,7 @@ func (h *Handlers) renderSourceDetail(
|
||||
}
|
||||
|
||||
data := map[string]any{
|
||||
"Webhook": webhook,
|
||||
tmplKeyWebhook: webhook,
|
||||
"Entrypoints": entrypoints,
|
||||
"Targets": targets,
|
||||
"Events": events,
|
||||
@@ -351,8 +352,8 @@ func (h *Handlers) HandleSourceEdit() http.HandlerFunc {
|
||||
}
|
||||
|
||||
data := map[string]any{
|
||||
"Webhook": webhook,
|
||||
"Error": "",
|
||||
tmplKeyWebhook: webhook,
|
||||
tmplKeyError: "",
|
||||
}
|
||||
|
||||
h.renderTemplate(w, r, "source_edit.html", data)
|
||||
@@ -415,8 +416,8 @@ func (h *Handlers) applyWebhookEdit(
|
||||
name := r.FormValue("name")
|
||||
if name == "" {
|
||||
data := map[string]any{
|
||||
"Webhook": *webhook,
|
||||
"Error": "Name is required",
|
||||
tmplKeyWebhook: *webhook,
|
||||
tmplKeyError: "Name is required",
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
@@ -589,7 +590,7 @@ func (h *Handlers) HandleSourceLogs() http.HandlerFunc {
|
||||
}
|
||||
|
||||
data := map[string]any{
|
||||
"Webhook": webhook,
|
||||
tmplKeyWebhook: webhook,
|
||||
"Events": evts,
|
||||
"Page": page,
|
||||
"TotalPages": totalPages,
|
||||
@@ -815,6 +816,7 @@ func (h *Handlers) processTargetCreate(
|
||||
targetType := database.TargetType(r.FormValue("type"))
|
||||
targetURL := r.FormValue("url")
|
||||
maxRetriesStr := r.FormValue("max_retries")
|
||||
expiry := r.FormValue("expiry")
|
||||
|
||||
if name == "" {
|
||||
http.Error(
|
||||
@@ -834,7 +836,7 @@ func (h *Handlers) processTargetCreate(
|
||||
}
|
||||
|
||||
configJSON, err := h.buildTargetConfig(
|
||||
w, r, targetType, targetURL,
|
||||
w, r, targetType, targetURL, expiry,
|
||||
)
|
||||
if err != nil {
|
||||
return
|
||||
@@ -892,18 +894,28 @@ func parseNonNegativeInt(s string) int {
|
||||
}
|
||||
|
||||
// buildTargetConfig builds the JSON config string for a target.
|
||||
// The expiry form value is read by the caller (which bounds the
|
||||
// request body) and applies to database targets only.
|
||||
func (h *Handlers) buildTargetConfig(
|
||||
w http.ResponseWriter,
|
||||
r *http.Request,
|
||||
targetType database.TargetType,
|
||||
targetURL string,
|
||||
targetURL, expiry string,
|
||||
) (string, error) {
|
||||
switch targetType {
|
||||
case database.TargetTypeHTTP:
|
||||
return h.buildHTTPTargetConfig(w, r, targetURL)
|
||||
return h.buildURLTargetConfig(
|
||||
w, r, targetURL, "url",
|
||||
"URL is required for HTTP targets",
|
||||
)
|
||||
case database.TargetTypeSlack:
|
||||
return h.buildSlackTargetConfig(w, r, targetURL)
|
||||
case database.TargetTypeDatabase, database.TargetTypeLog:
|
||||
return h.buildURLTargetConfig(
|
||||
w, r, targetURL, "webhookUrl",
|
||||
"Webhook URL is required for Slack targets",
|
||||
)
|
||||
case database.TargetTypeDatabase:
|
||||
return h.buildDatabaseTargetConfig(w, expiry)
|
||||
case database.TargetTypeLog:
|
||||
return "", nil
|
||||
default:
|
||||
http.Error(
|
||||
@@ -915,16 +927,18 @@ func (h *Handlers) buildTargetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
// buildHTTPTargetConfig builds config JSON for an HTTP target.
|
||||
func (h *Handlers) buildHTTPTargetConfig(
|
||||
// buildURLTargetConfig builds config JSON for a target whose
|
||||
// configuration is a single SSRF-validated URL stored under
|
||||
// configKey. missingMsg is the error shown when no URL is given.
|
||||
func (h *Handlers) buildURLTargetConfig(
|
||||
w http.ResponseWriter,
|
||||
r *http.Request,
|
||||
targetURL string,
|
||||
targetURL, configKey, missingMsg string,
|
||||
) (string, error) {
|
||||
if targetURL == "" {
|
||||
http.Error(
|
||||
w,
|
||||
"URL is required for HTTP targets",
|
||||
missingMsg,
|
||||
http.StatusBadRequest,
|
||||
)
|
||||
|
||||
@@ -949,7 +963,7 @@ func (h *Handlers) buildHTTPTargetConfig(
|
||||
return "", err
|
||||
}
|
||||
|
||||
cfg := map[string]any{"url": targetURL}
|
||||
cfg := map[string]any{configKey: targetURL}
|
||||
|
||||
configBytes, err := json.Marshal(cfg)
|
||||
if err != nil {
|
||||
@@ -964,41 +978,33 @@ func (h *Handlers) buildHTTPTargetConfig(
|
||||
return string(configBytes), nil
|
||||
}
|
||||
|
||||
// buildSlackTargetConfig builds config JSON for a Slack target.
|
||||
func (h *Handlers) buildSlackTargetConfig(
|
||||
// buildDatabaseTargetConfig builds config JSON for a database
|
||||
// (archive) target. The optional expiry (a form value read by
|
||||
// the caller, which bounds the request body) is validated here,
|
||||
// at creation time, so an unparseable value is rejected with a
|
||||
// 400 instead of failing every subsequent delivery. An empty
|
||||
// expiry yields an empty config (the keep-forever default).
|
||||
func (h *Handlers) buildDatabaseTargetConfig(
|
||||
w http.ResponseWriter,
|
||||
r *http.Request,
|
||||
targetURL string,
|
||||
expiry string,
|
||||
) (string, error) {
|
||||
if targetURL == "" {
|
||||
http.Error(
|
||||
w,
|
||||
"Webhook URL is required for Slack targets",
|
||||
http.StatusBadRequest,
|
||||
)
|
||||
|
||||
return "", errMissingURL
|
||||
expiry = strings.TrimSpace(expiry)
|
||||
if expiry == "" {
|
||||
return "", nil
|
||||
}
|
||||
|
||||
err := delivery.ValidateTargetURL(
|
||||
r.Context(), targetURL,
|
||||
)
|
||||
err := delivery.ValidateArchiveExpiry(expiry)
|
||||
if err != nil {
|
||||
h.log.Warn(
|
||||
"target URL blocked by SSRF protection",
|
||||
"url", targetURL,
|
||||
"error", err,
|
||||
)
|
||||
http.Error(
|
||||
w,
|
||||
"Invalid target URL: "+err.Error(),
|
||||
"Invalid archive expiry: "+err.Error(),
|
||||
http.StatusBadRequest,
|
||||
)
|
||||
|
||||
return "", err
|
||||
}
|
||||
|
||||
cfg := map[string]any{"webhookUrl": targetURL}
|
||||
cfg := map[string]any{"expiry": expiry}
|
||||
|
||||
configBytes, err := json.Marshal(cfg)
|
||||
if err != nil {
|
||||
|
||||
@@ -484,8 +484,13 @@ func metricsAuthMiddleware(
|
||||
return middleware.NewForTest(log, cfg, sessManager)
|
||||
}
|
||||
|
||||
func TestMetricsAuth_ValidCredentials(t *testing.T) {
|
||||
t.Parallel()
|
||||
// runMetricsAuthRequest sends a GET /metrics request with the
|
||||
// given basic-auth password through MetricsAuth and reports
|
||||
// whether the wrapped handler ran plus the recorded response.
|
||||
func runMetricsAuthRequest(
|
||||
t *testing.T, password string,
|
||||
) (bool, *httptest.ResponseRecorder) {
|
||||
t.Helper()
|
||||
|
||||
m := metricsAuthMiddleware(t)
|
||||
|
||||
@@ -503,12 +508,20 @@ func TestMetricsAuth_ValidCredentials(t *testing.T) {
|
||||
context.Background(),
|
||||
http.MethodGet, "/metrics", nil,
|
||||
)
|
||||
req.SetBasicAuth("admin", "secret")
|
||||
req.SetBasicAuth("admin", password)
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
|
||||
handler.ServeHTTP(w, req)
|
||||
|
||||
return called, w
|
||||
}
|
||||
|
||||
func TestMetricsAuth_ValidCredentials(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
called, w := runMetricsAuthRequest(t, "secret")
|
||||
|
||||
assert.True(
|
||||
t, called,
|
||||
"handler should be called with valid basic auth",
|
||||
@@ -519,27 +532,7 @@ func TestMetricsAuth_ValidCredentials(t *testing.T) {
|
||||
func TestMetricsAuth_InvalidCredentials(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
m := metricsAuthMiddleware(t)
|
||||
|
||||
var called bool
|
||||
|
||||
handler := m.MetricsAuth()(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
called = true
|
||||
|
||||
w.WriteHeader(http.StatusOK)
|
||||
},
|
||||
))
|
||||
|
||||
req := httptest.NewRequestWithContext(
|
||||
context.Background(),
|
||||
http.MethodGet, "/metrics", nil,
|
||||
)
|
||||
req.SetBasicAuth("admin", "wrong-password")
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
|
||||
handler.ServeHTTP(w, req)
|
||||
called, w := runMetricsAuthRequest(t, "wrong-password")
|
||||
|
||||
assert.False(
|
||||
t, called,
|
||||
|
||||
@@ -173,8 +173,18 @@ func TestSetUser_SetsAllFields(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
func TestGetUserID(t *testing.T) {
|
||||
t.Parallel()
|
||||
// testSessionGetter exercises a session string getter before and
|
||||
// after SetUser: it must report false with an empty value on a
|
||||
// fresh session, then true with the expected value once
|
||||
// SetUser(sess, "user-xyz", "bob") has run.
|
||||
func testSessionGetter(
|
||||
t *testing.T,
|
||||
get func(
|
||||
*session.Session, *sessions.Session,
|
||||
) (string, bool),
|
||||
expected string,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
s := testSession(t)
|
||||
|
||||
@@ -185,44 +195,46 @@ func TestGetUserID(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
|
||||
// Before setting user
|
||||
userID, ok := s.GetUserID(sess)
|
||||
val, ok := get(s, sess)
|
||||
assert.False(
|
||||
t, ok, "should return false when no user ID is set",
|
||||
t, ok, "should return false before SetUser",
|
||||
)
|
||||
assert.Empty(t, userID)
|
||||
assert.Empty(t, val)
|
||||
|
||||
// After setting user
|
||||
s.SetUser(sess, "user-xyz", "bob")
|
||||
|
||||
userID, ok = s.GetUserID(sess)
|
||||
val, ok = get(s, sess)
|
||||
assert.True(t, ok)
|
||||
assert.Equal(t, "user-xyz", userID)
|
||||
assert.Equal(t, expected, val)
|
||||
}
|
||||
|
||||
func TestGetUserID(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
testSessionGetter(
|
||||
t,
|
||||
func(
|
||||
s *session.Session, sess *sessions.Session,
|
||||
) (string, bool) {
|
||||
return s.GetUserID(sess)
|
||||
},
|
||||
"user-xyz",
|
||||
)
|
||||
}
|
||||
|
||||
func TestGetUsername(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
s := testSession(t)
|
||||
|
||||
req := httptest.NewRequestWithContext(
|
||||
context.Background(), http.MethodGet, "/", nil)
|
||||
|
||||
sess, err := s.Get(req)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Before setting user
|
||||
username, ok := s.GetUsername(sess)
|
||||
assert.False(
|
||||
t, ok, "should return false when no username is set",
|
||||
testSessionGetter(
|
||||
t,
|
||||
func(
|
||||
s *session.Session, sess *sessions.Session,
|
||||
) (string, bool) {
|
||||
return s.GetUsername(sess)
|
||||
},
|
||||
"bob",
|
||||
)
|
||||
assert.Empty(t, username)
|
||||
|
||||
// After setting user
|
||||
s.SetUser(sess, "user-xyz", "bob")
|
||||
|
||||
username, ok = s.GetUsername(sess)
|
||||
assert.True(t, ok)
|
||||
assert.Equal(t, "bob", username)
|
||||
}
|
||||
|
||||
// --- IsAuthenticated Tests ---
|
||||
|
||||
@@ -12,7 +12,12 @@ import (
|
||||
// middleware and handler tests to use real session functionality. The key
|
||||
// parameter is the raw 32-byte authentication key used for session encryption
|
||||
// and CSRF cookie signing.
|
||||
func NewForTest(store *sessions.CookieStore, cfg *config.Config, log *slog.Logger, key []byte) *Session {
|
||||
func NewForTest(
|
||||
store *sessions.CookieStore,
|
||||
cfg *config.Config,
|
||||
log *slog.Logger,
|
||||
key []byte,
|
||||
) *Session {
|
||||
return &Session{
|
||||
store: store,
|
||||
key: key,
|
||||
|
||||
@@ -10,11 +10,11 @@ set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
# Pinned versions, 2026-07-07. Never "latest"; exact versions only.
|
||||
GOLANGCI_LINT_VERSION="2.11.3"
|
||||
# sha256 of golangci-lint-2.11.3-linux-<arch>.tar.gz release archives
|
||||
GOLANGCI_LINT_SHA256_AMD64="87bb8cddbcc825d5778b64e8a91b46c0526b247f4e2f2904dea74ec7450475d1"
|
||||
GOLANGCI_LINT_SHA256_ARM64="ee3d95f301359e7d578e6d99c8ad5aeadbabc5a13009a30b2b0df11c8058afe9"
|
||||
# Pinned versions, 2026-08-07. Never "latest"; exact versions only.
|
||||
GOLANGCI_LINT_VERSION="2.12.2"
|
||||
# sha256 of golangci-lint-2.12.2-linux-<arch>.tar.gz release archives
|
||||
GOLANGCI_LINT_SHA256_AMD64="8df580d2670fed8fa984aac0507099af8df275e665215f5c7a2ae3943893a553"
|
||||
GOLANGCI_LINT_SHA256_ARM64="44cd40a8c76c86755375adfeea52cfd3533cb43d7bd647771e0ae065e166df3a"
|
||||
|
||||
PKGMGR=""
|
||||
SUDO=""
|
||||
|
||||
@@ -113,6 +113,10 @@
|
||||
<input type="url" name="url" placeholder="https://hooks.slack.com/services/..." :disabled="targetType !== 'slack'" class="input text-sm">
|
||||
<p class="text-xs text-gray-500 mt-1">Slack or Mattermost incoming webhook URL. Payloads are pretty-printed in code blocks.</p>
|
||||
</div>
|
||||
<div x-show="targetType === 'database'">
|
||||
<input type="text" name="expiry" placeholder="never" :disabled="targetType !== 'database'" class="input text-sm">
|
||||
<p class="text-xs text-gray-500 mt-1">Archive expiry: "never" (default) keeps rows forever, or a duration like "720h" prunes older rows.</p>
|
||||
</div>
|
||||
<button type="submit" class="btn-primary text-sm">Add Target</button>
|
||||
</form>
|
||||
</div>
|
||||
|
||||
Reference in New Issue
Block a user