mavend: split tick.go along the three concerns already in it (V-422)
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.
This commit is contained in:
@@ -10,22 +10,17 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/calendar"
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/delivery"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/morning"
|
||||
"github.com/kami/maven/internal/pattern"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/routine"
|
||||
"github.com/kami/maven/internal/store"
|
||||
@@ -315,613 +310,6 @@ func (t *tickLoop) repeatPhrase(rule string) (body, summary string) {
|
||||
return pn.Body, pn.Summary
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// detectPatterns runs the pattern detector proactively over every
|
||||
// action+object pair that has ever produced an event, independent of
|
||||
// whichever fact write (or channel) last touched it (Vikunja #43). This is
|
||||
// what makes pattern inference actually proactive: it fires on the daemon's
|
||||
// own schedule reading accumulated history, not only as a side effect of a
|
||||
// live voice turn.
|
||||
//
|
||||
// Idempotence and noise are handled by the store, not here — this function
|
||||
// is safe to call every tick:
|
||||
// - Same pattern, tick after tick: detectAndPropose's LookupProposedRoutine
|
||||
// check plus proposed_routines' UNIQUE(action, object) constraint (with
|
||||
// CreateProposedRoutine's ON CONFLICT DO NOTHING) mean a pair that
|
||||
// already has a row — in ANY status — produces no second row and no log
|
||||
// spam beyond the one line at genuine creation.
|
||||
// - A DISMISSED proposal must never come back. DismissProposedRoutine flips
|
||||
// status in place; the row is never deleted. So the same Lookup check
|
||||
// that stops a duplicate "proposed" also stops a "dismissed" one from
|
||||
// resurrecting — there is nothing tick-specific to get right here beyond
|
||||
// calling the same shared path the voice route already used.
|
||||
//
|
||||
// By default this only creates a row for the /routines page to show: it does
|
||||
// not notify, ring, or speak. Detection is not the same act as disturbing him
|
||||
// about it, and Maven is "not a nag, not autonomous" (CLAUDE.md). Announcing
|
||||
// is opt-in through the pattern_proposals config block — see announceProposal
|
||||
// for the restraints that apply even then. A proposal only starts producing
|
||||
// recurring nudges once he accepts it (fireAcceptedRoutines).
|
||||
func (t *tickLoop) detectPatterns(ctx context.Context, now time.Time, state loop.State) {
|
||||
pairs, err := t.store.DistinctEventPairs(ctx)
|
||||
if err != nil {
|
||||
log.Printf("tick: distinct event pairs: %v", err)
|
||||
return
|
||||
}
|
||||
announced := false
|
||||
for _, p := range pairs {
|
||||
r, _, err := detectAndPropose(ctx, t.store, p.Action, p.Object, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: detect pattern %s/%s: %v", p.Action, p.Object, err)
|
||||
continue
|
||||
}
|
||||
if r == nil {
|
||||
continue // no stable pattern, or already proposed/accepted/dismissed
|
||||
}
|
||||
log.Printf("tick: proposed routine: %s/%s every %.1f days", r.Action, r.Object, r.IntervalDays)
|
||||
// One announcement per tick at most, whatever the scan turned up. The
|
||||
// rest are on /routines; they are not lost, they are just not shouted.
|
||||
// Nor are they queued: the row now exists, so no later tick re-detects
|
||||
// them and they are never announced. See announceProposal.
|
||||
if announced {
|
||||
continue
|
||||
}
|
||||
announced = t.announceProposal(ctx, r, now, state)
|
||||
}
|
||||
}
|
||||
|
||||
// announceProposal offers a freshly inferred routine through the ordinary
|
||||
// care-delivery path, if announcing is switched on at all. Returns true when
|
||||
// something was actually sent.
|
||||
//
|
||||
// Everything here is restraint. The feature is off unless configured; when on
|
||||
// it is sev1 (the lowest severity, so quiet hours, away presence and snooze
|
||||
// all suppress it via loop.Gate exactly like a care nudge); it is spaced by
|
||||
// proposalCfg.Cooldown across every pair, not per pair; and a suppressed or
|
||||
// dropped announcement is NOT retried — the cooldown clock advances only on a
|
||||
// real send, but the proposal row already exists, so the next tick will not
|
||||
// re-detect it and nothing queues up behind it. A missed announcement means
|
||||
// he reads it on /routines instead, which is the whole point of the page.
|
||||
//
|
||||
// What the cooldown is and is not. detectAndPropose returns non-nil only for a
|
||||
// newly created row, so a pair gets exactly one chance to be spoken: the tick
|
||||
// that first proposes it. Combined with one announcement per tick, the first
|
||||
// tick over a populated history announces one pattern and permanently silences
|
||||
// every other pattern found in the same pass. That is the intent, not an
|
||||
// oversight — an inferred routine is not worth a second attempt at his
|
||||
// attention, and /routines lists all of them. So the cooldown does not drain a
|
||||
// backlog. It only spaces announcements of genuinely new pairs discovered on
|
||||
// later ticks. If it should ever become "one per day until each is mentioned",
|
||||
// that needs a queue rather than this counter.
|
||||
//
|
||||
// Cooldown gets its default here as well as in applyDefaults. That is
|
||||
// deliberate: a tickLoop assembled directly in a test never goes through Load,
|
||||
// and an unspaced announcer is not what those tests mean to exercise.
|
||||
//
|
||||
// The body is the detector's own literal Russian phrasing (pattern.PhraseRoutine
|
||||
// — "ты заправляешь поилку раз в 7 дней — напоминать?"), not LLM-generated, so
|
||||
// an inferred routine cannot arrive worded as something Maven never observed.
|
||||
func (t *tickLoop) announceProposal(ctx context.Context, r *pattern.ProposedRoutine, now time.Time, state loop.State) bool {
|
||||
if !t.proposalCfg.AnnounceProposals() {
|
||||
return false
|
||||
}
|
||||
cooldown := time.Duration(t.proposalCfg.Cooldown)
|
||||
if cooldown <= 0 {
|
||||
cooldown = config.DefaultProposalCooldown
|
||||
}
|
||||
if !t.lastProposalAt.IsZero() && now.Sub(t.lastProposalAt) < cooldown {
|
||||
return false
|
||||
}
|
||||
|
||||
rule := loop.Rule{Name: "proposal:" + r.Action + " " + r.Object, Severity: loop.Sev1}
|
||||
if !loop.Gate(state, rule) {
|
||||
return false
|
||||
}
|
||||
body := pattern.PhraseRoutine(r)
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{Rule: rule, Severity: rule.Severity, State: state},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
sent, err := t.dispatcher.DispatchNudge(ctx, pn, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: announce proposal %s/%s: %v", r.Action, r.Object, err)
|
||||
return false
|
||||
}
|
||||
if len(sent) == 0 {
|
||||
return false // routing dropped it — /routines still has it.
|
||||
}
|
||||
t.lastProposalAt = now
|
||||
return true
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
|
||||
// routinesFromConfig maps the config's routine blocks to the engine type.
|
||||
// Validation (cron parses, name/body present, severity defaulted) already ran
|
||||
// in config.Load, so this is a pure field copy.
|
||||
func routinesFromConfig(rc []config.RoutineConfig) []routine.Routine {
|
||||
if len(rc) == 0 {
|
||||
return nil
|
||||
}
|
||||
out := make([]routine.Routine, len(rc))
|
||||
for i, r := range rc {
|
||||
out[i] = routine.Routine{Name: r.Name, Cron: r.Cron, Body: r.Body, Severity: r.Severity}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// fireRoutines dispatches the routines whose cron schedule crossed since their
|
||||
// last fire. Each is delivered as a nudge through the normal routing table
|
||||
// (ChannelsFor(severity, presence)) with a "routine:"-prefixed rule name so it
|
||||
// can't collide with a care rule in the feedback autotuner. A dispatch failure
|
||||
// logs and continues — one bad send must not skip the rest, and routine.Due has
|
||||
// already advanced the last-fire time so a transient failure drops that fire
|
||||
// rather than replaying it every tick (a routine is clockwork, not an alarm —
|
||||
// no repeat-til-ack).
|
||||
func (t *tickLoop) fireRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
for _, r := range routine.Due(t.routines, t.routineLast, now) {
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{
|
||||
Rule: loop.Rule{Name: "routine:" + r.Name, Severity: loop.Severity(r.Severity)},
|
||||
Severity: loop.Severity(r.Severity),
|
||||
State: state,
|
||||
},
|
||||
Body: r.Body,
|
||||
Summary: r.Body,
|
||||
}
|
||||
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
||||
log.Printf("tick: dispatch routine %s: %v", r.Name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// fireAcceptedRoutines nudges about the routines the user accepted, once per
|
||||
// interval (Vikunja #366). Accepting used to create a single reminder, so a
|
||||
// non-weekly routine fired once and went quiet forever; the schedule lives in
|
||||
// the proposed_routines row now and the loop re-reads it every tick.
|
||||
//
|
||||
// A routine is a care-class nudge and goes through the restraint gate like any
|
||||
// other: quiet hours, away presence and snooze all suppress it. Reminders bypass
|
||||
// that gate; routines must not. A suppressed nudge is NOT marked fired, so it
|
||||
// goes out on the next tick that the gate allows — one nudge, held, not dropped
|
||||
// and not repeated.
|
||||
//
|
||||
// The body is literal text built from the detected action and object, not
|
||||
// LLM-phrased, so a routine can't hallucinate. It nudges; it never acts.
|
||||
func (t *tickLoop) fireAcceptedRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
rows, err := t.store.ListAcceptedRoutines(ctx)
|
||||
if err != nil {
|
||||
log.Printf("tick: list accepted routines: %v", err)
|
||||
return
|
||||
}
|
||||
accepted := make([]routine.Accepted, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
if r.AcceptedTs == nil {
|
||||
continue // accepted before the schedule column existed — no clock to start from.
|
||||
}
|
||||
accepted = append(accepted, routine.Accepted{
|
||||
ID: r.ID,
|
||||
Name: r.Action + " " + r.Object,
|
||||
IntervalDays: r.IntervalDays,
|
||||
Accepted: *r.AcceptedTs,
|
||||
LastFired: r.LastFiredTs,
|
||||
})
|
||||
}
|
||||
|
||||
for _, a := range routine.DueAccepted(accepted, now) {
|
||||
rule := loop.Rule{Name: "routine:" + a.Name, Severity: loop.Sev1}
|
||||
if !loop.Gate(state, rule) {
|
||||
continue
|
||||
}
|
||||
body := "пора: " + a.Name
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{Rule: rule, Severity: rule.Severity, State: state},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
sent, err := t.dispatcher.DispatchNudge(ctx, pn, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: dispatch accepted routine %d: %v", a.ID, err)
|
||||
continue
|
||||
}
|
||||
if len(sent) == 0 {
|
||||
continue // routing dropped it — leave it due.
|
||||
}
|
||||
if err := t.store.MarkRoutineFired(ctx, a.ID, now); err != nil {
|
||||
log.Printf("tick: mark routine %d fired: %v", a.ID, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// fireMorningRoutines checks each configured checklist against today's facts
|
||||
// and dispatches a nag listing exactly what's still missing, at most once per
|
||||
// routine per calendar day. Fact reads happen here (not in loop.Gatherer)
|
||||
// because the item↔fact-key mapping is morning-routine-specific, not a rule
|
||||
// concern — pulling it into the shared gather path would leak that mapping
|
||||
// into loop's "rules declare wanted keys" contract. Bodies are literal
|
||||
// operator text (item labels joined), not LLM-phrased, same rationale as
|
||||
// cron routines: deterministic, can't hallucinate a checklist item.
|
||||
func (t *tickLoop) fireMorningRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
if len(t.morningRoutines) == 0 {
|
||||
return
|
||||
}
|
||||
facts := t.gatherMorningFacts(ctx)
|
||||
|
||||
for _, cand := range morning.Due(t.morningRoutines, facts, t.morningLast, now) {
|
||||
body := morningNudgeBody(cand)
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{
|
||||
Rule: loop.Rule{Name: "morning:" + cand.Routine.Name, Severity: loop.Severity(cand.Routine.Severity)},
|
||||
Severity: loop.Severity(cand.Routine.Severity),
|
||||
State: state,
|
||||
},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
||||
log.Printf("tick: dispatch morning routine %s: %v", cand.Routine.Name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// morningNudgeBody words the one message a routine gets per day. Required
|
||||
// items are what she says was not done; optional ones follow, worded as
|
||||
// something he could still do rather than something he owes (Vikunja #473).
|
||||
// Operator text, not phrased by the model, for the same reason it always was:
|
||||
// a checklist item must not be invented.
|
||||
func morningNudgeBody(cand morning.Candidate) string {
|
||||
labels := func(items []morning.Item) string {
|
||||
out := make([]string, len(items))
|
||||
for i, it := range items {
|
||||
out[i] = it.Label
|
||||
}
|
||||
return strings.Join(out, ", ")
|
||||
}
|
||||
body := fmt.Sprintf("%s: не сделано — %s", cand.Routine.Name, labels(morning.Required(cand.Missing)))
|
||||
if opt := morning.OptionalOnly(cand.Missing); len(opt) > 0 {
|
||||
body += fmt.Sprintf(". если будет время — %s", labels(opt))
|
||||
}
|
||||
return body
|
||||
}
|
||||
|
||||
// gatherMorningFacts reads the latest fact for every item's fact_key across
|
||||
// all configured morning routines. Shared by fireMorningRoutines (nudge
|
||||
// decision) and morningStatus (read-only query) so the two paths can never
|
||||
// disagree about what evidence exists.
|
||||
func (t *tickLoop) gatherMorningFacts(ctx context.Context) map[string]store.Fact {
|
||||
keys := make(map[string]struct{})
|
||||
for _, r := range t.morningRoutines {
|
||||
for _, it := range r.Items {
|
||||
keys[it.FactKey] = struct{}{}
|
||||
}
|
||||
}
|
||||
facts := make(map[string]store.Fact, len(keys))
|
||||
for k := range keys {
|
||||
f, err := t.store.LatestFact(ctx, k)
|
||||
if err == nil {
|
||||
facts[k] = f
|
||||
continue
|
||||
}
|
||||
if err != store.ErrNoFact {
|
||||
log.Printf("tick: morning: latest fact %s: %v", k, err)
|
||||
}
|
||||
}
|
||||
return facts
|
||||
}
|
||||
|
||||
// morningStatus is the read-only "what's missing" query the web UI (and
|
||||
// eventually a voice query) calls. Pure recompute over the current facts —
|
||||
// no dedupe/nudge-time gating, unlike fireMorningRoutines: this answers
|
||||
// "state right now," not "should we nag."
|
||||
func (t *tickLoop) morningStatus(ctx context.Context, now time.Time) []ipc.MorningRoutineStatus {
|
||||
if len(t.morningRoutines) == 0 {
|
||||
return nil
|
||||
}
|
||||
facts := t.gatherMorningFacts(ctx)
|
||||
out := make([]ipc.MorningRoutineStatus, 0, len(t.morningRoutines))
|
||||
for _, r := range t.morningRoutines {
|
||||
st := morning.Evaluate(r, facts, now)
|
||||
done := make(map[string]bool, len(st.Completed))
|
||||
for _, it := range st.Completed {
|
||||
done[it.Key] = true
|
||||
}
|
||||
items := make([]ipc.MorningRoutineItem, len(r.Items))
|
||||
for i, it := range r.Items {
|
||||
items[i] = ipc.MorningRoutineItem{Key: it.Key, Label: it.Label, Done: done[it.Key]}
|
||||
}
|
||||
out = append(out, ipc.MorningRoutineStatus{
|
||||
Name: r.Name,
|
||||
Active: st.Active,
|
||||
WindowStart: r.WindowStart,
|
||||
WindowEnd: r.WindowEnd,
|
||||
Items: items,
|
||||
})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// dayPlan is the read-only "what does today hold" query (Vikunja #128). It is
|
||||
// the impure half of morning.BuildPlan: it reads the calendar events, the
|
||||
// pending reminders and the checklist facts, and the pure builder orders them.
|
||||
//
|
||||
// It never dispatches. Asking for the plan is a query like any other; the only
|
||||
// unprompted delivery in maven stays with the morning nudge and the
|
||||
// dispatcher's policy.
|
||||
func (t *tickLoop) dayPlan(ctx context.Context, now time.Time) ipc.DayPlan {
|
||||
y, m, d := now.Date()
|
||||
dayStart := time.Date(y, m, d, 0, 0, 0, 0, now.Location())
|
||||
dayEnd := dayStart.AddDate(0, 0, 1)
|
||||
|
||||
var events []morning.PlanEntry
|
||||
facts, err := t.store.CalendarEvents(ctx, dayStart, dayEnd)
|
||||
if err != nil {
|
||||
log.Printf("tick: day plan: calendar events: %v", err)
|
||||
}
|
||||
for _, f := range facts {
|
||||
events = append(events, morning.PlanEntry{
|
||||
At: f.Ts,
|
||||
// The plan prints the hour itself, so the "@ 14:00-14:30" tail the
|
||||
// fact value carries would say it twice.
|
||||
Text: calendar.FactSummary(f.Value),
|
||||
Kind: morning.PlanEvent,
|
||||
// Provenance below a calendar read (an ambient relay, #126) is
|
||||
// hedged rather than recited as fact.
|
||||
Uncertain: f.Confidence < 1.0,
|
||||
})
|
||||
}
|
||||
|
||||
var reminders []morning.PlanEntry
|
||||
rems, err := t.store.PendingReminders(ctx, dayStart, dayEnd)
|
||||
if err != nil {
|
||||
log.Printf("tick: day plan: pending reminders: %v", err)
|
||||
}
|
||||
for _, r := range rems {
|
||||
if r.Status != store.ReminderPending {
|
||||
continue
|
||||
}
|
||||
fire := r.NextFireTs
|
||||
if fire.IsZero() {
|
||||
fire = r.FireTs
|
||||
}
|
||||
reminders = append(reminders, morning.PlanEntry{
|
||||
At: fire,
|
||||
Text: r.Text(),
|
||||
Kind: morning.PlanReminder,
|
||||
})
|
||||
}
|
||||
|
||||
var checklistFacts map[string]store.Fact
|
||||
if len(t.morningRoutines) > 0 {
|
||||
checklistFacts = t.gatherMorningFacts(ctx)
|
||||
}
|
||||
plan := morning.BuildPlan(t.morningRoutines, checklistFacts, events, reminders, now)
|
||||
|
||||
out := ipc.DayPlan{Date: plan.Date, Spoken: plan.FormatRU()}
|
||||
out.Items = make([]ipc.DayPlanItem, len(plan.Items))
|
||||
for i, it := range plan.Items {
|
||||
out.Items[i] = ipc.DayPlanItem{
|
||||
At: it.At,
|
||||
Text: it.Text,
|
||||
Kind: string(it.Kind),
|
||||
Uncertain: it.Uncertain,
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// tune — the feedback auto-tuner's impure step. runs on a slow cadence
|
||||
// (autotuneInterval, see run) so it doesn't write a fact every tick. for each
|
||||
// rule:
|
||||
@@ -994,101 +382,3 @@ func (t *tickLoop) trace() *loop.TickTrace {
|
||||
defer t.mu.Unlock()
|
||||
return t.lastTrace
|
||||
}
|
||||
|
||||
// daemonAPI wraps a store-backed CoreAPI and overrides TickTrace with the
|
||||
// daemon's in-memory tick trace cache.
|
||||
type daemonAPI struct {
|
||||
ipc.CoreAPI
|
||||
getTrace func() *loop.TickTrace
|
||||
getMorningStatus func(ctx context.Context) []ipc.MorningRoutineStatus
|
||||
getDayPlan func(ctx context.Context) ipc.DayPlan
|
||||
chatFn func(ctx context.Context, conversation, text string) string
|
||||
getMCPServers func() []ipc.MCPServerStatus
|
||||
getEvents func(n int) []ipc.IntakeEvent
|
||||
}
|
||||
|
||||
// RecentEvents — the unified intake journal (Vikunja #283). Empty, not an
|
||||
// error, when no bus was wired: "nothing has arrived" and "the journal is off"
|
||||
// look the same to a reader on purpose, because neither is a fault and the
|
||||
// page renders both as an empty table.
|
||||
func (d *daemonAPI) RecentEvents(ctx context.Context, n int) ([]ipc.IntakeEvent, error) {
|
||||
if d.getEvents == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return d.getEvents(n), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) Chat(ctx context.Context, conversation, text string) (string, error) {
|
||||
if d.chatFn == nil {
|
||||
return "", errors.New("mavend: chat not available")
|
||||
}
|
||||
return d.chatFn(ctx, conversation, text), nil
|
||||
}
|
||||
|
||||
// MCPServers — the configured MCP servers and their health (Vikunja #251).
|
||||
// Empty, not an error, when the mcp block is absent: "not configured" is the
|
||||
// default state and the web surface renders it as such.
|
||||
func (d *daemonAPI) MCPServers(ctx context.Context) ([]ipc.MCPServerStatus, error) {
|
||||
if d.getMCPServers == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return d.getMCPServers(), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) TickTrace(ctx context.Context) (ipc.TickTrace, error) {
|
||||
trace := d.getTrace()
|
||||
if trace == nil {
|
||||
return ipc.TickTrace{}, nil
|
||||
}
|
||||
return toIPCTickTrace(*trace), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) MorningStatus(ctx context.Context) ([]ipc.MorningRoutineStatus, error) {
|
||||
if d.getMorningStatus == nil {
|
||||
return nil, errors.New("mavend: morning status not available")
|
||||
}
|
||||
return d.getMorningStatus(ctx), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) DayPlan(ctx context.Context) (ipc.DayPlan, error) {
|
||||
if d.getDayPlan == nil {
|
||||
return ipc.DayPlan{}, errors.New("mavend: day plan not available")
|
||||
}
|
||||
return d.getDayPlan(ctx), nil
|
||||
}
|
||||
|
||||
func toIPCTickTrace(t loop.TickTrace) ipc.TickTrace {
|
||||
rules := make([]ipc.RuleTrace, len(t.RuleTraces))
|
||||
for i, r := range t.RuleTraces {
|
||||
rules[i] = toIPCRuleTrace(r)
|
||||
}
|
||||
return ipc.TickTrace{
|
||||
Now: t.Now,
|
||||
Winner: t.Winner,
|
||||
Rules: rules,
|
||||
}
|
||||
}
|
||||
|
||||
func toIPCRuleTrace(r loop.RuleTrace) ipc.RuleTrace {
|
||||
return ipc.RuleTrace{
|
||||
RuleName: r.RuleName,
|
||||
Severity: int(r.Severity),
|
||||
PredicateResult: r.PredicateResult,
|
||||
GateResult: r.GateResult,
|
||||
GateBlockedBy: r.GateBlockedBy,
|
||||
GateDetail: toIPCGateDetail(r.GateDetail),
|
||||
WasSelected: r.WasSelected,
|
||||
LostTo: r.LostTo,
|
||||
}
|
||||
}
|
||||
|
||||
func toIPCGateDetail(d loop.GateDetail) ipc.GateDetail {
|
||||
return ipc.GateDetail{
|
||||
SnoozeUntil: d.SnoozeUntil,
|
||||
CooldownUntil: d.CooldownUntil,
|
||||
QuietHours: d.QuietHours,
|
||||
CalendarBusy: d.CalendarBusy,
|
||||
Presence: d.Presence,
|
||||
InertKeysMissing: d.InertKeysMissing,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
// mavend/tick_api.go — the daemonAPI read surface over the tick loop.
|
||||
//
|
||||
// Split out of tick.go, move-only (Vikunja #422). What mavweb asks the daemon
|
||||
// for, and the loop-to-ipc conversions those answers need.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
)
|
||||
|
||||
// daemonAPI wraps a store-backed CoreAPI and overrides TickTrace with the
|
||||
// daemon's in-memory tick trace cache.
|
||||
type daemonAPI struct {
|
||||
ipc.CoreAPI
|
||||
getTrace func() *loop.TickTrace
|
||||
getMorningStatus func(ctx context.Context) []ipc.MorningRoutineStatus
|
||||
getDayPlan func(ctx context.Context) ipc.DayPlan
|
||||
chatFn func(ctx context.Context, conversation, text string) string
|
||||
getMCPServers func() []ipc.MCPServerStatus
|
||||
getEvents func(n int) []ipc.IntakeEvent
|
||||
}
|
||||
|
||||
// RecentEvents — the unified intake journal (Vikunja #283). Empty, not an
|
||||
// error, when no bus was wired: "nothing has arrived" and "the journal is off"
|
||||
// look the same to a reader on purpose, because neither is a fault and the
|
||||
// page renders both as an empty table.
|
||||
func (d *daemonAPI) RecentEvents(ctx context.Context, n int) ([]ipc.IntakeEvent, error) {
|
||||
if d.getEvents == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return d.getEvents(n), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) Chat(ctx context.Context, conversation, text string) (string, error) {
|
||||
if d.chatFn == nil {
|
||||
return "", errors.New("mavend: chat not available")
|
||||
}
|
||||
return d.chatFn(ctx, conversation, text), nil
|
||||
}
|
||||
|
||||
// MCPServers — the configured MCP servers and their health (Vikunja #251).
|
||||
// Empty, not an error, when the mcp block is absent: "not configured" is the
|
||||
// default state and the web surface renders it as such.
|
||||
func (d *daemonAPI) MCPServers(ctx context.Context) ([]ipc.MCPServerStatus, error) {
|
||||
if d.getMCPServers == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return d.getMCPServers(), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) TickTrace(ctx context.Context) (ipc.TickTrace, error) {
|
||||
trace := d.getTrace()
|
||||
if trace == nil {
|
||||
return ipc.TickTrace{}, nil
|
||||
}
|
||||
return toIPCTickTrace(*trace), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) MorningStatus(ctx context.Context) ([]ipc.MorningRoutineStatus, error) {
|
||||
if d.getMorningStatus == nil {
|
||||
return nil, errors.New("mavend: morning status not available")
|
||||
}
|
||||
return d.getMorningStatus(ctx), nil
|
||||
}
|
||||
|
||||
func (d *daemonAPI) DayPlan(ctx context.Context) (ipc.DayPlan, error) {
|
||||
if d.getDayPlan == nil {
|
||||
return ipc.DayPlan{}, errors.New("mavend: day plan not available")
|
||||
}
|
||||
return d.getDayPlan(ctx), nil
|
||||
}
|
||||
|
||||
func toIPCTickTrace(t loop.TickTrace) ipc.TickTrace {
|
||||
rules := make([]ipc.RuleTrace, len(t.RuleTraces))
|
||||
for i, r := range t.RuleTraces {
|
||||
rules[i] = toIPCRuleTrace(r)
|
||||
}
|
||||
return ipc.TickTrace{
|
||||
Now: t.Now,
|
||||
Winner: t.Winner,
|
||||
Rules: rules,
|
||||
}
|
||||
}
|
||||
|
||||
func toIPCRuleTrace(r loop.RuleTrace) ipc.RuleTrace {
|
||||
return ipc.RuleTrace{
|
||||
RuleName: r.RuleName,
|
||||
Severity: int(r.Severity),
|
||||
PredicateResult: r.PredicateResult,
|
||||
GateResult: r.GateResult,
|
||||
GateBlockedBy: r.GateBlockedBy,
|
||||
GateDetail: toIPCGateDetail(r.GateDetail),
|
||||
WasSelected: r.WasSelected,
|
||||
LostTo: r.LostTo,
|
||||
}
|
||||
}
|
||||
|
||||
func toIPCGateDetail(d loop.GateDetail) ipc.GateDetail {
|
||||
return ipc.GateDetail{
|
||||
SnoozeUntil: d.SnoozeUntil,
|
||||
CooldownUntil: d.CooldownUntil,
|
||||
QuietHours: d.QuietHours,
|
||||
CalendarBusy: d.CalendarBusy,
|
||||
Presence: d.Presence,
|
||||
InertKeysMissing: d.InertKeysMissing,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,233 @@
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,197 @@
|
||||
// mavend/tick_morning.go — the morning checklist and the day plan.
|
||||
//
|
||||
// Split out of tick.go, move-only (Vikunja #422). One window per configured
|
||||
// routine, nudging once at the end for what is still open, plus the read
|
||||
// surfaces /morning renders.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/calendar"
|
||||
"github.com/kami/maven/internal/delivery"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/morning"
|
||||
"github.com/kami/maven/internal/store"
|
||||
)
|
||||
|
||||
// fireMorningRoutines checks each configured checklist against today's facts
|
||||
// and dispatches a nag listing exactly what's still missing, at most once per
|
||||
// routine per calendar day. Fact reads happen here (not in loop.Gatherer)
|
||||
// because the item↔fact-key mapping is morning-routine-specific, not a rule
|
||||
// concern — pulling it into the shared gather path would leak that mapping
|
||||
// into loop's "rules declare wanted keys" contract. Bodies are literal
|
||||
// operator text (item labels joined), not LLM-phrased, same rationale as
|
||||
// cron routines: deterministic, can't hallucinate a checklist item.
|
||||
func (t *tickLoop) fireMorningRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
if len(t.morningRoutines) == 0 {
|
||||
return
|
||||
}
|
||||
facts := t.gatherMorningFacts(ctx)
|
||||
|
||||
for _, cand := range morning.Due(t.morningRoutines, facts, t.morningLast, now) {
|
||||
body := morningNudgeBody(cand)
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{
|
||||
Rule: loop.Rule{Name: "morning:" + cand.Routine.Name, Severity: loop.Severity(cand.Routine.Severity)},
|
||||
Severity: loop.Severity(cand.Routine.Severity),
|
||||
State: state,
|
||||
},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
||||
log.Printf("tick: dispatch morning routine %s: %v", cand.Routine.Name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// morningNudgeBody words the one message a routine gets per day. Required
|
||||
// items are what she says was not done; optional ones follow, worded as
|
||||
// something he could still do rather than something he owes (Vikunja #473).
|
||||
// Operator text, not phrased by the model, for the same reason it always was:
|
||||
// a checklist item must not be invented.
|
||||
func morningNudgeBody(cand morning.Candidate) string {
|
||||
labels := func(items []morning.Item) string {
|
||||
out := make([]string, len(items))
|
||||
for i, it := range items {
|
||||
out[i] = it.Label
|
||||
}
|
||||
return strings.Join(out, ", ")
|
||||
}
|
||||
body := fmt.Sprintf("%s: не сделано — %s", cand.Routine.Name, labels(morning.Required(cand.Missing)))
|
||||
if opt := morning.OptionalOnly(cand.Missing); len(opt) > 0 {
|
||||
body += fmt.Sprintf(". если будет время — %s", labels(opt))
|
||||
}
|
||||
return body
|
||||
}
|
||||
|
||||
// gatherMorningFacts reads the latest fact for every item's fact_key across
|
||||
// all configured morning routines. Shared by fireMorningRoutines (nudge
|
||||
// decision) and morningStatus (read-only query) so the two paths can never
|
||||
// disagree about what evidence exists.
|
||||
func (t *tickLoop) gatherMorningFacts(ctx context.Context) map[string]store.Fact {
|
||||
keys := make(map[string]struct{})
|
||||
for _, r := range t.morningRoutines {
|
||||
for _, it := range r.Items {
|
||||
keys[it.FactKey] = struct{}{}
|
||||
}
|
||||
}
|
||||
facts := make(map[string]store.Fact, len(keys))
|
||||
for k := range keys {
|
||||
f, err := t.store.LatestFact(ctx, k)
|
||||
if err == nil {
|
||||
facts[k] = f
|
||||
continue
|
||||
}
|
||||
if err != store.ErrNoFact {
|
||||
log.Printf("tick: morning: latest fact %s: %v", k, err)
|
||||
}
|
||||
}
|
||||
return facts
|
||||
}
|
||||
|
||||
// morningStatus is the read-only "what's missing" query the web UI (and
|
||||
// eventually a voice query) calls. Pure recompute over the current facts —
|
||||
// no dedupe/nudge-time gating, unlike fireMorningRoutines: this answers
|
||||
// "state right now," not "should we nag."
|
||||
func (t *tickLoop) morningStatus(ctx context.Context, now time.Time) []ipc.MorningRoutineStatus {
|
||||
if len(t.morningRoutines) == 0 {
|
||||
return nil
|
||||
}
|
||||
facts := t.gatherMorningFacts(ctx)
|
||||
out := make([]ipc.MorningRoutineStatus, 0, len(t.morningRoutines))
|
||||
for _, r := range t.morningRoutines {
|
||||
st := morning.Evaluate(r, facts, now)
|
||||
done := make(map[string]bool, len(st.Completed))
|
||||
for _, it := range st.Completed {
|
||||
done[it.Key] = true
|
||||
}
|
||||
items := make([]ipc.MorningRoutineItem, len(r.Items))
|
||||
for i, it := range r.Items {
|
||||
items[i] = ipc.MorningRoutineItem{Key: it.Key, Label: it.Label, Done: done[it.Key]}
|
||||
}
|
||||
out = append(out, ipc.MorningRoutineStatus{
|
||||
Name: r.Name,
|
||||
Active: st.Active,
|
||||
WindowStart: r.WindowStart,
|
||||
WindowEnd: r.WindowEnd,
|
||||
Items: items,
|
||||
})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// dayPlan is the read-only "what does today hold" query (Vikunja #128). It is
|
||||
// the impure half of morning.BuildPlan: it reads the calendar events, the
|
||||
// pending reminders and the checklist facts, and the pure builder orders them.
|
||||
//
|
||||
// It never dispatches. Asking for the plan is a query like any other; the only
|
||||
// unprompted delivery in maven stays with the morning nudge and the
|
||||
// dispatcher's policy.
|
||||
func (t *tickLoop) dayPlan(ctx context.Context, now time.Time) ipc.DayPlan {
|
||||
y, m, d := now.Date()
|
||||
dayStart := time.Date(y, m, d, 0, 0, 0, 0, now.Location())
|
||||
dayEnd := dayStart.AddDate(0, 0, 1)
|
||||
|
||||
var events []morning.PlanEntry
|
||||
facts, err := t.store.CalendarEvents(ctx, dayStart, dayEnd)
|
||||
if err != nil {
|
||||
log.Printf("tick: day plan: calendar events: %v", err)
|
||||
}
|
||||
for _, f := range facts {
|
||||
events = append(events, morning.PlanEntry{
|
||||
At: f.Ts,
|
||||
// The plan prints the hour itself, so the "@ 14:00-14:30" tail the
|
||||
// fact value carries would say it twice.
|
||||
Text: calendar.FactSummary(f.Value),
|
||||
Kind: morning.PlanEvent,
|
||||
// Provenance below a calendar read (an ambient relay, #126) is
|
||||
// hedged rather than recited as fact.
|
||||
Uncertain: f.Confidence < 1.0,
|
||||
})
|
||||
}
|
||||
|
||||
var reminders []morning.PlanEntry
|
||||
rems, err := t.store.PendingReminders(ctx, dayStart, dayEnd)
|
||||
if err != nil {
|
||||
log.Printf("tick: day plan: pending reminders: %v", err)
|
||||
}
|
||||
for _, r := range rems {
|
||||
if r.Status != store.ReminderPending {
|
||||
continue
|
||||
}
|
||||
fire := r.NextFireTs
|
||||
if fire.IsZero() {
|
||||
fire = r.FireTs
|
||||
}
|
||||
reminders = append(reminders, morning.PlanEntry{
|
||||
At: fire,
|
||||
Text: r.Text(),
|
||||
Kind: morning.PlanReminder,
|
||||
})
|
||||
}
|
||||
|
||||
var checklistFacts map[string]store.Fact
|
||||
if len(t.morningRoutines) > 0 {
|
||||
checklistFacts = t.gatherMorningFacts(ctx)
|
||||
}
|
||||
plan := morning.BuildPlan(t.morningRoutines, checklistFacts, events, reminders, now)
|
||||
|
||||
out := ipc.DayPlan{Date: plan.Date, Spoken: plan.FormatRU()}
|
||||
out.Items = make([]ipc.DayPlanItem, len(plan.Items))
|
||||
for i, it := range plan.Items {
|
||||
out.Items[i] = ipc.DayPlanItem{
|
||||
At: it.At,
|
||||
Text: it.Text,
|
||||
Kind: string(it.Kind),
|
||||
Uncertain: it.Uncertain,
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -0,0 +1,234 @@
|
||||
// mavend/tick_routines.go — pattern detection and the routines that fire.
|
||||
//
|
||||
// Split out of tick.go, move-only (Vikunja #422). Configured routines, the
|
||||
// routines he accepted on /routines, and the tick-side detector that proposes
|
||||
// new ones.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
"github.com/kami/maven/internal/delivery"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/pattern"
|
||||
"github.com/kami/maven/internal/routine"
|
||||
)
|
||||
|
||||
// detectPatterns runs the pattern detector proactively over every
|
||||
// action+object pair that has ever produced an event, independent of
|
||||
// whichever fact write (or channel) last touched it (Vikunja #43). This is
|
||||
// what makes pattern inference actually proactive: it fires on the daemon's
|
||||
// own schedule reading accumulated history, not only as a side effect of a
|
||||
// live voice turn.
|
||||
//
|
||||
// Idempotence and noise are handled by the store, not here — this function
|
||||
// is safe to call every tick:
|
||||
// - Same pattern, tick after tick: detectAndPropose's LookupProposedRoutine
|
||||
// check plus proposed_routines' UNIQUE(action, object) constraint (with
|
||||
// CreateProposedRoutine's ON CONFLICT DO NOTHING) mean a pair that
|
||||
// already has a row — in ANY status — produces no second row and no log
|
||||
// spam beyond the one line at genuine creation.
|
||||
// - A DISMISSED proposal must never come back. DismissProposedRoutine flips
|
||||
// status in place; the row is never deleted. So the same Lookup check
|
||||
// that stops a duplicate "proposed" also stops a "dismissed" one from
|
||||
// resurrecting — there is nothing tick-specific to get right here beyond
|
||||
// calling the same shared path the voice route already used.
|
||||
//
|
||||
// By default this only creates a row for the /routines page to show: it does
|
||||
// not notify, ring, or speak. Detection is not the same act as disturbing him
|
||||
// about it, and Maven is "not a nag, not autonomous" (CLAUDE.md). Announcing
|
||||
// is opt-in through the pattern_proposals config block — see announceProposal
|
||||
// for the restraints that apply even then. A proposal only starts producing
|
||||
// recurring nudges once he accepts it (fireAcceptedRoutines).
|
||||
func (t *tickLoop) detectPatterns(ctx context.Context, now time.Time, state loop.State) {
|
||||
pairs, err := t.store.DistinctEventPairs(ctx)
|
||||
if err != nil {
|
||||
log.Printf("tick: distinct event pairs: %v", err)
|
||||
return
|
||||
}
|
||||
announced := false
|
||||
for _, p := range pairs {
|
||||
r, _, err := detectAndPropose(ctx, t.store, p.Action, p.Object, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: detect pattern %s/%s: %v", p.Action, p.Object, err)
|
||||
continue
|
||||
}
|
||||
if r == nil {
|
||||
continue // no stable pattern, or already proposed/accepted/dismissed
|
||||
}
|
||||
log.Printf("tick: proposed routine: %s/%s every %.1f days", r.Action, r.Object, r.IntervalDays)
|
||||
// One announcement per tick at most, whatever the scan turned up. The
|
||||
// rest are on /routines; they are not lost, they are just not shouted.
|
||||
// Nor are they queued: the row now exists, so no later tick re-detects
|
||||
// them and they are never announced. See announceProposal.
|
||||
if announced {
|
||||
continue
|
||||
}
|
||||
announced = t.announceProposal(ctx, r, now, state)
|
||||
}
|
||||
}
|
||||
|
||||
// announceProposal offers a freshly inferred routine through the ordinary
|
||||
// care-delivery path, if announcing is switched on at all. Returns true when
|
||||
// something was actually sent.
|
||||
//
|
||||
// Everything here is restraint. The feature is off unless configured; when on
|
||||
// it is sev1 (the lowest severity, so quiet hours, away presence and snooze
|
||||
// all suppress it via loop.Gate exactly like a care nudge); it is spaced by
|
||||
// proposalCfg.Cooldown across every pair, not per pair; and a suppressed or
|
||||
// dropped announcement is NOT retried — the cooldown clock advances only on a
|
||||
// real send, but the proposal row already exists, so the next tick will not
|
||||
// re-detect it and nothing queues up behind it. A missed announcement means
|
||||
// he reads it on /routines instead, which is the whole point of the page.
|
||||
//
|
||||
// What the cooldown is and is not. detectAndPropose returns non-nil only for a
|
||||
// newly created row, so a pair gets exactly one chance to be spoken: the tick
|
||||
// that first proposes it. Combined with one announcement per tick, the first
|
||||
// tick over a populated history announces one pattern and permanently silences
|
||||
// every other pattern found in the same pass. That is the intent, not an
|
||||
// oversight — an inferred routine is not worth a second attempt at his
|
||||
// attention, and /routines lists all of them. So the cooldown does not drain a
|
||||
// backlog. It only spaces announcements of genuinely new pairs discovered on
|
||||
// later ticks. If it should ever become "one per day until each is mentioned",
|
||||
// that needs a queue rather than this counter.
|
||||
//
|
||||
// Cooldown gets its default here as well as in applyDefaults. That is
|
||||
// deliberate: a tickLoop assembled directly in a test never goes through Load,
|
||||
// and an unspaced announcer is not what those tests mean to exercise.
|
||||
//
|
||||
// The body is the detector's own literal Russian phrasing (pattern.PhraseRoutine
|
||||
// — "ты заправляешь поилку раз в 7 дней — напоминать?"), not LLM-generated, so
|
||||
// an inferred routine cannot arrive worded as something Maven never observed.
|
||||
func (t *tickLoop) announceProposal(ctx context.Context, r *pattern.ProposedRoutine, now time.Time, state loop.State) bool {
|
||||
if !t.proposalCfg.AnnounceProposals() {
|
||||
return false
|
||||
}
|
||||
cooldown := time.Duration(t.proposalCfg.Cooldown)
|
||||
if cooldown <= 0 {
|
||||
cooldown = config.DefaultProposalCooldown
|
||||
}
|
||||
if !t.lastProposalAt.IsZero() && now.Sub(t.lastProposalAt) < cooldown {
|
||||
return false
|
||||
}
|
||||
|
||||
rule := loop.Rule{Name: "proposal:" + r.Action + " " + r.Object, Severity: loop.Sev1}
|
||||
if !loop.Gate(state, rule) {
|
||||
return false
|
||||
}
|
||||
body := pattern.PhraseRoutine(r)
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{Rule: rule, Severity: rule.Severity, State: state},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
sent, err := t.dispatcher.DispatchNudge(ctx, pn, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: announce proposal %s/%s: %v", r.Action, r.Object, err)
|
||||
return false
|
||||
}
|
||||
if len(sent) == 0 {
|
||||
return false // routing dropped it — /routines still has it.
|
||||
}
|
||||
t.lastProposalAt = now
|
||||
return true
|
||||
}
|
||||
|
||||
// routinesFromConfig maps the config's routine blocks to the engine type.
|
||||
// Validation (cron parses, name/body present, severity defaulted) already ran
|
||||
// in config.Load, so this is a pure field copy.
|
||||
func routinesFromConfig(rc []config.RoutineConfig) []routine.Routine {
|
||||
if len(rc) == 0 {
|
||||
return nil
|
||||
}
|
||||
out := make([]routine.Routine, len(rc))
|
||||
for i, r := range rc {
|
||||
out[i] = routine.Routine{Name: r.Name, Cron: r.Cron, Body: r.Body, Severity: r.Severity}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// fireRoutines dispatches the routines whose cron schedule crossed since their
|
||||
// last fire. Each is delivered as a nudge through the normal routing table
|
||||
// (ChannelsFor(severity, presence)) with a "routine:"-prefixed rule name so it
|
||||
// can't collide with a care rule in the feedback autotuner. A dispatch failure
|
||||
// logs and continues — one bad send must not skip the rest, and routine.Due has
|
||||
// already advanced the last-fire time so a transient failure drops that fire
|
||||
// rather than replaying it every tick (a routine is clockwork, not an alarm —
|
||||
// no repeat-til-ack).
|
||||
func (t *tickLoop) fireRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
for _, r := range routine.Due(t.routines, t.routineLast, now) {
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{
|
||||
Rule: loop.Rule{Name: "routine:" + r.Name, Severity: loop.Severity(r.Severity)},
|
||||
Severity: loop.Severity(r.Severity),
|
||||
State: state,
|
||||
},
|
||||
Body: r.Body,
|
||||
Summary: r.Body,
|
||||
}
|
||||
if _, err := t.dispatcher.DispatchNudge(ctx, pn, now); err != nil {
|
||||
log.Printf("tick: dispatch routine %s: %v", r.Name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// fireAcceptedRoutines nudges about the routines the user accepted, once per
|
||||
// interval (Vikunja #366). Accepting used to create a single reminder, so a
|
||||
// non-weekly routine fired once and went quiet forever; the schedule lives in
|
||||
// the proposed_routines row now and the loop re-reads it every tick.
|
||||
//
|
||||
// A routine is a care-class nudge and goes through the restraint gate like any
|
||||
// other: quiet hours, away presence and snooze all suppress it. Reminders bypass
|
||||
// that gate; routines must not. A suppressed nudge is NOT marked fired, so it
|
||||
// goes out on the next tick that the gate allows — one nudge, held, not dropped
|
||||
// and not repeated.
|
||||
//
|
||||
// The body is literal text built from the detected action and object, not
|
||||
// LLM-phrased, so a routine can't hallucinate. It nudges; it never acts.
|
||||
func (t *tickLoop) fireAcceptedRoutines(ctx context.Context, now time.Time, state loop.State) {
|
||||
rows, err := t.store.ListAcceptedRoutines(ctx)
|
||||
if err != nil {
|
||||
log.Printf("tick: list accepted routines: %v", err)
|
||||
return
|
||||
}
|
||||
accepted := make([]routine.Accepted, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
if r.AcceptedTs == nil {
|
||||
continue // accepted before the schedule column existed — no clock to start from.
|
||||
}
|
||||
accepted = append(accepted, routine.Accepted{
|
||||
ID: r.ID,
|
||||
Name: r.Action + " " + r.Object,
|
||||
IntervalDays: r.IntervalDays,
|
||||
Accepted: *r.AcceptedTs,
|
||||
LastFired: r.LastFiredTs,
|
||||
})
|
||||
}
|
||||
|
||||
for _, a := range routine.DueAccepted(accepted, now) {
|
||||
rule := loop.Rule{Name: "routine:" + a.Name, Severity: loop.Sev1}
|
||||
if !loop.Gate(state, rule) {
|
||||
continue
|
||||
}
|
||||
body := "пора: " + a.Name
|
||||
pn := delivery.PhrasedNudge{
|
||||
Candidate: loop.Candidate{Rule: rule, Severity: rule.Severity, State: state},
|
||||
Body: body,
|
||||
Summary: body,
|
||||
}
|
||||
sent, err := t.dispatcher.DispatchNudge(ctx, pn, now)
|
||||
if err != nil {
|
||||
log.Printf("tick: dispatch accepted routine %d: %v", a.ID, err)
|
||||
continue
|
||||
}
|
||||
if len(sent) == 0 {
|
||||
continue // routing dropped it — leave it due.
|
||||
}
|
||||
if err := t.store.MarkRoutineFired(ctx, a.ID, now); err != nil {
|
||||
log.Printf("tick: mark routine %d fired: %v", a.ID, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user