Files
Maven/cmd/mavend/main.go
T
kami 6c92f85d10 feat(ecosystem): compliant Praxis/Hexis integration + vendored build
Bring the Nexus/Praxis/Hexis integration in line with
MAVEN_ECOSYSTEM_ARCHITECTURE.md:

- Praxis over HTTP: drop the in-process praxis.db open (praxisstore/
  praxistools) and call praxisd's /api/v1/tools/* API via a new praxisClient.
  Honors the "no component reads another's DB" invariant (AC#12).
  PraxisConfig.DBPath -> URL.
- Hexis confirmation gate: mutating capabilities (ReadOnly=false) now park a
  bound pendingHexis confirmation and require a spoken "да" before executing;
  read-only run immediately (AC#7, no auto attention->action).
- Capability safety: >1 verb match is ambiguous -> ask instead of firing the
  first; ambiguous Nexus resolution asks for clarification (AC#2).
- Correlation IDs on Hexis execute, recorded in the cross-service trace.
- Bug: importance arrives as JSON float64 over HTTP, not int.
- Tests: confirm-gate, decline, read-only, and ambiguity paths.

Build: vendor/ bakes in the hexis client (replace-directed at a sibling repo
outside the Docker context); Dockerfile builds from vendor and no longer
`go mod download`s the unreachable replace paths.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-19 20:24:33 +04:00

578 lines
18 KiB
Go

// Package main is mavend — maven's daemon.
//
// "core = the only key-holder": one process holds the unlocked store + the
// trigger loop; modules are separate processes, key-free, fail-independent.
// the daemon wires Store → Gatherer → Tick → phraser → delivery, runs the 60s
// ticker, owns the cold-start unlock dance, and exposes the CoreAPI boundary
// over a unix socket for modules (router/delivery/poller/...) to call.
//
// Floor (this file): pluggable seams wired with the deterministic Stubs.
// - phraser Stub (no LLM)
// - voice sink wired via wireVoice: embedder/classifier seeded with ~10
// examples across 5 intents; stt + tts stubs in-process by default,
// remote module sockets when configured; TCP listener on voice.bind.
// The voice sink (voicesink.Sink via Sessions) is wired into the
// dispatcher — the reactive path (push-to-talk) AND proactive nudges
// (care-when-present, sev3/sev4 present) both route through the same
// stt→router→tts→client pipeline.
// - auth FloorEnrollment + nil Session — cold-start unlock assumed: today
// the store opens plain sqlite (sqlcipher deferred). the daemon runs
// "unlocked" — the locked-until-asserted dance lands with the Session
// verifier + ask-password transport (open spec item).
//
// Cold-start unlock (2026-07-06):
//
// When a passkey credential is enrolled AND no env key is set, the daemon
// starts in LOCKED mode: the IPC server runs but rejects all store methods
// except MethodAssertStepUp and MethodUnlock. A passkey assertion followed
// by MethodUnlock (with the same credential's public key) unwraps the at-rest
// AES-256 key from a wrapped blob on disk (HKDF-SHA256 + AES-GCM) and opens
// the encrypted store. After unlock, the daemon wires voice, loop, and
// delivery and runs normally.
//
// Fallback: when db_key_env is set (or no wrapped file exists), the daemon
// starts unlocked from the env key (pre-unlock behavior). Enrolling a passkey
// while unlocked calls MethodStoreEncryptionKey to wrap the env key and
// persist the wrapped blob — enabling cold-start unlock on the next boot
// after the env key is removed.
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"log"
"net"
"os"
"os/signal"
"sync"
"syscall"
"time"
"github.com/kami/maven/internal/auth"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/delivery"
"github.com/kami/maven/internal/delivery/ntfysink"
"github.com/kami/maven/internal/delivery/telegramsink"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/webauthn"
"github.com/kami/maven/internal/loop"
)
var errLocked = errors.New("mavend: daemon locked — complete passkey assertion first")
// daemonLock tracks whether the daemon is in locked (pre-unlock) mode.
// In locked mode, all CoreAPI methods return errLocked. The unlock path
// replaces the CoreAPI with the real store adapter and flips the flag.
type daemonLock struct {
mu sync.Mutex
locked bool
}
func newDaemonLock(locked bool) *daemonLock {
return &daemonLock{locked: locked}
}
func (l *daemonLock) isLocked() bool {
l.mu.Lock()
defer l.mu.Unlock()
return l.locked
}
func (l *daemonLock) unlock() {
l.mu.Lock()
defer l.mu.Unlock()
l.locked = false
}
func main() {
if err := run(os.Args[1:]); err != nil {
fmt.Fprintln(os.Stderr, "mavend:", err)
os.Exit(1)
}
}
// lockedAPI is a dummy CoreAPI used while the daemon is locked. Every method
// returns errLocked. The wire protocol's StoreAPI methods all go through the
// Server dispatch on CoreAPI, so returning errLocked from each is correct.
type lockedAPI struct{}
var _ ipc.CoreAPI = (*lockedAPI)(nil)
func (l *lockedAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
return 0, errLocked
}
func (l *lockedAPI) LatestFact(ctx context.Context, key string) (ipc.Fact, error) {
return ipc.Fact{}, errLocked
}
func (l *lockedAPI) LatestFactBySource(ctx context.Context, key, source string) (ipc.Fact, error) {
return ipc.Fact{}, errLocked
}
func (l *lockedAPI) Since(ctx context.Context, key string, now time.Time) (time.Duration, error) {
return 0, errLocked
}
func (l *lockedAPI) Presence(ctx context.Context) (ipc.Presence, error) {
return ipc.Presence{}, errLocked
}
func (l *lockedAPI) CreateReminder(ctx context.Context, fire time.Time, payload, cron string) (int64, error) {
return 0, errLocked
}
func (l *lockedAPI) MarkReminder(ctx context.Context, id int64, status string) error {
return errLocked
}
func (l *lockedAPI) ListReminders(ctx context.Context, n int) ([]ipc.Reminder, error) {
return nil, errLocked
}
func (l *lockedAPI) RecordNudge(ctx context.Context, rule, channel, message string, ts time.Time) (int64, error) {
return 0, errLocked
}
func (l *lockedAPI) ResolveNudge(ctx context.Context, id int64, outcome string, ts time.Time) error {
return errLocked
}
func (l *lockedAPI) RecentOutcomes(ctx context.Context, rule string, n int) ([]string, error) {
return nil, errLocked
}
func (l *lockedAPI) RecentFacts(ctx context.Context, n int) ([]ipc.Fact, error) {
return nil, errLocked
}
func (l *lockedAPI) CalendarEvents(ctx context.Context, from, to time.Time) ([]ipc.Fact, error) {
return nil, errLocked
}
func (l *lockedAPI) RecentNudges(ctx context.Context, n int) ([]ipc.Nudge, error) {
return nil, errLocked
}
func (l *lockedAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
return 0, errLocked
}
func (l *lockedAPI) QueryNotes(ctx context.Context, embedding []float32, k int) ([]ipc.Note, error) {
return nil, errLocked
}
func (l *lockedAPI) RecentNotes(ctx context.Context, n int) ([]ipc.Note, error) {
return nil, errLocked
}
func (l *lockedAPI) ProposeTool(ctx context.Context, name, utterance, scope string, ts time.Time) (bool, error) {
return false, errLocked
}
func (l *lockedAPI) EnableTool(ctx context.Context, name string, cmd []string, destructive bool, scope string, ts time.Time) error {
return errLocked
}
func (l *lockedAPI) DisableTool(ctx context.Context, name string) error { return errLocked }
func (l *lockedAPI) DeleteTool(ctx context.Context, name string) error { return errLocked }
func (l *lockedAPI) ListProposedRoutines(ctx context.Context) ([]ipc.ProposedRoutine, error) {
return nil, errLocked
}
func (l *lockedAPI) DismissProposedRoutine(ctx context.Context, id int64) error { return errLocked }
func (l *lockedAPI) LookupTool(ctx context.Context, name string) (ipc.Tool, error) {
return ipc.Tool{}, errLocked
}
func (l *lockedAPI) ListTools(ctx context.Context, status string) ([]ipc.Tool, error) {
return nil, errLocked
}
func (l *lockedAPI) RevertFact(ctx context.Context, key string) (int64, error) { return 0, errLocked }
func (l *lockedAPI) Chat(ctx context.Context, text string) (string, error) {
return "", errLocked
}
func (l *lockedAPI) TickTrace(ctx context.Context) (ipc.TickTrace, error) {
return ipc.TickTrace{}, errLocked
}
func run(args []string) error {
cfgPath := flag.String("config", defaultConfigPath(), "path to mavend JSON config")
wrappedKeyPath := flag.String("wrapped-key-file", "", "path to wrapped encryption key blob (enables cold-start unlock)")
flag.CommandLine.Parse(args)
cfg, err := config.Load(*cfgPath)
if err != nil {
return err
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
defer stop()
// ----- encryption key: env key → wrapped key → plaintext -----
// Priority: env key (from config/docker) > wrapped key (cold-start unlock) > plaintext (dev/CI).
// When a wrapped key file exists AND no env key is set, the daemon starts
// LOCKED and waits for a passkey assertion to unwrap it.
envKey, err := cfg.DBEncryptionKey()
if err != nil {
return err
}
wrappedExists := false
if *wrappedKeyPath != "" {
if _, err := os.Stat(*wrappedKeyPath); err == nil {
wrappedExists = true
} else if !errors.Is(err, os.ErrNotExist) {
return fmt.Errorf("check wrapped key: %w", err)
}
} else {
// default path alongside the config
defaultWrapped := cfg.DefaultWrappedKeyPath()
if _, err := os.Stat(defaultWrapped); err == nil {
*wrappedKeyPath = defaultWrapped
wrappedExists = true
}
}
locked := wrappedExists && envKey == nil
dl := newDaemonLock(locked)
var st *store.Store
var envKeyBytes []byte // kept for WrapKeyFn (enrollment wraps this key)
if !locked {
// Normal boot: env key or plaintext (dev/CI)
if envKey != nil {
envKeyBytes = make([]byte, len(envKey))
copy(envKeyBytes, envKey)
st, err = store.OpenEncrypted(ctx, cfg.DBPath, cfg.DBTmpfs, envKey)
} else {
st, err = store.Open(ctx, cfg.DBPath)
}
if err != nil {
return fmt.Errorf("open store: %w", err)
}
defer st.Close()
}
// ----- daemon components (only wired when unlocked) -----
// Pre-declare so the unlock path can wire them later.
var (
gatherer *loop.Gatherer
rules []loop.Rule
phr phraser.Phraser
voiceW *voiceWiring
dispatcher *delivery.Dispatcher
tl *tickLoop
coreAPI ipc.CoreAPI
eco *ecosystemWiring
)
if !locked {
rules = loop.DefaultRules()
gatherer = loop.NewGatherer(st, rules)
if cfg.QuietHours != nil {
gatherer.SetQuietHours(cfg.QuietHours.Start, cfg.QuietHours.End)
}
// phraser
phr = phraser.NewStub()
if cfg.Phraser != nil {
pc := phraser.Config{
ModelPath: cfg.Phraser.ModelPath,
BinPath: cfg.Phraser.BinPath,
Listen: cfg.Phraser.Listen,
NGpuLayers: cfg.Phraser.NGpuLayers,
NCtx: cfg.Phraser.NCtx,
Timeout: time.Duration(cfg.Phraser.Timeout),
Persona: personaFromCfg(cfg),
}
if pc.BinPath == "" {
pc.BinPath = "llama-server"
}
if pc.Listen == "" {
pc.Listen = "127.0.0.1:0"
}
if pc.NCtx <= 0 {
pc.NCtx = 2048
}
if pc.Timeout <= 0 {
pc.Timeout = 30 * time.Second
}
var err error
phr, err = phraser.NewLLMPhraser(ctx, pc)
if err != nil {
return fmt.Errorf("phraser: %w", err)
}
}
// ecosystem — nexus + hexis + praxis (all over HTTP; no direct DB access)
eco = wireEcosystem(cfg)
// voice
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
// delivery
var ntfy delivery.Sink
if cfg.Ntfy != nil {
s, err := ntfysink.New(*cfg.Ntfy)
if err != nil {
return fmt.Errorf("wire ntfy sink: %w", err)
}
ntfy = s
}
var telegram delivery.Sink
if cfg.Telegram != nil {
s, err := telegramsink.New(*cfg.Telegram)
if err != nil {
return fmt.Errorf("wire telegram sink: %w", err)
}
telegram = s
}
var voiceSink delivery.Sink
if voiceW != nil {
voiceSink = voiceW.voiceSink
}
dispatcher = delivery.NewDispatcher(delivery.Config{
Ntfy: ntfy,
Telegram: telegram,
Voice: voiceSink,
Ack: st,
Nudges: st,
Reminders: st,
})
// tick loop
tickInterval := time.Duration(cfg.TickInterval)
repeatInterval := time.Duration(cfg.RepeatInterval)
autotuneInterval := time.Duration(cfg.AutotuneInterval)
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines))
coreAPI = &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
getTrace: tl.trace,
}
if voiceW != nil && voiceW.handler != nil {
api := coreAPI.(*daemonAPI)
api.chatFn = voiceW.handler.handleText
}
} else {
// locked mode: dummy CoreAPI that returns errLocked for everything
coreAPI = &lockedAPI{}
}
// ----- IPC boundary (core ↔ modules) -----
srv, err := ipc.Listen(cfg.SocketPath, coreAPI)
if err != nil {
return fmt.Errorf("ipc listen: %w", err)
}
passkeySess := webauthn.NewPasskeySession(5 * time.Minute)
// Set Server.Check — in locked mode, block everything except unlock-path methods.
if locked {
srv.Check = func(ctx context.Context, m ipc.Method, _ json.RawMessage) error {
switch m {
case ipc.MethodAssertStepUp, ipc.MethodUnlock:
return nil // allowed in locked mode
default:
return errLocked
}
}
} else {
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
}
srv.StepUp = func(ctx context.Context) error { return passkeySess.Assert(ctx, auth.Scope{}) }
// WrapKeyFn — wraps the env key with a passkey credential public key and
// persists the wrapped blob. Only wired when the daemon has the key in
// memory (env key mode). Called by mavweb after passkey enrollment.
if envKeyBytes != nil {
srv.WrapKeyFn = func(ctx context.Context, publicKey []byte) error {
blob, err := webauthn.WrapKey(envKeyBytes, publicKey)
if err != nil {
return fmt.Errorf("wrap encryption key: %w", err)
}
wp := *wrappedKeyPath
if wp == "" {
wp = cfg.DefaultWrappedKeyPath()
}
if err := os.WriteFile(wp, blob, 0o600); err != nil {
return fmt.Errorf("write wrapped key: %w", err)
}
log.Printf("mavend: wrapped encryption key with passkey credential (%d bytes)", len(blob))
return nil
}
}
// UnlockFn — cold-start unlock: unwraps the encryption key from the wrapped
// blob using the passkey credential public key, opens the store, wires all
// daemon components, and replaces the locked API.
if locked {
srv.UnlockFn = func(ctx context.Context, publicKey []byte) error {
wp := *wrappedKeyPath
blob, err := os.ReadFile(wp)
if err != nil {
return fmt.Errorf("read wrapped key: %w", err)
}
key, err := webauthn.UnwrapKey(blob, publicKey)
if err != nil {
return fmt.Errorf("unwrap key: %w", err)
}
// Open the store with the unwrapped key.
st, err = store.OpenEncrypted(ctx, cfg.DBPath, cfg.DBTmpfs, key)
if err != nil {
return fmt.Errorf("unlock store: %w", err)
}
// Wire everything.
rules = loop.DefaultRules()
gatherer = loop.NewGatherer(st, rules)
if cfg.QuietHours != nil {
gatherer.SetQuietHours(cfg.QuietHours.Start, cfg.QuietHours.End)
}
phr = phraser.NewStub()
if cfg.Phraser != nil {
pc := phraser.Config{
ModelPath: cfg.Phraser.ModelPath,
BinPath: cfg.Phraser.BinPath,
Listen: cfg.Phraser.Listen,
NGpuLayers: cfg.Phraser.NGpuLayers,
NCtx: cfg.Phraser.NCtx,
Timeout: time.Duration(cfg.Phraser.Timeout),
Persona: personaFromCfg(cfg),
}
if pc.BinPath == "" {
pc.BinPath = "llama-server"
}
if pc.Listen == "" {
pc.Listen = "127.0.0.1:0"
}
if pc.NCtx <= 0 {
pc.NCtx = 2048
}
if pc.Timeout <= 0 {
pc.Timeout = 30 * time.Second
}
phr, err = phraser.NewLLMPhraser(ctx, pc)
if err != nil {
return fmt.Errorf("phraser: %w", err)
}
}
eco = wireEcosystem(cfg)
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
var ntfy delivery.Sink
if cfg.Ntfy != nil {
s, err := ntfysink.New(*cfg.Ntfy)
if err != nil {
return fmt.Errorf("wire ntfy sink: %w", err)
}
ntfy = s
}
var telegram delivery.Sink
if cfg.Telegram != nil {
s, err := telegramsink.New(*cfg.Telegram)
if err != nil {
return fmt.Errorf("wire telegram sink: %w", err)
}
telegram = s
}
var voiceSink delivery.Sink
if voiceW != nil {
voiceSink = voiceW.voiceSink
}
dispatcher = delivery.NewDispatcher(delivery.Config{
Ntfy: ntfy,
Telegram: telegram,
Voice: voiceSink,
Ack: st,
Nudges: st,
Reminders: st,
})
tickInterval := time.Duration(cfg.TickInterval)
repeatInterval := time.Duration(cfg.RepeatInterval)
autotuneInterval := time.Duration(cfg.AutotuneInterval)
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines))
// Swap the CoreAPI from lockedAPI to the real store adapter.
newAPI := &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
getTrace: tl.trace,
}
if voiceW != nil && voiceW.handler != nil {
newAPI.chatFn = voiceW.handler.handleText
}
srv.SetAPI(newAPI)
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
// Start voice server.
if voiceW != nil {
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("voice serve: %v", err)
}
}()
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
}
// Start tick loop.
go func() {
tl.run(ctx)
}()
dl.unlock()
log.Printf("mavend: unlocked via passkey assertion")
return nil
}
}
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
if err := srv.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("ipc serve: %v", err)
}
}()
log.Printf("mavend: ipc listening on %s", srv.Path())
if !locked && voiceW != nil {
wg.Add(1)
go func() {
defer wg.Done()
if err := voiceW.server.Serve(); err != nil && !errors.Is(err, net.ErrClosed) {
log.Printf("voice serve: %v", err)
}
}()
log.Printf("mavend: voice listening on %s", voiceW.server.Addr())
}
if !locked {
wg.Add(1)
go func() {
defer wg.Done()
tl.run(ctx)
}()
}
<-ctx.Done()
log.Printf("mavend: shutdown signal received")
if err := srv.Close(); err != nil {
log.Printf("ipc close: %v", err)
}
if voiceW != nil {
voiceW.close()
}
wg.Wait()
log.Printf("mavend: bye")
return nil
}
// personaFromCfg extracts the voice persona from the config, or returns ""
// when voice isn't configured. Used to pass a character prompt into the
// LLM phraser without requiring voice to be enabled.
func personaFromCfg(cfg *config.Config) string {
if cfg.Voice != nil {
return cfg.Voice.Persona
}
return ""
}