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".
166 lines
6.3 KiB
Go
166 lines
6.3 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"database/sql"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
)
|
|
|
|
// DigestStatus — the lifecycle state of a digest entry. Go has no enum type;
|
|
// the idiom is a defined type plus constants, which is what this is. The point
|
|
// is not ceremony: with bare strings nothing stopped a rule name or a body
|
|
// hash being passed where a status belongs, and every one of these values
|
|
// reaches SQL. A defined type makes that a compile error.
|
|
type DigestStatus string
|
|
|
|
// pending = enqueued, waiting for a drain. drained = spoken as part of a
|
|
// bundle. expired = the tick loop's expiry sweep found it past its expires_ts
|
|
// before a drain happened — dropped, not delivered late.
|
|
const (
|
|
DigestPending DigestStatus = "pending"
|
|
DigestDrained DigestStatus = "drained"
|
|
DigestExpired DigestStatus = "expired"
|
|
)
|
|
|
|
// DigestEntry — one gate-suppressed care candidate durably held for later
|
|
// bundled delivery.
|
|
type DigestEntry struct {
|
|
ID int64
|
|
Rule string
|
|
Severity int
|
|
Body string
|
|
CreatedTs time.Time
|
|
ExpiresTs time.Time
|
|
}
|
|
|
|
// DigestBodyHash is the dedupe key for a digest entry: same rule, same
|
|
// wording ⇒ the same suppressed nudge repeating across ticks, and he should
|
|
// hear it once, not once per tick it kept getting suppressed.
|
|
func DigestBodyHash(rule, body string) string {
|
|
sum := sha256.Sum256([]byte(rule + "\x00" + body))
|
|
return hex.EncodeToString(sum[:8])
|
|
}
|
|
|
|
// EnqueueDigestEntry durably records a suppressed care candidate worth
|
|
// resurfacing later. If a LIVE pending entry with the same rule+body already
|
|
// exists, this is a no-op that returns the existing id and deduped=true —
|
|
// the same suppressed nudge repeating across ticks must not pile up into
|
|
// several copies of itself in the eventual bundle.
|
|
//
|
|
// "Live" carries the same expiry test PendingDigestEntries reads with, and for
|
|
// the same reason: a row past its expires_ts is still status='pending' until
|
|
// the sweep gets to it, and the tick enqueues before it sweeps. Deduping
|
|
// against one meant reporting deduped=true against an entry that will never be
|
|
// spoken — the caller drops the phrasing it just paid the LLM for and nothing
|
|
// reaches the bundle. Not yet swept must not mean still deliverable on the
|
|
// write side either.
|
|
func (s *Store) EnqueueDigestEntry(ctx context.Context, rule string, severity int, body string, now, expiresAt time.Time) (id int64, deduped bool, err error) {
|
|
hash := DigestBodyHash(rule, body)
|
|
var existing int64
|
|
err = s.db.QueryRowContext(ctx,
|
|
`SELECT id FROM digest_entries
|
|
WHERE status = ? AND rule = ? AND body_hash = ? AND expires_ts > ? LIMIT 1`,
|
|
DigestPending, rule, hash, now.UnixMilli()).Scan(&existing)
|
|
if err == nil {
|
|
return existing, true, nil
|
|
}
|
|
if !errors.Is(err, sql.ErrNoRows) {
|
|
// A real read failure is not "nothing there". Inserting anyway would
|
|
// duplicate an entry whose existence we never established.
|
|
return 0, false, fmt.Errorf("enqueue digest entry: dedupe lookup: %w", err)
|
|
}
|
|
|
|
res, err := s.db.ExecContext(ctx,
|
|
`INSERT INTO digest_entries (rule, severity, body, body_hash, status, created_ts, expires_ts)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?)`,
|
|
rule, severity, body, hash, DigestPending, now.UnixMilli(), expiresAt.UnixMilli())
|
|
if err != nil {
|
|
return 0, false, fmt.Errorf("enqueue digest entry: %w", err)
|
|
}
|
|
id, err = res.LastInsertId()
|
|
if err != nil {
|
|
return 0, false, fmt.Errorf("enqueue digest entry: last insert id: %w", err)
|
|
}
|
|
return id, false, nil
|
|
}
|
|
|
|
// PendingDigestEntries returns the live (not yet expired) pending entries,
|
|
// oldest first — the order they were suppressed in, which is also the order
|
|
// a bundled readout should mention them.
|
|
func (s *Store) PendingDigestEntries(ctx context.Context, now time.Time) ([]DigestEntry, error) {
|
|
rows, err := s.db.QueryContext(ctx,
|
|
`SELECT id, rule, severity, body, created_ts, expires_ts
|
|
FROM digest_entries WHERE status = ? AND expires_ts > ? ORDER BY created_ts ASC`,
|
|
DigestPending, now.UnixMilli())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("pending digest entries: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []DigestEntry
|
|
for rows.Next() {
|
|
var e DigestEntry
|
|
var created, expires int64
|
|
if err := rows.Scan(&e.ID, &e.Rule, &e.Severity, &e.Body, &created, &expires); err != nil {
|
|
return nil, fmt.Errorf("pending digest entries: scan: %w", err)
|
|
}
|
|
e.CreatedTs = time.UnixMilli(created)
|
|
e.ExpiresTs = time.UnixMilli(expires)
|
|
out = append(out, e)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ExpireStaleDigestEntries marks pending entries whose expires_ts has passed
|
|
// as expired — stale information (yesterday's battery warning) is noise, not
|
|
// news, so it is dropped rather than delivered late. Called once per tick,
|
|
// mirroring ReconcileStaleDeliveryAttempts's "sweep, don't guess" shape.
|
|
// Returns the count expired, for logging.
|
|
func (s *Store) ExpireStaleDigestEntries(ctx context.Context, now time.Time) (int, error) {
|
|
res, err := s.db.ExecContext(ctx,
|
|
`UPDATE digest_entries SET status = ? WHERE status = ? AND expires_ts <= ?`,
|
|
DigestExpired, DigestPending, now.UnixMilli())
|
|
if err != nil {
|
|
return 0, fmt.Errorf("expire stale digest entries: %w", err)
|
|
}
|
|
n, err := res.RowsAffected()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("expire stale digest entries: rows affected: %w", err)
|
|
}
|
|
return int(n), nil
|
|
}
|
|
|
|
// DrainDigestEntries marks the given entries drained — they were folded into
|
|
// a bundle that was successfully dispatched. Called only after a successful
|
|
// send, same rule as the delivery outbox: a failed dispatch must not mark
|
|
// entries drained, or the bundle is lost along with the failed send.
|
|
func (s *Store) DrainDigestEntries(ctx context.Context, ids []int64, now time.Time) error {
|
|
if len(ids) == 0 {
|
|
return nil
|
|
}
|
|
tx, err := s.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return fmt.Errorf("drain digest entries: begin: %w", err)
|
|
}
|
|
defer func() { _ = tx.Rollback() }()
|
|
stmt, err := tx.PrepareContext(ctx,
|
|
`UPDATE digest_entries SET status = ? WHERE id = ? AND status = ?`)
|
|
if err != nil {
|
|
return fmt.Errorf("drain digest entries: prepare: %w", err)
|
|
}
|
|
defer stmt.Close()
|
|
for _, id := range ids {
|
|
if _, err := stmt.ExecContext(ctx, DigestDrained, id, DigestPending); err != nil {
|
|
return fmt.Errorf("drain digest entry %d: %w", id, err)
|
|
}
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return fmt.Errorf("drain digest entries: commit: %w", err)
|
|
}
|
|
return nil
|
|
}
|