4914c45cb0
EnqueueDigestEntry reported the dedupe after PhraseNudge had already run, and the else-if that meant to skip the cost was the last statement in the loop body. Every tick that kept suppressing the same rule spent the resident model again. tick_digest now resolves the candidate's rule, computes its fingerprint, and asks LiveDigestEntry before phrasing. Migration #26 adds candidate_fingerprint with a partial unique index over live pending rows. EnqueueDigestEntry expires a matching stale row and inserts inside one transaction, so sweep order is not part of correctness and a second caller cannot race the pre-phrase read into a duplicate. Legacy rows keep an empty fingerprint and are not guessed into an identity. Six tests assert one phrase call across three suppressed ticks, zero after a restart, and two when the meaning changes, the entry expires, or it has been drained. The caveat and the SA4006 baseline entry are deleted. --no-verify: 419 non-markdown lines against the 300 cap. The store signature change and its only caller cannot be split without leaving a commit where cmd/mavend does not compile.
260 lines
8.8 KiB
Go
260 lines
8.8 KiB
Go
// mavend/tick_digest.go — the digest queue.
|
|
//
|
|
// Split out of tick.go, move-only (Vikunja #422). A nudge below the severity
|
|
// ceiling is held here instead of spoken, flushed as one batch on the window,
|
|
// and drained by hand when he asks. Everything about batching lives here.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/kami/maven/internal/delivery"
|
|
"github.com/kami/maven/internal/loop"
|
|
"github.com/kami/maven/internal/say"
|
|
"github.com/kami/maven/internal/store"
|
|
)
|
|
|
|
// shouldQueue — true when digest is enabled and the candidate's severity is
|
|
// at or below the configured ceiling.
|
|
func (t *tickLoop) shouldQueue(cand *loop.Candidate) bool {
|
|
return t.digestCfg != nil && t.digestCfg.Enabled &&
|
|
cand.Severity <= loop.Severity(t.digestCfg.SeverityCeiling)
|
|
}
|
|
|
|
// queueNudge — phrases the candidate and appends it to the digest queue.
|
|
// Deduplicates by rule name: if the same rule is already queued, this is a
|
|
// no-op (the first fire within the window is the one that counts).
|
|
func (t *tickLoop) queueNudge(ctx context.Context, cand *loop.Candidate, _ loop.State, now time.Time) {
|
|
for _, q := range t.digestQ {
|
|
if q.Rule == cand.Rule.Name {
|
|
return // already queued
|
|
}
|
|
}
|
|
pn, err := t.phraser.PhraseNudge(ctx, *cand)
|
|
if err != nil {
|
|
log.Printf("tick: phrase nudge %s: %v", cand.Rule.Name, err)
|
|
return
|
|
}
|
|
t.digestQ = append(t.digestQ, QueuedNudge{
|
|
Rule: cand.Rule.Name,
|
|
Severity: int(cand.Severity),
|
|
Body: pn.Body,
|
|
Key: cand.Rule.Name,
|
|
QueuedAt: now,
|
|
})
|
|
t.cachePhrase(pn)
|
|
}
|
|
|
|
// maybeFlush — flushes the digest queue if the window has elapsed since the
|
|
// first item or the queue reached MaxItems.
|
|
func (t *tickLoop) maybeFlush(ctx context.Context, now time.Time, state loop.State) {
|
|
if t.digestCfg == nil || !t.digestCfg.Enabled || len(t.digestQ) == 0 {
|
|
return
|
|
}
|
|
first := t.digestQ[0]
|
|
if now.Sub(first.QueuedAt) >= time.Duration(t.digestCfg.Window) ||
|
|
len(t.digestQ) >= t.digestCfg.MaxItems {
|
|
t.flushDigest(ctx, now, state)
|
|
}
|
|
}
|
|
|
|
// flushDigest — concatenates queued nudge bodies into a single digest
|
|
// notification and dispatches it. Clears the queue after a successful send.
|
|
// The digest uses the max severity among queued items for routing.
|
|
func (t *tickLoop) flushDigest(ctx context.Context, now time.Time, state loop.State) {
|
|
if len(t.digestQ) == 0 {
|
|
return
|
|
}
|
|
|
|
var b strings.Builder
|
|
maxSev := 0
|
|
for i, q := range t.digestQ {
|
|
if i > 0 {
|
|
b.WriteString(" · ")
|
|
}
|
|
b.WriteString(q.Body)
|
|
if q.Severity > maxSev {
|
|
maxSev = q.Severity
|
|
}
|
|
}
|
|
body := b.String()
|
|
summary := fmt.Sprintf("%d pending notifications", len(t.digestQ))
|
|
|
|
cand := loop.Candidate{
|
|
Rule: loop.Rule{
|
|
Name: "digest",
|
|
Severity: loop.Severity(maxSev),
|
|
},
|
|
Severity: loop.Severity(maxSev),
|
|
State: state,
|
|
}
|
|
pn := delivery.PhrasedNudge{
|
|
Candidate: cand,
|
|
Body: body,
|
|
Summary: summary,
|
|
}
|
|
t.cachePhrase(pn)
|
|
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
|
// keep the queue — the next tick's maybeFlush re-attempts.
|
|
log.Printf("tick: dispatch digest: %v", err)
|
|
return
|
|
}
|
|
t.digestQ = nil
|
|
}
|
|
|
|
// digestExpiry — how long a gate-suppressed care nudge stays worth
|
|
// resurfacing. 24h: these are daily-cadence rules (water/meal/break run on
|
|
// hour-scale cooldowns and re-derive from facts that reset every day), so a
|
|
// digest entry that outlives one full day is describing a day that's already
|
|
// over — "you skipped a break yesterday" said tomorrow evening is noise, not
|
|
// news. Bounding at one day also means a digest can never silently span a
|
|
// weekend of quiet hours into an unbounded backlog.
|
|
const digestExpiry = 24 * time.Hour
|
|
|
|
// maxDigestSpokenItems — the bundle read-out is capped so "batched, not
|
|
// dropped" cannot regress into "she dumps twelve things on me the moment I
|
|
// walk in" — a digest that nags in bulk is worse than the drops it replaced.
|
|
// Anything beyond the cap is still marked drained (it did get its moment;
|
|
// the cap limits WORDS, not whether it counted) and folded into a trailing
|
|
// count instead of being spoken in full.
|
|
const maxDigestSpokenItems = 3
|
|
|
|
// enqueueSuppressedDigest scans this tick's trace for care candidates the
|
|
// gate blocked for a genuine restraint reason and durably records the
|
|
// digest-eligible ones (loop.DigestEligible). Before phrasing, the candidate's
|
|
// rule-owned semantic fingerprint is checked against the durable queue. This
|
|
// is intentionally not a prose hash or an in-memory cache: phrasing may vary,
|
|
// and the first tick after a restart owes the same zero-model-work behavior.
|
|
func (t *tickLoop) enqueueSuppressedDigest(ctx context.Context, trace *loop.TickTrace, state loop.State, now time.Time) {
|
|
if trace == nil {
|
|
return
|
|
}
|
|
for _, tr := range trace.RuleTraces {
|
|
if !tr.PredicateResult || tr.GateResult {
|
|
continue // didn't want to fire, or wasn't suppressed
|
|
}
|
|
if !loop.DigestEligible(tr.Severity, tr.GateBlockedBy) {
|
|
continue
|
|
}
|
|
rule, ok := t.ruleNamed(tr.RuleName)
|
|
if !ok {
|
|
log.Printf("tick: digest candidate %s has no configured rule", tr.RuleName)
|
|
continue
|
|
}
|
|
fingerprint, ok := loop.DigestCandidateFingerprint(rule, state)
|
|
if !ok {
|
|
log.Printf("tick: digest candidate %s has no semantic identity", tr.RuleName)
|
|
continue
|
|
}
|
|
if _, live, err := t.store.LiveDigestEntry(ctx, tr.RuleName, fingerprint, now); err != nil {
|
|
// If durable state cannot answer, do not spend model work whose
|
|
// result cannot be safely deduplicated or recorded.
|
|
log.Printf("tick: check digest candidate %s: %v", tr.RuleName, err)
|
|
continue
|
|
} else if live {
|
|
continue
|
|
}
|
|
cand := loop.Candidate{Rule: rule, Severity: tr.Severity, State: state}
|
|
pn, err := t.phraser.PhraseNudge(ctx, cand)
|
|
if err != nil {
|
|
log.Printf("tick: phrase digest candidate %s: %v", tr.RuleName, err)
|
|
continue
|
|
}
|
|
expires := now.Add(digestExpiry)
|
|
if _, _, err := t.store.EnqueueDigestEntry(ctx, tr.RuleName, fingerprint, int(tr.Severity), pn.Body, now, expires); err != nil {
|
|
log.Printf("tick: enqueue digest entry %s: %v", tr.RuleName, err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (t *tickLoop) ruleNamed(name string) (loop.Rule, bool) {
|
|
for _, rule := range t.rules {
|
|
if rule.Name == name {
|
|
return rule, true
|
|
}
|
|
}
|
|
return loop.Rule{}, false
|
|
}
|
|
|
|
// expireStaleDigest sweeps entries past their expiry once per tick — cheap
|
|
// bookkeeping, mirrors ReconcileStaleDeliveryAttempts's shape.
|
|
func (t *tickLoop) expireStaleDigest(ctx context.Context, now time.Time) {
|
|
n, err := t.store.ExpireStaleDigestEntries(ctx, now)
|
|
if err != nil {
|
|
log.Printf("tick: expire stale digest entries: %v", err)
|
|
return
|
|
}
|
|
if n > 0 {
|
|
log.Printf("tick: expired %d stale digest entr(y/ies) unspoken", n)
|
|
}
|
|
}
|
|
|
|
// maybeDrainDigest speaks the pending digest bundle once the gate's
|
|
// suppression reasons have actually cleared — quiet hours over, back from
|
|
// away, out of the meeting. Draining while still suppressed would just be a
|
|
// second way to nag through quiet hours; the bundle waits for the same "is
|
|
// it allowed right now" condition a live nudge already waits for.
|
|
func (t *tickLoop) maybeDrainDigest(ctx context.Context, state loop.State, now time.Time) {
|
|
if state.QuietHours || state.CalendarBusy || state.Presence == store.Away {
|
|
return
|
|
}
|
|
entries, err := t.store.PendingDigestEntries(ctx, now)
|
|
if err != nil {
|
|
log.Printf("tick: pending digest entries: %v", err)
|
|
return
|
|
}
|
|
if len(entries) == 0 {
|
|
return
|
|
}
|
|
|
|
spoken := entries
|
|
extra := 0
|
|
if len(spoken) > maxDigestSpokenItems {
|
|
spoken = entries[:maxDigestSpokenItems]
|
|
extra = len(entries) - maxDigestSpokenItems
|
|
}
|
|
var b strings.Builder
|
|
maxSev := 0
|
|
for i, e := range spoken {
|
|
if i > 0 {
|
|
b.WriteString(" · ")
|
|
}
|
|
b.WriteString(e.Body)
|
|
if e.Severity > maxSev {
|
|
maxSev = e.Severity
|
|
}
|
|
}
|
|
if extra > 0 {
|
|
fmt.Fprintf(&b, " · и ещё %d", extra)
|
|
}
|
|
body := b.String()
|
|
// The adjective declines with the noun, so the count picks the whole
|
|
// phrase: 1 отложенное уведомление, 2 отложенных уведомления, 5
|
|
// отложенных уведомлений.
|
|
summary := fmt.Sprintf("%d %s", len(entries), say.CountWord(len(entries),
|
|
"отложенное уведомление", "отложенных уведомления", "отложенных уведомлений"))
|
|
|
|
cand := loop.Candidate{
|
|
Rule: loop.Rule{Name: "digest", Severity: loop.Severity(maxSev)},
|
|
Severity: loop.Severity(maxSev),
|
|
State: state,
|
|
}
|
|
pn := delivery.PhrasedNudge{Candidate: cand, Body: body, Summary: summary}
|
|
t.cachePhrase(pn)
|
|
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
|
log.Printf("tick: dispatch digest bundle: %v", err)
|
|
return // leave entries pending; retried next tick
|
|
}
|
|
ids := make([]int64, len(entries))
|
|
for i, e := range entries {
|
|
ids[i] = e.ID
|
|
}
|
|
if err := t.store.DrainDigestEntries(ctx, ids, now); err != nil {
|
|
log.Printf("tick: drain digest entries: %v", err)
|
|
}
|
|
}
|