Files
Maven/internal/ipc/server.go
T
kami 708a69375f tasks: key derived captures by external id and record who resolved
A task extracted from mail deduped on the live-norm index only, so once he
finished it the row left the live set and the next poll of the same immutable
message re-extracted it as a fresh candidate. mavmaild is a read-only reader
and marks nothing read, so that repeats forever. Derived rows now carry an
ext_id built from the message uid and the extracted span, unique across every
status, while voice keeps live-only norm dedupe because saying an errand again
is the recurrence signal. A derived source can no longer capture straight to
open, and saying a task out loud that Maven had only proposed promotes the
candidate instead of answering that it is already in the list.

SetTaskStatus was classified AuthRead. Resolving a task is not additive, it
erases work off his list, so it is a write, and the row now records the caller
that moved it. ListTasks was unbounded. The list-query matcher claimed any
utterance with "что мне делать", including "с чем мне помочь", and the urgency
stripper matched inside words.

Found in review of #60.
2026-08-01 14:16:39 +04:00

1205 lines
40 KiB
Go

package ipc
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"net"
"os"
"sync"
"sync/atomic"
"time"
"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) 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) 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) 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) 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, 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
// own goroutine but share the single CoreAPI (and therefore the single store
// writer — store is single-connection, SetMaxOpenConns(1), so serialization is
// already guaranteed at the db; the Server adds no locking of its own).
type Server struct {
api atomic.Value // stores CoreAPI
path string
ln net.Listener
wg sync.WaitGroup
done chan struct{}
accept sync.Mutex // guards wg.Add vs Close's wg.Wait sequence
// Check — optional authorization hook. dispatch runs it BEFORE method
// dispatch, with the raw params, so the auth layer can make verdicts
// that depend on the call's shape (e.g. WriteFact's source). A non-nil
// error aborts the call; the wire code is codeForbidden when the error
// satisfies errors.Is(ErrForbidden), else codeInternal.
//
// Nil ⇒ today's auth floor: any same-uid caller (the 0600 socket perms)
// is authorized, identical to pre-auth behavior. The daemon sets this to
// auth.Gate.Check once the auth layer is constructed; there is no module
// change to gain or lose the seam.
Check CheckFunc
// StepUp — optional handler for MethodAssertStepUp. When a real Session
// (PasskeySession) is wired, the daemon sets this to session.Assert so a
// module (mavweb) can assert a user-verification gesture over IPC. Nil ⇒
// MethodAssertStepUp returns ErrUnknownMethod (same as pre-stepup floor).
StepUp StepUpFunc
// WrapKeyFn — wraps the in-memory store encryption key under the passkey
// PRF secret (HKDF-AESGCM) and writes the wrapped blob to disk.
// Set by the daemon; nil ⇒ MethodStoreEncryptionKey returns ErrUnknownMethod.
WrapKeyFn WrapKeyFunc
// IngestMailFn — extracts task candidates from one fetched message. Set by
// the daemon only when an email block is configured AND there is a
// llama-server to extract with; nil ⇒ MethodIngestMail returns
// ErrUnknownMethod, so a mail reader pointed at a core that is not
// configured for mail is refused rather than silently ignored.
//
// Like StepUp/WrapKeyFn/UnlockFn this bypasses CoreAPI: it is not a store
// operation, it needs the resident model, and it must not become a method
// every CoreAPI implementation has to carry.
IngestMailFn IngestMailFunc
// SwapModelFn / ModelStatusFn — the on-the-fly resident model swap (Vikunja
// #250) and its read side. Set by the daemon only when phraser.swap_models
// lists at least one model AND the phraser owns a llama-server; nil ⇒ both
// methods answer ErrUnknownMethod, which is what "off unless configured"
// looks like at the wire.
//
// They bypass CoreAPI for the same reason IngestMailFn does: this is not a
// store operation, it needs the daemon's llama-server, and no other CoreAPI
// implementation should have to carry it. MethodSwapModel is AuthStepUp in
// internal/auth — owner-triggered, never an act and never a timer.
SwapModelFn SwapModelFunc
ModelStatusFn ModelStatusFunc
// DescribeImageFn — looks at one image (Vikunja #252). Set by the daemon only
// when a media store is configured AND vision is enabled with a local
// endpoint; nil ⇒ MethodDescribeImage answers ErrUnknownMethod, so a surface
// cannot make Maven accept a photo by merely sending one.
//
// It bypasses CoreAPI for the same reason IngestMailFn does: it needs a blob
// store and a vision server, neither of which is a store operation, and no
// other CoreAPI implementation should have to carry it.
DescribeImageFn DescribeImageFunc
// Capture* — the meeting recorder (Vikunja #253). Set by the daemon only
// when a media store is configured AND capture.enabled is true; nil ⇒ all
// four methods answer ErrUnknownMethod. That is the load-bearing default for
// this capability: on an unconfigured box there is no wire path that begins a
// recording, so nothing can be recorded by accident, by a bug in a surface,
// or by a model deciding it would be helpful.
//
// They bypass CoreAPI because a recorder needs a blob store, an STT worker
// and a llama-server, none of which is a store operation.
CaptureStartFn CaptureStartFunc
CaptureAppendFn CaptureAppendFunc
CaptureStopFn CaptureStopFunc
CaptureStatusFn CaptureStatusFunc
// Speaker* — voice identification (Vikunja #255). Set by the daemon only
// when a speaker block is configured; nil ⇒ all three methods answer
// ErrUnknownMethod, so on an unconfigured box no wire path enrols a voice.
EnrollSpeakerFn EnrollSpeakerFunc
ListSpeakersFn ListSpeakersFunc
ForgetSpeakerFn ForgetSpeakerFunc
// UnlockFn — unwraps the store encryption key from the wrapped blob using
// the passkey PRF secret, opens the encrypted store, and wires
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
// in locked mode; nil ⇒ MethodUnlock returns ErrUnknownMethod.
UnlockFn UnlockFunc
// now is injected so tests can drive time; the loop already works in
// absolute ts supplied by callers, so this isn't load-bearing for live ops.
}
// WrapKeyFunc — wraps the store encryption key under the passkey-derived
// secret (a 32-byte WebAuthn PRF output) and persists the wrapped blob.
type WrapKeyFunc func(ctx context.Context, secret []byte) error
// UnlockFunc — unwraps the store encryption key using the passkey-derived
// secret and completes daemon initialization.
type UnlockFunc func(ctx context.Context, secret []byte) error
// SwapModelFunc — loads another resident model in place of the live one.
type SwapModelFunc func(ctx context.Context, req SwapModelReq) (SwapModelResp, error)
// ModelStatusFunc — reports the resident model and the swap allowlist.
type ModelStatusFunc func(ctx context.Context) (ModelStatusResp, error)
// IngestMailFunc — core-side mail extraction. Returns what was captured.
type IngestMailFunc func(ctx context.Context, req IngestMailReq) (IngestMailResp, error)
// DescribeImageFunc — core-side image intake + description.
type DescribeImageFunc func(ctx context.Context, req DescribeImageReq) (DescribeImageResp, error)
// CaptureStartFunc / CaptureAppendFunc / CaptureStopFunc / CaptureStatusFunc —
// the four core-side halves of the meeting recorder.
type CaptureStartFunc func(ctx context.Context, req CaptureStartReq) (CaptureStartResp, error)
type CaptureAppendFunc func(ctx context.Context, req CaptureAppendReq) (CaptureAppendResp, error)
type CaptureStopFunc func(ctx context.Context, req CaptureStopReq) (CaptureStopResp, error)
type CaptureStatusFunc func(ctx context.Context) (CaptureStatusResp, error)
// EnrollSpeakerFunc / ListSpeakersFunc / ForgetSpeakerFunc — the core-side
// halves of voice enrolment.
type EnrollSpeakerFunc func(ctx context.Context, req EnrollSpeakerReq) (EnrollSpeakerResp, error)
type ListSpeakersFunc func(ctx context.Context) (ListSpeakersResp, error)
type ForgetSpeakerFunc func(ctx context.Context, req ForgetSpeakerReq) error
// CheckFunc — the auth hook signature. Wired by the daemon (auth.Gate.Check
// satisfies this); dispatch calls it once per request after param-unmarshal
// independence (it gets the raw params, may unmarshal what it needs — ipc
// already unmarshals for the typed call separately). Keeping Check on raw
// params means ipc doesn't need to know each method's authority shape, and
// auth doesn't need to leak implementation into ipc.
type CheckFunc func(ctx context.Context, m Method, params json.RawMessage) error
// StepUpFunc — records a user-verification gesture. Set by the daemon when
// a real Session is wired (PasskeySession); nil means not available.
// MethodAssertStepUp dispatch calls this instead of going through CoreAPI.
type StepUpFunc func(ctx context.Context) error
// Listen creates a Server bound to path. path's parent dir must exist and be
// 0700 (we chmod it if we own it); the socket file itself is created 0600 so
// only the same unix user can connect — the current "auth floor", same radius
// as wg at the network boundary. Removing a stale socket at path first lets
// the daemon restart cleanly.
func Listen(path string, api CoreAPI) (*Server, error) {
_ = os.Remove(path) // stale socket from a crashed daemon; ignore missing
if err := os.MkdirAll(parentDir(path), 0o700); err != nil {
return nil, fmt.Errorf("ipc: mkdir socket dir: %w", err)
}
// umask could widen the perms on socket creation; tighten then chmod to
// be explicit. 0600 ⇒ read+write by owner only.
oldMask := unix.Umask(0o077)
ln, err := net.Listen("unix", path)
unix.Umask(oldMask)
if err != nil {
return nil, fmt.Errorf("ipc: listen %s: %w", path, err)
}
if err := os.Chmod(path, 0o600); err != nil {
_ = ln.Close()
_ = os.Remove(path)
return nil, fmt.Errorf("ipc: chmod socket: %w", err)
}
s := &Server{
path: path,
ln: ln,
done: make(chan struct{}),
}
s.api.Store(api)
return s, nil
}
// Serve accepts connections until the listener closes. Each connection is
// served in its own goroutine; a panicking handler or a malformed frame tears
// down only that conn, not the server (a misbehaving module can't kill core).
func (s *Server) Serve() error {
for {
c, err := s.ln.Accept()
if err != nil {
select {
case <-s.done:
return nil // graceful Close
default:
return fmt.Errorf("ipc: accept: %w", err)
}
}
// wg.Add under accept mutex so Close's wg.Wait (also under accept) sees
// a consistent counter — a connection accepted just before Close closes
// the listener must be tracked before Wait starts.
s.accept.Lock()
s.wg.Add(1)
s.accept.Unlock()
go func(c net.Conn) {
defer s.wg.Done()
defer c.Close()
s.serveConn(c)
}(c)
}
}
func (s *Server) serveConn(c net.Conn) {
caller, callerOK := peerCaller(c)
ctx := context.Background()
if callerOK {
ctx = WithCaller(ctx, caller)
}
for {
var req Request
if err := readFrame(c, &req); err != nil {
return // EOF or malformed ⇒ end this conn; nothing to recover
}
// redispatch expects the framework's recover so one bad call can't
// take the goroutine (and therefore the conn) with it.
result, err := s.safeDispatch(ctx, req)
resp := Response{}
if err != nil {
resp.Error = rpcErr(err)
} else {
resp.Result = result
}
if err := writeFrame(c, resp); err != nil {
return
}
}
}
func (s *Server) safeDispatch(ctx context.Context, req Request) (result json.RawMessage, err error) {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("ipc: panic dispatching %s: %v", req.Method, r)
}
}()
return s.dispatch(ctx, req)
}
// handlerFunc — one table entry's shape: unmarshal req.Params (if it wants
// any), call the matching CoreAPI method against the api passed in, marshal
// the result. api is a parameter, not a closed-over field, precisely so a
// table built once at package init never pins a stale CoreAPI — see the note
// on methodTable below about SetAPI.
type handlerFunc func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error)
// withParams adapts a (typed params, typed result) CoreAPI call into a
// handlerFunc: unmarshal into P, call fn, marshal R. On error the result is
// dropped (marshalResult's output is never read when err != nil — see
// serveConn) so every entry can uniformly return early on error without
// re-deriving what the pre-table per-arm code used to return in that case.
func withParams[P any, R any](fn func(ctx context.Context, api CoreAPI, p P) (R, error)) handlerFunc {
return func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error) {
var p P
if err := unmarshalParams(raw, &p); err != nil {
return nil, err
}
r, err := fn(ctx, api, p)
if err != nil {
return nil, err
}
return marshalResult(r), nil
}
}
// withParamsVoid is withParams for the error-only methods (mark/resolve/
// enable/disable/...): params in, no result out, wire reply is always null.
func withParamsVoid[P any](fn func(ctx context.Context, api CoreAPI, p P) error) handlerFunc {
return func(ctx context.Context, api CoreAPI, raw json.RawMessage) (json.RawMessage, error) {
var p P
if err := unmarshalParams(raw, &p); err != nil {
return nil, err
}
return marshalResult(nil), fn(ctx, api, p)
}
}
// withoutParams is withParams for the handful of methods that take no
// params at all (Presence, TickTrace, MorningStatus, ListProposedRoutines).
// It does NOT call unmarshalParams — matching the pre-table arms, which
// never touched req.Params for these four methods.
func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error)) handlerFunc {
return func(ctx context.Context, api CoreAPI, _ json.RawMessage) (json.RawMessage, error) {
r, err := fn(ctx, api)
if err != nil {
return nil, err
}
return marshalResult(r), nil
}
}
// methodTable — one entry per CoreAPI-backed method. Built once at package
// init, not per-Server and not per-dispatch: entries close over nothing but
// the CoreAPI method being called, and dispatch passes in the *current*
// api (loaded fresh via s.api.Load() every call, same as before the table
// existed) as an argument — so SetAPI's runtime swap (the unlock transition)
// is still honored on the very next request with no extra plumbing here.
//
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock,
// MethodIngestMail, MethodSwapModel, MethodModelStatus,
// MethodDescribeImage and the four MethodCapture* methods are NOT in this
// table: they bypass CoreAPI entirely (s.StepUp / s.WrapKeyFn / s.UnlockFn /
// s.IngestMailFn / s.DescribeImageFn / s.Capture*Fn), so dispatch
// special-cases them before consulting the table.
var methodTable = map[Method]handlerFunc{
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
id, err := api.WriteFact(ctx, p)
return idResp{ID: id}, err
}),
MethodLatestFact: withParams(func(ctx context.Context, api CoreAPI, p keyReq) (Fact, error) {
return api.LatestFact(ctx, p.Key)
}),
MethodLatestFactBySource: withParams(func(ctx context.Context, api CoreAPI, p keySourceReq) (Fact, error) {
return api.LatestFactBySource(ctx, p.Key, p.Source)
}),
MethodSince: withParams(func(ctx context.Context, api CoreAPI, p sinceReq) (sinceResp, error) {
d, err := api.Since(ctx, p.Key, p.Now)
return sinceResp{Dur: d}, err
}),
MethodPresence: withoutParams(func(ctx context.Context, api CoreAPI) (Presence, error) {
return api.Presence(ctx)
}),
MethodCreateReminder: withParams(func(ctx context.Context, api CoreAPI, p createReminderReq) (idResp, error) {
id, err := api.CreateReminder(ctx, p.Fire, p.Payload, p.Cron)
return idResp{ID: id}, err
}),
MethodMarkReminder: withParamsVoid(func(ctx context.Context, api CoreAPI, p markReminderReq) error {
return api.MarkReminder(ctx, p.ID, p.Status)
}),
MethodListReminders: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Reminder, error) {
out, err := api.ListReminders(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []Reminder{}
}
return out, nil
}),
MethodRecordNudge: withParams(func(ctx context.Context, api CoreAPI, p recordNudgeReq) (idResp, error) {
id, err := api.RecordNudge(ctx, p.Rule, p.Channel, p.Message, p.Ts)
return idResp{ID: id}, err
}),
MethodResolveNudge: withParamsVoid(func(ctx context.Context, api CoreAPI, p resolveNudgeReq) error {
return api.ResolveNudge(ctx, p.ID, p.Outcome, p.Ts)
}),
MethodRecentOutcomes: withParams(func(ctx context.Context, api CoreAPI, p outcomesReq) ([]string, error) {
out, err := api.RecentOutcomes(ctx, p.Rule, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []string{} // stable non-null on the wire
}
return out, nil
}),
MethodRecentFacts: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Fact, error) {
out, err := api.RecentFacts(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []Fact{}
}
return out, nil
}),
MethodCalendarEvents: withParams(func(ctx context.Context, api CoreAPI, p calendarEventsReq) ([]Fact, error) {
out, err := api.CalendarEvents(ctx, p.From, p.To)
if err != nil {
return nil, err
}
if out == nil {
out = []Fact{}
}
return out, nil
}),
MethodRecentNudges: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Nudge, error) {
out, err := api.RecentNudges(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []Nudge{}
}
return out, nil
}),
MethodWriteNote: withParams(func(ctx context.Context, api CoreAPI, p writeNoteReq) (idResp, error) {
id, err := api.WriteNote(ctx, p.Ts, p.Text, p.Embedding, p.Source)
return idResp{ID: id}, err
}),
MethodQueryNotes: withParams(func(ctx context.Context, api CoreAPI, p queryNotesReq) ([]Note, error) {
out, err := api.QueryNotes(ctx, p.Embedding, p.K)
if err != nil {
return nil, err
}
if out == nil {
out = []Note{}
}
return out, nil
}),
MethodRecentNotes: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]Note, error) {
out, err := api.RecentNotes(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []Note{}
}
return out, nil
}),
MethodProposeTool: withParams(func(ctx context.Context, api CoreAPI, p proposeToolReq) (proposeToolResp, error) {
ok, err := api.ProposeTool(ctx, p.Name, p.Utterance, p.Scope, p.Ts)
return proposeToolResp{Proposed: ok}, err
}),
MethodEnableTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p enableToolReq) error {
return api.EnableTool(ctx, p.Name, p.Cmd, p.Destructive, p.Scope, p.Ts)
}),
MethodDisableTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p disableToolReq) error {
return api.DisableTool(ctx, p.Name)
}),
MethodLookupTool: withParams(func(ctx context.Context, api CoreAPI, p lookupToolReq) (Tool, error) {
return api.LookupTool(ctx, p.Name)
}),
MethodListTools: withParams(func(ctx context.Context, api CoreAPI, p listToolsReq) (listToolsResp, error) {
out, err := api.ListTools(ctx, p.Status)
if err != nil {
return listToolsResp{}, err
}
if out == nil {
out = []Tool{}
}
return listToolsResp{Tools: out}, nil
}),
// MethodDeleteTool shares disableToolReq — both take just a tool name.
MethodDeleteTool: withParamsVoid(func(ctx context.Context, api CoreAPI, p disableToolReq) error {
return api.DeleteTool(ctx, p.Name)
}),
MethodCaptureTask: withParams(func(ctx context.Context, api CoreAPI, p CaptureTaskReq) (CaptureTaskResp, error) {
return api.CaptureTask(ctx, p)
}),
MethodListTasks: withParams(func(ctx context.Context, api CoreAPI, p listTasksReq) (listTasksResp, error) {
out, err := api.ListTasks(ctx, p.Status)
if err != nil {
return listTasksResp{}, err
}
if out == nil {
out = []Task{}
}
return listTasksResp{Tasks: out}, nil
}),
MethodSetTaskStatus: withParamsVoid(func(ctx context.Context, api CoreAPI, p setTaskStatusReq) error {
return api.SetTaskStatus(ctx, p.ID, p.Status, p.Ts, p.By)
}),
MethodListProposedRoutines: withoutParams(func(ctx context.Context, api CoreAPI) (listProposedRoutinesResp, error) {
out, err := api.ListProposedRoutines(ctx)
if err != nil {
return listProposedRoutinesResp{}, err
}
if out == nil {
out = []ProposedRoutine{}
}
return listProposedRoutinesResp{Routines: out}, nil
}),
MethodDismissProposedRoutine: withParamsVoid(func(ctx context.Context, api CoreAPI, p dismissProposedRoutineReq) error {
return api.DismissProposedRoutine(ctx, p.ID)
}),
MethodAcceptProposedRoutine: withParamsVoid(func(ctx context.Context, api CoreAPI, p acceptProposedRoutineReq) error {
return api.AcceptProposedRoutine(ctx, p.ID)
}),
MethodRevertFact: withParams(func(ctx context.Context, api CoreAPI, p revertReq) (map[string]int64, error) {
newID, err := api.RevertFact(ctx, p.Key)
if err != nil {
return nil, err
}
return map[string]int64{"new_id": newID}, nil
}),
MethodChat: withParams(func(ctx context.Context, api CoreAPI, p chatReq) (chatResp, error) {
reply, err := api.Chat(ctx, p.Text)
return chatResp{Reply: reply}, err
}),
MethodTickTrace: withoutParams(func(ctx context.Context, api CoreAPI) (TickTrace, error) {
return api.TickTrace(ctx)
}),
// MorningStatus intentionally has no nil→[]T{} normalization here — the
// pre-table arm marshaled api.MorningStatus's result as-is (a nil slice
// serializes as JSON null), and this preserves that exact wire shape.
MethodDayPlan: withoutParams(func(ctx context.Context, api CoreAPI) (DayPlan, error) {
return api.DayPlan(ctx)
}),
MethodMorningStatus: withoutParams(func(ctx context.Context, api CoreAPI) ([]MorningRoutineStatus, error) {
return api.MorningStatus(ctx)
}),
MethodRecentEvents: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]IntakeEvent, error) {
out, err := api.RecentEvents(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []IntakeEvent{}
}
return out, nil
}),
MethodMCPServers: withoutParams(func(ctx context.Context, api CoreAPI) ([]MCPServerStatus, error) {
out, err := api.MCPServers(ctx)
if err != nil {
return nil, err
}
if out == nil {
out = []MCPServerStatus{}
}
return out, nil
}),
}
// dispatch unmarshals params for req.Method and calls the matching CoreAPI
// method. Unknown method ⇒ ErrUnknownMethod; a malformed params payload ⇒
// ErrBadParams with the underlying text (local, server-side, not shipped to
// the module except as a generic message via rpcErr).
//
// Authorization runs ONCE at the top: if Server.Check is set, we call it with
// the raw params before any method-specific unmarshal; auth unmarshals fields
// it cares about (WriteFact's source, etc.) itself. A nil Check is the floor
// and is invisible at the wire — pre-auth Server behavior is unchanged.
func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, error) {
api := s.api.Load().(CoreAPI)
if s.Check != nil {
if err := s.Check(ctx, req.Method, req.Params); err != nil {
return nil, err
}
}
// These bypass CoreAPI entirely — they drive Server fields set
// directly by the daemon (StepUp / WrapKeyFn / UnlockFn), not store
// state, so they can never be table entries keyed on a CoreAPI method.
switch req.Method {
case MethodAssertStepUp:
if s.StepUp != nil {
return marshalResult(nil), s.StepUp(ctx)
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodStoreEncryptionKey:
if s.WrapKeyFn != nil {
var p storeEncryptionKeyReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
return marshalResult(nil), s.WrapKeyFn(ctx, p.Secret)
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodUnlock:
if s.UnlockFn != nil {
var p unlockReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
return marshalResult(nil), s.UnlockFn(ctx, p.Secret)
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodIngestMail:
if s.IngestMailFn != nil {
var p IngestMailReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.IngestMailFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodSwapModel:
if s.SwapModelFn != nil {
var p SwapModelReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.SwapModelFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodDescribeImage:
if s.DescribeImageFn != nil {
var p DescribeImageReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.DescribeImageFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStart:
if s.CaptureStartFn != nil {
var p CaptureStartReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureStartFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureAppend:
if s.CaptureAppendFn != nil {
var p CaptureAppendReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureAppendFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStop:
if s.CaptureStopFn != nil {
var p CaptureStopReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.CaptureStopFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodCaptureStatus:
if s.CaptureStatusFn != nil {
resp, err := s.CaptureStatusFn(ctx)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodEnrollSpeaker:
if s.EnrollSpeakerFn != nil {
var p EnrollSpeakerReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
resp, err := s.EnrollSpeakerFn(ctx, p)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodListSpeakers:
if s.ListSpeakersFn != nil {
resp, err := s.ListSpeakersFn(ctx)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodForgetSpeaker:
if s.ForgetSpeakerFn != nil {
var p ForgetSpeakerReq
if err := unmarshalParams(req.Params, &p); err != nil {
return nil, err
}
if err := s.ForgetSpeakerFn(ctx, p); err != nil {
return nil, err
}
return marshalResult(nil), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
case MethodModelStatus:
if s.ModelStatusFn != nil {
resp, err := s.ModelStatusFn(ctx)
if err != nil {
return nil, err
}
return marshalResult(resp), nil
}
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
}
h, ok := methodTable[req.Method]
if !ok {
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
}
return h(ctx, api, req.Params)
}
func unmarshalParams(raw json.RawMessage, v any) error {
if len(raw) == 0 {
raw = []byte("null")
}
if err := json.Unmarshal(raw, v); err != nil {
return fmt.Errorf("%w: %v", ErrBadParams, err)
}
return nil
}
func marshalResult(v any) json.RawMessage {
if v == nil {
return json.RawMessage("null")
}
b, _ := json.Marshal(v)
return b
}
// Close stops accepting and waits for in-flight connections to drain. The
// socket file is removed so a restart can rebind cleanly. Idempotent.
func (s *Server) Close() error {
select {
case <-s.done:
return nil
default:
close(s.done)
}
err := s.ln.Close()
// Under accept lock: after the listener closes, no new Accept can complete,
// so no new wg.Add will be called. The Wait is safe to observe the wg
// counter because any in-flight Accept that already got a conn either
// already called wg.Add (before releasing the lock) or will see the closed
// listener error and not call wg.Add at all.
s.accept.Lock()
s.wg.Wait()
s.accept.Unlock()
_ = os.Remove(s.path)
return err
}
// Path returns the filesystem path of the listening socket.
func (s *Server) Path() string { return s.path }
// SetAPI atomically replaces the CoreAPI the server dispatches to. Used by
// the daemon's unlock path: in locked mode a dummy API returns errors for all
// store methods; after unlock, the real store API is swapped in. Safe to call
// while the server is serving (dispatch loads api once per request via atomic).
func (s *Server) SetAPI(api CoreAPI) { s.api.Store(api) }
func parentDir(p string) string {
if i := lastIndexByte(p, '/'); i >= 0 {
if i == 0 {
return "/"
}
return p[:i]
}
return "."
}
func lastIndexByte(s string, b byte) int {
for i := len(s) - 1; i >= 0; i-- {
if s[i] == b {
return i
}
}
return -1
}
// peerCaller — read SO_PEERCRED off a unix conn to identify the connecting
// process. Returns ok=false on a non-unix conn or a platform without
// SO_PEERCRED; the caller then proceeds without a Caller (the socket perms
// already proved same-user). Linux only today; on other platforms this floors
// to "unknown caller" rather than failing — the wire still works.
func peerCaller(c net.Conn) (Caller, bool) {
uc, ok := c.(*net.UnixConn)
if !ok {
return Caller{}, false
}
raw, err := uc.SyscallConn()
if err != nil {
return Caller{}, false
}
var cred *unix.Ucred
ctrlErr := raw.Control(func(fd uintptr) {
cred, err = unix.GetsockoptUcred(int(fd), unix.SOL_SOCKET, unix.SO_PEERCRED)
})
if ctrlErr != nil || err != nil || cred == nil {
return Caller{}, false
}
return Caller{Uid: int32(cred.Uid), Pid: int32(cred.Pid)}, true
}