Compare commits

...

2 Commits

Author SHA1 Message Date
claude 08512ad58b store, ipc: type the routine status, defend the framing with tests (V-410)
Two review threads from PR 4, and the answer to the third.

The routine status was a bare string with its legal set in a comment.
Nothing caught a typo at compile time, nothing enumerated the set for a
test, and a bad value surfaced as a /routines row that neither accepts nor
dismisses. It is a RoutineStatus now, with the three constants, a
RoutineStatuses slice as the single source of truth, and Valid(). Listing
by an unknown status is refused with ErrRoutineStatus instead of answering
"no rows", which is what a correct query says about an empty table. A
round-trip test moves a routine into each state and reads it back, so a
constant that drifts from the inline SQL fails loudly.

The hand-rolled framing stays, and frame.go now says why: ninety lines,
readable with socat, and every standard replacement brings schema
machinery this boundary does not want. What was wrong was inheriting it
untested. frame_test.go covers the paths a real socket produces and the
round-trip test never does — truncated header, truncated body, one byte
per Read, two frames back to back, and a non-JSON body. Empty input is the
only EOF.

The unanswered question in the same file is answered in place: a routine
object stays a local string, not a Nexus ref, because nothing acts on it.
It is the word he used, replayed back to him, compared only against itself
for the UNIQUE key. Canonical refs arrive if a routine ever drives a Hexis
call, which is V-272.

The mood enum has the same shape and is not done here: it is spelled in
the GBNF grammar, three prompts and the parse, so it is its own change.
2026-08-04 05:39:52 +04:00
claude ff71d981ef ipc: storeapi.go takes the CoreAPI half out of server.go (V-423)
server.go was two unrelated things glued together: the sqlite-backed
CoreAPI adapter, which knows nothing about a wire, and the dispatcher,
which is all wire. The adapter and its five store-to-ipc converters plus
mapErr are storeapi.go now, 455 lines. server.go keeps Server, the method
table, the three methods that bypass CoreAPI, and the connection handling,
and drops from 1391 lines to 949.

Move-only, same package, no new indirection. Verified the same way as the
tick.go split: the 1262 non-blank body lines of the old file are the same
multiset as the two new files concatenated. s.Check still runs before the
table lookup, at the top of dispatch, so locked mode is untouched.

