20184874b2
Backlog item #3 (20-07-2026-BACKLOG.md). A morning routine is a checklist for a daily window: several items, each evidenced by a fact key, completed in any order, checked once near the end of the window. Modelling it as four independent reminder timers would stack into exactly the kind of noise Maven is supposed not to produce, so the engine nags at most once per day per routine and only for what is actually still missing. internal/morning follows the established pure-engine pattern (loop, routine, pattern): no store, no clock of its own. Evaluate answers "what's still missing" at any point; Due decides whether to nag. The impurity — reading facts under the store lock, holding the last-nudge map across ticks — stays in the tick driver, which calls Due each tick exactly as it does for loop.Rule and routine.Routine. Completion evidence is a fact key's latest non-voided value timestamped inside today's window, so manual ("выпил воды", voice-tapped) and inferred (another daemon writing the same key) are indistinguishable and both count. Weekdays scopes which days a routine applies to, so weekday/weekend variants are two routine rows rather than a special case in the engine. Exposed read-only: a MorningStatus RPC over ipc, and a /morning page in mavweb built on the same server-rendered shape as /trace — no live-update loop, since checklist state moves on the scale of minutes. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01X5JApcrCRVGmqrxnhynSik
469 lines
16 KiB
Go
469 lines
16 KiB
Go
package ipc
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Client — the module side of the boundary. Wraps a unix-socket connection
|
|
// and satisfies CoreAPI, so a module imports ipc, holds a CoreAPI, and is
|
|
// agnostic to whether it's been wired in-process (tests / daemon-embedded)
|
|
// or over this socket (full topology). The swappability is the seam auth
|
|
// will insert into without touching module code.
|
|
//
|
|
// One Client ⇒ one conn ⇒ one concurrent request at a time. A module that
|
|
// wants parallel requests opens one Client per goroutine; the store is the
|
|
// bottleneck anyway (single writer), so pipelining buys nothing here and a
|
|
// per-Client lock keeps frame interleaving impossible by construction.
|
|
type Client struct {
|
|
conn net.Conn
|
|
path string // kept so a dropped conn can be re-dialed (core restart)
|
|
mu sync.Mutex
|
|
}
|
|
|
|
// errWriteLost marks a conn drop while sending the request frame: the request
|
|
// never reached the server (or the server never saw a complete frame), so
|
|
// retrying is always safe regardless of method — nothing was applied to
|
|
// retry twice.
|
|
var errWriteLost = errors.New("ipc: connection lost before request sent")
|
|
|
|
// errReadLost marks a conn drop while waiting for the reply: the request was
|
|
// sent and may have already been applied server-side before the connection
|
|
// died (core restart mid-request, crash after commit but before reply, etc).
|
|
// Retrying here can double-apply a mutation, which ECOSYSTEM-SPEC.md's
|
|
// "never retry an unknown outcome" invariant forbids. call() only auto-retries
|
|
// this for read-only methods (idempotent by construction); a mutation method
|
|
// returns ErrAmbiguousOutcome instead so the caller can decide — the outcome
|
|
// genuinely is unknown, not safely retryable and not safely reported as failed.
|
|
var errReadLost = errors.New("ipc: connection lost awaiting reply")
|
|
|
|
// ErrAmbiguousOutcome is returned when a mutation's request may or may not
|
|
// have been applied server-side (the connection dropped after the request was
|
|
// sent, before the reply arrived). Callers must not blindly retry — the retry
|
|
// itself could double-apply. Surface this to the user/operator rather than
|
|
// silently treating it as either success or failure.
|
|
var ErrAmbiguousOutcome = errors.New("ipc: mutation outcome unknown (connection lost awaiting reply)")
|
|
|
|
// readOnlyMethods are safe to retry on an ambiguous (post-send) connection
|
|
// loss: replaying a read cannot double-apply anything. Every method not
|
|
// listed here is treated as a mutation for retry purposes — being
|
|
// conservative (refusing to retry) is the safe default for a method added
|
|
// here by omission.
|
|
var readOnlyMethods = map[Method]bool{
|
|
MethodLatestFact: true,
|
|
MethodLatestFactBySource: true,
|
|
MethodSince: true,
|
|
MethodPresence: true,
|
|
MethodListReminders: true,
|
|
MethodRecentOutcomes: true,
|
|
MethodRecentFacts: true,
|
|
MethodCalendarEvents: true,
|
|
MethodRecentNudges: true,
|
|
MethodQueryNotes: true,
|
|
MethodRecentNotes: true,
|
|
MethodLookupTool: true,
|
|
MethodListTools: true,
|
|
MethodListProposedRoutines: true,
|
|
MethodTickTrace: true,
|
|
MethodMorningStatus: true,
|
|
}
|
|
|
|
// Dial connects to a core socket at path and returns a Client. The module
|
|
// owns its Client lifecycle; Close on shutdown.
|
|
func Dial(path string) (*Client, error) {
|
|
c, err := net.Dial("unix", path)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("ipc: dial %s: %w", path, err)
|
|
}
|
|
return &Client{conn: c, path: path}, nil
|
|
}
|
|
|
|
func (c *Client) Close() error {
|
|
if c.conn == nil {
|
|
return nil
|
|
}
|
|
return c.conn.Close()
|
|
}
|
|
|
|
// DialWait is Dial with patience: it retries with capped backoff until the
|
|
// socket is reachable or timeout elapses. Core loads models on boot and may
|
|
// come up after its modules (compose depends_on orders container start, not
|
|
// socket readiness), so a module that Dial'd once would crash-loop on a cold
|
|
// start. Every core-dialing module should use this instead of Dial. Mid-life
|
|
// core restarts are handled separately by the Client's own redial-on-drop.
|
|
func DialWait(path string, timeout time.Duration) (*Client, error) {
|
|
deadline := time.Now().Add(timeout)
|
|
delay := 200 * time.Millisecond
|
|
for {
|
|
c, err := Dial(path)
|
|
if err == nil {
|
|
return c, nil
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return nil, err
|
|
}
|
|
time.Sleep(delay)
|
|
if delay < 2*time.Second {
|
|
delay *= 2
|
|
}
|
|
}
|
|
}
|
|
|
|
// call — the single request/response engine. Serialized by c.mu so a frame
|
|
// and its reply always pair up; no interleaving to disambiguate. A wire
|
|
// RpcError is rehydrated into the matching package sentinel (errors.Is works
|
|
// the same as the in-process path — the boundary is transparent to callers).
|
|
func (c *Client) call(ctx context.Context, m Method, params, result any) error {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
|
|
// Honor ctx cancellation by closing the conn — a half-sent frame would
|
|
// desync the stream; tearing down is the clean recovery. A fresh Dial
|
|
// is the module's responsibility on the next call (modules are long-lived
|
|
// processes; a dropped conn is recoverable, not fatal).
|
|
select {
|
|
case <-ctx.Done():
|
|
c.drop()
|
|
return ctx.Err()
|
|
default:
|
|
}
|
|
|
|
var raw json.RawMessage
|
|
if params != nil {
|
|
b, err := json.Marshal(params)
|
|
if err != nil {
|
|
return fmt.Errorf("ipc: marshal params: %w", err)
|
|
}
|
|
raw = b
|
|
}
|
|
|
|
var resp Response
|
|
err := c.roundtrip(m, raw, &resp)
|
|
switch {
|
|
case errors.Is(err, errWriteLost):
|
|
// The request never left; a duplicate send can't double-apply.
|
|
// Redial (roundtrip re-dials on a nil conn) and retry exactly once.
|
|
err = c.roundtrip(m, raw, &resp)
|
|
case errors.Is(err, errReadLost):
|
|
if readOnlyMethods[m] {
|
|
// A duplicate read can't double-apply either — safe to replay.
|
|
err = c.roundtrip(m, raw, &resp)
|
|
} else {
|
|
// The mutation may have already committed server-side. Do not
|
|
// retry: report the ambiguity instead of guessing.
|
|
return fmt.Errorf("%w: %v", ErrAmbiguousOutcome, err)
|
|
}
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if resp.Error != nil {
|
|
return hydrate(resp.Error)
|
|
}
|
|
if result == nil {
|
|
return nil
|
|
}
|
|
// "null" body into a pointer is valid (sets the zero value); marshal a
|
|
// RawMessage directly to avoid extra encode/decode churn.
|
|
return json.Unmarshal(resp.Result, result)
|
|
}
|
|
|
|
// roundtrip sends one request and reads its reply on c.conn, lazily (re)dialing
|
|
// if the conn is nil (fresh Client or a prior drop). A write-phase (or dial)
|
|
// failure is wrapped in errWriteLost (always safe to retry); a read-phase
|
|
// failure is wrapped in errReadLost (ambiguous — call() only retries it for
|
|
// read-only methods). Either way a failed conn is dropped so the next call
|
|
// re-dials clean. Caller holds c.mu.
|
|
func (c *Client) roundtrip(m Method, raw json.RawMessage, resp *Response) error {
|
|
if c.conn == nil {
|
|
conn, err := net.Dial("unix", c.path)
|
|
if err != nil {
|
|
return fmt.Errorf("%w: dial %s: %v", errWriteLost, c.path, err)
|
|
}
|
|
c.conn = conn
|
|
}
|
|
if err := writeFrame(c.conn, Request{Method: m, Params: raw}); err != nil {
|
|
c.drop()
|
|
return fmt.Errorf("%w: %v", errWriteLost, err)
|
|
}
|
|
if err := readFrame(c.conn, resp); err != nil {
|
|
c.drop()
|
|
return fmt.Errorf("%w: %v", errReadLost, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// drop closes and forgets the current conn so the next call re-dials.
|
|
func (c *Client) drop() {
|
|
if c.conn != nil {
|
|
_ = c.conn.Close()
|
|
c.conn = nil
|
|
}
|
|
}
|
|
|
|
// hydrate rehydrates a wire RpcError into the matching package sentinel. The
|
|
// code↔sentinel table is the only place the wire "knows" about errors; keep it
|
|
// in sync with codeOf in wire.go.
|
|
func hydrate(e *RpcError) error {
|
|
switch e.Code {
|
|
case codeNoFact:
|
|
return fmt.Errorf("%w: %s", ErrNoFact, e.Message)
|
|
case codeConfidence:
|
|
return fmt.Errorf("%w: %s", ErrConfidence, e.Message)
|
|
case codeVoidsMissing:
|
|
return fmt.Errorf("%w: %s", ErrVoidsMissing, e.Message)
|
|
case codeNudgeNotFound:
|
|
return fmt.Errorf("%w: %s", ErrNudgeNotFound, e.Message)
|
|
case codeNudgeOutcome:
|
|
return fmt.Errorf("%w: %s", ErrNudgeOutcome, e.Message)
|
|
case codeReminderMissing:
|
|
return fmt.Errorf("%w: %s", ErrReminderNotFound, e.Message)
|
|
case codeReminderState:
|
|
return fmt.Errorf("%w: %s", ErrReminderState, e.Message)
|
|
case codeToolNotFound:
|
|
return fmt.Errorf("%w: %s", ErrToolNotFound, e.Message)
|
|
case codeUnknownMethod:
|
|
return fmt.Errorf("%w: %s", ErrUnknownMethod, e.Message)
|
|
case codeBadParams:
|
|
return fmt.Errorf("%w: %s", ErrBadParams, e.Message)
|
|
case codeForbidden:
|
|
return fmt.Errorf("%w: %s", ErrForbidden, e.Message)
|
|
default:
|
|
return errors.New(e.Error())
|
|
}
|
|
}
|
|
|
|
// CoreAPI implementation on *Client. Each method is a thin call() shim; the
|
|
// shape mirrors the CoreAPI interface 1:1 so the embedded-doc intent (module
|
|
// holds a CoreAPI, transport-agnostic) reads straight off the signatures.
|
|
|
|
func (c *Client) WriteFact(ctx context.Context, req WriteFactReq) (int64, error) {
|
|
var r idResp
|
|
if err := c.call(ctx, MethodWriteFact, req, &r); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.ID, nil
|
|
}
|
|
|
|
func (c *Client) LatestFact(ctx context.Context, key string) (Fact, error) {
|
|
var f Fact
|
|
if err := c.call(ctx, MethodLatestFact, keyReq{Key: key}, &f); err != nil {
|
|
return Fact{}, err
|
|
}
|
|
return f, nil
|
|
}
|
|
|
|
func (c *Client) LatestFactBySource(ctx context.Context, key, source string) (Fact, error) {
|
|
var f Fact
|
|
if err := c.call(ctx, MethodLatestFactBySource, keySourceReq{Key: key, Source: source}, &f); err != nil {
|
|
return Fact{}, err
|
|
}
|
|
return f, nil
|
|
}
|
|
|
|
func (c *Client) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
|
|
var r sinceResp
|
|
if err := c.call(ctx, MethodSince, sinceReq{Key: key, Now: now}, &r); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.Dur, nil
|
|
}
|
|
|
|
func (c *Client) Presence(ctx context.Context) (Presence, error) {
|
|
var p Presence
|
|
if err := c.call(ctx, MethodPresence, nil, &p); err != nil {
|
|
return Presence{}, err
|
|
}
|
|
return p, nil
|
|
}
|
|
|
|
func (c *Client) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
|
|
var r idResp
|
|
if err := c.call(ctx, MethodCreateReminder, createReminderReq{Fire: fire, Payload: payload, Cron: cron}, &r); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.ID, nil
|
|
}
|
|
|
|
func (c *Client) MarkReminder(ctx context.Context, id int64, status string) error {
|
|
return c.call(ctx, MethodMarkReminder, markReminderReq{ID: id, Status: status}, nil)
|
|
}
|
|
|
|
func (c *Client) ListReminders(ctx context.Context, n int) ([]Reminder, error) {
|
|
var out []Reminder
|
|
if err := c.call(ctx, MethodListReminders, nReq{N: n}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
|
|
var r idResp
|
|
if err := c.call(ctx, MethodRecordNudge, recordNudgeReq{Rule: rule, Channel: channel, Message: message, Ts: ts}, &r); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.ID, nil
|
|
}
|
|
|
|
func (c *Client) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
|
|
return c.call(ctx, MethodResolveNudge, resolveNudgeReq{ID: id, Outcome: outcome, Ts: ts}, nil)
|
|
}
|
|
|
|
func (c *Client) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
|
|
var out []string
|
|
if err := c.call(ctx, MethodRecentOutcomes, outcomesReq{Rule: rule, N: n}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) RecentFacts(ctx context.Context, n int) ([]Fact, error) {
|
|
var out []Fact
|
|
if err := c.call(ctx, MethodRecentFacts, nReq{N: n}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) CalendarEvents(ctx context.Context, from, to time.Time) ([]Fact, error) {
|
|
var out []Fact
|
|
if err := c.call(ctx, MethodCalendarEvents, calendarEventsReq{From: from, To: to}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) RecentNudges(ctx context.Context, n int) ([]Nudge, error) {
|
|
var out []Nudge
|
|
if err := c.call(ctx, MethodRecentNudges, nReq{N: n}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
|
|
var r idResp
|
|
if err := c.call(ctx, MethodWriteNote, writeNoteReq{Ts: ts, Text: text, Embedding: embedding, Source: source}, &r); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.ID, nil
|
|
}
|
|
|
|
func (c *Client) QueryNotes(ctx context.Context, embedding []float32, k int) ([]Note, error) {
|
|
var out []Note
|
|
if err := c.call(ctx, MethodQueryNotes, queryNotesReq{Embedding: embedding, K: k}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) RecentNotes(ctx context.Context, n int) ([]Note, error) {
|
|
var out []Note
|
|
if err := c.call(ctx, MethodRecentNotes, nReq{N: n}, &out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (c *Client) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
|
|
var r proposeToolResp
|
|
if err := c.call(ctx, MethodProposeTool, proposeToolReq{Name: name, Scope: scope, Utterance: utterance, Ts: ts}, &r); err != nil {
|
|
return false, err
|
|
}
|
|
return r.Proposed, nil
|
|
}
|
|
|
|
func (c *Client) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
|
|
return c.call(ctx, MethodEnableTool, enableToolReq{Name: name, Scope: scope, Cmd: cmd, Destructive: destructive, Ts: ts}, nil)
|
|
}
|
|
|
|
func (c *Client) DisableTool(ctx context.Context, name string) error {
|
|
return c.call(ctx, MethodDisableTool, disableToolReq{Name: name}, nil)
|
|
}
|
|
|
|
func (c *Client) AssertStepUp(ctx context.Context) error {
|
|
return c.call(ctx, MethodAssertStepUp, nil, nil)
|
|
}
|
|
|
|
func (c *Client) StoreEncryptionKey(ctx context.Context, publicKey []byte) error {
|
|
return c.call(ctx, MethodStoreEncryptionKey, storeEncryptionKeyReq{PublicKey: publicKey}, nil)
|
|
}
|
|
|
|
func (c *Client) Unlock(ctx context.Context, publicKey []byte) error {
|
|
return c.call(ctx, MethodUnlock, unlockReq{PublicKey: publicKey}, nil)
|
|
}
|
|
|
|
func (c *Client) LookupTool(ctx context.Context, name string) (Tool, error) {
|
|
var t Tool
|
|
if err := c.call(ctx, MethodLookupTool, lookupToolReq{Name: name}, &t); err != nil {
|
|
return Tool{}, err
|
|
}
|
|
return t, nil
|
|
}
|
|
|
|
func (c *Client) ListTools(ctx context.Context, status string) ([]Tool, error) {
|
|
var r listToolsResp
|
|
if err := c.call(ctx, MethodListTools, listToolsReq{Status: status}, &r); err != nil {
|
|
return nil, err
|
|
}
|
|
return r.Tools, nil
|
|
}
|
|
|
|
func (c *Client) DeleteTool(ctx context.Context, name string) error {
|
|
return c.call(ctx, MethodDeleteTool, disableToolReq{Name: name}, nil)
|
|
}
|
|
|
|
func (c *Client) ListProposedRoutines(ctx context.Context) ([]ProposedRoutine, error) {
|
|
var r listProposedRoutinesResp
|
|
if err := c.call(ctx, MethodListProposedRoutines, nil, &r); err != nil {
|
|
return nil, err
|
|
}
|
|
return r.Routines, nil
|
|
}
|
|
|
|
func (c *Client) DismissProposedRoutine(ctx context.Context, id int64) error {
|
|
return c.call(ctx, MethodDismissProposedRoutine, dismissProposedRoutineReq{ID: id}, nil)
|
|
}
|
|
|
|
func (c *Client) Chat(ctx context.Context, text string) (string, error) {
|
|
var r chatResp
|
|
if err := c.call(ctx, MethodChat, chatReq{Text: text}, &r); err != nil {
|
|
return "", err
|
|
}
|
|
return r.Reply, nil
|
|
}
|
|
|
|
func (c *Client) TickTrace(ctx context.Context) (TickTrace, error) {
|
|
var t TickTrace
|
|
if err := c.call(ctx, MethodTickTrace, nil, &t); err != nil {
|
|
return TickTrace{}, err
|
|
}
|
|
return t, nil
|
|
}
|
|
|
|
func (c *Client) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
|
|
var s []MorningRoutineStatus
|
|
if err := c.call(ctx, MethodMorningStatus, nil, &s); err != nil {
|
|
return nil, err
|
|
}
|
|
return s, nil
|
|
}
|
|
|
|
func (c *Client) RevertFact(ctx context.Context, key string) (int64, error) {
|
|
var result struct {
|
|
NewID int64 `json:"new_id"`
|
|
}
|
|
if err := c.call(ctx, MethodRevertFact, map[string]string{"key": key}, &result); err != nil {
|
|
return 0, err
|
|
}
|
|
return result.NewID, nil
|
|
}
|
|
|
|
// Compile-time check: *Client satisfies CoreAPI.
|
|
var _ CoreAPI = (*Client)(nil)
|