Files
Maven/internal/ipc/client.go
T
kami 20184874b2 Add the morning routine engine — a daily checklist, not four timers
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
2026-07-30 23:49:10 +04:00

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)