Files
Maven/internal/ipc/ipc_test.go
T
kami d52f60c54e maven: fix test mocks for CalendarEvents interface (verification)
- Add CalendarEvents method to recordingAPI in auth_test.go
- Add CalendarEvents method to fakeCore in handlers_test.go

Co-Authored-By: opencode <opencode@anthropic.com>
2026-07-06 04:20:16 +04:00

427 lines
14 KiB
Go

package ipc
import (
"bytes"
"context"
"encoding/binary"
"encoding/json"
"errors"
"io"
"net"
"os"
"path/filepath"
"testing"
"time"
"github.com/kami/maven/internal/store"
)
// tmpSocket — a socket path under a 0700 temp dir, unique per test.
func tmpSocket(t *testing.T) string {
t.Helper()
dir := t.TempDir()
return filepath.Join(dir, "maven.sock")
}
// newServerWithStore spins a real store + Server + Client so the boundary
// is exercised exactly as the daemon wires it. Returns the api (for direct
// in-process expectations) and a client going through the socket.
func newServerWithStore(t *testing.T) (CoreAPI, *Server, *Client, *store.Store) {
t.Helper()
dir := t.TempDir()
s, err := store.Open(context.Background(), filepath.Join(dir, "maven.db"))
if err != nil {
t.Fatalf("open store: %v", err)
}
t.Cleanup(func() { _ = s.Close() })
api := NewStoreAPI(s)
srv, err := Listen(tmpSocket(t), api)
if err != nil {
t.Fatalf("listen: %v", err)
}
done := make(chan struct{})
go func() {
_ = srv.Serve()
close(done)
}()
t.Cleanup(func() {
_ = srv.Close()
<-done
})
cli, err := Dial(srv.Path())
if err != nil {
t.Fatalf("dial: %v", err)
}
t.Cleanup(func() { _ = cli.Close() })
return api, srv, cli, s
}
// TestDialWait_WaitsForLateServer — a module may start before core's socket
// exists (core loads models first). DialWait must keep retrying until the
// socket appears rather than fail on the first attempt.
func TestDialWait_WaitsForLateServer(t *testing.T) {
dir := t.TempDir()
sock := filepath.Join(dir, "maven.sock")
srvCh := make(chan *Server, 1)
errCh := make(chan error, 1)
// bring the server up only after DialWait is already retrying
go func() {
time.Sleep(300 * time.Millisecond)
s, err := store.Open(context.Background(), filepath.Join(dir, "maven.db"))
if err != nil {
errCh <- err
return
}
t.Cleanup(func() { _ = s.Close() })
srv, err := Listen(sock, NewStoreAPI(s))
if err != nil {
errCh <- err
return
}
go func() { _ = srv.Serve() }()
srvCh <- srv
}()
cli, err := DialWait(sock, 5*time.Second)
if err != nil {
t.Fatalf("DialWait should connect once the server appears: %v", err)
}
select {
case err := <-errCh:
t.Fatalf("server setup failed: %v", err)
case srv := <-srvCh:
_ = cli.Close() // close client first so srv.Close's handler wait returns
_ = srv.Close()
}
}
// TestClient_ReconnectsAfterServerRestart — a long-lived module (e.g. mavweb)
// must survive a core restart. The first server is closed and a new one is
// brought up on the SAME socket path (as a daemon restart does); the client's
// cached conn is now dead. The next call must transparently re-dial and succeed
// instead of failing forever with "broken pipe".
func TestClient_ReconnectsAfterServerRestart(t *testing.T) {
dir := t.TempDir()
sock := filepath.Join(dir, "maven.sock")
serve := func() *Server {
s, err := store.Open(context.Background(), filepath.Join(dir, "maven.db"))
if err != nil {
t.Fatalf("open store: %v", err)
}
t.Cleanup(func() { _ = s.Close() })
srv, err := Listen(sock, NewStoreAPI(s))
if err != nil {
t.Fatalf("listen: %v", err)
}
go func() { _ = srv.Serve() }()
return srv
}
srv1 := serve()
cli, err := Dial(sock)
if err != nil {
t.Fatalf("dial: %v", err)
}
t.Cleanup(func() { _ = cli.Close() })
// works against the first server
if _, err := cli.Presence(context.Background()); err != nil {
t.Fatalf("call before restart: %v", err)
}
// Simulate a core restart: the client's conn dies (as it would when the
// daemon process exits), then a fresh server binds the SAME path. Closing
// the client side first also lets srv1's handler goroutine see EOF and
// exit, so srv1.Close()'s wg.Wait() returns instead of blocking on a
// parked reader.
cli.conn.Close()
_ = srv1.Close()
srv2 := serve()
// The cached conn is dead — the call must transparently re-dial and succeed.
if _, err := cli.Presence(context.Background()); err != nil {
t.Fatalf("call after restart should have re-dialed, got: %v", err)
}
// Teardown order matters: Server.Close waits for its handler goroutine,
// which is parked reading the (now live, re-dialed) client conn. Close the
// client first so the handler sees EOF and Close returns instead of hanging.
_ = cli.Close()
_ = srv2.Close()
}
// TestFrame_Roundtrip — JSON over a length prefix survives the loop, and the
// prefix itself encodes the length exactly. The framing is the only thing
// keeping a module's request paired with core's reply; it's worth a direct test.
func TestFrame_Roundtrip(t *testing.T) {
var buf bytes.Buffer
type payload struct {
Msg string `json:"m"`
N int `json:"n"`
}
want := payload{Msg: "hello", N: 42}
if err := writeFrame(&buf, want); err != nil {
t.Fatalf("writeFrame: %v", err)
}
// header length must equal the JSON body length that follows.
var hdr [4]byte
if _, err := io.ReadFull(&buf, hdr[:]); err != nil {
t.Fatalf("read hdr: %v", err)
}
bodyLen := binary.BigEndian.Uint32(hdr[:])
if int(bodyLen) != buf.Len() {
t.Fatalf("prefix length %d != body %d", bodyLen, buf.Len())
}
// readFrame consumes header+body together; recombine so it sees a whole frame.
full := append(hdr[:], buf.Bytes()...)
var got payload
if err := readFrame(bytes.NewReader(full), &got); err != nil {
t.Fatalf("readFrame: %v", err)
}
if got != want {
t.Fatalf("roundtrip mismatch: got %+v want %+v", got, want)
}
}
func TestFrame_TooLarge(t *testing.T) {
// Encode-side guard refuses to ship anything bigger than maxFrame; the
// socket never sees it. Defense against a confused peer, not a real path.
big := make([]byte, maxFrame+1)
if err := writeFrame(io.Discard, big); !errors.Is(err, ErrFrameTooLarge) {
t.Fatalf("writeFrame: got %v, want ErrFrameTooLarge", err)
}
// Decode-side guard refuses a header claiming a too-large body; the conn
// is now desynced (length read, body not), but readFrame doesn't have to
// recover — the caller closes it.
var hdr [4]byte
binary.BigEndian.PutUint32(hdr[:], maxFrame+1)
if err := readFrame(bytes.NewReader(hdr[:]), nil); !errors.Is(err, ErrFrameTooLarge) {
t.Fatalf("readFrame: got %v, want ErrFrameTooLarge", err)
}
}
// TestSocket_Perms — the auth floor. 0600 ⇒ only the same unix user can
// connect. If this regresses to world-readable, every user on the box is a
// module; that's the entire auth model today, so assert it.
func TestSocket_Perms(t *testing.T) {
_, srv, _, _ := newServerWithStore(t)
fi, err := os.Stat(srv.Path())
if err != nil {
t.Fatalf("stat socket: %v", err)
}
mode := fi.Mode().Perm()
if mode != 0o600 {
t.Fatalf("socket perm = %#o, want 0600", mode)
}
}
// TestStoreAPI_Direct — the in-process adapter path (no socket) maps store
// sentinels to ipc sentinels. The boundary's contract is that error identity
// is the same on both sides; this pins it for the daemon-embedded modules
// (router, delivery today) that never go over the wire.
func TestStoreAPI_Direct(t *testing.T) {
api, _, _, _ := newServerWithStore(t)
ctx := context.Background()
// missing key ⇒ ErrNoFact
if _, err := api.LatestFact(ctx, "nope"); !errors.Is(err, ErrNoFact) {
t.Fatalf("LatestFact missing: got %v, want ErrNoFact", err)
}
// bad confidence ⇒ ErrConfidence
if _, err := api.WriteFact(ctx, WriteFactReq{
Ts: time.Now(), Kind: "self", Key: "water", Value: "1", Source: "tap:water", Confidence: 0,
}); !errors.Is(err, ErrConfidence) {
t.Fatalf("WriteFact conf=0: got %v, want ErrConfidence", err)
}
// since on missing key ⇒ ErrNoFact
if _, err := api.Since(ctx, "nope", time.Now()); !errors.Is(err, ErrNoFact) {
t.Fatalf("Since missing: got %v, want ErrNoFact", err)
}
// reminder idempotency: invalid status ⇒ ErrReminderState
if err := api.MarkReminder(ctx, 99999, "weird"); !errors.Is(err, ErrReminderState) {
t.Fatalf("MarkReminder weird: got %v, want ErrReminderState", err)
}
// resolve nonexistent nudge ⇒ ErrNudgeNotFound
if err := api.ResolveNudge(ctx, 99999, "acted", time.Now()); !errors.Is(err, ErrNudgeNotFound) {
t.Fatalf("ResolveNudge none: got %v, want ErrNudgeNotFound", err)
}
}
// TestClient_E2E — full socket round trip against a real store. Drives every
// method end-to-end and asserts sentinel identity survives the wire. This is
// the test that catches the boundary bugs: param shape mismatch, sentinel
// code drift, dto mapping, framing interleaving.
func TestClient_E2E(t *testing.T) {
_, _, cli, _ := newServerWithStore(t)
ctx := context.Background()
now := time.Now().UTC().Truncate(time.Millisecond)
// write a tap (self, confidence 1.0) and read it back.
id, err := cli.WriteFact(ctx, WriteFactReq{
Ts: now, Kind: "self", Key: "water", Value: "1", Source: "tap:water", Confidence: 1.0,
})
if err != nil {
t.Fatalf("WriteFact: %v", err)
}
if id <= 0 {
t.Fatalf("WriteFact returned id %d", id)
}
f, err := cli.LatestFact(ctx, "water")
if err != nil {
t.Fatalf("LatestFact: %v", err)
}
if f.Key != "water" || f.Value != "1" || f.Source != "tap:water" || f.Confidence != 1.0 {
t.Fatalf("LatestFact mismatch: %+v", f)
}
if !f.Ts.Equal(now) {
t.Fatalf("Ts roundtrip: got %v want %v", f.Ts, now)
}
// provenance scope: a foreign source doesn't see the tap value.
if _, err := cli.LatestFactBySource(ctx, "water", "poll:evil"); !errors.Is(err, ErrNoFact) {
t.Fatalf("LatestFactBySource foreign: got %v, want ErrNoFact", err)
}
if _, err := cli.LatestFactBySource(ctx, "water", "tap:water"); err != nil {
t.Fatalf("LatestFactBySource own: %v", err)
}
// since: ~0 elapsed since "now".
d, err := cli.Since(ctx, "water", now.Add(time.Second))
if err != nil {
t.Fatalf("Since: %v", err)
}
if d != time.Second {
t.Fatalf("Since dur = %v, want 1s", d)
}
// since missing ⇒ ErrNoFact over the wire.
if _, err := cli.Since(ctx, "nope", now); !errors.Is(err, ErrNoFact) {
t.Fatalf("Since missing: got %v, want ErrNoFact", err)
}
// presence cold-start ⇒ away, score 0.
pres, err := cli.Presence(ctx)
if err != nil {
t.Fatalf("Presence: %v", err)
}
if pres.Bucket != Away || pres.Score != 0 {
t.Fatalf("Presence cold-start = %+v, want away/0", pres)
}
// reminder lifecycle: create → mark fired → re-mark ⇒ ErrReminderState.
rid, err := cli.CreateReminder(ctx, now.Add(time.Hour), `{"text":"wake me 7"}`, "")
if err != nil {
t.Fatalf("CreateReminder: %v", err)
}
if err := cli.MarkReminder(ctx, rid, "fired"); err != nil {
t.Fatalf("MarkReminder fired: %v", err)
}
if err := cli.MarkReminder(ctx, rid, "fired"); !errors.Is(err, ErrReminderState) {
t.Fatalf("MarkReminder twice: got %v, want ErrReminderState", err)
}
// nudge lifecycle: record → resolve acted → resolve again ⇒ ErrNudgeOutcome.
nid, err := cli.RecordNudge(ctx, "water", "voice", "drink", now)
if err != nil {
t.Fatalf("RecordNudge: %v", err)
}
if err := cli.ResolveNudge(ctx, nid, "acted", now); err != nil {
t.Fatalf("ResolveNudge acted: %v", err)
}
if err := cli.ResolveNudge(ctx, nid, "ignored", now); !errors.Is(err, ErrNudgeOutcome) {
t.Fatalf("ResolveNudge twice: got %v, want ErrNudgeOutcome", err)
}
// feedback loop read: RecentOutcomes returns the resolved outcome.
out, err := cli.RecentOutcomes(ctx, "water", 5)
if err != nil {
t.Fatalf("RecentOutcomes: %v", err)
}
if len(out) != 1 || out[0] != "acted" {
t.Fatalf("RecentOutcomes = %v, want [acted]", out)
}
// empty result over the wire is a stable [] not null (server coerces).
if got, err := cli.RecentOutcomes(ctx, "never_fired_rule", 5); err != nil || len(got) != 0 {
t.Fatalf("RecentOutcomes empty = %v err=%v, want []", got, err)
}
}
// TestCaller_Peercred — when the client dials, core sees a Caller with the
// test process's own uid via SO_PEERCRED. This is the seam auth scopes on;
// asserting it's populated today means the future auth layer has its input.
func TestCaller_Peercred(t *testing.T) {
_, srv, cli, _ := newServerWithStore(t)
ctx := context.Background()
// round-trip any call; the server annotates ctx with a Caller on accept.
if _, err := cli.LatestFact(ctx, "nope"); err != nil && !errors.Is(err, ErrNoFact) {
t.Fatalf("LatestFact: %v", err)
}
// introspect the server's view: re-accept a conn manually and read creds.
uc, err := dialRaw(srv.Path())
if err != nil {
t.Fatalf("dialRaw: %v", err)
}
defer uc.Close()
c, ok := peerCaller(uc)
if !ok {
t.Skip("SO_PEERCRED unavailable on this platform; skipping")
}
if c.Uid != int32(os.Getuid()) {
t.Fatalf("peercred uid = %d, want %d", c.Uid, os.Getuid())
}
}
// TestDispatch_UnknownMethod — an unknown method over the wire comes back as
// ErrUnknownMethod, not a panic or a dropped conn. The server must stay up
// for the next (legitimate) request on the same conn.
func TestDispatch_UnknownMethod(t *testing.T) {
_, srv, _, _ := newServerWithStore(t)
uc, err := dialRaw(srv.Path())
if err != nil {
t.Fatalf("dialRaw: %v", err)
}
defer uc.Close()
// send garbage method on the raw conn, read back its error, then send a
// real method on the SAME conn to confirm the server survived.
if err := writeFrame(uc, Request{Method: Method("definitely_not_a_method")}); err != nil {
t.Fatalf("writeFrame: %v", err)
}
var resp Response
if err := readFrame(uc, &resp); err != nil {
t.Fatalf("readFrame: %v", err)
}
if resp.Error == nil || !errors.Is(hydrate(resp.Error), ErrUnknownMethod) {
t.Fatalf("unknown method response = %+v, want ErrUnknownMethod", resp.Error)
}
// same conn, legit follow-up: prove the goroutine is still alive.
if err := writeFrame(uc, Request{Method: MethodLatestFact, Params: mustJSON(keyReq{Key: "nope"})}); err != nil {
t.Fatalf("writeFrame follow-up: %v", err)
}
if err := readFrame(uc, &resp); err != nil {
t.Fatalf("readFrame follow-up: %v", err)
}
if resp.Error == nil || !errors.Is(hydrate(resp.Error), ErrNoFact) {
t.Fatalf("follow-up response = %+v, want ErrNoFact", resp.Error)
}
}
// dialRaw — a bare unix conn for tests that want to script the wire directly
// (send an unknown method, follow up on the same conn, inspect framing).
func dialRaw(path string) (net.Conn, error) {
return net.Dial("unix", path)
}
func mustJSON(v any) []byte {
b, err := json.Marshal(v)
if err != nil {
panic(err)
}
return b
}