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.WriteFact(ctx, req.Ts, store.FactKind(req.Kind), req.Key, 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) 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) 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 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 with a passkey // credential public key (HKDF-AESGCM) and writes the wrapped blob to disk. // Set by the daemon; nil ⇒ MethodStoreEncryptionKey returns ErrUnknownMethod. WrapKeyFn WrapKeyFunc // UnlockFn — unwraps the store encryption key from the wrapped blob using // the passkey credential public key, 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 with the given credential // public key and persists the wrapped blob. type WrapKeyFunc func(ctx context.Context, publicKey []byte) error // UnlockFunc — unwraps the store encryption key using the given credential // public key and completes daemon initialization. type UnlockFunc func(ctx context.Context, publicKey []byte) 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) } // 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 } } switch req.Method { case MethodWriteFact: var p WriteFactReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } id, err := api.WriteFact(ctx, p) return marshalResult(idResp{ID: id}), err case MethodLatestFact: var p keyReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } f, err := api.LatestFact(ctx, p.Key) if err != nil { return nil, err } return marshalResult(f), nil case MethodLatestFactBySource: var p keySourceReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } f, err := api.LatestFactBySource(ctx, p.Key, p.Source) if err != nil { return nil, err } return marshalResult(f), nil case MethodSince: var p sinceReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } d, err := api.Since(ctx, p.Key, p.Now) if err != nil { return nil, err } return marshalResult(sinceResp{Dur: d}), nil case MethodPresence: pres, err := api.Presence(ctx) if err != nil { return nil, err } return marshalResult(pres), nil case MethodCreateReminder: var p createReminderReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } id, err := api.CreateReminder(ctx, p.Fire, p.Payload, p.Cron) return marshalResult(idResp{ID: id}), err case MethodMarkReminder: var p markReminderReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } err := api.MarkReminder(ctx, p.ID, p.Status) return marshalResult(nil), err case MethodListReminders: var p nReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.ListReminders(ctx, p.N) if err != nil { return nil, err } if out == nil { out = []Reminder{} } return marshalResult(out), nil case MethodRecordNudge: var p recordNudgeReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } id, err := api.RecordNudge(ctx, p.Rule, p.Channel, p.Message, p.Ts) return marshalResult(idResp{ID: id}), err case MethodResolveNudge: var p resolveNudgeReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } err := api.ResolveNudge(ctx, p.ID, p.Outcome, p.Ts) return marshalResult(nil), err case MethodRecentOutcomes: var p outcomesReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } 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 marshalResult(out), nil case MethodRecentFacts: var p nReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.RecentFacts(ctx, p.N) if err != nil { return nil, err } if out == nil { out = []Fact{} } return marshalResult(out), nil case MethodCalendarEvents: var p calendarEventsReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.CalendarEvents(ctx, p.From, p.To) if err != nil { return nil, err } if out == nil { out = []Fact{} } return marshalResult(out), nil case MethodRecentNudges: var p nReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.RecentNudges(ctx, p.N) if err != nil { return nil, err } if out == nil { out = []Nudge{} } return marshalResult(out), nil case MethodWriteNote: var p writeNoteReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } id, err := api.WriteNote(ctx, p.Ts, p.Text, p.Embedding, p.Source) return marshalResult(idResp{ID: id}), err case MethodQueryNotes: var p queryNotesReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.QueryNotes(ctx, p.Embedding, p.K) if err != nil { return nil, err } if out == nil { out = []Note{} } return marshalResult(out), nil case MethodRecentNotes: var p nReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.RecentNotes(ctx, p.N) if err != nil { return nil, err } if out == nil { out = []Note{} } return marshalResult(out), nil case MethodProposeTool: var p proposeToolReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } ok, err := api.ProposeTool(ctx, p.Name, p.Utterance, p.Scope, p.Ts) if err != nil { return nil, err } return marshalResult(proposeToolResp{Proposed: ok}), nil case MethodEnableTool: var p enableToolReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } return marshalResult(nil), api.EnableTool(ctx, p.Name, p.Cmd, p.Destructive, p.Scope, p.Ts) case MethodDisableTool: var p disableToolReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } return marshalResult(nil), api.DisableTool(ctx, p.Name) case MethodLookupTool: var p lookupToolReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } t, err := api.LookupTool(ctx, p.Name) if err != nil { return nil, err } return marshalResult(t), nil case MethodListTools: var p listToolsReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } out, err := api.ListTools(ctx, p.Status) if err != nil { return nil, err } if out == nil { out = []Tool{} } return marshalResult(listToolsResp{Tools: out}), nil case MethodDeleteTool: var p disableToolReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } return marshalResult(nil), api.DeleteTool(ctx, p.Name) case MethodListProposedRoutines: out, err := api.ListProposedRoutines(ctx) if err != nil { return nil, err } if out == nil { out = []ProposedRoutine{} } return marshalResult(listProposedRoutinesResp{Routines: out}), nil case MethodDismissProposedRoutine: var p dismissProposedRoutineReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } return marshalResult(nil), api.DismissProposedRoutine(ctx, p.ID) case MethodRevertFact: var p struct { Key string `json:"key"` } if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } newID, err := api.RevertFact(ctx, p.Key) if err != nil { return nil, err } return marshalResult(map[string]int64{"new_id": newID}), nil case MethodChat: var p chatReq if err := unmarshalParams(req.Params, &p); err != nil { return nil, err } reply, err := api.Chat(ctx, p.Text) if err != nil { return nil, err } return marshalResult(chatResp{Reply: reply}), nil case MethodTickTrace: t, err := api.TickTrace(ctx) if err != nil { return nil, err } return marshalResult(t), nil 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.PublicKey) } 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.PublicKey) } return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method) default: return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method) } } 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 }