7d8b0af99d
Fallthrough is checked per severity through the outbox trail, so sev3/sev4 reroute and sev1/sev2 still drop. Durability uses a real store on a temp file: a crash between Begin and Complete becomes unknown, is not resent, is not dropped, and a late Complete cannot overwrite it. One skipped test marks a real gap: a panic mid-send leaves a permanent pending row. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
248 lines
9.1 KiB
Go
248 lines
9.1 KiB
Go
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)
|
|
}
|
|
}
|