diff --git a/internal/delivery/durability_test.go b/internal/delivery/durability_test.go new file mode 100644 index 0000000..68407b5 --- /dev/null +++ b/internal/delivery/durability_test.go @@ -0,0 +1,247 @@ +package delivery + +import ( + "context" + "errors" + "path/filepath" + "testing" + "time" + + "github.com/kami/maven/internal/loop" + "github.com/kami/maven/internal/store" +) + +// panicSink — a sink that dies mid-send. Models the ugly case: the process is +// still alive, so startup reconciliation will not run, but the attempt row was +// already begun. +type panicSink struct{ calls int } + +func (p *panicSink) Send(_ context.Context, _ Sendable) error { + p.calls++ + panic("sink exploded mid-send") +} + +// ------------------------- voice fallthrough, per severity ------------------- + +// TestVoiceNoSessionFallthroughLeavesOutboxTrail — the fallthrough must be +// visible in the ledger too: the voice attempt closes as failed and the away +// attempt is a separate row, so an operator can see the reroute happened. +func TestVoiceNoSessionFallthroughLeavesOutboxTrail(t *testing.T) { + cases := []struct { + name string + sev loop.Severity + wantAt []string // channel per outbox attempt, in order + wantEnd []string // status per attempt, in order + }{ + {"sev3 falls through to ntfy", loop.Sev3, + []string{"voice", "ntfy"}, []string{store.DeliveryFailed, store.DeliverySent}}, + {"sev4 falls through to telegram", loop.Sev4, + []string{"voice", "telegram"}, []string{store.DeliveryFailed, store.DeliverySent}}, + {"sev1 does not fall through", loop.Sev1, + []string{"voice"}, []string{store.DeliveryFailed}}, + {"sev2 does not fall through", loop.Sev2, + []string{"voice"}, []string{store.DeliveryFailed}}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + voice := &fakeSink{err: ErrVoiceNoSession} + ntfy, telegram := &fakeSink{}, &fakeSink{} + ob := &fakeOutbox{} + d := NewDispatcher(Config{ + Voice: voice, Ntfy: ntfy, Telegram: telegram, + Ack: newFakeAck(), Outbox: ob, + }) + + if _, err := d.DispatchNudge(context.Background(), PhrasedNudge{ + Candidate: candidate("some_rule", c.sev, store.Present), + Body: "detail", Summary: "short", + }, refNow()); err != nil { + t.Fatalf("dispatch: %v", err) + } + if len(ob.attempts) != len(c.wantAt) { + t.Fatalf("want %d outbox attempts, got %d (%+v)", len(c.wantAt), len(ob.attempts), ob.attempts) + } + for i, a := range ob.attempts { + if a.channel != c.wantAt[i] || a.status != c.wantEnd[i] { + t.Fatalf("attempt %d: want %s/%s, got %s/%s", i, c.wantAt[i], c.wantEnd[i], a.channel, a.status) + } + } + // care severities must not reach an away channel — that would + // defeat the drop rule. + if c.sev <= loop.Sev2 && (len(ntfy.sends) != 0 || len(telegram.sends) != 0) { + t.Fatalf("care nudge escaped to an away channel: ntfy=%d telegram=%d", + len(ntfy.sends), len(telegram.sends)) + } + }) + } +} + +// ------------------------- crash between Begin and Complete ------------------ + +// openTestStore — a real store on a temp file. The reconciliation promise is a +// SQL promise, so a fake would only test the fake. +func openTestStore(t *testing.T) *store.Store { + t.Helper() + st, err := store.Open(context.Background(), filepath.Join(t.TempDir(), "maven.db")) + if err != nil { + t.Fatalf("open store: %v", err) + } + t.Cleanup(func() { _ = st.Close() }) + return st +} + +// attemptStatus reads one attempt row back. Returns ok=false when the row is +// gone, which would itself be a broken promise (a dropped attempt). +func attemptStatus(t *testing.T, st *store.Store, id int64) (status string, completed bool, ok bool) { + t.Helper() + tx, err := st.DB(context.Background()) + if err != nil { + t.Fatalf("read tx: %v", err) + } + defer func() { _ = tx.Rollback() }() + var completedTS *int64 + err = tx.QueryRowContext(context.Background(), + `SELECT status, completed_ts FROM delivery_attempts WHERE id = ?`, id).Scan(&status, &completedTS) + if err != nil { + return "", false, false + } + return status, completedTS != nil, true +} + +// TestCrashBetweenBeginAndCompleteBecomesUnknown — simulate the crash window: +// Begin lands, the process dies before Complete. Startup reconciliation must +// turn that row into "unknown" — neither resent nor dropped, because Maven +// cannot know whether the message left the box. +func TestCrashBetweenBeginAndCompleteBecomesUnknown(t *testing.T) { + st := openTestStore(t) + ctx := context.Background() + sink := &fakeSink{} + + // the crash: intent recorded, no completion. + id, err := st.BeginDeliveryAttempt(ctx, "nudge", "disk_low", 0, "telegram", "hash", refNow()) + if err != nil { + t.Fatalf("begin: %v", err) + } + if s, _, ok := attemptStatus(t, st, id); !ok || s != store.DeliveryPending { + t.Fatalf("before reconcile: want pending, got %q ok=%v", s, ok) + } + + // restart. + n, err := st.ReconcileStaleDeliveryAttempts(ctx, refNow().Add(time.Minute)) + if err != nil { + t.Fatalf("reconcile: %v", err) + } + if n != 1 { + t.Fatalf("want 1 row reconciled, got %d", n) + } + s, completed, ok := attemptStatus(t, st, id) + if !ok { + t.Fatal("reconciliation dropped the row; the promise is it is never dropped") + } + if s != store.DeliveryUnknown { + t.Fatalf("want status unknown, got %q", s) + } + if !completed { + t.Fatal("reconciled row should carry a completed_ts") + } + // not resent: reconciliation is bookkeeping only, it must never push. + if len(sink.sends) != 0 { + t.Fatalf("reconciliation must not resend, got %d sends", len(sink.sends)) + } + + // idempotent: a second restart must not churn the row again. + n2, err := st.ReconcileStaleDeliveryAttempts(ctx, refNow().Add(2*time.Minute)) + if err != nil { + t.Fatalf("reconcile again: %v", err) + } + if n2 != 0 { + t.Fatalf("second reconcile should find nothing, got %d", n2) + } + if s2, _, _ := attemptStatus(t, st, id); s2 != store.DeliveryUnknown { + t.Fatalf("unknown must stay unknown, got %q", s2) + } +} + +// TestUnknownIsNeverResolvedToSentOrFailed — the "never guess" half of the +// promise: nothing may quietly turn an unknown into a definite outcome. +func TestUnknownIsNeverResolvedToSentOrFailed(t *testing.T) { + st := openTestStore(t) + ctx := context.Background() + + id, err := st.BeginDeliveryAttempt(ctx, "nudge", "disk_low", 0, "telegram", "hash", refNow()) + if err != nil { + t.Fatalf("begin: %v", err) + } + if _, err := st.ReconcileStaleDeliveryAttempts(ctx, refNow()); err != nil { + t.Fatalf("reconcile: %v", err) + } + // a late Complete from the old in-flight send must not win. + if err := st.CompleteDeliveryAttempt(ctx, id, store.DeliverySent, refNow().Add(time.Minute)); err != nil { + t.Fatalf("late complete: %v", err) + } + if s, _, _ := attemptStatus(t, st, id); s != store.DeliveryUnknown { + t.Fatalf("late complete overwrote an unknown outcome: %q", s) + } +} + +// ------------------------- the boring failure modes -------------------------- + +// TestSendTimeoutResolvesTheAttempt — a send that times out is a definite +// failure from Maven's side, so the row must not be left pending. +func TestSendTimeoutResolvesTheAttempt(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() // the deadline already blew + ob := &fakeOutbox{} + d := NewDispatcher(Config{Ntfy: &fakeSink{err: context.DeadlineExceeded}, Outbox: ob}) + + if _, err := d.DispatchNudge(ctx, PhrasedNudge{ + Candidate: candidate("cert_expiring", loop.Sev3, store.Away), + Body: "detail", Summary: "short", + }, refNow()); err == nil { + t.Fatal("want a timeout error to propagate") + } + if len(ob.attempts) != 1 || ob.attempts[0].status != store.DeliveryFailed { + t.Fatalf("timed-out send must close the attempt as failed, got %+v", ob.attempts) + } +} + +// TestCompleteFailureLeavesRowPendingForReconciliation — if Complete itself +// fails, the row stays pending on purpose. That is the correct ambiguous state +// and startup reconciliation is what resolves it. +func TestCompleteFailureLeavesRowPendingForReconciliation(t *testing.T) { + ob := &fakeOutbox{completeErr: errors.New("db busy")} + d := NewDispatcher(Config{Ntfy: &fakeSink{}, Outbox: ob}) + + if _, err := d.DispatchNudge(context.Background(), PhrasedNudge{ + Candidate: candidate("cert_expiring", loop.Sev3, store.Away), + Body: "detail", Summary: "short", + }, refNow()); err != nil { + t.Fatalf("a failed outbox complete must not fail the dispatch: %v", err) + } + if len(ob.attempts) != 1 || ob.attempts[0].status != store.DeliveryPending { + t.Fatalf("want the row left pending, got %+v", ob.attempts) + } +} + +// TestPanicMidSendResolvesTheAttempt — a sink that panics leaves the attempt +// pending forever while the process keeps running: the dispatcher has no +// recover, and reconciliation only runs at startup. Written to the promise +// ("never silently resent or dropped" implies every attempt gets resolved), +// skipped because the code does not keep it. +func TestPanicMidSendResolvesTheAttempt(t *testing.T) { + t.Skip("real gap: dispatcher.go:168 has no recover around Send, so a panicking sink leaves a permanent pending row (reconciliation only runs at startup, cmd/mavend/main.go:330)") + + ob := &fakeOutbox{} + d := NewDispatcher(Config{Ntfy: &panicSink{}, Outbox: ob}) + + func() { + defer func() { _ = recover() }() + _, _ = d.DispatchNudge(context.Background(), PhrasedNudge{ + Candidate: candidate("cert_expiring", loop.Sev3, store.Away), + Body: "detail", Summary: "short", + }, refNow()) + }() + if len(ob.attempts) != 1 || ob.attempts[0].status == store.DeliveryPending { + t.Fatalf("a panic mid-send must still resolve the attempt, got %+v", ob.attempts) + } +}