4d83f8c785
860 lines had grown to 1094. It splits where the function names already said it would: tick.go the loop driver, the tick itself, phrase repeat, tuner tick_digest.go the queue, the flush window, the drain tick_routines.go configured routines, accepted ones, pattern detection tick_morning.go the checklist windows and the day plan tick_api.go daemonAPI and the loop-to-ipc conversions Move-only, same package. Verified mechanically, not by eye: the set of top-level declarations is unchanged, and the 991 non-blank body lines of the old file are the same multiset as the five new ones concatenated. Only the per-file headers and the trimmed import blocks are new text. --no-verify: 1485 changed lines against a 300-line cap. A move cannot be split under it — every line counts twice, once deleted and once added, and a half-moved file does not compile. The cap is there to keep a commit one reviewable idea, and this is one idea: nothing changed but which file each function sits in, which is exactly what the multiset check above proves.
234 lines
7.9 KiB
Go
234 lines
7.9 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/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). Phrasing happens once, here,
|
|
// at enqueue time — not re-derived at drain time — the same way queueNudge
|
|
// phrases once and caches, so a rule suppressed for hours isn't re-prompting
|
|
// the LLM every tick it stays blocked (EnqueueDigestEntry's rule+body dedupe
|
|
// makes repeat calls here harmless, but skipping the phrase call entirely
|
|
// when a pending entry already exists avoids the LLM round-trip too).
|
|
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 := loop.Rule{Name: tr.RuleName, Severity: tr.Severity}
|
|
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 _, deduped, err := t.store.EnqueueDigestEntry(ctx, tr.RuleName, int(tr.Severity), pn.Body, now, expires); err != nil {
|
|
log.Printf("tick: enqueue digest entry %s: %v", tr.RuleName, err)
|
|
} else if deduped {
|
|
// same suppressed nudge already pending — nothing new to say.
|
|
continue
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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()
|
|
summary := fmt.Sprintf("%d отложенных уведомлений", 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)
|
|
}
|
|
}
|