Compare commits
1 Commits
next
...
a5bc29f433
| Author | SHA1 | Date | |
|---|---|---|---|
| a5bc29f433 |
25
README.md
25
README.md
@@ -1234,16 +1234,6 @@ webhooker solves this by acting as a durable intermediary:
|
|||||||
backoff. Every delivery attempt is logged with status codes, response
|
backoff. Every delivery attempt is logged with status codes, response
|
||||||
bodies, and timing.
|
bodies, and timing.
|
||||||
|
|
||||||
**That guarantee is at-least-once, not exactly-once.** When a send
|
|
||||||
reaches its target but the write recording that outcome fails, the
|
|
||||||
delivery is deliberately left in a recoverable state rather than
|
|
||||||
marked done — losing a delivery is the worse failure — so the
|
|
||||||
pending sweep picks it up about fifteen minutes later, or the next
|
|
||||||
restart does, and the target receives a payload it already got.
|
|
||||||
webhooker adds no delivery identifier of its own to an outbound
|
|
||||||
request, so **make your receiver idempotent** against whatever the
|
|
||||||
payload itself carries.
|
|
||||||
|
|
||||||
3. **Observability** — Full request/response logging for every webhook
|
3. **Observability** — Full request/response logging for every webhook
|
||||||
received and every delivery attempted. Prometheus metrics expose
|
received and every delivery attempted. Prometheus metrics expose
|
||||||
volume, latency, and error rates. The web UI provides real-time
|
volume, latency, and error rates. The web UI provides real-time
|
||||||
@@ -2626,15 +2616,12 @@ abuse limit later; they are tracked as future work.
|
|||||||
| `POST` | `/source/{id}/edit` | Edit webhook submission |
|
| `POST` | `/source/{id}/edit` | Edit webhook submission |
|
||||||
| `POST` | `/source/{id}/delete` | Delete webhook |
|
| `POST` | `/source/{id}/delete` | Delete webhook |
|
||||||
| `GET` | `/source/{id}/logs` | Webhook event logs |
|
| `GET` | `/source/{id}/logs` | Webhook event logs |
|
||||||
| `GET` | `/source/{id}/logs/{eventID}/body` | Download an event's full stored body. The log page renders each body only up to its cap, so this is the only route that serves a whole one; it is offered wherever a body is shown truncated |
|
|
||||||
| `POST` | `/source/{id}/deliveries/{deliveryID}/replay` | Replay a finished delivery: creates a new delivery for the same event against the target's current configuration (30 per minute per bucket, then `429`) |
|
| `POST` | `/source/{id}/deliveries/{deliveryID}/replay` | Replay a finished delivery: creates a new delivery for the same event against the target's current configuration (30 per minute per bucket, then `429`) |
|
||||||
| `POST` | `/source/{id}/events/{eventID}/resubmit` | Resubmit a stored event: creates a new event copying it and fans that out to every currently active target (30 per minute per bucket, then `429`) |
|
| `POST` | `/source/{id}/events/{eventID}/resubmit` | Resubmit a stored event: creates a new event copying it and fans that out to every currently active target (30 per minute per bucket, then `429`) |
|
||||||
| `POST` | `/source/{id}/entrypoints` | Add entrypoint to webhook |
|
| `POST` | `/source/{id}/entrypoints` | Add entrypoint to webhook |
|
||||||
| `POST` | `/source/{id}/entrypoints/{entrypointID}/delete` | Delete an entrypoint |
|
| `POST` | `/source/{id}/entrypoints/{entrypointID}/delete` | Delete an entrypoint |
|
||||||
| `POST` | `/source/{id}/entrypoints/{entrypointID}/toggle` | Enable or disable an entrypoint |
|
| `POST` | `/source/{id}/entrypoints/{entrypointID}/toggle` | Enable or disable an entrypoint |
|
||||||
| `POST` | `/source/{id}/targets` | Add target to webhook |
|
| `POST` | `/source/{id}/targets` | Add target to webhook |
|
||||||
| `GET` | `/source/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked |
|
|
||||||
| `POST` | `/source/{id}/targets/{targetID}/edit` | Edit target submission |
|
|
||||||
| `POST` | `/source/{id}/targets/{targetID}/delete` | Delete a target |
|
| `POST` | `/source/{id}/targets/{targetID}/delete` | Delete a target |
|
||||||
| `POST` | `/source/{id}/targets/{targetID}/toggle` | Enable or disable a target |
|
| `POST` | `/source/{id}/targets/{targetID}/toggle` | Enable or disable a target |
|
||||||
|
|
||||||
@@ -2673,8 +2660,6 @@ webhooker/
|
|||||||
├── internal/
|
├── internal/
|
||||||
│ ├── banner/
|
│ ├── banner/
|
||||||
│ │ └── banner.go # Ruled block for the one credential shown in the clear
|
│ │ └── banner.go # Ruled block for the one credential shown in the clear
|
||||||
│ ├── ciscript/
|
|
||||||
│ │ └── doc.go # Tests for the CI shell scripts in script/; no runtime code
|
|
||||||
│ ├── resetpw/
|
│ ├── resetpw/
|
||||||
│ │ └── resetpw.go # `webhooker resetpw`: set an account's password, stopped deployments only
|
│ │ └── resetpw.go # `webhooker resetpw`: set an account's password, stopped deployments only
|
||||||
│ ├── config/
|
│ ├── config/
|
||||||
@@ -2744,17 +2729,13 @@ webhooker/
|
|||||||
│ │ ├── ratelimit.go # Per-IP rate limiting middleware (go-chi/httprate)
|
│ │ ├── ratelimit.go # Per-IP rate limiting middleware (go-chi/httprate)
|
||||||
│ │ ├── loginguard.go # Login failure counters and the Argon2id verification semaphore
|
│ │ ├── loginguard.go # Login failure counters and the Argon2id verification semaphore
|
||||||
│ │ └── testing.go # NewForTest: Middleware without the fx lifecycle
|
│ │ └── testing.go # NewForTest: Middleware without the fx lifecycle
|
||||||
│ ├── reqtls/
|
|
||||||
│ │ └── reqtls.go # IsTLS: the one TLS predicate, r.TLS or X-Forwarded-Proto
|
|
||||||
│ ├── server/
|
│ ├── server/
|
||||||
│ │ ├── server.go # Server struct, fx lifecycle, signal handling
|
│ │ ├── server.go # Server struct, fx lifecycle, signal handling
|
||||||
│ │ ├── http.go # HTTP server setup with timeouts
|
│ │ ├── http.go # HTTP server setup with timeouts
|
||||||
│ │ └── routes.go # All route definitions
|
│ │ └── routes.go # All route definitions
|
||||||
│ ├── session/
|
│ └── session/
|
||||||
│ │ ├── session.go # Cookie-based session management
|
│ ├── session.go # Cookie-based session management
|
||||||
│ │ └── testing.go # NewForTest: Session without the fx lifecycle
|
│ └── testing.go # NewForTest: Session without the fx lifecycle
|
||||||
│ └── versionscript/
|
|
||||||
│ └── doc.go # Tests for script/version and the build files that use it
|
|
||||||
├── static/
|
├── static/
|
||||||
│ ├── static.go # //go:embed directive
|
│ ├── static.go # //go:embed directive
|
||||||
│ ├── css/input.css # Tailwind input, source for tailwind.css (make css)
|
│ ├── css/input.css # Tailwind input, source for tailwind.css (make css)
|
||||||
|
|||||||
41
TODO.md
41
TODO.md
@@ -18,27 +18,18 @@ Issue branches do NOT touch this file — the manager maintains it on
|
|||||||
|
|
||||||
# Status
|
# Status
|
||||||
|
|
||||||
The milestone (https://git.eeqj.de/sneak/webhooker/milestone/9) is the
|
1.0.0 is open, with work remaining. The milestone
|
||||||
authoritative list, and the only place to read a count or a state of
|
(https://git.eeqj.de/sneak/webhooker/milestone/9) is the authoritative
|
||||||
play from. This file records where the project is, not what is in
|
list, and the only place to read a count or a state of play from. This
|
||||||
flight: a sentence whose truth depends on a branch being unmerged is
|
file records where the project is, not what is in flight: a sentence
|
||||||
wrong the moment it merges, and this file has been wrong that way
|
whose truth depends on a branch being unmerged is wrong the moment it
|
||||||
before.
|
merges, and this file has been wrong that way before.
|
||||||
|
|
||||||
The durability defect that held the tag has landed
|
The tag is held on a durability defect
|
||||||
(https://git.eeqj.de/sneak/webhooker/issues/256, commit `8d64259`).
|
(https://git.eeqj.de/sneak/webhooker/issues/256): a concurrent reader
|
||||||
Every SQLite handle opens with WAL journaling and a busy timeout, a
|
of a per-webhook event database strands delivered webhooks at
|
||||||
bookkeeping write that fails leaves its delivery in a recoverable
|
`pending`, and the next restart re-delivers them. That issue gates
|
||||||
state rather than a lying one, and recovery skips a delivery that
|
`v1.0.0`, and is where the fix's own state is tracked.
|
||||||
already has a successful result row. Final pre-tag verification
|
|
||||||
exercised it and confirmed it holds. Whatever the milestone still
|
|
||||||
shows open is what remains before `v1.0.0`.
|
|
||||||
|
|
||||||
Delivery is at-least-once by design, not by accident: a send whose
|
|
||||||
result row does not land is attempted again, so a receiver can see a
|
|
||||||
duplicate. That is deliberate — the alternative is a silent lost
|
|
||||||
delivery — and the README says so under Rationale. It is not a defect
|
|
||||||
to re-file.
|
|
||||||
|
|
||||||
One caveat on reading a green check: a docs-only commit deliberately
|
One caveat on reading a green check: a docs-only commit deliberately
|
||||||
replays from the layer cache
|
replays from the layer cache
|
||||||
@@ -48,11 +39,11 @@ commit invalidates the `COPY` layer and genuinely executes.
|
|||||||
|
|
||||||
# Next Step
|
# Next Step
|
||||||
|
|
||||||
Clear the rest of the open 1.0.0 milestone
|
Land https://git.eeqj.de/sneak/webhooker/issues/256, then clear the
|
||||||
(https://git.eeqj.de/sneak/webhooker/milestone/9) and tag `v1.0.0`.
|
rest of the open 1.0.0 milestone and tag `v1.0.0`. Merging `next` into
|
||||||
Merging `next` into `main` is a separate act from tagging and waits on
|
`main` is a separate act from tagging and waits on neither of those:
|
||||||
neither of those: `next` is kept mergeable at all times, which is the
|
`next` is kept mergeable at all times, which is the point of the
|
||||||
point of the branch.
|
branch.
|
||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ package delivery
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -440,7 +439,7 @@ func (e *Engine) processNewTask(
|
|||||||
|
|
||||||
event := buildEventFromTask(task)
|
event := buildEventFromTask(task)
|
||||||
|
|
||||||
event, err = e.hydrateEvent(
|
event, err = e.resolveEventBody(
|
||||||
webhookDB, event, task,
|
webhookDB, event, task,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -506,13 +505,9 @@ func (e *Engine) processRetryTask(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if e.abandonRetryForMissingTarget(webhookDB, d, task) {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
event := buildEventFromTask(task)
|
event := buildEventFromTask(task)
|
||||||
|
|
||||||
event, err = e.hydrateEvent(
|
event, err = e.resolveEventBody(
|
||||||
webhookDB, event, task,
|
webhookDB, event, task,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -534,64 +529,6 @@ func (e *Engine) processRetryTask(
|
|||||||
e.processDelivery(ctx, webhookDB, d, task)
|
e.processDelivery(ctx, webhookDB, d, task)
|
||||||
}
|
}
|
||||||
|
|
||||||
// abandonRetryForMissingTarget stops a retry chain whose target has
|
|
||||||
// been deleted, and reports whether it did.
|
|
||||||
//
|
|
||||||
// A scheduled retry lives in memory as a time.AfterFunc holding the
|
|
||||||
// target's configuration as it was when the chain began, and nothing
|
|
||||||
// else on this path reads the target row. Without this check a
|
|
||||||
// deletion stops nothing: the timer keeps firing and keeps sending to
|
|
||||||
// the destination the operator removed, for the whole remaining
|
|
||||||
// backoff chain. Terminalising in the recovery and sweep paths alone
|
|
||||||
// is not enough, because those only see the delivery once nothing
|
|
||||||
// holds it in memory — which is to say after a restart.
|
|
||||||
//
|
|
||||||
// The worker already owns this delivery, so the terminal write happens
|
|
||||||
// here directly, exactly as a target's own Deliver fails one. Claiming
|
|
||||||
// it again through the recovery gate would only fail against the
|
|
||||||
// reference the worker itself is holding.
|
|
||||||
//
|
|
||||||
// A lookup that fails for any other reason is not a deletion — it is
|
|
||||||
// the main database being unreadable — and the delivery goes ahead as
|
|
||||||
// it did before. A guard that terminally failed deliveries on a
|
|
||||||
// transient fault would be worse than the bug it fixes.
|
|
||||||
func (e *Engine) abandonRetryForMissingTarget(
|
|
||||||
webhookDB *gorm.DB,
|
|
||||||
d *database.Delivery,
|
|
||||||
task *Task,
|
|
||||||
) bool {
|
|
||||||
_, err := e.loadTarget(task.TargetID)
|
|
||||||
if err == nil {
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
if !errors.Is(err, gorm.ErrRecordNotFound) {
|
|
||||||
e.log.Warn(
|
|
||||||
"could not confirm the target of a retrying "+
|
|
||||||
"delivery still exists; attempting anyway",
|
|
||||||
"delivery_id", task.DeliveryID,
|
|
||||||
"target_id", task.TargetID,
|
|
||||||
"error", err,
|
|
||||||
)
|
|
||||||
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
targetType, reason := e.missingTargetReason(task.TargetID)
|
|
||||||
|
|
||||||
e.log.Warn(
|
|
||||||
"abandoning scheduled retry: target is gone",
|
|
||||||
"webhook_id", task.WebhookID,
|
|
||||||
"delivery_id", task.DeliveryID,
|
|
||||||
"target_id", task.TargetID,
|
|
||||||
"target_type", targetType,
|
|
||||||
)
|
|
||||||
|
|
||||||
e.failDelivery(webhookDB, d, targetType, reason)
|
|
||||||
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
|
|
||||||
func (e *Engine) recoverInFlight(ctx context.Context) {
|
func (e *Engine) recoverInFlight(ctx context.Context) {
|
||||||
var webhookIDs []string
|
var webhookIDs []string
|
||||||
|
|
||||||
@@ -696,20 +633,6 @@ func (e *Engine) recoverSingleRetry(
|
|||||||
) {
|
) {
|
||||||
target, err := e.loadTarget(d.TargetID)
|
target, err := e.loadTarget(d.TargetID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// A target that is merely gone is an operator action with a
|
|
||||||
// terminal answer. Any other failure is the main database
|
|
||||||
// refusing to read, which is transient and must leave the
|
|
||||||
// delivery alone: failing every retrying delivery of every
|
|
||||||
// webhook on one bad read would be a far larger fault than
|
|
||||||
// the strand it is meant to clear.
|
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
||||||
e.failMissingTargetRetry(
|
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
e.log.Error(
|
e.log.Error(
|
||||||
"failed to load target for retrying "+
|
"failed to load target for retrying "+
|
||||||
"delivery recovery",
|
"delivery recovery",
|
||||||
@@ -1105,16 +1028,6 @@ func (e *Engine) sweepSingleRetry(
|
|||||||
) {
|
) {
|
||||||
target, err := e.loadTarget(d.TargetID)
|
target, err := e.loadTarget(d.TargetID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Deleted is terminal, unreadable is not; see
|
|
||||||
// recoverSingleRetry.
|
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
||||||
e.failMissingTargetRetry(
|
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
e.log.Error(
|
e.log.Error(
|
||||||
"retry sweep: failed to load target",
|
"retry sweep: failed to load target",
|
||||||
"delivery_id", d.ID,
|
"delivery_id", d.ID,
|
||||||
@@ -1221,113 +1134,6 @@ func (e *Engine) failUnretryableRetry(
|
|||||||
target.Type,
|
target.Type,
|
||||||
)
|
)
|
||||||
|
|
||||||
e.failDelivery(webhookDB, d, target.Type, reason)
|
|
||||||
}
|
|
||||||
|
|
||||||
// failMissingTargetRetry terminally fails an orphaned retrying
|
|
||||||
// delivery whose target row is gone. Both restart recovery and the
|
|
||||||
// periodic sweep call it, so the transition exists once.
|
|
||||||
//
|
|
||||||
// Until it existed both paths logged the failed lookup and returned,
|
|
||||||
// which left the delivery retrying for the life of the database and
|
|
||||||
// the sweep repeating the same error every minute forever. Failing it
|
|
||||||
// with a recorded reason is the treatment the other orphaned-retry
|
|
||||||
// cases already get, so all of them read alike in the event log.
|
|
||||||
//
|
|
||||||
// Logged at warn rather than error: a deleted target is an operator
|
|
||||||
// action, not a system fault.
|
|
||||||
func (e *Engine) failMissingTargetRetry(
|
|
||||||
webhookDB *gorm.DB,
|
|
||||||
webhookID string,
|
|
||||||
d *database.Delivery,
|
|
||||||
) {
|
|
||||||
// Terminal, and reached from the recovery paths, so it takes
|
|
||||||
// ownership like every other write they make.
|
|
||||||
if !e.inflight.retainIdle(d.ID) {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
defer e.inflight.release(d.ID)
|
|
||||||
|
|
||||||
targetType, reason := e.missingTargetReason(d.TargetID)
|
|
||||||
|
|
||||||
e.log.Warn(
|
|
||||||
"failing orphaned retrying delivery: "+
|
|
||||||
"its target no longer exists",
|
|
||||||
"webhook_id", webhookID,
|
|
||||||
"delivery_id", d.ID,
|
|
||||||
"target_id", d.TargetID,
|
|
||||||
"target_type", targetType,
|
|
||||||
)
|
|
||||||
|
|
||||||
e.failDelivery(webhookDB, d, targetType, reason)
|
|
||||||
}
|
|
||||||
|
|
||||||
// missingTargetReason describes a target id that no longer resolves,
|
|
||||||
// and returns the type of the deleted row where there still is one.
|
|
||||||
//
|
|
||||||
// The lookup is Unscoped because deletes are soft: the row survives
|
|
||||||
// with deleted_at set, invisible to loadTarget's default scope.
|
|
||||||
// Reading it is what separates "you deleted this target" from "this id
|
|
||||||
// never named a row" — different things to whoever reads the event
|
|
||||||
// log, and only the first is something an operator did. The widened
|
|
||||||
// scope is deliberately confined to this terminal path: the engine's
|
|
||||||
// normal target loading must go on refusing a deleted target, or
|
|
||||||
// deleting one would stop nothing.
|
|
||||||
//
|
|
||||||
// The type comes back so the caller can label the delivery's status
|
|
||||||
// transition with it. Where the row is gone entirely there is no type
|
|
||||||
// to give, and updateDeliveryStatus leaves the counter alone rather
|
|
||||||
// than opening a series named by the empty string.
|
|
||||||
func (e *Engine) missingTargetReason(
|
|
||||||
targetID string,
|
|
||||||
) (database.TargetType, string) {
|
|
||||||
var target database.Target
|
|
||||||
|
|
||||||
err := e.database.DB().Unscoped().
|
|
||||||
First(&target, "id = ?", targetID).Error
|
|
||||||
if err != nil {
|
|
||||||
return "", fmt.Sprintf(
|
|
||||||
"target %s no longer exists; the delivery "+
|
|
||||||
"cannot be retried and has been failed "+
|
|
||||||
"terminally",
|
|
||||||
targetID,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
return target.Type, fmt.Sprintf(
|
|
||||||
"target %q (type %s) was deleted; the delivery "+
|
|
||||||
"cannot be retried and has been failed terminally",
|
|
||||||
target.Name, target.Type,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// failDelivery records why a delivery is over and then marks it
|
|
||||||
// failed. The caller must already own the delivery: every call site is
|
|
||||||
// either a worker holding the reference runTask took, or a recovery
|
|
||||||
// path that took one through retainIdle.
|
|
||||||
//
|
|
||||||
// The result row is written first and a failure to write it stops the
|
|
||||||
// transition, which is what keeps a delivery from ending failed with
|
|
||||||
// an empty event log — the state that leaves an operator with nothing
|
|
||||||
// but a server log line to work out what happened. A delivery whose
|
|
||||||
// reason could not be recorded stays in the non-terminal state it
|
|
||||||
// already holds, where the sweep will find it again; see
|
|
||||||
// bookkeepingFailed.
|
|
||||||
//
|
|
||||||
// The target type is a parameter rather than read off d because the
|
|
||||||
// orphaned-retry callers deliberately hold a delivery loaded without
|
|
||||||
// its Target relation: populating d.Target would make GORM's
|
|
||||||
// SaveBeforeAssociations upsert the whole target row — plaintext
|
|
||||||
// config, which for a slack target is the credential — into the
|
|
||||||
// per-webhook event database. See
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/206.
|
|
||||||
func (e *Engine) failDelivery(
|
|
||||||
webhookDB *gorm.DB,
|
|
||||||
d *database.Delivery,
|
|
||||||
targetType database.TargetType,
|
|
||||||
reason string,
|
|
||||||
) {
|
|
||||||
err := e.recordResult(
|
err := e.recordResult(
|
||||||
webhookDB,
|
webhookDB,
|
||||||
d,
|
d,
|
||||||
@@ -1344,8 +1150,14 @@ func (e *Engine) failDelivery(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The type is passed rather than assigned onto d: the delivery
|
||||||
|
// is loaded here without its target relation, and populating
|
||||||
|
// d.Target would make GORM's SaveBeforeAssociations upsert the
|
||||||
|
// whole target row — plaintext config, which for a slack target
|
||||||
|
// is the credential — into the per-webhook event database. See
|
||||||
|
// https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||||
e.settleStatus(
|
e.settleStatus(
|
||||||
webhookDB, d, targetType,
|
webhookDB, d, target.Type,
|
||||||
database.DeliveryStatusFailed,
|
database.DeliveryStatusFailed,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -1366,19 +1178,9 @@ func (e *Engine) processDelivery(
|
|||||||
"type", d.Target.Type,
|
"type", d.Target.Type,
|
||||||
)
|
)
|
||||||
|
|
||||||
// The reason is recorded, not just logged. This branch used
|
e.settleStatus(
|
||||||
// to fail the delivery with no DeliveryResult at all, which
|
|
||||||
// showed in the event log as "failed, no attempts recorded
|
|
||||||
// yet" and left one server log line as the only account of
|
|
||||||
// why anywhere.
|
|
||||||
e.failDelivery(
|
|
||||||
webhookDB, d, d.Target.Type,
|
webhookDB, d, d.Target.Type,
|
||||||
fmt.Sprintf(
|
database.DeliveryStatusFailed,
|
||||||
"unknown target type %q: this build has no "+
|
|
||||||
"delivery implementation for it, so no "+
|
|
||||||
"attempt was made",
|
|
||||||
d.Target.Type,
|
|
||||||
),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
return
|
return
|
||||||
@@ -1547,11 +1349,6 @@ func truncate(s string, maxLen int) string {
|
|||||||
|
|
||||||
// --- Helper functions ---
|
// --- Helper functions ---
|
||||||
|
|
||||||
// buildEventFromTask reconstructs the event a Task describes, as far
|
|
||||||
// as the Task itself goes. The fields it cannot fill — the body when
|
|
||||||
// it was too large to inline, and the receipt time, which no Task
|
|
||||||
// carries — come from the stored row in hydrateEvent, which every
|
|
||||||
// caller of this function runs next.
|
|
||||||
func buildEventFromTask(task *Task) database.Event {
|
func buildEventFromTask(task *Task) database.Event {
|
||||||
event := database.Event{
|
event := database.Event{
|
||||||
EntrypointID: task.EntrypointID,
|
EntrypointID: task.EntrypointID,
|
||||||
@@ -1579,67 +1376,29 @@ func buildTargetFromTask(task *Task) database.Target {
|
|||||||
return target
|
return target
|
||||||
}
|
}
|
||||||
|
|
||||||
// hydrateEvent fills in the event fields a Task does not carry, by
|
func (e *Engine) resolveEventBody(
|
||||||
// reading the stored event row.
|
|
||||||
//
|
|
||||||
// CreatedAt is the event's receipt time and lives only in that row.
|
|
||||||
// The Slack target renders it into every message it sends, so an
|
|
||||||
// unhydrated event puts the zero time in front of a human on every
|
|
||||||
// notification the product delivers. See
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/257.
|
|
||||||
//
|
|
||||||
// The body comes from the same row when the Task did not inline it,
|
|
||||||
// which is the case for a body at or above MaxInlineBodySize.
|
|
||||||
//
|
|
||||||
// A read failure is fatal to the delivery only when the body depended
|
|
||||||
// on it. When the Task inlined the body, the delivery has everything
|
|
||||||
// it needs to be sent and goes ahead with the timestamp unset: the row
|
|
||||||
// can be gone under a retention reap while a queued delivery still
|
|
||||||
// holds its body, and dropping a deliverable event to protect one
|
|
||||||
// metadata field would be a worse failure than the one it prevents.
|
|
||||||
func (e *Engine) hydrateEvent(
|
|
||||||
webhookDB *gorm.DB,
|
webhookDB *gorm.DB,
|
||||||
event database.Event,
|
event database.Event,
|
||||||
task *Task,
|
task *Task,
|
||||||
) (database.Event, error) {
|
) (database.Event, error) {
|
||||||
columns := []string{"created_at"}
|
if task.Body != nil {
|
||||||
|
|
||||||
if task.Body == nil {
|
|
||||||
columns = append(columns, "body")
|
|
||||||
}
|
|
||||||
|
|
||||||
var dbEvent database.Event
|
|
||||||
|
|
||||||
err := webhookDB.Select(columns).
|
|
||||||
First(&dbEvent, "id = ?", task.EventID).Error
|
|
||||||
if err != nil {
|
|
||||||
if task.Body == nil {
|
|
||||||
return event, fmt.Errorf(
|
|
||||||
"fetching event body: %w", err,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
e.log.Warn(
|
|
||||||
"could not read the stored event; delivering "+
|
|
||||||
"the inlined body without its receipt time",
|
|
||||||
"event_id", task.EventID,
|
|
||||||
"delivery_id", task.DeliveryID,
|
|
||||||
"error", err,
|
|
||||||
)
|
|
||||||
|
|
||||||
event.Body = *task.Body
|
event.Body = *task.Body
|
||||||
|
|
||||||
return event, nil
|
return event, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
event.CreatedAt = dbEvent.CreatedAt
|
var dbEvent database.Event
|
||||||
|
|
||||||
if task.Body != nil {
|
err := webhookDB.Select("body").
|
||||||
event.Body = *task.Body
|
First(&dbEvent, "id = ?", task.EventID).Error
|
||||||
} else {
|
if err != nil {
|
||||||
event.Body = dbEvent.Body
|
return event, fmt.Errorf(
|
||||||
|
"fetching event body: %w", err,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
event.Body = dbEvent.Body
|
||||||
|
|
||||||
return event, nil
|
return event, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -377,17 +377,6 @@ func TestProcessRetryTask_SuccessfulRetry(t *testing.T) {
|
|||||||
|
|
||||||
bodyStr := event.Body
|
bodyStr := event.Body
|
||||||
cfg := iHTTPConfig(ts.URL)
|
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(
|
task := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"retry-target", cfg, 5, 2, &bodyStr,
|
"retry-target", cfg, 5, 2, &bodyStr,
|
||||||
@@ -467,12 +456,6 @@ func TestProcessRetryTask_LargeBody_FetchFromDB(
|
|||||||
)
|
)
|
||||||
|
|
||||||
cfg := iHTTPConfig(ts.URL)
|
cfg := iHTTPConfig(ts.URL)
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "retry-large",
|
|
||||||
database.TargetTypeHTTP, cfg, 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
task := iTask(
|
task := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"retry-large", cfg, 5, 2, nil,
|
"retry-large", cfg, 5, 2, nil,
|
||||||
@@ -575,12 +558,6 @@ func TestWorkerLifecycle_ProcessesRetryChannel(
|
|||||||
|
|
||||||
bodyStr := event.Body
|
bodyStr := event.Body
|
||||||
cfg := iHTTPConfig(ts.URL)
|
cfg := iHTTPConfig(ts.URL)
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "retry-chan-test",
|
|
||||||
database.TargetTypeHTTP, cfg, 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
task := iTask(
|
task := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"retry-chan-test", cfg, 5, 2, &bodyStr,
|
"retry-chan-test", cfg, 5, 2, &bodyStr,
|
||||||
|
|||||||
@@ -96,15 +96,7 @@ func TestEventDBHoldsNoTargetRows(t *testing.T) {
|
|||||||
)
|
)
|
||||||
assertNoTargetRows(t, dbPath)
|
assertNoTargetRows(t, dbPath)
|
||||||
|
|
||||||
// A retry. Its target exists in the main database, because the
|
// A retry.
|
||||||
// engine confirms a scheduled retry's target has not been
|
|
||||||
// deleted before running it; see
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/107.
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "leaky-target",
|
|
||||||
database.TargetTypeHTTP, cfg, 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
rd := iSeedDelivery(
|
rd := iSeedDelivery(
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
t, s.WebhookDB, event.ID, targetID,
|
||||||
database.DeliveryStatusRetrying,
|
database.DeliveryStatusRetrying,
|
||||||
|
|||||||
@@ -1,442 +0,0 @@
|
|||||||
package delivery_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"encoding/json"
|
|
||||||
"io"
|
|
||||||
"net/http"
|
|
||||||
"net/http/httptest"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/google/uuid"
|
|
||||||
"github.com/stretchr/testify/assert"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
"gorm.io/gorm"
|
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
|
||||||
)
|
|
||||||
|
|
||||||
// tsEventCreatedAt is the receipt time seeded on the events these
|
|
||||||
// tests deliver. It is far enough from both the zero time and from
|
|
||||||
// now that neither can be mistaken for it.
|
|
||||||
func tsEventCreatedAt() time.Time {
|
|
||||||
return time.Date(
|
|
||||||
2026, time.March, 4, 5, 6, 7, 0, time.UTC,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// tsZeroStamp is what a Slack message renders when the event handed
|
|
||||||
// to FormatSlackMessage carries no CreatedAt.
|
|
||||||
const tsZeroStamp = "*Timestamp:* `0001-01-01T00:00:00Z`"
|
|
||||||
|
|
||||||
// tsEventBody is the body seeded on every event in this file. It is
|
|
||||||
// small enough that a Task can inline it.
|
|
||||||
const tsEventBody = `{"hello":"world"}`
|
|
||||||
|
|
||||||
// tsUndeliverableHook stands in for a Slack incoming webhook on the
|
|
||||||
// tests that never send: the config parser requires a URL, but no
|
|
||||||
// request is made.
|
|
||||||
const tsUndeliverableHook = "https://hooks.slack.com/services/T/B/x"
|
|
||||||
|
|
||||||
// tsSink is a stand-in Slack incoming webhook that records the raw
|
|
||||||
// body posted to it.
|
|
||||||
type tsSink struct {
|
|
||||||
*httptest.Server
|
|
||||||
|
|
||||||
bodies chan []byte
|
|
||||||
}
|
|
||||||
|
|
||||||
func newTSSink(t *testing.T) *tsSink {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
s := &tsSink{bodies: make(chan []byte, 8)}
|
|
||||||
|
|
||||||
s.Server = httptest.NewServer(http.HandlerFunc(
|
|
||||||
func(w http.ResponseWriter, r *http.Request) {
|
|
||||||
body, _ := io.ReadAll(r.Body)
|
|
||||||
|
|
||||||
select {
|
|
||||||
case s.bodies <- body:
|
|
||||||
default:
|
|
||||||
}
|
|
||||||
|
|
||||||
w.WriteHeader(http.StatusOK)
|
|
||||||
},
|
|
||||||
))
|
|
||||||
|
|
||||||
t.Cleanup(s.Close)
|
|
||||||
|
|
||||||
return s
|
|
||||||
}
|
|
||||||
|
|
||||||
// text returns the Slack message text from the single payload the
|
|
||||||
// sink received.
|
|
||||||
func (s *tsSink) text(t *testing.T) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case raw := <-s.bodies:
|
|
||||||
t.Logf("raw slack payload: %s", raw)
|
|
||||||
|
|
||||||
var payload struct {
|
|
||||||
Text string `json:"text"`
|
|
||||||
}
|
|
||||||
|
|
||||||
require.NoError(t, json.Unmarshal(raw, &payload))
|
|
||||||
|
|
||||||
return payload.Text
|
|
||||||
case <-time.After(5 * time.Second):
|
|
||||||
t.Fatal("slack sink received no payload")
|
|
||||||
|
|
||||||
return ""
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func tsSlackConfig(t *testing.T, url string) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
data, err := json.Marshal(
|
|
||||||
delivery.SlackTargetConfig{WebhookURL: url},
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
return string(data)
|
|
||||||
}
|
|
||||||
|
|
||||||
// tsSeedEvent writes an event whose CreatedAt is tsEventCreatedAt
|
|
||||||
// rather than the write time, so an assertion on the rendered
|
|
||||||
// timestamp cannot pass by accident against "roughly now".
|
|
||||||
func tsSeedEvent(
|
|
||||||
t *testing.T, db *gorm.DB, webhookID string,
|
|
||||||
) database.Event {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
event := database.Event{
|
|
||||||
WebhookID: webhookID,
|
|
||||||
EntrypointID: uuid.New().String(),
|
|
||||||
Method: http.MethodPost,
|
|
||||||
Headers: `{}`,
|
|
||||||
Body: tsEventBody,
|
|
||||||
ContentType: "application/json",
|
|
||||||
}
|
|
||||||
event.ID = uuid.New().String()
|
|
||||||
event.CreatedAt = tsEventCreatedAt()
|
|
||||||
event.UpdatedAt = tsEventCreatedAt()
|
|
||||||
|
|
||||||
require.NoError(t, db.Create(&event).Error)
|
|
||||||
|
|
||||||
var stored database.Event
|
|
||||||
|
|
||||||
require.NoError(t,
|
|
||||||
db.First(&stored, "id = ?", event.ID).Error,
|
|
||||||
)
|
|
||||||
require.Equal(t,
|
|
||||||
tsEventCreatedAt().UTC(), stored.CreatedAt.UTC(),
|
|
||||||
"seeded created_at did not round-trip",
|
|
||||||
)
|
|
||||||
|
|
||||||
return event
|
|
||||||
}
|
|
||||||
|
|
||||||
// tsSeedTarget writes the slack target row into the main database.
|
|
||||||
// The retry path confirms the target still exists before sending.
|
|
||||||
func tsSeedTarget(
|
|
||||||
t *testing.T, mainDB *gorm.DB, webhookID, config string,
|
|
||||||
) database.Target {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
target := database.Target{
|
|
||||||
WebhookID: webhookID,
|
|
||||||
Name: "slack-sink",
|
|
||||||
Type: database.TargetTypeSlack,
|
|
||||||
Config: config,
|
|
||||||
Active: true,
|
|
||||||
}
|
|
||||||
|
|
||||||
require.NoError(t, mainDB.Create(&target).Error)
|
|
||||||
|
|
||||||
return target
|
|
||||||
}
|
|
||||||
|
|
||||||
func tsTask(
|
|
||||||
d database.Delivery,
|
|
||||||
event database.Event,
|
|
||||||
webhookID string,
|
|
||||||
target database.Target,
|
|
||||||
attemptNum int,
|
|
||||||
body *string,
|
|
||||||
) delivery.Task {
|
|
||||||
return delivery.Task{
|
|
||||||
DeliveryID: d.ID,
|
|
||||||
EventID: event.ID,
|
|
||||||
WebhookID: webhookID,
|
|
||||||
EntrypointID: event.EntrypointID,
|
|
||||||
TargetID: target.ID,
|
|
||||||
TargetName: target.Name,
|
|
||||||
TargetType: database.TargetTypeSlack,
|
|
||||||
TargetConfig: target.Config,
|
|
||||||
MaxRetries: 0,
|
|
||||||
Method: event.Method,
|
|
||||||
Headers: event.Headers,
|
|
||||||
ContentType: event.ContentType,
|
|
||||||
Body: body,
|
|
||||||
AttemptNum: attemptNum,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func tsAssertRealTimestamp(t *testing.T, text string) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
assert.NotContains(t, text, tsZeroStamp,
|
|
||||||
"slack message carries the zero timestamp",
|
|
||||||
)
|
|
||||||
assert.Contains(t, text,
|
|
||||||
"*Timestamp:* `"+
|
|
||||||
tsEventCreatedAt().UTC().Format(time.RFC3339)+"`",
|
|
||||||
"slack message does not carry the event's receipt time",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// tsCase is one end-to-end delivery of a seeded event to a slack
|
|
||||||
// sink, over whichever engine path `process` names.
|
|
||||||
type tsCase struct {
|
|
||||||
// status is the delivery row's status before the engine runs.
|
|
||||||
// The retry path refuses a delivery that is not retrying.
|
|
||||||
status database.DeliveryStatus
|
|
||||||
|
|
||||||
// inlineBody mirrors a Task built for a body under
|
|
||||||
// MaxInlineBodySize. When false the engine reads the body back
|
|
||||||
// from the stored row.
|
|
||||||
inlineBody bool
|
|
||||||
|
|
||||||
attemptNum int
|
|
||||||
|
|
||||||
process func(
|
|
||||||
ctx context.Context, e *delivery.Engine, task *delivery.Task,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// run delivers one event through the named path and returns the
|
|
||||||
// Slack message text the sink received.
|
|
||||||
func (c tsCase) run(t *testing.T) (iSetup, database.Delivery, string) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
sink := newTSSink(t)
|
|
||||||
|
|
||||||
cfg := tsSlackConfig(t, sink.URL)
|
|
||||||
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
|
|
||||||
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, target.ID, c.status,
|
|
||||||
)
|
|
||||||
|
|
||||||
var body *string
|
|
||||||
|
|
||||||
if c.inlineBody {
|
|
||||||
bodyStr := event.Body
|
|
||||||
body = &bodyStr
|
|
||||||
}
|
|
||||||
|
|
||||||
task := tsTask(
|
|
||||||
d, event, s.WebhookID, target, c.attemptNum, body,
|
|
||||||
)
|
|
||||||
|
|
||||||
c.process(context.TODO(), s.Engine, &task)
|
|
||||||
|
|
||||||
return s, d, sink.text(t)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestSlackFirstAttemptCarriesEventTimestamp covers the path an
|
|
||||||
// event takes on its first delivery: the task comes from the
|
|
||||||
// receiver and the engine reconstructs the event from it.
|
|
||||||
func TestSlackFirstAttemptCarriesEventTimestamp(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s, d, text := tsCase{
|
|
||||||
status: database.DeliveryStatusPending,
|
|
||||||
inlineBody: true,
|
|
||||||
attemptNum: 1,
|
|
||||||
process: func(
|
|
||||||
ctx context.Context,
|
|
||||||
e *delivery.Engine,
|
|
||||||
task *delivery.Task,
|
|
||||||
) {
|
|
||||||
e.ExportProcessNewTask(ctx, task)
|
|
||||||
},
|
|
||||||
}.run(t)
|
|
||||||
|
|
||||||
tsAssertRealTimestamp(t, text)
|
|
||||||
|
|
||||||
iAssertStatus(t, s.WebhookDB, d.ID,
|
|
||||||
database.DeliveryStatusDelivered,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestSlackFirstAttemptLargeBodyCarriesEventTimestamp covers the
|
|
||||||
// first-attempt path for an event whose body exceeded
|
|
||||||
// MaxInlineBodySize, so the task carries no body and the engine
|
|
||||||
// reads it back from the stored row.
|
|
||||||
func TestSlackFirstAttemptLargeBodyCarriesEventTimestamp(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
_, _, text := tsCase{
|
|
||||||
status: database.DeliveryStatusPending,
|
|
||||||
inlineBody: false,
|
|
||||||
attemptNum: 1,
|
|
||||||
process: func(
|
|
||||||
ctx context.Context,
|
|
||||||
e *delivery.Engine,
|
|
||||||
task *delivery.Task,
|
|
||||||
) {
|
|
||||||
e.ExportProcessNewTask(ctx, task)
|
|
||||||
},
|
|
||||||
}.run(t)
|
|
||||||
|
|
||||||
tsAssertRealTimestamp(t, text)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestSlackRetryCarriesEventTimestamp covers the retry path, which
|
|
||||||
// reconstructs the event from the same task the first attempt used.
|
|
||||||
func TestSlackRetryCarriesEventTimestamp(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s, d, text := tsCase{
|
|
||||||
status: database.DeliveryStatusRetrying,
|
|
||||||
inlineBody: true,
|
|
||||||
attemptNum: 2,
|
|
||||||
process: func(
|
|
||||||
ctx context.Context,
|
|
||||||
e *delivery.Engine,
|
|
||||||
task *delivery.Task,
|
|
||||||
) {
|
|
||||||
e.ExportProcessRetryTask(ctx, task)
|
|
||||||
},
|
|
||||||
}.run(t)
|
|
||||||
|
|
||||||
tsAssertRealTimestamp(t, text)
|
|
||||||
|
|
||||||
iAssertStatus(t, s.WebhookDB, d.ID,
|
|
||||||
database.DeliveryStatusDelivered,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFormatSlackMessageOverTaskReconstructedEvent asserts on the
|
|
||||||
// formatted message directly, over the event the delivery paths
|
|
||||||
// reconstruct from a Task. It is the unit-level guard under the
|
|
||||||
// end-to-end tests: revert the CreatedAt population in hydrateEvent
|
|
||||||
// and this fails on the zero timestamp.
|
|
||||||
func TestFormatSlackMessageOverTaskReconstructedEvent(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
cfg := tsSlackConfig(t, tsUndeliverableHook)
|
|
||||||
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
|
|
||||||
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, target.ID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
|
|
||||||
bodyStr := event.Body
|
|
||||||
task := tsTask(d, event, s.WebhookID, target, 1, &bodyStr)
|
|
||||||
|
|
||||||
rebuilt, err := s.Engine.ExportEventForTask(
|
|
||||||
s.WebhookDB, &task,
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
assert.False(t, rebuilt.CreatedAt.IsZero(),
|
|
||||||
"reconstructed event carries the zero time",
|
|
||||||
)
|
|
||||||
assert.Equal(t,
|
|
||||||
tsEventCreatedAt().UTC(), rebuilt.CreatedAt.UTC(),
|
|
||||||
)
|
|
||||||
|
|
||||||
tsAssertRealTimestamp(
|
|
||||||
t, delivery.FormatSlackMessage(&rebuilt),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFormatSlackMessageZeroTimestamp asserts the rendering choice
|
|
||||||
// directly, without going through the engine: a zero CreatedAt (the
|
|
||||||
// shape a reaped-row fallback produces) renders as "unknown" rather
|
|
||||||
// than the year-1 zero time, while a real CreatedAt still renders as
|
|
||||||
// RFC3339.
|
|
||||||
func TestFormatSlackMessageZeroTimestamp(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
zeroEvent := database.Event{
|
|
||||||
Method: http.MethodPost,
|
|
||||||
ContentType: testContentType,
|
|
||||||
Body: tsEventBody,
|
|
||||||
}
|
|
||||||
|
|
||||||
zeroText := delivery.FormatSlackMessage(&zeroEvent)
|
|
||||||
|
|
||||||
assert.NotContains(t, zeroText, "0001-01-01",
|
|
||||||
"slack message carries the zero-time year",
|
|
||||||
)
|
|
||||||
assert.Contains(t, zeroText, "*Timestamp:* `unknown`",
|
|
||||||
"slack message does not mark an unset receipt time as unknown",
|
|
||||||
)
|
|
||||||
|
|
||||||
nonZeroEvent := zeroEvent
|
|
||||||
nonZeroEvent.CreatedAt = tsEventCreatedAt()
|
|
||||||
|
|
||||||
nonZeroText := delivery.FormatSlackMessage(&nonZeroEvent)
|
|
||||||
|
|
||||||
assert.Contains(t, nonZeroText,
|
|
||||||
"*Timestamp:* `"+
|
|
||||||
tsEventCreatedAt().UTC().Format(time.RFC3339)+"`",
|
|
||||||
"slack message does not render a real receipt time as RFC3339",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEventReconstructionSurvivesAReapedRow pins the fallback: an
|
|
||||||
// event row reaped by retention while its delivery still holds the
|
|
||||||
// body inline is still delivered, with the receipt time unset,
|
|
||||||
// rather than dropped.
|
|
||||||
func TestEventReconstructionSurvivesAReapedRow(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
cfg := tsSlackConfig(t, tsUndeliverableHook)
|
|
||||||
target := tsSeedTarget(t, s.MainDB, s.WebhookID, cfg)
|
|
||||||
event := tsSeedEvent(t, s.WebhookDB, s.WebhookID)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, target.ID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
|
|
||||||
bodyStr := event.Body
|
|
||||||
task := tsTask(d, event, s.WebhookID, target, 1, &bodyStr)
|
|
||||||
|
|
||||||
require.NoError(t, s.WebhookDB.Unscoped().Delete(
|
|
||||||
&database.Event{}, "id = ?", event.ID,
|
|
||||||
).Error)
|
|
||||||
|
|
||||||
rebuilt, err := s.Engine.ExportEventForTask(
|
|
||||||
s.WebhookDB, &task,
|
|
||||||
)
|
|
||||||
require.NoError(t, err)
|
|
||||||
assert.Equal(t, bodyStr, rebuilt.Body)
|
|
||||||
assert.True(t, rebuilt.CreatedAt.IsZero())
|
|
||||||
|
|
||||||
// A task with no inlined body has nothing left to deliver, so
|
|
||||||
// the same reaped row is an error there.
|
|
||||||
noBody := task
|
|
||||||
noBody.Body = nil
|
|
||||||
|
|
||||||
_, err = s.Engine.ExportEventForTask(s.WebhookDB, &noBody)
|
|
||||||
require.Error(t, err)
|
|
||||||
}
|
|
||||||
@@ -151,16 +151,6 @@ func (e *Engine) ExportProcessRetryTask(
|
|||||||
e.processRetryTask(ctx, task)
|
e.processRetryTask(ctx, task)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ExportEventForTask exposes the event reconstruction the delivery
|
|
||||||
// paths run: buildEventFromTask followed by hydrateEvent.
|
|
||||||
func (e *Engine) ExportEventForTask(
|
|
||||||
webhookDB *gorm.DB, task *Task,
|
|
||||||
) (database.Event, error) {
|
|
||||||
return e.hydrateEvent(
|
|
||||||
webhookDB, buildEventFromTask(task), task,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ExportProcessDelivery exposes processDelivery.
|
// ExportProcessDelivery exposes processDelivery.
|
||||||
func (e *Engine) ExportProcessDelivery(
|
func (e *Engine) ExportProcessDelivery(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
|
|||||||
@@ -229,13 +229,6 @@ func mExhaustRetries(t *testing.T, s iSetup) {
|
|||||||
body := event.Body
|
body := event.Body
|
||||||
cfg := iHTTPConfig(ts.URL)
|
cfg := iHTTPConfig(ts.URL)
|
||||||
|
|
||||||
// The retry below is only run if its target still exists; see
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/107.
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "metrics-fail",
|
|
||||||
database.TargetTypeHTTP, cfg, 2,
|
|
||||||
)
|
|
||||||
|
|
||||||
first := iTask(
|
first := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"metrics-fail", cfg, 2, 1, &body,
|
"metrics-fail", cfg, 2, 1, &body,
|
||||||
@@ -296,13 +289,6 @@ func TestDeliveryMetrics_CircuitBreakerGauge(t *testing.T) {
|
|||||||
// rather than the budget is what stops the delivery.
|
// rather than the budget is what stops the delivery.
|
||||||
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
||||||
|
|
||||||
// The retries below are only run if their target still exists;
|
|
||||||
// see https://git.eeqj.de/sneak/webhooker/issues/107.
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "metrics-trip",
|
|
||||||
database.TargetTypeHTTP, cfg, maxRetries,
|
|
||||||
)
|
|
||||||
|
|
||||||
first := iTask(
|
first := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"metrics-trip", cfg, maxRetries, 1, &body,
|
"metrics-trip", cfg, maxRetries, 1, &body,
|
||||||
@@ -367,11 +353,6 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
|
|||||||
cfg := iHTTPConfig(ts.URL)
|
cfg := iHTTPConfig(ts.URL)
|
||||||
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "metrics-blocked",
|
|
||||||
database.TargetTypeHTTP, cfg, maxRetries,
|
|
||||||
)
|
|
||||||
|
|
||||||
first := iTask(
|
first := iTask(
|
||||||
d, event, s.WebhookID, targetID,
|
d, event, s.WebhookID, targetID,
|
||||||
"metrics-blocked", cfg, maxRetries, 1, &body,
|
"metrics-blocked", cfg, maxRetries, 1, &body,
|
||||||
|
|||||||
@@ -231,15 +231,10 @@ func FormatSlackMessage(
|
|||||||
event.ContentType,
|
event.ContentType,
|
||||||
)
|
)
|
||||||
|
|
||||||
timestamp := "unknown"
|
|
||||||
if !event.CreatedAt.IsZero() {
|
|
||||||
timestamp = event.CreatedAt.UTC().Format(time.RFC3339)
|
|
||||||
}
|
|
||||||
|
|
||||||
fmt.Fprintf(
|
fmt.Fprintf(
|
||||||
&b,
|
&b,
|
||||||
"*Timestamp:* `%s`\n",
|
"*Timestamp:* `%s`\n",
|
||||||
timestamp,
|
event.CreatedAt.UTC().Format(time.RFC3339),
|
||||||
)
|
)
|
||||||
|
|
||||||
fmt.Fprintf(
|
fmt.Fprintf(
|
||||||
|
|||||||
@@ -1,531 +0,0 @@
|
|||||||
package delivery_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"net/http"
|
|
||||||
"net/http/httptest"
|
|
||||||
"sync/atomic"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
"github.com/google/uuid"
|
|
||||||
"github.com/stretchr/testify/assert"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
"sneak.berlin/go/webhooker/internal/database"
|
|
||||||
"sneak.berlin/go/webhooker/internal/delivery"
|
|
||||||
)
|
|
||||||
|
|
||||||
// The two terminal-state gaps of
|
|
||||||
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
|
||||||
// with nothing in its event log to say why, and a retrying delivery
|
|
||||||
// whose target was deleted, which used to keep sending and then never
|
|
||||||
// terminalise.
|
|
||||||
|
|
||||||
// tUnknownType is a target type no build implements. It stands in for
|
|
||||||
// a target whose type was written by a build that knew a type this one
|
|
||||||
// does not.
|
|
||||||
const tUnknownType = database.TargetType("pubsub")
|
|
||||||
|
|
||||||
// tSeedDeletedTarget creates a target, a retrying delivery against it
|
|
||||||
// with one recorded failed attempt, and then deletes the target the
|
|
||||||
// way the source page does.
|
|
||||||
//
|
|
||||||
// It asserts the delete is soft, because that is the whole reason the
|
|
||||||
// engine could not tell a deleted target from a target id that never
|
|
||||||
// named a row: the surviving row is invisible to a scoped read.
|
|
||||||
func tSeedDeletedTarget(
|
|
||||||
t *testing.T,
|
|
||||||
s iSetup,
|
|
||||||
name, url string,
|
|
||||||
) string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, name,
|
|
||||||
database.TargetTypeHTTP, iHTTPConfig(url), 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"target":"deleted"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
|
|
||||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
|
||||||
|
|
||||||
require.NoError(t, s.MainDB.Delete(
|
|
||||||
&database.Target{}, "id = ?", targetID,
|
|
||||||
).Error)
|
|
||||||
|
|
||||||
var scoped, unscoped int64
|
|
||||||
|
|
||||||
require.NoError(t, s.MainDB.
|
|
||||||
Model(&database.Target{}).
|
|
||||||
Where("id = ?", targetID).
|
|
||||||
Count(&scoped).Error)
|
|
||||||
|
|
||||||
require.NoError(t, s.MainDB.Unscoped().
|
|
||||||
Model(&database.Target{}).
|
|
||||||
Where("id = ?", targetID).
|
|
||||||
Count(&unscoped).Error)
|
|
||||||
|
|
||||||
require.Zero(t, scoped,
|
|
||||||
"the deleted target is still visible to a scoped read",
|
|
||||||
)
|
|
||||||
require.Equal(t, int64(1), unscoped,
|
|
||||||
"the delete was hard, so this test proves nothing about "+
|
|
||||||
"the soft-delete case it exists for",
|
|
||||||
)
|
|
||||||
|
|
||||||
return d.ID
|
|
||||||
}
|
|
||||||
|
|
||||||
// tLastResult returns a delivery's final recorded attempt, asserting
|
|
||||||
// the expected number of them.
|
|
||||||
func tLastResult(
|
|
||||||
t *testing.T,
|
|
||||||
s iSetup,
|
|
||||||
deliveryID string,
|
|
||||||
want int,
|
|
||||||
) database.DeliveryResult {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
results := iResults(t, s.WebhookDB, deliveryID)
|
|
||||||
require.Len(t, results, want)
|
|
||||||
|
|
||||||
return results[want-1]
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- 1. A failure with nothing recorded ---
|
|
||||||
|
|
||||||
func TestProcessDelivery_UnknownTargetType_RecordsWhy(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"unknown":"type"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
seeded := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
|
||||||
database.DeliveryStatusPending,
|
|
||||||
)
|
|
||||||
|
|
||||||
target := database.Target{
|
|
||||||
Name: "mystery",
|
|
||||||
Type: tUnknownType,
|
|
||||||
Config: iHTTPConfig("http://example.com/hook"),
|
|
||||||
}
|
|
||||||
target.ID = targetID
|
|
||||||
|
|
||||||
d := database.Delivery{
|
|
||||||
EventID: event.ID,
|
|
||||||
TargetID: targetID,
|
|
||||||
Status: database.DeliveryStatusPending,
|
|
||||||
Event: event,
|
|
||||||
Target: target,
|
|
||||||
}
|
|
||||||
d.ID = seeded.ID
|
|
||||||
|
|
||||||
body := event.Body
|
|
||||||
task := iTask(
|
|
||||||
seeded, event, s.WebhookID, targetID, "mystery",
|
|
||||||
target.Config, 0, 1, &body,
|
|
||||||
)
|
|
||||||
task.TargetType = tUnknownType
|
|
||||||
|
|
||||||
s.Engine.ExportProcessDelivery(
|
|
||||||
context.Background(), s.WebhookDB, &d, &task,
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, d.ID, database.DeliveryStatusFailed,
|
|
||||||
)
|
|
||||||
|
|
||||||
last := tLastResult(t, s, d.ID, 1)
|
|
||||||
|
|
||||||
assert.False(t, last.Success)
|
|
||||||
assert.Equal(t, 1, last.AttemptNum)
|
|
||||||
assert.Contains(t, last.Error, string(tUnknownType),
|
|
||||||
"the recorded reason does not name the offending type",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- 2. A retrying delivery whose target is gone ---
|
|
||||||
|
|
||||||
func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
iCreateWebhook(
|
|
||||||
t, s.MainDB, s.WebhookID, "deleted-target-recovery",
|
|
||||||
)
|
|
||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
|
||||||
t, s, "gone-on-recovery", "http://example.com/hook",
|
|
||||||
)
|
|
||||||
|
|
||||||
s.Engine.ExportRecoverWebhookDeliveries(
|
|
||||||
context.Background(), s.WebhookID,
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, deliveryID,
|
|
||||||
database.DeliveryStatusFailed,
|
|
||||||
)
|
|
||||||
|
|
||||||
last := tLastResult(t, s, deliveryID, 2)
|
|
||||||
|
|
||||||
assert.False(t, last.Success)
|
|
||||||
assert.Equal(t, 2, last.AttemptNum)
|
|
||||||
assert.Contains(t, last.Error, "gone-on-recovery")
|
|
||||||
assert.Contains(t, last.Error, "was deleted")
|
|
||||||
|
|
||||||
assert.Empty(t, s.Engine.ExportRetryCh(),
|
|
||||||
"a delivery whose target is gone was rescheduled",
|
|
||||||
)
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld(),
|
|
||||||
"the terminal path leaked its ownership reference",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
iCreateWebhook(
|
|
||||||
t, s.MainDB, s.WebhookID, "deleted-target-sweep",
|
|
||||||
)
|
|
||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
|
||||||
t, s, "gone-on-sweep", "http://example.com/hook",
|
|
||||||
)
|
|
||||||
|
|
||||||
// Twice, because the bug was an error the sweep repeated every
|
|
||||||
// minute for the life of the database: the second sweep must
|
|
||||||
// find nothing left to do.
|
|
||||||
s.Engine.ExportSweepWebhookRetries(
|
|
||||||
context.Background(), s.WebhookID,
|
|
||||||
)
|
|
||||||
s.Engine.ExportSweepWebhookRetries(
|
|
||||||
context.Background(), s.WebhookID,
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, deliveryID,
|
|
||||||
database.DeliveryStatusFailed,
|
|
||||||
)
|
|
||||||
|
|
||||||
last := tLastResult(t, s, deliveryID, 2)
|
|
||||||
|
|
||||||
assert.Contains(t, last.Error, "gone-on-sweep")
|
|
||||||
assert.Contains(t, last.Error, "was deleted")
|
|
||||||
|
|
||||||
assert.Empty(t, s.Engine.ExportRetryCh())
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestSweepSingleRetry_TargetNeverExisted covers the other half of the
|
|
||||||
// soft-delete distinction: an id with no row at all, deleted or
|
|
||||||
// otherwise, must not be reported as something the operator deleted.
|
|
||||||
func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
iCreateWebhook(
|
|
||||||
t, s.MainDB, s.WebhookID, "target-never-existed",
|
|
||||||
)
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"target":"absent"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
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,
|
|
||||||
)
|
|
||||||
|
|
||||||
last := tLastResult(t, s, d.ID, 2)
|
|
||||||
|
|
||||||
assert.Contains(t, last.Error, targetID)
|
|
||||||
assert.Contains(t, last.Error, "no longer exists")
|
|
||||||
assert.NotContains(t, last.Error, "was deleted",
|
|
||||||
"an id that never named a row was reported as a deletion",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFailMissingTargetRetry_WritesNoTargetRow holds the new terminal
|
|
||||||
// path to the same rule as the existing one: no target row, and so no
|
|
||||||
// plaintext target config, may be written into the per-webhook event
|
|
||||||
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
|
||||||
func TestFailMissingTargetRetry_WritesNoTargetRow(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
iCreateWebhook(
|
|
||||||
t, s.MainDB, s.WebhookID, "no-target-row-deleted",
|
|
||||||
)
|
|
||||||
|
|
||||||
hookURL := "https://hooks.slack.com/services/T00/B00/x"
|
|
||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
|
||||||
t, s, "credential-bearing", hookURL,
|
|
||||||
)
|
|
||||||
|
|
||||||
s.Engine.ExportSweepWebhookRetries(
|
|
||||||
context.Background(), s.WebhookID,
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, deliveryID,
|
|
||||||
database.DeliveryStatusFailed,
|
|
||||||
)
|
|
||||||
|
|
||||||
var configs []string
|
|
||||||
|
|
||||||
require.NoError(t, s.WebhookDB.
|
|
||||||
Table("targets").
|
|
||||||
Pluck("config", &configs).Error)
|
|
||||||
|
|
||||||
assert.Empty(t, configs,
|
|
||||||
"the deleted-target terminal path wrote a target row "+
|
|
||||||
"into the per-webhook event database",
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- 3. The scheduled retry chain ---
|
|
||||||
|
|
||||||
// tRetryChainSetup wires a counting sink and a retrying delivery
|
|
||||||
// against a live target pointing at it, and returns the task a
|
|
||||||
// scheduled retry would carry — config and all, snapshotted as
|
|
||||||
// ScheduleRetry snapshots it.
|
|
||||||
func tRetryChainSetup(
|
|
||||||
t *testing.T,
|
|
||||||
s iSetup,
|
|
||||||
name string,
|
|
||||||
hits *atomic.Int64,
|
|
||||||
) (delivery.Task, string) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ts := httptest.NewServer(http.HandlerFunc(
|
|
||||||
func(w http.ResponseWriter, _ *http.Request) {
|
|
||||||
hits.Add(1)
|
|
||||||
w.WriteHeader(http.StatusOK)
|
|
||||||
},
|
|
||||||
))
|
|
||||||
t.Cleanup(ts.Close)
|
|
||||||
|
|
||||||
iCreateWebhook(t, s.MainDB, s.WebhookID, name)
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
cfg := iHTTPConfig(ts.URL)
|
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, name,
|
|
||||||
database.TargetTypeHTTP, cfg, 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"chain":"retry"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
|
|
||||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
|
||||||
|
|
||||||
body := event.Body
|
|
||||||
|
|
||||||
return iTask(
|
|
||||||
d, event, s.WebhookID, targetID, name, cfg, 5, 2, &body,
|
|
||||||
), targetID
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestProcessRetryTask_TargetDeleted_MakesNoAttempt is the half the
|
|
||||||
// deployability audit found worse than filed: terminalising on
|
|
||||||
// recovery and sweep alone leaves the already-scheduled timer chain
|
|
||||||
// running, and it holds the target's configuration from before the
|
|
||||||
// deletion, so it goes on sending to a destination that was removed.
|
|
||||||
func TestProcessRetryTask_TargetDeleted_MakesNoAttempt(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
var hits atomic.Int64
|
|
||||||
|
|
||||||
task, targetID := tRetryChainSetup(
|
|
||||||
t, s, "gone-mid-chain", &hits,
|
|
||||||
)
|
|
||||||
|
|
||||||
require.NoError(t, s.MainDB.Delete(
|
|
||||||
&database.Target{}, "id = ?", targetID,
|
|
||||||
).Error)
|
|
||||||
|
|
||||||
s.Engine.ExportProcessRetryTask(
|
|
||||||
context.Background(), &task,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Zero(t, hits.Load(),
|
|
||||||
"a scheduled retry fired at a target the operator "+
|
|
||||||
"had already deleted",
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, task.DeliveryID,
|
|
||||||
database.DeliveryStatusFailed,
|
|
||||||
)
|
|
||||||
|
|
||||||
last := tLastResult(t, s, task.DeliveryID, 2)
|
|
||||||
|
|
||||||
assert.False(t, last.Success)
|
|
||||||
assert.Contains(t, last.Error, "was deleted")
|
|
||||||
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestProcessRetryTask_TargetPresent_StillDelivers is the guard's
|
|
||||||
// mutation check: a liveness check that refused every retry would pass
|
|
||||||
// the test above and break every retry there is.
|
|
||||||
func TestProcessRetryTask_TargetPresent_StillDelivers(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
var hits atomic.Int64
|
|
||||||
|
|
||||||
task, _ := tRetryChainSetup(t, s, "still-there", &hits)
|
|
||||||
|
|
||||||
s.Engine.ExportProcessRetryTask(
|
|
||||||
context.Background(), &task,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Equal(t, int64(1), hits.Load())
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, task.DeliveryID,
|
|
||||||
database.DeliveryStatusDelivered,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestProcessRetryTask_TargetUnreadable_StillDelivers pins the other
|
|
||||||
// half of the guard: only a target that is confirmed gone stops a
|
|
||||||
// retry. A main database that cannot be read is a transient fault, and
|
|
||||||
// a guard that abandoned deliveries on one would be a worse bug than
|
|
||||||
// the one it fixes.
|
|
||||||
func TestProcessRetryTask_TargetUnreadable_StillDelivers(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
var hits atomic.Int64
|
|
||||||
|
|
||||||
task, _ := tRetryChainSetup(t, s, "unreadable-main", &hits)
|
|
||||||
|
|
||||||
sqlDB, err := s.MainDB.DB()
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, sqlDB.Close())
|
|
||||||
|
|
||||||
s.Engine.ExportProcessRetryTask(
|
|
||||||
context.Background(), &task,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Equal(t, int64(1), hits.Load(),
|
|
||||||
"a retry was abandoned because the main database "+
|
|
||||||
"could not be read, not because its target was gone",
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, task.DeliveryID,
|
|
||||||
database.DeliveryStatusDelivered,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone is the
|
|
||||||
// same rule on the recovery path. A read failure that is not
|
|
||||||
// "record not found" must leave every retrying delivery of every
|
|
||||||
// webhook exactly as it was.
|
|
||||||
func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
|
|
||||||
t *testing.T,
|
|
||||||
) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
s := newISetup(t)
|
|
||||||
|
|
||||||
iCreateWebhook(
|
|
||||||
t, s.MainDB, s.WebhookID, "unreadable-on-recovery",
|
|
||||||
)
|
|
||||||
|
|
||||||
targetID := uuid.New().String()
|
|
||||||
|
|
||||||
iCreateTarget(
|
|
||||||
t, s.MainDB, targetID, s.WebhookID, "healthy",
|
|
||||||
database.TargetTypeHTTP,
|
|
||||||
iHTTPConfig("http://example.com/hook"), 5,
|
|
||||||
)
|
|
||||||
|
|
||||||
event := iSeedEvent(
|
|
||||||
t, s.WebhookDB, s.WebhookID, `{"still":"retrying"}`,
|
|
||||||
)
|
|
||||||
|
|
||||||
d := iSeedDelivery(
|
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
|
|
||||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
|
||||||
|
|
||||||
sqlDB, err := s.MainDB.DB()
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, sqlDB.Close())
|
|
||||||
|
|
||||||
s.Engine.ExportRecoverRetryingDeliveries(
|
|
||||||
s.WebhookDB, s.WebhookID,
|
|
||||||
)
|
|
||||||
|
|
||||||
iAssertStatus(
|
|
||||||
t, s.WebhookDB, d.ID,
|
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Len(t, iResults(t, s.WebhookDB, d.ID), 1,
|
|
||||||
"an unreadable main database produced a terminal "+
|
|
||||||
"failure row",
|
|
||||||
)
|
|
||||||
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user