Files
Maven/cmd/mavend/main.go
T
kami 062d4252ef Tell her what she can actually do
The context block now lists her real capabilities: reminders, notes and
facts (write and recall), and the calendar — all three are code paths in
mavend today. Weather, telegram and shell acts are listed only when the
config actually has them, because offering something she cannot do is
worse than staying quiet about it.

Also drops the pronouns from the optional name/city line. The block's
own "ты" is Maven, so "тебя зовут" read as her name and "его" would have
shown her the third-person form she must never use about him. They are
plain labels now.

Vikunja #394.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CGeSZxh1DCtRxmFVSYVGvJ
2026-07-31 15:57:58 +04:00

631 lines
20 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/loop"
"github.com/kami/maven/internal/persona"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/store"
"github.com/kami/maven/internal/webauthn"
)
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) AcceptProposedRoutine(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 (l *lockedAPI) MorningStatus(ctx context.Context) ([]ipc.MorningRoutineStatus, error) {
return nil, 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)")
reembed := flag.Bool("reembed", false, "re-embed every stored note and fact with the configured embedder, then serve normally (run once after an embedder swap)")
flag.CommandLine.Parse(args)
reembedOnStart = *reembed
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
factWorker *factEnrichmentWorker
)
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),
ContextBlock: contextBlockFn(cfg, time.Now),
}
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
}
// A crashed prior run may have left "pending" delivery attempts (send
// may have landed externally, then the process died before recording
// it) — reconcile them to "unknown" before the tick loop resumes
// sending, so nothing auto-resends into that ambiguity.
if _, err := st.ReconcileStaleDeliveryAttempts(context.Background(), time.Now()); err != nil {
log.Printf("delivery outbox reconcile: %v", err)
}
dispatcher = delivery.NewDispatcher(delivery.Config{
Ntfy: ntfy,
Telegram: telegram,
Voice: voiceSink,
Ack: st,
Nudges: st,
Reminders: st,
Outbox: 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), config.MorningRoutinesFromConfig(cfg.MorningRoutines))
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
coreAPI = &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
}
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),
ContextBlock: contextBlockFn(cfg, time.Now),
}
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
}
if _, err := st.ReconcileStaleDeliveryAttempts(context.Background(), time.Now()); err != nil {
log.Printf("delivery outbox reconcile: %v", err)
}
dispatcher = delivery.NewDispatcher(delivery.Config{
Ntfy: ntfy,
Telegram: telegram,
Voice: voiceSink,
Ack: st,
Nudges: st,
Reminders: st,
Outbox: 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), config.MorningRoutinesFromConfig(cfg.MorningRoutines))
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
// Swap the CoreAPI from lockedAPI to the real store adapter.
newAPI := &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
}
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)
}()
// Start fact-entity enrichment worker.
go func() {
factWorker.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)
}()
wg.Add(1)
go func() {
defer wg.Done()
factWorker.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
}
// personaFacts reads the optional, deployment-specific facts (his name, his
// city, the free-text persona string) out of the config. Everything here may
// be empty — the context block is correct without any of it.
func personaFacts(cfg *config.Config) persona.Facts {
f := persona.Facts{
// Telegram lives outside the voice block, so it counts either way.
Telegram: cfg.Telegram != nil && cfg.Telegram.BotToken != "" && cfg.Telegram.ChatID != "",
}
if cfg.Voice == nil {
return f
}
f.OwnerName = cfg.Voice.OwnerName
f.City = cfg.Voice.City
f.Static = cfg.Voice.Persona
// Same test wireVoice uses to pick the real provider over the stub.
f.Weather = cfg.Voice.Weather != nil && cfg.Voice.Weather.Provider == "open-meteo"
f.Tools = len(cfg.Voice.Tools) > 0
return f
}
// contextBlockFn returns the per-turn renderer of the shared context block.
// Per turn, not once at startup, because the block states the current time.
func contextBlockFn(cfg *config.Config, now func() time.Time) func() string {
f := personaFacts(cfg)
return func() string { return f.Block(now()) }
}