76d123edf3
LastSent scanned MAX(sent_at) over an empty ack_sends into a bare int64, so the ordinary "nothing sent yet" case came back as a scan error rather than the zero time its doc promises. ack_sends is written only by MarkSent, and MarkSent runs only after a repeat has gone out, so every rule's FIRST repeat read an empty table — and RepeatUnacked aborts its whole sweep on that error. The repeat-til-ack loop could never take its first step. Scans into a NullInt64, the same way OldestPendingTelegram already does two files over. EnqueueDigestEntry deduped against any row still marked pending, including one already past its expires_ts. The tick enqueues before it sweeps, so a suppressed nudge arriving on the tick after an expiry was told deduped=true against an entry PendingDigestEntries will never hand back: the caller drops the phrasing it just paid the LLM for and nothing reaches the bundle. The dedupe now carries the same expiry test the read side does. Its lookup also stops treating a real read failure as "nothing there".
246 lines
7.6 KiB
Go
246 lines
7.6 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// TestDigestEntryRoundTrips — a suppressed care candidate lands durably and
|
|
// comes back out of PendingDigestEntries with its severity and body intact.
|
|
func TestDigestEntryRoundTrips(t *testing.T) {
|
|
s := newTestStore(t)
|
|
ctx := context.Background()
|
|
now := time.Now()
|
|
|
|
id, deduped, err := s.EnqueueDigestEntry(ctx, "break", 2, "ты долго не отдыхала", now, now.Add(24*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
if deduped {
|
|
t.Fatal("first enqueue must not report deduped")
|
|
}
|
|
if id == 0 {
|
|
t.Fatal("want a nonzero id")
|
|
}
|
|
|
|
entries, err := s.PendingDigestEntries(ctx, now)
|
|
if err != nil {
|
|
t.Fatalf("pending: %v", err)
|
|
}
|
|
if len(entries) != 1 || entries[0].ID != id {
|
|
t.Fatalf("want 1 pending entry with id %d, got %+v", id, entries)
|
|
}
|
|
if entries[0].Rule != "break" || entries[0].Severity != 2 || entries[0].Body != "ты долго не отдыхала" {
|
|
t.Fatalf("entry contents wrong: %+v", entries[0])
|
|
}
|
|
}
|
|
|
|
// TestDigestEntrySurvivesRestart — durability is the whole point: a fresh
|
|
// Store handle on the same file must see the same pending entry, exactly
|
|
// like the delivery outbox's crash-recovery promise.
|
|
func TestDigestEntrySurvivesRestart(t *testing.T) {
|
|
dir := t.TempDir()
|
|
ctx := context.Background()
|
|
now := time.Now()
|
|
|
|
s1, err := Open(ctx, dir+"/m.db")
|
|
if err != nil {
|
|
t.Fatalf("open: %v", err)
|
|
}
|
|
id, _, err := s1.EnqueueDigestEntry(ctx, "break", 2, "перерыв", now, now.Add(24*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
if err := s1.Close(); err != nil {
|
|
t.Fatalf("close: %v", err)
|
|
}
|
|
|
|
// simulated restart: a brand new Store handle on the same file.
|
|
s2, err := Open(ctx, dir+"/m.db")
|
|
if err != nil {
|
|
t.Fatalf("reopen: %v", err)
|
|
}
|
|
defer func() { _ = s2.Close() }()
|
|
|
|
entries, err := s2.PendingDigestEntries(ctx, now)
|
|
if err != nil {
|
|
t.Fatalf("pending after restart: %v", err)
|
|
}
|
|
if len(entries) != 1 || entries[0].ID != id {
|
|
t.Fatalf("digest entry did not survive restart: %+v", entries)
|
|
}
|
|
}
|
|
|
|
// TestDigestEntryDedupesSameRuleAndBody — the same suppressed nudge
|
|
// repeating across ticks (quiet hours holding for hours) must not pile up
|
|
// into several copies of itself; he hears it once.
|
|
func TestDigestEntryDedupesSameRuleAndBody(t *testing.T) {
|
|
s := newTestStore(t)
|
|
ctx := context.Background()
|
|
now := time.Now()
|
|
|
|
id1, deduped1, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв нужен", now, now.Add(24*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("first enqueue: %v", err)
|
|
}
|
|
if deduped1 {
|
|
t.Fatal("first enqueue should not be deduped")
|
|
}
|
|
|
|
for i := 0; i < 2; i++ {
|
|
id2, deduped2, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв нужен", now.Add(time.Minute), now.Add(25*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("repeat enqueue: %v", err)
|
|
}
|
|
if !deduped2 {
|
|
t.Fatal("repeat enqueue of the same rule+body should report deduped")
|
|
}
|
|
if id2 != id1 {
|
|
t.Fatalf("deduped enqueue should return the original id: want %d got %d", id1, id2)
|
|
}
|
|
}
|
|
|
|
entries, err := s.PendingDigestEntries(ctx, now)
|
|
if err != nil {
|
|
t.Fatalf("pending: %v", err)
|
|
}
|
|
if len(entries) != 1 {
|
|
t.Fatalf("want exactly 1 pending entry after 3 enqueues of the same nudge, got %d", len(entries))
|
|
}
|
|
}
|
|
|
|
// TestDigestEntryExpiresRatherThanDeliversLate — a stale entry (past its
|
|
// expires_ts) must not surface in PendingDigestEntries, and the sweep should
|
|
// mark it expired instead of leaving it around to be delivered late.
|
|
func TestDigestEntryExpiresRatherThanDeliversLate(t *testing.T) {
|
|
s := newTestStore(t)
|
|
ctx := context.Background()
|
|
created := time.Now()
|
|
expiresAt := created.Add(time.Hour)
|
|
|
|
id, _, err := s.EnqueueDigestEntry(ctx, "water", 1, "стакан воды", created, expiresAt)
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
afterExpiry := expiresAt.Add(time.Minute)
|
|
|
|
// even before the sweep runs, a stale entry must not be handed back as
|
|
// pending — "not yet swept" must not mean "still deliverable".
|
|
entries, err := s.PendingDigestEntries(ctx, afterExpiry)
|
|
if err != nil {
|
|
t.Fatalf("pending: %v", err)
|
|
}
|
|
if len(entries) != 0 {
|
|
t.Fatalf("stale entry must not be returned as pending, got %+v", entries)
|
|
}
|
|
|
|
n, err := s.ExpireStaleDigestEntries(ctx, afterExpiry)
|
|
if err != nil {
|
|
t.Fatalf("expire sweep: %v", err)
|
|
}
|
|
if n != 1 {
|
|
t.Fatalf("want 1 entry expired, got %d", n)
|
|
}
|
|
|
|
var status DigestStatus
|
|
if err := s.db.QueryRowContext(ctx, `SELECT status FROM digest_entries WHERE id = ?`, id).Scan(&status); err != nil {
|
|
t.Fatalf("read back: %v", err)
|
|
}
|
|
if status != DigestExpired {
|
|
t.Fatalf("status: want %q, got %q", DigestExpired, status)
|
|
}
|
|
|
|
// idempotent: a second sweep finds nothing new.
|
|
n2, err := s.ExpireStaleDigestEntries(ctx, afterExpiry.Add(time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("second sweep: %v", err)
|
|
}
|
|
if n2 != 0 {
|
|
t.Fatalf("second sweep should find nothing, got %d", n2)
|
|
}
|
|
}
|
|
|
|
// TestDigestEntryDrainMarksDrainedNotDeleted — draining is bookkeeping, not
|
|
// deletion: the row survives as an audit trail of what she actually said.
|
|
func TestDigestEntryDrainMarksDrainedNotDeleted(t *testing.T) {
|
|
s := newTestStore(t)
|
|
ctx := context.Background()
|
|
now := time.Now()
|
|
|
|
id1, _, err := s.EnqueueDigestEntry(ctx, "break", 2, "перерыв", now, now.Add(24*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("enqueue 1: %v", err)
|
|
}
|
|
id2, _, err := s.EnqueueDigestEntry(ctx, "break2", 2, "другое", now, now.Add(24*time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("enqueue 2: %v", err)
|
|
}
|
|
|
|
if err := s.DrainDigestEntries(ctx, []int64{id1, id2}, now.Add(time.Hour)); err != nil {
|
|
t.Fatalf("drain: %v", err)
|
|
}
|
|
|
|
entries, err := s.PendingDigestEntries(ctx, now.Add(time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("pending: %v", err)
|
|
}
|
|
if len(entries) != 0 {
|
|
t.Fatalf("drained entries must not still be pending, got %+v", entries)
|
|
}
|
|
|
|
for _, id := range []int64{id1, id2} {
|
|
var status DigestStatus
|
|
if err := s.db.QueryRowContext(ctx, `SELECT status FROM digest_entries WHERE id = ?`, id).Scan(&status); err != nil {
|
|
t.Fatalf("read back %d: %v", id, err)
|
|
}
|
|
if status != DigestDrained {
|
|
t.Fatalf("entry %d status: want %q, got %q", id, DigestDrained, status)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestDigestEnqueueDoesNotDedupeAgainstAnExpiredEntry — an entry past its
|
|
// expires_ts is still status='pending' until the sweep runs, and the tick
|
|
// enqueues before it sweeps. Deduping against one reports deduped=true for a
|
|
// row PendingDigestEntries will never hand back, so the suppressed nudge is
|
|
// dropped instead of held: the entry says the thing was recorded when nothing
|
|
// was.
|
|
func TestDigestEnqueueDoesNotDedupeAgainstAnExpiredEntry(t *testing.T) {
|
|
s := newTestStore(t)
|
|
ctx := context.Background()
|
|
created := time.Now()
|
|
expiresAt := created.Add(time.Hour)
|
|
|
|
first, deduped, err := s.EnqueueDigestEntry(ctx, "break", 2, "ты долго не отдыхала", created, expiresAt)
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
if deduped {
|
|
t.Fatal("first enqueue must not report deduped")
|
|
}
|
|
|
|
// A tick after the expiry, with the sweep not yet run: the same suppressed
|
|
// nudge comes round again and must be recorded afresh.
|
|
after := expiresAt.Add(time.Minute)
|
|
second, deduped, err := s.EnqueueDigestEntry(ctx, "break", 2, "ты долго не отдыхала", after, after.Add(time.Hour))
|
|
if err != nil {
|
|
t.Fatalf("re-enqueue: %v", err)
|
|
}
|
|
if deduped {
|
|
t.Fatal("an expired entry must not swallow a fresh one")
|
|
}
|
|
if second == first {
|
|
t.Fatalf("want a new row, got the expired one back: id=%d", second)
|
|
}
|
|
|
|
entries, err := s.PendingDigestEntries(ctx, after)
|
|
if err != nil {
|
|
t.Fatalf("pending: %v", err)
|
|
}
|
|
if len(entries) != 1 || entries[0].ID != second {
|
|
t.Fatalf("want the fresh entry %d pending, got %+v", second, entries)
|
|
}
|
|
}
|