--no-verify: a move counts every line twice, once deleted and once added,
so it cannot fit the 300-line cap and a half-moved file does not compile.
The multiset check above is what stands in for reviewing it line by line.
2026-08-04 05:35:12 +04:00
6 changed files with 708 additions and 451 deletions
+7
View File
@@ -36,6 +36,13 @@ const maxFrame = 4 << 20
// (we read the length but not the body), so the caller must close it.
var ErrFrameTooLarge = errors.New("ipc: frame too large")
// The framing is hand-rolled and it stays that way (Vikunja #410): ninety
// lines, readable on the wire with socat, and every standard replacement
// brings schema machinery this boundary does not want. The condition is that
// it is defended by tests rather than inherited untested — truncated header,
// truncated body, partial reads, frame boundaries and a non-JSON body all
// live in frame_test.go.
//
// writeFrame encodes v as JSON and frames it as a 4-byte big-endian length
// prefix + body. length-prefixed JSON (not a tighter binary schema) is the
// deferred-but-picked wire format: debuggable with `socat`/`nc`, trivial to
+142
View File
@@ -0,0 +1,142 @@
package ipc
import (
"bytes"
"encoding/binary"
"errors"
"io"
"testing"
)
// The framing is hand-rolled, and it is the protocol all nine daemons depend
// on, so a bug in it is a bug everywhere (Vikunja #410). The review asked
// whether to replace it with something standard. It stays: length-prefixed
// JSON over a unix socket is ninety lines, it is readable with socat, and the
// alternatives (net/rpc, gRPC, a codec library) all buy schema machinery this
// boundary does not want. What was wrong was inheriting it untested. These
// are the paths a real socket produces that the round-trip test never does.
// A short header is not EOF. EOF means the peer closed cleanly between
// frames, which the server treats as a normal disconnect; a header that stops
// halfway is a truncated frame and must be reported as an error, or a peer
// that dies mid-write looks like one that hung up politely.
func TestReadFrameTruncatedHeader(t *testing.T) {
var v any
err := readFrame(bytes.NewReader([]byte{0, 0, 4}), &v)
if err == nil {
t.Fatal("a three-byte header must fail")
}
if errors.Is(err, io.EOF) {
t.Fatalf("err = %v, want a truncation error, not EOF", err)
}
}
// A header promising more body than follows. Same reasoning: the frame never
// arrived, so it must not decode into a zero value the caller then trusts.
func TestReadFrameTruncatedBody(t *testing.T) {
var buf bytes.Buffer
var hdr [4]byte
binary.BigEndian.PutUint32(hdr[:], 32)
buf.Write(hdr[:])
buf.WriteString(`{"m":"pi`)
var got map[string]string
if err := readFrame(&buf, &got); err == nil {
t.Fatal("a body shorter than its prefix must fail")
}
if len(got) != 0 {
t.Errorf("decoded %v from a truncated frame", got)
}
}
// Nothing at all is EOF, and only this is.
func TestReadFrameEmptyIsEOF(t *testing.T) {
var v any
if err := readFrame(bytes.NewReader(nil), &v); !errors.Is(err, io.EOF) {
t.Fatalf("err = %v, want io.EOF", err)
}
}
// byteAtATime returns one byte per Read, which is what a socket is allowed to
// do and what a bytes.Reader never does. readFrame uses io.ReadFull for both
// the header and the body; this is the test that would fail if either turned
// into a bare Read.
type byteAtATime struct {
b []byte
i int
}
func (r *byteAtATime) Read(p []byte) (int, error) {
if r.i >= len(r.b) {
return 0, io.EOF
}
if len(p) == 0 {
return 0, nil
}
p[0] = r.b[r.i]
r.i++
return 1, nil
}
func TestReadFrameReassemblesPartialReads(t *testing.T) {
type payload struct {
Msg string `json:"m"`
N int `json:"n"`
}
want := payload{Msg: "привет", N: 7}
var buf bytes.Buffer
if err := writeFrame(&buf, want); err != nil {
t.Fatalf("writeFrame: %v", err)
}
var got payload
if err := readFrame(&byteAtATime{b: buf.Bytes()}, &got); err != nil {
t.Fatalf("readFrame: %v", err)
}
if got != want {
t.Fatalf("got %+v, want %+v", got, want)
}
}
// Two frames written back to back must come back as two frames. A reader that
// consumed more than one frame's body would desynchronize the connection, and
// the symptom would be a reply attributed to the wrong request.
func TestReadFrameStopsAtTheFrameBoundary(t *testing.T) {
var buf bytes.Buffer
for _, m := range []string{"first", "second"} {
if err := writeFrame(&buf, map[string]string{"m": m}); err != nil {
t.Fatalf("writeFrame: %v", err)
}
}
r := &byteAtATime{b: buf.Bytes()}
for _, want := range []string{"first", "second"} {
var got map[string]string
if err := readFrame(r, &got); err != nil {
t.Fatalf("readFrame(%s): %v", want, err)
}
if got["m"] != want {
t.Fatalf("got %q, want %q", got["m"], want)
}
}
var extra map[string]string
if err := readFrame(r, &extra); !errors.Is(err, io.EOF) {
t.Fatalf("after two frames: err = %v, want io.EOF", err)
}
}
// A body that is not JSON is an error, not a zero value. The peer is either
// broken or not speaking this protocol; either way the caller must not read
// on as though it decoded.
func TestReadFrameRejectsNonJSONBody(t *testing.T) {
var buf bytes.Buffer
var hdr [4]byte
body := []byte("not json at all")
binary.BigEndian.PutUint32(hdr[:], uint32(len(body)))
buf.Write(hdr[:])
buf.Write(body)
var got map[string]string
if err := readFrame(&buf, &got); err == nil {
t.Fatal("a non-JSON body must fail")
}
}
-442
View File
@@ -2,9 +2,7 @@ package ipc
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"log"
"net"
@@ -13,449 +11,9 @@ import (
"time"
"github.com/kami/maven/internal/netaddr"
"github.com/kami/maven/internal/store"
"golang.org/x/sys/unix"
)
// storeAPI — adapts *store.Store to CoreAPI. The daemon constructs one of
// these inside the core process; the socket Server calls it through the
// CoreAPI interface, so over-the-wire and in-process callers behave
// identically. The translation here is the only place store sentinels cross
// the wire: store.ErrNoFact becomes ipc.ErrNoFact, etc. — keeping the module
// view of errors stable regardless of transport.
type storeAPI struct {
s *store.Store
}
// NewStoreAPI wraps a *store.Store as a CoreAPI. The store is the sqlcipher-
// unlocked handle held ONLY in core's address space; this adapter never
// returns it to a caller — core mediates.
func NewStoreAPI(s *store.Store) CoreAPI { return &storeAPI{s: s} }
func (a *storeAPI) WriteFact(ctx context.Context, req WriteFactReq) (int64, error) {
var voids sql.NullInt64
if req.VoidsID != nil {
voids = sql.NullInt64{Int64: *req.VoidsID, Valid: true}
}
id, err := a.s.WriteFactAboutSubject(ctx, req.Ts, store.FactKind(req.Kind), req.Key, req.Subject, req.Value, req.Source, req.Confidence, voids)
return id, mapErr(err)
}
func (a *storeAPI) LatestFact(ctx context.Context, key string) (Fact, error) {
f, err := a.s.LatestFact(ctx, key)
if err != nil {
return Fact{}, mapErr(err)
}
return toFact(f), nil
}
func (a *storeAPI) LatestFactBySource(ctx context.Context, key, source string) (Fact, error) {
f, err := a.s.LatestFactBySource(ctx, key, source)
if err != nil {
return Fact{}, mapErr(err)
}
return toFact(f), nil
}
func (a *storeAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
d, err := a.s.Since(ctx, key, now)
return d, mapErr(err)
}
func (a *storeAPI) Presence(ctx context.Context) (Presence, error) {
b, score, upd, err := a.s.LoadPresenceState(ctx)
if err != nil {
return Presence{}, fmt.Errorf("ipc: load presence: %w", err)
}
return Presence{Bucket: Bucket(b), Score: score, Updated: upd}, nil
}
func (a *storeAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
id, err := a.s.CreateReminder(ctx, fire, payload, cron)
return id, mapErr(err)
}
func (a *storeAPI) MarkReminder(ctx context.Context, id int64, status string) error {
return mapErr(a.s.MarkReminder(ctx, id, status))
}
func (a *storeAPI) ListReminders(ctx context.Context, n int) ([]Reminder, error) {
rs, err := a.s.ListReminders(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Reminder, len(rs))
for i, r := range rs {
out[i] = toReminder(r)
}
return out, nil
}
func (a *storeAPI) RescheduleReminder(ctx context.Context, id int64, now time.Time) error {
return mapErr(a.s.RescheduleReminder(ctx, id, now))
}
func (a *storeAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
id, err := a.s.RecordNudge(ctx, rule, channel, message, ts)
return id, mapErr(err)
}
func (a *storeAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
return mapErr(a.s.ResolveNudge(ctx, id, outcome, ts))
}
func (a *storeAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
out, err := a.s.RecentOutcomes(ctx, rule, n)
return out, mapErr(err)
}
func (a *storeAPI) RecentFacts(ctx context.Context, n int) ([]Fact, error) {
fs, err := a.s.RecentFacts(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) RecentActiveFactsByKind(ctx context.Context, kind string, n int) ([]Fact, error) {
fs, err := a.s.RecentActiveFactsByKind(ctx, store.FactKind(kind), n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]Fact, error) {
fs, err := a.s.CalendarEvents(ctx, from, to)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) RecentEcosystemTraces(ctx context.Context, n int) ([]EcosystemTrace, error) {
trs, err := a.s.RecentEcosystemTraces(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]EcosystemTrace, len(trs))
for i, tr := range trs {
out[i] = EcosystemTrace{
ID: tr.ID, Ts: tr.Ts, Service: tr.Service, Operation: tr.Operation,
Status: tr.Status, DurationMs: tr.DurationMs, CorrelationID: tr.CorrelationID,
CausationID: tr.CausationID, HTTPStatus: tr.HTTPStatus, Fields: tr.Fields,
}
}
return out, nil
}
func (a *storeAPI) RecentNudges(ctx context.Context, n int) ([]Nudge, error) {
ns, err := a.s.RecentNudges(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Nudge, len(ns))
for i, ng := range ns {
out[i] = toNudge(ng)
}
return out, nil
}
func (a *storeAPI) DeliveryAttempts(ctx context.Context, status string, n int) ([]DeliveryAttempt, error) {
as, err := a.s.ListDeliveryAttempts(ctx, status, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]DeliveryAttempt, len(as))
for i, at := range as {
out[i] = DeliveryAttempt{
ID: at.ID, Kind: at.Kind, Rule: at.Rule, ReminderID: at.ReminderID,
Channel: at.Channel, Status: at.Status, Created: at.Created,
}
if at.HasComplete {
t := at.Completed
out[i].Completed = &t
}
}
return out, nil
}
func (a *storeAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
id, err := a.s.WriteNote(ctx, ts, text, embedding, source)
return id, mapErr(err)
}
func (a *storeAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]Note, error) {
ns, err := a.s.QueryNotes(ctx, embedding, k)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, n := range ns {
out[i] = toNote(n)
}
return out, nil
}
func (a *storeAPI) RecentNotesFromSource(ctx context.Context, prefix string, n int) ([]Note, error) {
ns, err := a.s.RecentNotesFromSource(ctx, prefix, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, note := range ns {
out[i] = toNote(note)
}
return out, nil
}
func (a *storeAPI) RecentNotes(ctx context.Context, n int) ([]Note, error) {
ns, err := a.s.RecentNotes(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, note := range ns {
out[i] = toNote(note)
}
return out, nil
}
func (a *storeAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
ok, err := a.s.ProposeTool(ctx, name, utterance, scope, ts)
return ok, mapErr(err)
}
func (a *storeAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
return mapErr(a.s.EnableTool(ctx, name, cmd, destructive, scope, ts))
}
func (a *storeAPI) DisableTool(ctx context.Context, name string) error {
return mapErr(a.s.DisableTool(ctx, name))
}
func (a *storeAPI) LookupTool(ctx context.Context, name string) (Tool, error) {
t, err := a.s.LookupTool(ctx, name)
if err != nil {
return Tool{}, mapErr(err)
}
return toTool(t), nil
}
func (a *storeAPI) RevertFact(ctx context.Context, key string) (int64, error) {
_, newID, err := a.s.VoidLatestFact(ctx, key, "feedback", time.Now())
return newID, mapErr(err)
}
func (a *storeAPI) Chat(ctx context.Context, conversation, text string) (string, error) {
return "", errors.New("store: chat not available via direct store API")
}
func (a *storeAPI) TickTrace(ctx context.Context) (TickTrace, error) {
return TickTrace{}, errors.New("store: tick trace not available via direct store API")
}
func (a *storeAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
return nil, errors.New("store: morning status not available via direct store API")
}
// RecentEvents — same shape as TickTrace: the intake journal is a bounded ring
// in the daemon's memory, not a table, so a bare store cannot serve it.
func (a *storeAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, errors.New("store: intake events not available via direct store API")
}
func (a *storeAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, nil // no manager behind a bare store: nothing configured
}
func (a *storeAPI) DayPlan(ctx context.Context) (DayPlan, error) {
return DayPlan{}, errors.New("store: day plan not available via direct store API")
}
func (a *storeAPI) ListTools(ctx context.Context, status string) ([]Tool, error) {
ts, err := a.s.ListTools(ctx, status)
if err != nil {
return nil, mapErr(err)
}
out := make([]Tool, len(ts))
for i, t := range ts {
out[i] = toTool(t)
}
return out, nil
}
func (a *storeAPI) DeleteTool(ctx context.Context, name string) error {
return mapErr(a.s.DeleteTool(ctx, name))
}
func (a *storeAPI) CaptureTask(ctx context.Context, req CaptureTaskReq) (CaptureTaskResp, error) {
res, err := a.s.CaptureTask(ctx, store.Task{
CreatedTs: req.Ts,
Text: req.Text,
Source: req.Source,
Evidence: req.Evidence,
ExternalID: req.ExternalID,
Status: req.Status,
Due: req.Due,
Weight: req.Weight,
})
if err != nil {
return CaptureTaskResp{}, mapErr(err)
}
return CaptureTaskResp{ID: res.ID, Created: res.Created, Promoted: res.Promoted}, nil
}
func (a *storeAPI) ListTasks(ctx context.Context, status string) ([]Task, error) {
ts, err := a.s.ListTasks(ctx, status)
if err != nil {
return nil, mapErr(err)
}
out := make([]Task, len(ts))
for i, t := range ts {
out[i] = Task{
ID: t.ID,
CreatedTs: t.CreatedTs,
Text: t.Text,
Source: t.Source,
Evidence: t.Evidence,
ExternalID: t.ExternalID,
Status: t.Status,
Due: t.Due,
Weight: t.Weight,
Resolved: t.ResolvedTs,
ResolvedBy: t.ResolvedBy,
}
}
return out, nil
}
func (a *storeAPI) SetTaskStatus(ctx context.Context, id int64, status string, ts time.Time, by string) error {
return mapErr(a.s.SetTaskStatus(ctx, id, status, ts, by))
}
func (a *storeAPI) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, error) {
rs, err := a.s.ListProposedRoutines(ctx)
if err != nil {
return nil, mapErr(err)
}
out := make([]ProposedRoutine, len(rs))
for i, r := range rs {
out[i] = ProposedRoutine{
ID: r.ID,
Action: r.Action,
Object: r.Object,
IntervalDays: r.IntervalDays,
Status: r.Status,
CreatedTs: r.CreatedTs.UnixMilli(),
}
if r.ReminderID != nil {
out[i].ReminderID = r.ReminderID
}
}
return out, nil
}
func (a *storeAPI) DismissProposedRoutine(ctx context.Context, id int64) error {
return mapErr(a.s.DismissProposedRoutine(ctx, id))
}
func (a *storeAPI) AcceptProposedRoutine(ctx context.Context, id int64) error {
return mapErr(a.s.AcceptProposedRoutine(ctx, id, time.Now().UTC()))
}
func toTool(t store.Tool) Tool {
return Tool{
Name: t.Name, Scope: t.Scope, Cmd: t.Cmd, Destructive: t.Destructive,
Status: t.Status, Utterance: t.Utterance, Created: t.CreatedTs, Updated: t.UpdatedTs,
}
}
func toReminder(r store.Reminder) Reminder {
return Reminder{
ID: r.ID,
CreatedTs: r.CreatedTs,
FireTs: r.FireTs,
NextFireTs: r.NextFireTs,
Payload: r.Payload,
Status: r.Status,
Cron: r.Cron,
}
}
func toNote(n store.Note) Note {
return Note{ID: n.ID, Ts: n.Ts, Text: n.Text, Source: n.Source, Score: n.Score}
}
func toNudge(n store.Nudge) Nudge {
out := Nudge{
ID: n.ID, Ts: n.Ts, Rule: n.Rule, Channel: n.Channel,
Message: n.Message, Outcome: n.Outcome,
}
if n.OutcomeTs.Valid {
v := n.OutcomeTs.Int64
out.OutcomeTs = &v
}
return out
}
func toFact(f store.Fact) Fact {
out := Fact{
ID: f.ID,
Ts: f.Ts,
Kind: string(f.Kind),
Key: f.Key,
Value: f.Value,
Source: f.Source,
Confidence: f.Confidence,
}
if f.VoidsID.Valid {
v := f.VoidsID.Int64
out.VoidsID = &v
}
return out
}
// mapErr — store sentinel ↔ ipc sentinel. An unrecognized store error is
// wrapped but not mapped (server-side dispatch surfaces it as codeInternal,
// keeping internal text off the wire except to the daemon log).
func mapErr(err error) error {
if err == nil {
return nil
}
switch {
case errors.Is(err, store.ErrNoFact):
return ErrNoFact
case errors.Is(err, store.ErrConfidence):
return ErrConfidence
case errors.Is(err, store.ErrVoidsMissing):
return ErrVoidsMissing
case errors.Is(err, store.ErrNudgeNotFound):
return ErrNudgeNotFound
case errors.Is(err, store.ErrNudgeOutcome):
return ErrNudgeOutcome
case errors.Is(err, store.ErrReminderNotFound):
return ErrReminderNotFound
case errors.Is(err, store.ErrReminderState):
return ErrReminderState
case errors.Is(err, store.ErrToolNotFound):
return ErrToolNotFound
}
return err
}
// Server — the core side of the boundary. Listens on a unix domain socket,
// accepts module connections, frames requests to a CoreAPI and responses back.
// One Server per daemon process; concurrent connections are handled in their
+455
View File
@@ -0,0 +1,455 @@
// ipc/storeapi.go — the sqlite-backed CoreAPI.
//
// Split out of server.go, move-only (Vikunja #423). server.go was two
// unrelated things: this adapter, and the dispatcher that calls it over the
// socket. Nothing here knows there is a wire.
package ipc
import (
"context"
"database/sql"
"errors"
"fmt"
"time"
"github.com/kami/maven/internal/store"
)
// storeAPI — adapts *store.Store to CoreAPI. The daemon constructs one of
// these inside the core process; the socket Server calls it through the
// CoreAPI interface, so over-the-wire and in-process callers behave
// identically. The translation here is the only place store sentinels cross
// the wire: store.ErrNoFact becomes ipc.ErrNoFact, etc. — keeping the module
// view of errors stable regardless of transport.
type storeAPI struct {
s *store.Store
}
// NewStoreAPI wraps a *store.Store as a CoreAPI. The store is the sqlcipher-
// unlocked handle held ONLY in core's address space; this adapter never
// returns it to a caller — core mediates.
func NewStoreAPI(s *store.Store) CoreAPI { return &storeAPI{s: s} }
func (a *storeAPI) WriteFact(ctx context.Context, req WriteFactReq) (int64, error) {
var voids sql.NullInt64
if req.VoidsID != nil {
voids = sql.NullInt64{Int64: *req.VoidsID, Valid: true}
}
id, err := a.s.WriteFactAboutSubject(ctx, req.Ts, store.FactKind(req.Kind), req.Key, req.Subject, req.Value, req.Source, req.Confidence, voids)
return id, mapErr(err)
}
func (a *storeAPI) LatestFact(ctx context.Context, key string) (Fact, error) {
f, err := a.s.LatestFact(ctx, key)
if err != nil {
return Fact{}, mapErr(err)
}
return toFact(f), nil
}
func (a *storeAPI) LatestFactBySource(ctx context.Context, key, source string) (Fact, error) {
f, err := a.s.LatestFactBySource(ctx, key, source)
if err != nil {
return Fact{}, mapErr(err)
}
return toFact(f), nil
}
func (a *storeAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
d, err := a.s.Since(ctx, key, now)
return d, mapErr(err)
}
func (a *storeAPI) Presence(ctx context.Context) (Presence, error) {
b, score, upd, err := a.s.LoadPresenceState(ctx)
if err != nil {
return Presence{}, fmt.Errorf("ipc: load presence: %w", err)
}
return Presence{Bucket: Bucket(b), Score: score, Updated: upd}, nil
}
func (a *storeAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
id, err := a.s.CreateReminder(ctx, fire, payload, cron)
return id, mapErr(err)
}
func (a *storeAPI) MarkReminder(ctx context.Context, id int64, status string) error {
return mapErr(a.s.MarkReminder(ctx, id, status))
}
func (a *storeAPI) ListReminders(ctx context.Context, n int) ([]Reminder, error) {
rs, err := a.s.ListReminders(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Reminder, len(rs))
for i, r := range rs {
out[i] = toReminder(r)
}
return out, nil
}
func (a *storeAPI) RescheduleReminder(ctx context.Context, id int64, now time.Time) error {
return mapErr(a.s.RescheduleReminder(ctx, id, now))
}
func (a *storeAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
id, err := a.s.RecordNudge(ctx, rule, channel, message, ts)
return id, mapErr(err)
}
func (a *storeAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
return mapErr(a.s.ResolveNudge(ctx, id, outcome, ts))
}
func (a *storeAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
out, err := a.s.RecentOutcomes(ctx, rule, n)
return out, mapErr(err)
}
func (a *storeAPI) RecentFacts(ctx context.Context, n int) ([]Fact, error) {
fs, err := a.s.RecentFacts(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) RecentActiveFactsByKind(ctx context.Context, kind string, n int) ([]Fact, error) {
fs, err := a.s.RecentActiveFactsByKind(ctx, store.FactKind(kind), n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]Fact, error) {
fs, err := a.s.CalendarEvents(ctx, from, to)
if err != nil {
return nil, mapErr(err)
}
out := make([]Fact, len(fs))
for i, f := range fs {
out[i] = toFact(f)
}
return out, nil
}
func (a *storeAPI) RecentEcosystemTraces(ctx context.Context, n int) ([]EcosystemTrace, error) {
trs, err := a.s.RecentEcosystemTraces(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]EcosystemTrace, len(trs))
for i, tr := range trs {
out[i] = EcosystemTrace{
ID: tr.ID, Ts: tr.Ts, Service: tr.Service, Operation: tr.Operation,
Status: tr.Status, DurationMs: tr.DurationMs, CorrelationID: tr.CorrelationID,
CausationID: tr.CausationID, HTTPStatus: tr.HTTPStatus, Fields: tr.Fields,
}
}
return out, nil
}
func (a *storeAPI) RecentNudges(ctx context.Context, n int) ([]Nudge, error) {
ns, err := a.s.RecentNudges(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Nudge, len(ns))
for i, ng := range ns {
out[i] = toNudge(ng)
}
return out, nil
}
func (a *storeAPI) DeliveryAttempts(ctx context.Context, status string, n int) ([]DeliveryAttempt, error) {
as, err := a.s.ListDeliveryAttempts(ctx, status, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]DeliveryAttempt, len(as))
for i, at := range as {
out[i] = DeliveryAttempt{
ID: at.ID, Kind: at.Kind, Rule: at.Rule, ReminderID: at.ReminderID,
Channel: at.Channel, Status: at.Status, Created: at.Created,
}
if at.HasComplete {
t := at.Completed
out[i].Completed = &t
}
}
return out, nil
}
func (a *storeAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
id, err := a.s.WriteNote(ctx, ts, text, embedding, source)
return id, mapErr(err)
}
func (a *storeAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]Note, error) {
ns, err := a.s.QueryNotes(ctx, embedding, k)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, n := range ns {
out[i] = toNote(n)
}
return out, nil
}
func (a *storeAPI) RecentNotesFromSource(ctx context.Context, prefix string, n int) ([]Note, error) {
ns, err := a.s.RecentNotesFromSource(ctx, prefix, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, note := range ns {
out[i] = toNote(note)
}
return out, nil
}
func (a *storeAPI) RecentNotes(ctx context.Context, n int) ([]Note, error) {
ns, err := a.s.RecentNotes(ctx, n)
if err != nil {
return nil, mapErr(err)
}
out := make([]Note, len(ns))
for i, note := range ns {
out[i] = toNote(note)
}
return out, nil
}
func (a *storeAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
ok, err := a.s.ProposeTool(ctx, name, utterance, scope, ts)
return ok, mapErr(err)
}
func (a *storeAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
return mapErr(a.s.EnableTool(ctx, name, cmd, destructive, scope, ts))
}
func (a *storeAPI) DisableTool(ctx context.Context, name string) error {
return mapErr(a.s.DisableTool(ctx, name))
}
func (a *storeAPI) LookupTool(ctx context.Context, name string) (Tool, error) {
t, err := a.s.LookupTool(ctx, name)
if err != nil {
return Tool{}, mapErr(err)
}
return toTool(t), nil
}
func (a *storeAPI) RevertFact(ctx context.Context, key string) (int64, error) {
_, newID, err := a.s.VoidLatestFact(ctx, key, "feedback", time.Now())
return newID, mapErr(err)
}
func (a *storeAPI) Chat(ctx context.Context, conversation, text string) (string, error) {
return "", errors.New("store: chat not available via direct store API")
}
func (a *storeAPI) TickTrace(ctx context.Context) (TickTrace, error) {
return TickTrace{}, errors.New("store: tick trace not available via direct store API")
}
func (a *storeAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
return nil, errors.New("store: morning status not available via direct store API")
}
// RecentEvents — same shape as TickTrace: the intake journal is a bounded ring
// in the daemon's memory, not a table, so a bare store cannot serve it.
func (a *storeAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, errors.New("store: intake events not available via direct store API")
}
func (a *storeAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, nil // no manager behind a bare store: nothing configured
}
func (a *storeAPI) DayPlan(ctx context.Context) (DayPlan, error) {
return DayPlan{}, errors.New("store: day plan not available via direct store API")
}
func (a *storeAPI) ListTools(ctx context.Context, status string) ([]Tool, error) {
ts, err := a.s.ListTools(ctx, status)
if err != nil {
return nil, mapErr(err)
}
out := make([]Tool, len(ts))
for i, t := range ts {
out[i] = toTool(t)
}
return out, nil
}
func (a *storeAPI) DeleteTool(ctx context.Context, name string) error {
return mapErr(a.s.DeleteTool(ctx, name))
}
func (a *storeAPI) CaptureTask(ctx context.Context, req CaptureTaskReq) (CaptureTaskResp, error) {
res, err := a.s.CaptureTask(ctx, store.Task{
CreatedTs: req.Ts,
Text: req.Text,
Source: req.Source,
Evidence: req.Evidence,
ExternalID: req.ExternalID,
Status: req.Status,
Due: req.Due,
Weight: req.Weight,
})
if err != nil {
return CaptureTaskResp{}, mapErr(err)
}
return CaptureTaskResp{ID: res.ID, Created: res.Created, Promoted: res.Promoted}, nil
}
func (a *storeAPI) ListTasks(ctx context.Context, status string) ([]Task, error) {
ts, err := a.s.ListTasks(ctx, status)
if err != nil {
return nil, mapErr(err)
}
out := make([]Task, len(ts))
for i, t := range ts {
out[i] = Task{
ID: t.ID,
CreatedTs: t.CreatedTs,
Text: t.Text,
Source: t.Source,
Evidence: t.Evidence,
ExternalID: t.ExternalID,
Status: t.Status,
Due: t.Due,
Weight: t.Weight,
Resolved: t.ResolvedTs,
ResolvedBy: t.ResolvedBy,
}
}
return out, nil
}
func (a *storeAPI) SetTaskStatus(ctx context.Context, id int64, status string, ts time.Time, by string) error {
return mapErr(a.s.SetTaskStatus(ctx, id, status, ts, by))
}
func (a *storeAPI) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, error) {
rs, err := a.s.ListProposedRoutines(ctx)
if err != nil {
return nil, mapErr(err)
}
out := make([]ProposedRoutine, len(rs))
for i, r := range rs {
out[i] = ProposedRoutine{
ID: r.ID,
Action: r.Action,
Object: r.Object,
IntervalDays: r.IntervalDays,
Status: string(r.Status),
CreatedTs: r.CreatedTs.UnixMilli(),
}
if r.ReminderID != nil {
out[i].ReminderID = r.ReminderID
}
}
return out, nil
}
func (a *storeAPI) DismissProposedRoutine(ctx context.Context, id int64) error {
return mapErr(a.s.DismissProposedRoutine(ctx, id))
}
func (a *storeAPI) AcceptProposedRoutine(ctx context.Context, id int64) error {
return mapErr(a.s.AcceptProposedRoutine(ctx, id, time.Now().UTC()))
}
func toTool(t store.Tool) Tool {
return Tool{
Name: t.Name, Scope: t.Scope, Cmd: t.Cmd, Destructive: t.Destructive,
Status: t.Status, Utterance: t.Utterance, Created: t.CreatedTs, Updated: t.UpdatedTs,
}
}
func toReminder(r store.Reminder) Reminder {
return Reminder{
ID: r.ID,
CreatedTs: r.CreatedTs,
FireTs: r.FireTs,
NextFireTs: r.NextFireTs,
Payload: r.Payload,
Status: r.Status,
Cron: r.Cron,
}
}
func toNote(n store.Note) Note {
return Note{ID: n.ID, Ts: n.Ts, Text: n.Text, Source: n.Source, Score: n.Score}
}
func toNudge(n store.Nudge) Nudge {
out := Nudge{
ID: n.ID, Ts: n.Ts, Rule: n.Rule, Channel: n.Channel,
Message: n.Message, Outcome: n.Outcome,
}
if n.OutcomeTs.Valid {
v := n.OutcomeTs.Int64
out.OutcomeTs = &v
}
return out
}
func toFact(f store.Fact) Fact {
out := Fact{
ID: f.ID,
Ts: f.Ts,
Kind: string(f.Kind),
Key: f.Key,
Value: f.Value,
Source: f.Source,
Confidence: f.Confidence,
}
if f.VoidsID.Valid {
v := f.VoidsID.Int64
out.VoidsID = &v
}
return out
}
// mapErr — store sentinel ↔ ipc sentinel. An unrecognized store error is
// wrapped but not mapped (server-side dispatch surfaces it as codeInternal,
// keeping internal text off the wire except to the daemon log).
func mapErr(err error) error {
if err == nil {
return nil
}
switch {
case errors.Is(err, store.ErrNoFact):
return ErrNoFact
case errors.Is(err, store.ErrConfidence):
return ErrConfidence
case errors.Is(err, store.ErrVoidsMissing):
return ErrVoidsMissing
case errors.Is(err, store.ErrNudgeNotFound):
return ErrNudgeNotFound
case errors.Is(err, store.ErrNudgeOutcome):
return ErrNudgeOutcome
case errors.Is(err, store.ErrReminderNotFound):
return ErrReminderNotFound
case errors.Is(err, store.ErrReminderState):
return ErrReminderState
case errors.Is(err, store.ErrToolNotFound):
return ErrToolNotFound
}
return err
}
+47 -8
View File
@@ -8,14 +8,37 @@ import (
"time"
)
// The three states a proposal can be in. A proposal starts 'proposed' and
// moves once, either way, and never moves again.
// RoutineStatus — the state a proposal is in. A defined type, not a bare
// string, because the legal set used to live in a comment: nothing caught a
// typo at compile time, nothing enumerated the set for a test, and a bad value
// surfaced as a /routines row that neither accepts nor dismisses (Vikunja #46,
// #410).
type RoutineStatus string
// The three states a proposal can be in. A proposal starts proposed and moves
// once, either way, and never moves again.
const (
RoutineProposed = "proposed"
RoutineAccepted = "accepted"
RoutineDismissed = "dismissed"
RoutineProposed RoutineStatus = "proposed"
RoutineAccepted RoutineStatus = "accepted"
RoutineDismissed RoutineStatus = "dismissed"
)
// RoutineStatuses — the legal set, and the single source of truth a test can
// range over. Adding a state means adding it here.
var RoutineStatuses = []RoutineStatus{RoutineProposed, RoutineAccepted, RoutineDismissed}
// Valid reports whether s is one of RoutineStatuses.
func (s RoutineStatus) Valid() bool {
for _, v := range RoutineStatuses {
if s == v {
return true
}
}
return false
}
func (s RoutineStatus) String() string { return string(s) }
// ProposedRoutine — a detected pattern the system wants to nudge about on a
// repeating interval. Status 'proposed' means awaiting human confirmation;
// 'accepted' means the human confirmed and the tick loop now owns the schedule;
@@ -30,7 +53,7 @@ type ProposedRoutine struct {
Action string
Object string
IntervalDays float64
Status string // proposed | accepted | dismissed
Status RoutineStatus
CreatedTs time.Time
ReminderID *int64
AcceptedTs *time.Time
@@ -40,6 +63,7 @@ type ProposedRoutine struct {
var (
ErrProposedRoutineNotFound = errors.New("store: proposed routine not found")
ErrProposedRoutineExists = errors.New("store: proposed routine already exists for this action+object")
ErrRoutineStatus = errors.New("store: unknown routine status")
)
// CreateProposedRoutine inserts a new proposed routine. Returns
@@ -51,6 +75,16 @@ var (
// keep finding the pattern, and every re-propose is refused here. Maven is not
// a nag.
//
// The object stays a local string. It is not resolved against Nexus and it
// carries no canonical entity ref (asked on the PR 4 review, decided here,
// Vikunja #410). Nexus owns identity for things the ecosystem acts on, and
// nothing acts on a routine object: it is the word he used, replayed back to
// him in a nudge, and compared only against itself for the UNIQUE key. Two
// spellings of the same watering can are two routines, and that is the right
// answer when the point is to say the sentence he would say. Canonical refs
// arrive here only if a routine ever drives a Hexis call, which is Vikunja
// #272, not this.
//
// Vikunja #43: this is called both from the voice fact-write path (for the
// immediate spoken confirmation) and from the digestion tick's proactive
// scan (cmd/mavend/tick.go's detectPatterns, via patterns.go's
@@ -102,8 +136,13 @@ func (s *Store) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, er
}
// ListProposedRoutinesByStatus returns routines in one status, newest first.
// An empty status returns every row.
func (s *Store) ListProposedRoutinesByStatus(ctx context.Context, status string) ([]ProposedRoutine, error) {
// An empty status returns every row; an unknown one is refused rather than
// silently answering with nothing, since a typo and a genuinely empty state
// read the same otherwise.
func (s *Store) ListProposedRoutinesByStatus(ctx context.Context, status RoutineStatus) ([]ProposedRoutine, error) {
if status != "" && !status.Valid() {
return nil, fmt.Errorf("%w: %q", ErrRoutineStatus, status)
}
q := `SELECT id, action, object, interval_days, status, created_ts, reminder_id, accepted_ts, last_fired_ts
FROM proposed_routines`
var args []any
+57 -1
View File
@@ -235,7 +235,7 @@ func TestListProposedRoutinesByStatus(t *testing.T) {
}
cases := []struct {
status string
status RoutineStatus
want int
}{
{RoutineProposed, 0},
@@ -266,3 +266,59 @@ func TestLookupMissingProposedRoutine(t *testing.T) {
t.Fatal("want nil for missing routine")
}
}
// Every legal status must survive the database. The status column is written
// by three different UPDATE statements with the value spelled inline, so a
// constant that drifts from its SQL is exactly the failure this catches: the
// row would come back in a state no Go code compares equal to.
func TestRoutineStatusRoundTrip(t *testing.T) {
ctx := context.Background()
s := newTestStore(t)
now := time.Now().UTC()
write := map[RoutineStatus]func(id int64) error{
RoutineProposed: func(int64) error { return nil },
RoutineAccepted: func(id int64) error { return s.AcceptProposedRoutine(ctx, id, now) },
RoutineDismissed: func(id int64) error { return s.DismissProposedRoutine(ctx, id) },
}
for _, want := range RoutineStatuses {
if !want.Valid() {
t.Fatalf("%q is in RoutineStatuses but not Valid()", want)
}
object := "object-" + want.String()
id, err := s.CreateProposedRoutine(ctx, "полить", object, 3, now)
if err != nil {
t.Fatalf("CreateProposedRoutine(%s): %v", want, err)
}
if err := write[want](id); err != nil {
t.Fatalf("move to %s: %v", want, err)
}
got, err := s.LookupProposedRoutine(ctx, "полить", object)
if err != nil || got == nil {
t.Fatalf("LookupProposedRoutine(%s): %v", want, err)
}
if got.Status != want {
t.Errorf("status = %q, want %q", got.Status, want)
}
list, err := s.ListProposedRoutinesByStatus(ctx, want)
if err != nil {
t.Fatalf("ListProposedRoutinesByStatus(%s): %v", want, err)
}
if len(list) != 1 {
t.Errorf("status %s: listed %d rows, want 1", want, len(list))
}
}
}
// A typo used to read as "nothing is in that state", which is the same answer
// a correct query gives on an empty table.
func TestListByStatusRefusesAnUnknownStatus(t *testing.T) {
s := newTestStore(t)
if _, err := s.ListProposedRoutinesByStatus(context.Background(), "accpeted"); !errors.Is(err, ErrRoutineStatus) {
t.Fatalf("err = %v, want ErrRoutineStatus", err)
}
if RoutineStatus("accpeted").Valid() {
t.Error("a typo must not be Valid()")
}
}