diff --git a/CHANGELOG.md b/CHANGELOG.md index 08ad711a..445e7aa3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -128,6 +128,33 @@ property, about 80 bytes per tool, but the profile surcharge measured the schemas without it. It now measures the schemas as that connection is served them. (#515) +- **A shared connection refuses a write under an identity nobody declared.** + When several conversations shared one `plumb serve`, a state-changing call + carrying a per-call identity that no `session_start` on that connection had + declared (a model typing `plumb_agent: "my-session"`, or a client sending + `_meta` it never announced) was admitted. It got a fresh per-agent state + seeded from the connection's workspace, so a relative write landed in whatever + checkout the connection held, often another agent's. Now such a call is + refused, and the refusal says to call `session_start` first, with the + `workspace` the agent works in so that declaring does not leave it on another + agent's checkout. Reads are never refused. A connection used by one + conversation only (a main thread and its own subagents, as in the Claude Code + CLI) needs no declaration. A subagent stamped + `/` is admitted on its conversation's declaration and + works in its conversation's workspace rather than the connection's. It starts + there, and it follows when its conversation later moves itself to another + workspace, unless the subagent chose a workspace of its own. So a subagent of + an agent working in a worktree no longer writes into the main checkout. If + such a subagent moves the connection's pin (`scope: "connection"`), it stays + on its conversation's workspace, and `session_start` now says so rather than + naming the connection's new root as where its relative paths go. + Declarations are saved under the proxy session and restored when it + reconnects, including after a daemon restart. One is reclaimed only when the + idle reaper finds its `plumb serve` disconnected and the declaration older + than `[session] persist_state_ttl_minutes` (24 hours by default), and each + state-changing call refreshes it at most once per quarter of that (at most an + hour). With `persist_state` off, or after such a reclaim, an agent is refused + once and declares again with `session_start` (#513). - **A connection that closes mid-attach no longer leaks a language server, and a burst of roots notifications settles on the newest roots.** When a client reported a workspace change (`notifications/roots/list_changed`) diff --git a/docs/architecture.md b/docs/architecture.md index 086601b2..1d7a8b54 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -325,7 +325,7 @@ Rung 1 outranks client roots because it is the workspace the caller chose, not w **Single-workspace-per-connection contract.** Once a connection has attached a workspace, every path-bearing tool refuses paths outside the connection's allowed roots with a `workspace boundary violation` error (the allowed set is the workspace plus `extra_roots` read-write and `read_roots` / Go dependency roots read-only). `rename_symbol` also boundary-checks each output URI before applying. To switch projects, call `session_start` with an explicit `workspace`: a deliberate `workspace` arg re-pins the connection (re-attaching LSP/topology/quality/config) rather than being refused — clients may reuse one `plumb serve` across chats, so a fresh chat is not a fresh connection. The pin is **sticky once explicit** (issue #182): after a pin was set by `session_start`, a conflicting re-pin to a different project is refused unless the caller passes `force: true`, and a `roots/list_changed` that drops the pinned root can no longer move it — a peer agent multiplexed over the same `plumb serve` connection must not silently steal another agent's workspace (the refusal error names the `force: true` remediation, so a deliberate switch still self-heals; a roots- or auto-attach-held pin is not sticky, so the first explicit pin always lands). A connection that hits a violation is marked `Health: blocked` for the TUI, as is a refused sticky re-pin (cleared by the next successful explicit `session_start`, same-root or forced). `git`'s `repo` arg defaults to the pinned workspace when omitted. -**Shared connections key the pin per logical agent.** When a client identifies its agents (per-call `_meta` or the runtime-stamped argument on every call; `session_start.session_id` identifies only that call), each gets its own shard: its own pin, read/write trackers, undo store, rate budget, LSP routing and workspace-wide diagnostics/symbol aggregates. Three consequences are refusals a caller can meet. A shard seeded from the connection's pin is correctable only within the same tree; an unrelated workspace is refused, and the refusal carries `kind: pin_refused` with `details.scope = "agent"` — where `force: true` moves that agent's shard alone. A shard that has never chosen a workspace and is refused one it declared is RECORDED, and that agent's path-bearing calls are then refused by name until the declaration lands, so a relative path or a defaulted `git` repository cannot resolve inside the seeded root (another conversation's workspace); `session_start` is deliberately not behind that gate, so the remedy stays reachable. A refusal at `details.scope = "connection"` is a different animal: there `force: true` moves the pin every agent on the connection resolves against, so a client must surface it rather than retry it automatically. +**Shared connections key the pin per logical agent.** When a client identifies its agents (per-call `_meta` or the runtime-stamped argument on every call; `session_start.session_id` identifies only that call), each gets its own shard: its own pin, read/write trackers, undo store, rate budget, LSP routing and workspace-wide diagnostics/symbol aggregates. A shard is served to a state-changing call only when the identity's conversation half has declared itself through a successful `session_start` on the connection (a hook-stamped `/` rides its conversation); an undeclared identity is refused with that remedy instead of being handed a shard seeded from the connection's root (#513), except while every identity on the connection shares one conversation. A subagent's shard is seeded from its conversation's chosen root (its shard, or its persisted pin after a restart), follows its conversation when the conversation re-pins itself, is not dragged by a connection move while it sits on its conversation's chosen root, and inherits its conversation's refused declaration; a subagent that chose a root of its own keeps it. Whether a connection move drags a shard is one predicate (`followsConnectionLocked`), shared by the move and by `session_start`'s report of where the mover now resolves, so the two agree on which shards follow; on a connection's first pin there is no previous root to drag from, so a fresh shard is left at no workspace (#567). Declarations are persisted under the proxy session and restored on reconnect, including across a daemon restart. The idle reaper reclaims one only when it finds the serve disconnected and the declaration older than `persist_state_ttl_minutes`; the conversation's admitted state-changing calls refresh it at most once per min(TTL/4, 1 h). One so reclaimed, or any with `persist_state` off, must declare again. Three consequences are refusals a caller can meet. A shard seeded from the connection's pin is correctable only within the same tree; an unrelated workspace is refused, and the refusal carries `kind: pin_refused` with `details.scope = "agent"` — where `force: true` moves that agent's shard alone. A shard that has never chosen a workspace and is refused one it declared is RECORDED, and that agent's path-bearing calls are then refused by name until the declaration lands, so a relative path or a defaulted `git` repository cannot resolve inside the seeded root (another conversation's workspace); `session_start` is deliberately not behind that gate, so the remedy stays reachable. A refusal at `details.scope = "connection"` is a different animal: there `force: true` moves the pin every agent on the connection resolves against, so a client must surface it rather than retry it automatically. ## Persistence layout diff --git a/docs/threat-model.md b/docs/threat-model.md index 09ee4530..9c2b1345 100644 --- a/docs/threat-model.md +++ b/docs/threat-model.md @@ -241,7 +241,16 @@ shares one pin unless each call carries a logical-agent identity — per-call Code's PreToolUse hook). `session_start.session_id` identifies only that call. A client that sends none still shares one pin, and its anonymous state-changing calls are refused once two identities have been seen — see -[Known gaps](#known-gaps). +[Known gaps](#known-gaps). A per-call identity whose conversation never +declared itself through `session_start` on the connection is refused the same +way rather than given a fresh shard of the connection's root (#513), unless +every identity on the connection belongs to one conversation. Declarations +persist under the proxy session and survive a daemon restart that the serve +reconnects across. The idle reaper reclaims one only when it finds the serve +disconnected and the declaration older than `persist_state_ttl_minutes`; +state-changing calls refresh it at most once per min(TTL/4, 1 h). With +`persist_state` off, or after such a reclaim, the agent is refused once and must +call `session_start` again. ### A2 — Path escape via alias or traversal diff --git a/docs/tools.md b/docs/tools.md index 7a2d1128..e3b571ab 100644 --- a/docs/tools.md +++ b/docs/tools.md @@ -121,7 +121,14 @@ key, or `plumb_agent`, placed by a client **runtime** as a top-level key inside `plumb hooks install claude-code` PreToolUse hook is the emitter — it stamps every `mcp__plumb__*` call and also fills `session_id` on this tool); and `session_id` itself, which identifies this call only: on a connection other agents -share, a later write without a per-call identity is refused. Pass a stable value per agent: +share, a later write without a per-call identity is refused. It is also the +declaration a per-call identity needs there: a state-changing call whose +conversation half no successful `session_start` on the connection has declared +(its `session_id`, or the per-call identity it ran under) is refused with that +remedy, plus `workspace` so the declared agent is not left on the connection's +root, unless every identity on the connection belongs to one conversation. A +subagent stamped `/` rides its conversation's declaration +and works in its conversation's workspace, following it when it re-pins. Pass a stable value per agent: the conversation id for a main thread, `/` for a subagent. The session **record** — the name mail is addressed to, `plumb mail --external-id`, name inheritance on resume — is linked to the conversation half, diff --git a/internal/cli/boundary_policy.go b/internal/cli/boundary_policy.go index cf78c78d..ed155cf9 100644 --- a/internal/cli/boundary_policy.go +++ b/internal/cli/boundary_policy.go @@ -111,12 +111,17 @@ func (s *connSession) checkBoundaryFor(ctx context.Context, path string, want to // this can never block the call that clears it. func (s *connSession) declarationRefusedErr(ctx context.Context) error { id := mcp.LogicalAgentFromCtx(ctx) - p, ok := s.pendingDeclarationFor(id) + p, inherited, ok := s.pendingDeclarationForCall(ctx) if !ok { return nil } + whose := "its session_start" + if inherited { + // A subagent inherits its conversation's refused declaration (#513). + whose = fmt.Sprintf("its conversation %q's session_start", linkageIDOf(id)) + } return toolerror.Wrap( - fmt.Errorf("refusing this call for logical agent %q: its session_start naming %s was refused, so the agent is still on %s — a workspace it never chose, and possibly another conversation's. A path-bearing call here would quietly operate on that project. Re-issue session_start with workspace + session_id + force: true (on a shared connection force moves only THIS agent's shard), then retry", id, p.requested, p.sittingOn), + fmt.Errorf("refusing this call for logical agent %q: %s naming %s was refused, so the agent is still on %s — a workspace it never chose, and possibly another conversation's. A path-bearing call here would quietly operate on that project. Re-issue session_start with workspace + session_id + force: true (on a shared connection force moves only THIS agent's shard), then retry", id, whose, p.requested, p.sittingOn), toolerror.KindWorkspaceBoundary, toolerror.ClassPassForce, toolerror.WithTool("session_start"), diff --git a/internal/cli/conn_agent_shard.go b/internal/cli/conn_agent_shard.go index a867735e..3d1c82be 100644 --- a/internal/cli/conn_agent_shard.go +++ b/internal/cli/conn_agent_shard.go @@ -56,6 +56,9 @@ type agentShard struct { // "did this agent choose its root?" must accept either — see the // declaration-refusal marker in repinAgent. restored bool + // parentSeeded: seeded from its conversation's PERSISTED pin (parent not in + // memory), so a connection move must not drag it (#513). + parentSeeded bool // rosterID is the session.Info registered for THIS agent, so the workspace // it actually works in lists it (issue #472). Empty until the agent holds a @@ -111,6 +114,12 @@ func (s *connSession) shardFor(ctx context.Context) *agentShard { writeLimiter: tools.NewRateLimiter(s.store.Current().Edits.RateLimitPerMinute, time.Minute), pinOrigin: v.pinOrigin, } + // A hook-stamped subagent starts where its CONVERSATION chose to work, not + // where the connection happens to sit (issue #513 review). Seeding it from + // the connection pin sent a subagent of a parent that had re-pinned itself + // to a worktree into whichever checkout the connection held — another + // agent's — the exact misroute the declaration gate exists to prevent. + s.seedFromParentLocked(sh) // Restore a pin this agent persisted before the restart (PLAN-286): it takes // precedence over the connection's current pin. A pin that no longer verifies // is ignored, so the shard keeps the connection's root rather than resurrecting @@ -174,7 +183,7 @@ func (s *connSession) buildAgentPolicy(root, language string) *tools.PathPolicy // back to the connection's pin when the connection is not shared (or the call is // unattributed). workspace() stays the ctx-less default for background goroutines. func (s *connSession) workspaceFor(ctx context.Context) string { - if _, pending := s.pendingDeclarationFor(mcp.LogicalAgentFromCtx(ctx)); pending { + if _, _, pending := s.pendingDeclarationForCall(ctx); pending { // A refused declaration leaves nothing trustworthy to anchor to: the // shard's root is one this agent explicitly tried to leave. "" makes the // implicit resolvers — relative paths, git's default repository, @@ -302,11 +311,15 @@ func (s *connSession) repinAgent(ctx context.Context, root, language string, ori // agent then blocks on shardsMu behind it — all waiting on one agent's disk // I/O. Registered before the unlock defer so LIFO runs it AFTER sh.mu is // released, and it re-takes the lock itself. - var syncRoot, syncLang string + var syncRoot, syncLang, movedFrom string defer func() { if syncRoot != "" { s.syncAgentRoster(sh, syncRoot, syncLang) } + // After sh.mu is released: it takes shardsMu, then each shard's mu. + if movedFrom != "" { + s.followParentShard(sh.id, movedFrom) + } }() sh.mu.Lock() defer sh.mu.Unlock() @@ -402,6 +415,7 @@ func (s *connSession) repinAgent(ctx context.Context, root, language string, ori // used to wipe an agent's dirty-guard writes and undo history (PLAN-428). // The connection path keeps them under the same rule. if root != prev { + movedFrom = prev sh.readTracker.Reset() sh.writeTracker.Reset() sh.undoStore.Reset() @@ -443,68 +457,6 @@ func (s *connSession) seedShardOnLink(linkage string) { sh.readTracker.Hydrate(s.readTracker.Records()) } -// followConnectionShards re-seeds every shard that never chose a workspace of -// its own (!selfPinned) from the connection's NEW pin, after the connection -// itself moved away from prevRoot. A shard is seeded from the connection pin at -// first use — shardFor caches it BEFORE repinAgent can refuse, so one refused -// ask left the agent cached at a root whose sticky seed then refused the -// agent's next, entirely legitimate call (the exact PLAN-398 reproduction), -// while a fresh agent asking the same thing succeeded: the fresh shard seeded -// from the CURRENT pin, the stale one had not followed. Re-seeding here restores -// the invariant "a seeded shard sits where the connection sits" without -// touching shards whose agent deliberately pinned elsewhere — per-agent -// isolation means the connection's move cannot drag an agent that chose its own -// root. Runs OUTSIDE the connection mutate lane, in the documented lock order -// (shardsMu before sh.mu, s.mu innermost), so the per-tool-call hot path's lock -// pattern is unchanged; the writes mirror repinAgent's success path, held under -// one sh.mu acquisition each. -// -// Returns the ids of the agents whose shards followed, so session_start can -// tell the caller how many other agents its connection move took with it -// (issue #517). -func (s *connSession) followConnectionShards(prevRoot string) (followed []string) { - if prevRoot == "" { - return nil - } - v := s.view() - if v.acquiredRoot == "" || v.acquiredRoot == prevRoot { - return nil - } - s.shardsMu.Lock() - defer s.shardsMu.Unlock() - for _, sh := range s.shards { - sh.mu.Lock() - if sh.selfPinned || sh.root != prevRoot { - sh.mu.Unlock() - continue - } - sh.root = v.acquiredRoot - sh.language = v.acquiredLanguage - sh.pinOrigin = v.pinOrigin - // The connection landed on the root this agent asked for, so what its - // refusal was about is now simply true; holding the gate would refuse - // calls that are safe again. - if p, ok := s.pendingDeclarationFor(sh.id); ok && p.requested == sh.root { - s.clearDeclarationRefused(sh.id) - } - sh.policy = s.buildAgentPolicy(sh.root, sh.language) - sh.readTracker.Reset() - sh.writeTracker.Reset() - sh.undoStore.Reset() - // Copy what the calls below need while the lock is still held. - // shardsMu does not exclude repinAgent — that takes sh.mu alone — so - // reading sh.root/sh.language after the unlock would race a concurrent - // per-agent re-pin and could persist a root this shard no longer has. - // Same rule persistReadShard states: the shard's root is read under sh.mu. - root, language := sh.root, sh.language - sh.mu.Unlock() - followed = append(followed, sh.id) - s.rehydrateReadsForAgent(sh, root) - s.persistPinForAgent(sh, root, language, v.pinOrigin) - } - return followed -} - // persistReadShard mirrors a per-agent recorded read to the durable store, keyed // by (proxy session ID, logical-agent ID, workspace) so a shared connection's // per-agent reads survive a daemon restart. The shard's root is read under diff --git a/internal/cli/conn_agent_shard_follow.go b/internal/cli/conn_agent_shard_follow.go new file mode 100644 index 00000000..bc9c5574 --- /dev/null +++ b/internal/cli/conn_agent_shard_follow.go @@ -0,0 +1,98 @@ +package cli + +// conn_agent_shard_follow.go — which per-agent shards a move of the +// CONNECTION's pin takes with it, and the one predicate that decides it. +// +// Split from conn_agent_shard.go, which owns creating and resolving a shard and +// was over the file-size cap, so the move (followConnectionShards) and the +// report of it (connScopeCallerRoot) share followsConnectionLocked from one +// place. + +// followConnectionShards re-seeds every shard that never chose a workspace of +// its own (!selfPinned) from the connection's NEW pin, after the connection +// itself moved away from prevRoot. A shard is seeded from the connection pin at +// first use — shardFor caches it BEFORE repinAgent can refuse, so one refused +// ask left the agent cached at a root whose sticky seed then refused the +// agent's next, entirely legitimate call (the exact PLAN-398 reproduction), +// while a fresh agent asking the same thing succeeded: the fresh shard seeded +// from the CURRENT pin, the stale one had not followed. Re-seeding here restores +// the invariant "a seeded shard sits where the connection sits" without +// touching shards whose agent deliberately pinned elsewhere — per-agent +// isolation means the connection's move cannot drag an agent that chose its own +// root. Runs OUTSIDE the connection mutate lane, in the documented lock order +// (shardsMu before sh.mu, s.mu innermost), so the per-tool-call hot path's lock +// pattern is unchanged; the writes mirror repinAgent's success path, held under +// one sh.mu acquisition each. +// +// Returns the ids of the agents whose shards followed, so session_start can +// tell the caller how many other agents its connection move took with it +// (issue #517). +func (s *connSession) followConnectionShards(prevRoot string) (followed []string) { + if prevRoot == "" { + return nil + } + v := s.view() + if v.acquiredRoot == "" || v.acquiredRoot == prevRoot { + return nil + } + s.shardsMu.Lock() + defer s.shardsMu.Unlock() + for id, sh := range s.shards { + // A subagent on its conversation's CHOSEN root stays with it (#513); the + // parent's flag is read BEFORE the child is locked (never nest two mus). + parentChose := s.parentChoseLocked(id) + sh.mu.Lock() + if !followsConnectionLocked(sh, parentChose) || sh.root != prevRoot { + sh.mu.Unlock() + continue + } + sh.root = v.acquiredRoot + sh.language = v.acquiredLanguage + sh.pinOrigin = v.pinOrigin + // The connection landed on the root this agent asked for, so what its + // refusal was about is now simply true; holding the gate would refuse + // calls that are safe again. + if p, ok := s.pendingDeclarationFor(sh.id); ok && p.requested == sh.root { + s.clearDeclarationRefused(sh.id) + } + sh.policy = s.buildAgentPolicy(sh.root, sh.language) + sh.readTracker.Reset() + sh.writeTracker.Reset() + sh.undoStore.Reset() + // Copy what the calls below need while the lock is still held. + // shardsMu does not exclude repinAgent — that takes sh.mu alone — so + // reading sh.root/sh.language after the unlock would race a concurrent + // per-agent re-pin and could persist a root this shard no longer has. + // Same rule persistReadShard states: the shard's root is read under sh.mu. + root, language := sh.root, sh.language + sh.mu.Unlock() + followed = append(followed, sh.id) + s.rehydrateReadsForAgent(sh, root) + s.persistPinForAgent(sh, root, language, v.pinOrigin) + } + return followed +} + +// followsConnectionLocked is the one answer to "does a move of the connection's +// pin take this shard with it?". followConnectionShards asks it to decide whom +// to drag, and connScopeCallerRoot asks it to report where the caller of a +// connection-scoped move resolves afterwards. The two used to decide separately +// and disagreed about a subagent on its conversation's chosen root: the move +// left it there, but the report told it that it now worked in the connection's +// new root (review of #535 merged with #533). +// +// A shard follows unless its agent chose a root (selfPinned), or it is a +// subagent anchored to its conversation's chosen root, either seeded from the +// conversation's persisted pin (parentSeeded) or with the conversation's +// in-memory shard having chosen (parentChose). +// +// parentChose must be parentChoseLocked(sh.id), read under shardsMu BEFORE sh.mu +// is taken, so no path holds two shards' locks at once. Caller holds sh.mu. +// +// restored is deliberately not consulted. followConnectionShards drags only a +// shard that sits on the connection's previous root, and a shard restored from +// its own persisted pin is dragged from there like a seeded one. Whether it +// should be is the persisted-pin question #527 owns, so it is left as it was. +func followsConnectionLocked(sh *agentShard, parentChose bool) bool { + return !sh.selfPinned && !sh.parentSeeded && !parentChose +} diff --git a/internal/cli/conn_agent_shard_parent.go b/internal/cli/conn_agent_shard_parent.go new file mode 100644 index 00000000..29d7c127 --- /dev/null +++ b/internal/cli/conn_agent_shard_parent.go @@ -0,0 +1,156 @@ +package cli + +// conn_agent_shard_parent.go — what a hook-stamped subagent `/` +// inherits from its conversation `` (issue #513 review): the root its +// conversation chose, and the conversation's refused declaration. +// +// A subagent rarely calls session_start; it is admitted on its conversation's +// declaration, so the workspace it resolves against must be its conversation's +// too. Split from conn_agent_shard.go to keep that file under the size cap. + +import ( + "context" + + "github.com/plumbkit/plumb/internal/mcp" +) + +// seedFromParentLocked seeds a NEW subagent shard sh from its conversation's +// root, when the conversation chose one: its in-memory shard is self-pinned or +// restored from its own persisted pin, or, when that shard is not in memory +// (after a daemon restart, before the parent's next call), its persisted +// per-agent pin still verifies. Otherwise sh keeps the connection seed, which is +// then also the conversation's root. A no-op for an id that is its own linkage. +// +// Caller holds s.shardsMu; the parent's mu is taken beneath it, which is the +// documented lock order (shardsMu before sh.mu). +func (s *connSession) seedFromParentLocked(sh *agentShard) { + linkage := linkageIDOf(sh.id) + if linkage == sh.id || linkage == "" { + return + } + if parent, ok := s.shards[linkage]; ok { + parent.mu.RLock() + chose := parent.selfPinned || parent.restored + root, language, origin := parent.root, parent.language, parent.pinOrigin + parent.mu.RUnlock() + if chose { + // No parentSeeded here: a parent in memory that chose keeps saying + // so (selfPinned and restored are never cleared), and + // followConnectionShards asks it directly. + sh.root, sh.language, sh.pinOrigin = root, language, origin + } + return + } + root, language, origin, ok := s.loadPinForAgent(linkage) + if !ok { + return + } + if resolved, _, intact := s.restoreRootIntact(root); intact { + sh.root, sh.language, sh.pinOrigin = resolved, language, origin + // The parent is not in memory to be asked, so the child remembers. + sh.parentSeeded = true + } +} + +// followParentShard re-seeds, after conversation parentID moved its own shard +// off prevRoot, every subagent `/*` shard that never chose or +// restored a root of its own, wherever it currently sits. Seeding at creation alone left a subagent +// that had made any call before its parent re-pinned (a background subagent, a +// continued one) on the old root — another conversation's checkout — for good. +// The counterpart of followConnectionShards, for the parent's move. +// +// Lock order: shardsMu, then the parent's mu (read and RELEASED), then each +// child's mu in turn. No shard's mu is held while another's is taken. Called +// from repinAgent's post-unlock defer, so the parent's mu is not held here. +func (s *connSession) followParentShard(parentID, prevRoot string) { + if prevRoot == "" || linkageIDOf(parentID) != parentID { + return + } + s.shardsMu.Lock() + defer s.shardsMu.Unlock() + parent, ok := s.shards[parentID] + if !ok { + return + } + parent.mu.RLock() + root, language, origin := parent.root, parent.language, parent.pinOrigin + parent.mu.RUnlock() + for id, sh := range s.shards { + if id == parentID || linkageIDOf(id) != parentID { + continue + } + sh.mu.Lock() + // No root comparison: a subagent that never chose a root always follows + // its conversation. Comparing against prevRoot was NOT equivalent — + // a concurrent connection move could drag the subagent off prevRoot + // between the parent's move and this follow (followConnectionShards + // read parentChose before the parent committed), and the check would + // then strand it on the connection's new root (review of #535). + if sh.selfPinned || sh.restored { + sh.mu.Unlock() + continue + } + sh.root, sh.language, sh.pinOrigin = root, language, origin + sh.policy = s.buildAgentPolicy(root, language) + sh.readTracker.Reset() + sh.writeTracker.Reset() + sh.undoStore.Reset() + sh.mu.Unlock() + s.rehydrateReadsForAgent(sh, root) + } +} + +// parentChoseLocked reports whether id is a subagent whose conversation's +// in-memory shard chose or restored its root. Caller holds shardsMu and NO +// shard's mu; the parent's mu is taken and released here. +func (s *connSession) parentChoseLocked(id string) bool { + linkage := linkageIDOf(id) + if linkage == id { + return false + } + parent, ok := s.shards[linkage] + if !ok { + return false + } + parent.mu.RLock() + defer parent.mu.RUnlock() + return parent.selfPinned || parent.restored +} + +// pendingDeclarationForCall is pendingDeclarationFor as a CALL sees it: the +// agent's own refused declaration, or else — for a subagent that has not chosen +// a root of its own — its conversation's. Without the fallback a subagent of a +// refused parent sat on the connection seed, the very root its conversation had +// just been refused off, and resolved relative paths and git's default +// repository inside it. The second result reports that the marker is the +// conversation's, so the refusal can say so. +func (s *connSession) pendingDeclarationForCall(ctx context.Context) (p pendingDeclaration, inherited, ok bool) { + id := mcp.LogicalAgentFromCtx(ctx) + if p, ok := s.pendingDeclarationFor(id); ok { + return p, false, true + } + linkage := linkageIDOf(id) + if linkage == id { + return pendingDeclaration{}, false, false + } + p, ok = s.pendingDeclarationFor(linkage) + if !ok || s.agentChoseRoot(id) { + return pendingDeclaration{}, false, false + } + return p, true, true +} + +// agentChoseRoot reports whether id's shard holds a root the agent chose itself +// (self-pinned, or restored from its own persisted pin). Takes shardsMu then the +// shard's mu, the documented order; callers hold neither. +func (s *connSession) agentChoseRoot(id string) bool { + s.shardsMu.Lock() + sh, ok := s.shards[id] + s.shardsMu.Unlock() + if !ok { + return false + } + sh.mu.RLock() + defer sh.mu.RUnlock() + return sh.selfPinned || sh.restored +} diff --git a/internal/cli/conn_agent_shard_parent_test.go b/internal/cli/conn_agent_shard_parent_test.go new file mode 100644 index 00000000..0e214b35 --- /dev/null +++ b/internal/cli/conn_agent_shard_parent_test.go @@ -0,0 +1,315 @@ +package cli + +import ( + "context" + "fmt" + "strings" + "sync" + "testing" + "time" + + "github.com/plumbkit/plumb/internal/config" + "github.com/plumbkit/plumb/internal/mcp" + "github.com/plumbkit/plumb/internal/sessionstate" +) + +func idCtx(id string) context.Context { + return mcp.WithLogicalAgent(context.Background(), id) +} + +// newSharedConn builds a connection pinned to rootMain with the given agents +// declared at attach, so it is shared. +func newSharedConn(t *testing.T, s *connSession, rootMain string, agents ...string) { + t.Helper() + s.attachWorkspace(context.Background(), "file://"+rootMain) + if got := s.workspace(); got != rootMain { + t.Fatalf("precondition: connection pinned to %q, want %q", got, rootMain) + } + for _, a := range agents { + s.recordLogicalAgentAttach(a) + } +} + +// Issue #513 review B1: a hook-stamped subagent is admitted on its +// conversation's declaration, so it must work where its conversation chose to. +// Seeded from the connection pin instead, a subagent of a parent that had +// re-pinned itself to a worktree wrote into the other agent's checkout. +func TestSubagentInheritsItsConversationsChosenRoot(t *testing.T) { + t.Setenv("XDG_DATA_HOME", t.TempDir()) + s := newConnSession(context.Background(), detectTestPool(), nil, config.NewStore(config.Defaults()), nil, nil, newSharedBudgets()) + t.Cleanup(s.close) + rootMain, worktree := freshTempDir(t), freshTempDir(t) + mustGitDir(t, rootMain) + mustGitDir(t, worktree) + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+worktree, "", true, false); err != nil { + t.Fatalf("conv-y re-pin: %v", err) + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != worktree { + t.Errorf("conv-y's subagent resolves to %q, want its conversation's worktree %q", got, worktree) + } + // Control: a subagent whose conversation never chose a root sits where the + // connection (and so its conversation) sits. + if got := s.workspaceFor(idCtx("conv-x/sub")); got != rootMain { + t.Errorf("conv-x's subagent resolves to %q, want the connection's %q", got, rootMain) + } +} + +// After a daemon restart the parent's shard is not in memory until its next +// call; its persisted per-agent pin is the evidence of where it chose to work. +func TestSubagentInheritsItsConversationsPersistedRoot(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-parent-pin" + rootMain, worktree := freshTempDir(t), freshTempDir(t) + mustGitDir(t, rootMain) + mustGitDir(t, worktree) + if err := ss.UpsertPinForAgent(proxyID, "conv-y", worktree, "", sessionstate.PinSourceSessionStart); err != nil { + t.Fatalf("persist parent pin: %v", err) + } + s := newPersistSession(t, store, ss, proxyID) + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + + s.shardsMu.Lock() + _, parentLive := s.shards["conv-y"] + s.shardsMu.Unlock() + if parentLive { + t.Fatal("precondition: the parent's shard must not be in memory yet") + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != worktree { + t.Errorf("after a restart conv-y's subagent resolves to %q, want its conversation's persisted %q", got, worktree) + } + if got := s.workspaceFor(idCtx("conv-x/sub")); got != rootMain { + t.Errorf("control: conv-x's subagent resolves to %q, want the connection's %q", got, rootMain) + } +} + +// A subagent of a conversation whose declaration was REFUSED sits on the very +// root its conversation was refused off; it inherits the refusal until it +// chooses a root of its own. +func TestSubagentInheritsItsConversationsRefusedDeclaration(t *testing.T) { + t.Setenv("XDG_DATA_HOME", t.TempDir()) + s := newConnSession(context.Background(), detectTestPool(), nil, config.NewStore(config.Defaults()), nil, nil, newSharedBudgets()) + t.Cleanup(s.close) + rootMain, requested, own := freshTempDir(t), freshTempDir(t), freshTempDir(t) + for _, d := range []string{rootMain, requested, own} { + mustGitDir(t, d) + } + newSharedConn(t, s, rootMain, "conv-p", "conv-q") + s.markDeclarationRefused("conv-p", requested, rootMain) + + if got := s.workspaceFor(idCtx("conv-p/sub")); got != "" { + t.Errorf("a subagent of a refused conversation resolves to %q; it must resolve nowhere", got) + } + err := s.declarationRefusedErr(idCtx("conv-p/sub")) + if err == nil || !strings.Contains(err.Error(), `conversation "conv-p"`) { + t.Errorf("a subagent of a refused conversation must be refused naming it, got %v", err) + } + // Controls: another conversation's subagent is untouched. + if got := s.workspaceFor(idCtx("conv-q/sub")); got != rootMain { + t.Errorf("conv-q's subagent resolves to %q, want %q", got, rootMain) + } + if err := s.declarationRefusedErr(idCtx("conv-q/sub")); err != nil { + t.Errorf("conv-q's subagent refused: %v", err) + } + // A subagent that chose its own root is no longer held by its parent. + if _, err := s.repinWorkspace(idCtx("conv-p/sub"), "file://"+own, "", true, false); err != nil { + t.Fatalf("subagent re-pin: %v", err) + } + if got := s.workspaceFor(idCtx("conv-p/sub")); got != own { + t.Errorf("a subagent that chose %q resolves to %q", own, got) + } + if err := s.declarationRefusedErr(idCtx("conv-p/sub")); err != nil { + t.Errorf("a subagent that chose its own root is still refused: %v", err) + } +} + +func newThreeRootConn(t *testing.T) (s *connSession, rootMain, worktree, other string) { + t.Helper() + t.Setenv("XDG_DATA_HOME", t.TempDir()) + s = newConnSession(context.Background(), detectTestPool(), nil, config.NewStore(config.Defaults()), nil, nil, newSharedBudgets()) + t.Cleanup(s.close) + rootMain, worktree, other = freshTempDir(t), freshTempDir(t), freshTempDir(t) + for _, d := range []string{rootMain, worktree, other} { + mustGitDir(t, d) + } + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + return s, rootMain, worktree, other +} + +// Round-2 review of #535, N1: a subagent whose shard already exists when its +// conversation re-pins itself (a background subagent, a continued one) must +// follow — seeding only at creation left it in the other agent's checkout. +// A subagent that chose a root of its own stays where it chose. +func TestSubagentFollowsItsConversationsLaterRepin(t *testing.T) { + s, rootMain, worktree, other := newThreeRootConn(t) + if got := s.workspaceFor(idCtx("conv-y/sub")); got != rootMain { + t.Fatalf("precondition: the subagent starts on the connection root, got %q", got) + } + // Another conversation's subagent, on the same root BEFORE the move. + _ = s.workspaceFor(idCtx("conv-x/sub")) + if _, err := s.repinWorkspace(idCtx("conv-y/own"), "file://"+other, "", true, false); err != nil { + t.Fatalf("subagent's own re-pin: %v", err) + } + // A subagent that explicitly chose the root it was seeded on: no root + // moves, but it is now its own choice, so it must not follow either. + if _, err := s.repinWorkspace(idCtx("conv-y/here"), "file://"+rootMain, "", false, false); err != nil { + t.Fatalf("subagent's same-root declaration: %v", err) + } + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+worktree, "", true, false); err != nil { + t.Fatalf("conv-y re-pin: %v", err) + } + if got := s.workspaceFor(idCtx("conv-y/here")); got != rootMain { + t.Errorf("a subagent that declared %q followed its conversation to %q", rootMain, got) + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != worktree { + t.Errorf("a subagent created before its conversation's re-pin resolves to %q, want %q", got, worktree) + } + if got := s.workspaceFor(idCtx("conv-y/own")); got != other { + t.Errorf("a subagent that chose %q was moved to %q", other, got) + } + if got := s.workspaceFor(idCtx("conv-x/sub")); got != rootMain { + t.Errorf("another conversation's subagent moved to %q", got) + } + // And it keeps following: the conversation moves again. + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+rootMain, "", true, false); err != nil { + t.Fatalf("conv-y second re-pin: %v", err) + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != rootMain { + t.Errorf("after the second re-pin the subagent resolves to %q, want %q", got, rootMain) + } +} + +// N1b: a connection-scope move must not drag a subagent sitting on its +// conversation's CHOSEN root, even when that root equals the connection's. +func TestConnectionMoveDoesNotDragAParentSeededSubagent(t *testing.T) { + s, rootMain, worktree, other := newThreeRootConn(t) + // conv-y chooses rootMain (away and back), so it is self-pinned there. + for _, r := range []string{worktree, rootMain} { + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+r, "", true, false); err != nil { + t.Fatalf("conv-y re-pin to %s: %v", r, err) + } + } + _ = s.workspaceFor(idCtx("conv-y/sub")) // created from the parent's chosen root + _ = s.workspaceFor(idCtx("conv-x/seed")) // control: seeded from the connection + // A subagent created while ITS conversation sat on the seed; the + // conversation then declares that same root (no move, so nothing re-seeds + // the subagent) and so chose it. + _ = s.workspaceFor(idCtx("conv-z/early")) + s.recordLogicalAgentAttach("conv-z") + if _, err := s.repinWorkspace(idCtx("conv-z"), "file://"+rootMain, "", false, false); err != nil { + t.Fatalf("conv-z same-root declaration: %v", err) + } + if _, err := s.repinWorkspace(idCtx("conv-x"), "file://"+other, "", true, true); err != nil { + t.Fatalf("connection move: %v", err) + } + if got := s.workspace(); got != other { + t.Fatalf("precondition: the connection moved to %q, want %q", got, other) + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != rootMain { + t.Errorf("a connection move dragged conv-y's subagent to %q; its conversation is on %q", got, rootMain) + } + if got := s.workspaceFor(idCtx("conv-z/early")); got != rootMain { + t.Errorf("a connection move dragged conv-z's subagent to %q; its conversation chose %q", got, rootMain) + } + if got := s.workspaceFor(idCtx("conv-x/seed")); got != other { + t.Errorf("control: a connection-seeded subagent must follow the connection, got %q", got) + } +} + +// The same after a restart: the parent's shard is not in memory, its persisted +// pin seeded the subagent, and a connection move must not drag it. +func TestConnectionMoveDoesNotDragASubagentSeededFromAPersistedParent(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-parent-pin-move" + rootMain, other := freshTempDir(t), freshTempDir(t) + mustGitDir(t, rootMain) + mustGitDir(t, other) + if err := ss.UpsertPinForAgent(proxyID, "conv-y", rootMain, "", sessionstate.PinSourceSessionStart); err != nil { + t.Fatalf("persist parent pin: %v", err) + } + s := newPersistSession(t, store, ss, proxyID) + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + _ = s.workspaceFor(idCtx("conv-y/sub")) + if _, err := s.repinWorkspace(idCtx("conv-x"), "file://"+other, "", true, true); err != nil { + t.Fatalf("connection move: %v", err) + } + if got := s.workspaceFor(idCtx("conv-y/sub")); got != rootMain { + t.Errorf("a connection move dragged a subagent seeded from its conversation's persisted pin to %q", got) + } +} + +// A -race and deadlock guard for the follow paths: a conversation re-pins back +// and forth while its subagents make first calls and peers walk the shards. +func TestConcurrentParentRepinAndSubagentFirstCalls(t *testing.T) { + s, rootMain, worktree, _ := newThreeRootConn(t) + var wg sync.WaitGroup + wg.Add(3) + go func() { + defer wg.Done() + for i := range 40 { + r := worktree + if i%2 == 1 { + r = rootMain + } + _, _ = s.repinWorkspace(idCtx("conv-y"), "file://"+r, "", true, false) + } + }() + go func() { + defer wg.Done() + for i := range 400 { + ctx := idCtx(fmt.Sprintf("conv-y/sub-%d", i)) + _ = s.workspaceFor(ctx) + _ = s.declarationRefusedErr(ctx) + _ = s.checkBoundaryFor(ctx, "x.go", 0) + } + }() + go func() { + defer wg.Done() + for range 400 { + _ = s.pinnedPolicyGuard("/tmp/x") + _ = s.workspaceFor(idCtx("conv-x")) + } + }() + done := make(chan struct{}) + go func() { wg.Wait(); close(done) }() + select { + case <-done: + case <-time.After(90 * time.Second): + t.Fatal("deadlock: concurrent parent re-pins and subagent first calls did not finish") + } + // Settled: the conversation ended on rootMain (an even count of moves), and + // every subagent that never chose a root sits with it. + want := s.workspaceFor(idCtx("conv-y")) + for _, i := range []int{0, 199, 399} { + if got := s.workspaceFor(idCtx(fmt.Sprintf("conv-y/sub-%d", i))); got != want { + t.Errorf("conv-y/sub-%d resolves to %q, its conversation to %q", i, got, want) + } + } +} + +// A subagent whose root was restored from its OWN persisted pin chose it, and +// does not follow its conversation's later move. +func TestRestoredSubagentDoesNotFollowItsConversation(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-restored-sub" + rootMain, worktree := freshTempDir(t), freshTempDir(t) + mustGitDir(t, rootMain) + mustGitDir(t, worktree) + if err := ss.UpsertPinForAgent(proxyID, "conv-y/r", rootMain, "", sessionstate.PinSourceSessionStart); err != nil { + t.Fatalf("persist subagent pin: %v", err) + } + s := newPersistSession(t, store, ss, proxyID) + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + _ = s.workspaceFor(idCtx("conv-y/r")) + _ = s.workspaceFor(idCtx("conv-y/s")) // control: seeded, follows + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+worktree, "", true, false); err != nil { + t.Fatalf("conv-y re-pin: %v", err) + } + if got := s.workspaceFor(idCtx("conv-y/r")); got != rootMain { + t.Errorf("a subagent restored onto %q followed its conversation to %q", rootMain, got) + } + if got := s.workspaceFor(idCtx("conv-y/s")); got != worktree { + t.Errorf("control: a seeded subagent must follow, got %q", got) + } +} diff --git a/internal/cli/conn_attach.go b/internal/cli/conn_attach.go index 4bbb5700..b1abc57f 100644 --- a/internal/cli/conn_attach.go +++ b/internal/cli/conn_attach.go @@ -203,7 +203,10 @@ func explicitOrAutoAttach(explicit, autoAttach bool) bool { // onBeforeTool resolves the workspace root from the tool arguments when the // session has no primary workspace yet. Applies auto-attach and auto-attach- // persist when configured. -func (s *connSession) onBeforeTool(toolCtx context.Context, _ string, args json.RawMessage) { +func (s *connSession) onBeforeTool(toolCtx context.Context, name string, args json.RawMessage) { + // Runs only for ADMITTED calls: the refusal hook short-circuits dispatch + // before OnBeforeTool, so a refused call refreshes nothing. + s.refreshDeclaration(toolCtx, name) // Before the attached-already short circuit: an ALREADY attached session is // exactly the one whose primary can be stale after a live `enable-lsp`. // Generation-gated, so this is one atomic load on the steady-state path. diff --git a/internal/cli/conn_logical_agent.go b/internal/cli/conn_logical_agent.go index e13778e1..ec4558ff 100644 --- a/internal/cli/conn_logical_agent.go +++ b/internal/cli/conn_logical_agent.go @@ -22,6 +22,7 @@ import ( "github.com/plumbkit/plumb/internal/mcp" "github.com/plumbkit/plumb/internal/session" + "github.com/plumbkit/plumb/internal/stats" "github.com/plumbkit/plumb/internal/tools" ) @@ -36,6 +37,19 @@ type logicalAgentState struct { // declared for the connection's life, so a later re-check cannot un-see it // and flip the shared flag back off. seen map[string]struct{} + // declared is the set of LINKAGES (linkageIDOf) that declared themselves + // through session_start on this connection — its session_id argument, or + // the per-call identity a successful session_start ran under. It is the + // evidence refuse asks for before it lets a per-call identity route a + // state-changing call to a per-agent shard (issue #513): an identity nobody + // declared would otherwise get a fresh shard seeded from the connection's + // root, and a relative write would land in whatever workspace the connection + // holds. Keyed on the linkage so a hook-stamped subagent `/`, + // which usually never calls session_start itself, rides its parent's + // declaration. Like seen, it only grows. The value is when this process + // last refreshed the durable row (zero: never, e.g. restored after a + // restart); see refreshDue. + declared map[string]time.Time } // record commits an identity and reports the connection's shared STATE and, @@ -159,6 +173,11 @@ func shortIDPrefix(s string) string { func (l *logicalAgentState) sharedWith(id string) bool { l.mu.Lock() defer l.mu.Unlock() + return l.sharedWithLocked(id) +} + +// sharedWithLocked is sharedWith for a caller already holding l.mu. +func (l *logicalAgentState) sharedWithLocked(id string) bool { if len(l.seen) > 1 { return true } @@ -179,10 +198,29 @@ func (l *logicalAgentState) armed(id string) bool { return l.sharedWith(id) } +// refusalKind is why the ceiling refuses a call, or refusalNone. +type refusalKind int + +const ( + refusalNone refusalKind = iota + // refusalAnonymous: a shared connection and no per-call identity. + refusalAnonymous + // refusalUndeclared: a per-call identity whose linkage never declared + // itself through session_start on this connection (issue #513). + refusalUndeclared +) + // refuse reports whether a call declaring callID must be refused on this -// connection: the connection is shared (two or more distinct IDs observed) and -// the call is unattributable (no per-call ID). A non-shared connection needs no -// ID — the connection itself is the identity. +// connection. See refusal for the rule. +func (l *logicalAgentState) refuse(callID string) bool { + return l.refusal(callID) != refusalNone +} + +// refusal decides the fail-closed ceiling for a call declaring callID. +// +// Anonymous: refused once the connection is shared (two or more distinct IDs +// observed). A non-shared connection needs no ID — the connection itself is the +// identity. // // PLAN-394 removed the attach-time fallback from this decision. Before it, an // anonymous call was admitted whenever ANY session_start had attached — and @@ -191,22 +229,58 @@ func (l *logicalAgentState) armed(id string) bool { // force-pin, in the peer's project. Admitting a call on the strength of an // identity it did not present is attribution by guesswork; on a shared // connection only a presented ID admits a state-changing call. -func (l *logicalAgentState) refuse(callID string) bool { +// +// Identified: refused when the connection is shared COUNTING THE CALLER — the +// exact predicate shardFor routes on, so this refuses precisely the calls that +// would be served a per-agent shard — and the id's linkage was never declared +// through session_start (issue #513). Presenting an id is not the same as +// having declared one: an id nothing declared (a model typing `plumb_agent`, +// a client sending _meta it never announced) got a fresh shard seeded from the +// connection's root, so its relative write landed in whichever workspace the +// connection held — the misroute an anonymous call is refused for. The linkage, +// not the full id, is what must be declared, because a hook-stamped subagent +// `/` rarely calls session_start and is vouched for by ``. +// session_start itself is not state-changing, so the remedy stays reachable. +// Exempt: a connection where every observed identity, the caller's included, +// shares one linkage (see oneConversationLocked's call site). +func (l *logicalAgentState) refusal(callID string) refusalKind { l.mu.Lock() defer l.mu.Unlock() - if len(l.seen) <= 1 { - return false + if !l.sharedWithLocked(callID) { + return refusalNone } // No exemption for a connection where no caller has stamped yet. That // exemption (PLAN-440) let every anonymous write through on Claude // desktop's connector, whose stamp was being dropped in transit, and a // worktree edit landed in another agent's checkout (2026-09-30). An // unattributable write is refused; the refusal names the remedy. - return callID == "" + if callID == "" { + return refusalAnonymous + } + linkage := linkageIDOf(callID) + if _, ok := l.declared[linkage]; ok { + return refusalNone + } + // One conversation is not a shared connection in the sense the gate + // guards: a main thread that never called session_start and its own + // hook-stamped subagents all answer to the same linkage, so there is no + // OTHER conversation's workspace a call could be misrouted into. This is + // the Claude Code CLI's everyday topology; refusing it would lock the main + // thread out the moment its first subagent made any call. An identity with + // a different linkage — an invented one included — breaks the condition, + // and from then on every undeclared linkage is refused. + if l.oneConversationLocked(linkage) { + return refusalNone + } + return refusalUndeclared } -// recordLogicalAgentAttach records a session_start.session_id identity. -func (s *connSession) recordLogicalAgentAttach(id string) { s.recordLogicalAgent(id) } +// recordLogicalAgentAttach records a session_start.session_id identity, which +// is also a declaration of its linkage. +func (s *connSession) recordLogicalAgentAttach(id string) { + s.declareLogicalAgent(id) + s.recordLogicalAgent(id) +} // recordLogicalAgentCall records a per-call identity (_meta or the stamp // argument). @@ -215,14 +289,16 @@ func (s *connSession) recordLogicalAgentCall(id string) { } // recordCall commits an identity that arrived on the PER-CALL channel; -// recordAttach one that arrived at attach time. Both commit it the same way; -// the two names keep call sites and tests explicit about the channel. +// recordAttach one that arrived at attach time. Both commit it to seen the same +// way; only the attach channel is session_start, so only it also declares the +// linkage (issue #513). func (l *logicalAgentState) recordCall(id string) (shared, transition bool) { return l.record(id) } // recordAttach commits an attach-time identity. See recordCall. func (l *logicalAgentState) recordAttach(id string) (shared, transition bool) { + l.declare(id) return l.record(id) } @@ -312,13 +388,14 @@ func (s *connSession) markSharedConnectionDetected() { return } info.Health = "shared_connection_detected" - info.HealthMessage = "multiple logical agents share this connection; per-agent state is isolated, and a state-changing call carrying no identity is refused — " + sharedIdentityRemedy + info.HealthMessage = "multiple logical agents share this connection; per-agent state is isolated, and a state-changing call carrying no identity, or one no session_start declared, is refused — " + sharedIdentityRemedy }) } // refuseSharedStateChange is the fail-closed ceiling. It refuses a mutating // tool call that arrives on a shared connection without a trustworthy -// logical-agent identity, naming the supported topology and its remedy. Read +// logical-agent identity — none at all, or one whose linkage no session_start +// declared (issue #513) — naming the supported topology and its remedy. Read // calls are never refused: sharing read-only state is safe, and the acceptance // contract is about state-changing operations resetting a peer's pin, trackers, // rate budget, undo state or language. @@ -326,8 +403,11 @@ func (s *connSession) refuseSharedStateChange(_ context.Context, name, logicalAg if !slices.Contains(tools.StateChangingToolNames(), name) { return nil } - if !s.logicalAgents.refuse(logicalAgent) { + switch s.logicalAgents.refusal(logicalAgent) { + case refusalNone: return nil + case refusalUndeclared: + return undeclaredIdentityErr(name, logicalAgent) } // [collab] allow_unidentified_writes used to lift this refusal. It is no // longer honoured: with it set, an unattributable write resolved through @@ -340,6 +420,26 @@ func (s *connSession) refuseSharedStateChange(_ context.Context, name, logicalAg return fmt.Errorf("shared connection: %s is a state-changing call with no logical-agent identity, so it cannot be attributed to one of the agents multiplexing this connection, and plumb will not guess whose workspace it belongs to — %s%s", name, sharedIdentityRemedy, retired) } +// undeclaredIdentityErr is the refusal for a per-call identity whose linkage no +// session_start on this connection declared (issue #513). The session_id it +// names is the LINKAGE — for a hook-stamped subagent `/` that is +// its conversation's id, not the stamp — so the wording says which. +// +// It also asks for `workspace`. A session_start that only declares leaves the +// caller on whatever root its shard was seeded with, usually the connection's, +// which on a shared connection may be another agent's checkout: the declaration +// would then admit exactly the misrouted write this refusal stopped. +func undeclaredIdentityErr(name, logicalAgent string) error { + stamp := stats.SanitiseAgentID(logicalAgent) + linkage := linkageIDOf(logicalAgent) + which := "the identity you are stamping" + if linkage != logicalAgent { + which = "your conversation's id, the part of your stamp before `/`; a subagent is covered once its conversation has declared itself" + } + return fmt.Errorf("shared connection: %s carries the logical-agent identity %q, but no session_start on this connection has declared it, so plumb cannot tell which workspace it belongs to and will not guess — call session_start with session_id %q (%s) and workspace set to the absolute path of the project you are working in, then retry. Without workspace, your relative paths may resolve against the connection's root, which can be another agent's checkout", + name, stamp, stats.SanitiseAgentID(linkage), which) +} + // sharedIdentityRemedy names the ways an agent on a shared connection gets an // identity, cheapest first. Identity comes before topology on purpose // (PLAN-417): the previous wording led with "one plumb serve per agent", the diff --git a/internal/cli/conn_logical_agent_declared.go b/internal/cli/conn_logical_agent_declared.go new file mode 100644 index 00000000..f9891439 --- /dev/null +++ b/internal/cli/conn_logical_agent_declared.go @@ -0,0 +1,238 @@ +package cli + +// conn_logical_agent_declared.go — which conversations have DECLARED +// themselves on a connection through session_start, as distinct from which +// identities have merely been observed on it (issue #513). +// +// The fail-closed ceiling (conn_logical_agent.go) asks both questions. Seen +// decides whether the connection is shared; declared decides whether a +// per-call identity on a shared connection may be served a per-agent shard. +// Split from conn_logical_agent.go by responsibility and to keep it under the +// file-size cap. + +import ( + "context" + "slices" + "strings" + "time" + + "github.com/plumbkit/plumb/internal/mcp" + "github.com/plumbkit/plumb/internal/tools" +) + +// declare commits id's linkage as declared through session_start, and reports +// the linkage it committed ("" for a blank id). An existing entry keeps its +// refresh time. +func (l *logicalAgentState) declare(id string) string { + linkage := linkageIDOf(strings.TrimSpace(id)) + if linkage == "" { + return "" + } + l.mu.Lock() + defer l.mu.Unlock() + if l.declared == nil { + l.declared = make(map[string]time.Time) + } + if _, ok := l.declared[linkage]; !ok { + l.declared[linkage] = time.Time{} + } + return linkage +} + +// restoreDeclared commits linkages recovered from durable evidence after a +// daemon restart. Additive, like seed: a durable view is a lower bound. +func (l *logicalAgentState) restoreDeclared(linkages []string) { + for _, linkage := range linkages { + l.declare(linkage) + } +} + +// oneConversationLocked reports whether every identity observed on the +// connection answers to linkage. Caller holds l.mu. +func (l *logicalAgentState) oneConversationLocked(linkage string) bool { + for id := range l.seen { + if linkageIDOf(id) != linkage { + return false + } + } + return true +} + +// markRefreshed records that linkage's durable row was written at now. +func (l *logicalAgentState) markRefreshed(linkage string, now time.Time) { + l.mu.Lock() + defer l.mu.Unlock() + if _, ok := l.declared[linkage]; ok { + l.declared[linkage] = now + } +} + +// refreshDue reports whether linkage is declared and its durable row was last +// refreshed at least every ago, and if so claims the refresh (stamps now), so +// concurrent calls do not all write. It also returns the stamp it replaced, so +// a failed write can hand the slot back (releaseRefresh). An undeclared linkage +// is never due. +func (l *logicalAgentState) refreshDue(linkage string, now time.Time, every time.Duration) (prev time.Time, due bool) { + l.mu.Lock() + defer l.mu.Unlock() + last, ok := l.declared[linkage] + if !ok || now.Sub(last) < every { + return time.Time{}, false + } + l.declared[linkage] = now + return last, true +} + +// releaseRefresh hands back a refresh slot claimed at claimed, restoring prev, +// so a write that failed is retried on the next call rather than an interval +// later. A no-op if another call has claimed the slot since. +func (l *logicalAgentState) releaseRefresh(linkage string, claimed, prev time.Time) { + l.mu.Lock() + defer l.mu.Unlock() + if cur, ok := l.declared[linkage]; ok && cur.Equal(claimed) { + l.declared[linkage] = prev + } +} + +// declareLogicalAgent commits id's linkage as declared through session_start, +// in memory and durably under the proxy session (see restoreDeclaredLinkages). +// Only session_start's success path may call it: a refused call must leave no +// declaration behind. +func (s *connSession) declareLogicalAgent(id string) { + linkage := s.logicalAgents.declare(id) + if linkage == "" || s.sessionState == nil { + return + } + v := s.view() + if !v.session.PersistState || v.proxySessionID == "" { + return + } + if err := s.sessionState.RecordDeclaredLinkage(v.proxySessionID, linkage); err != nil { + s.log().Debug("daemon: recording the declared linkage failed", "linkage", logicalAgentLabel(linkage), "err", err) + return + } + s.logicalAgents.markRefreshed(linkage, time.Now()) +} + +// refreshDeclaration keeps a working conversation's durable declaration young. +// +// A declared_linkage row is otherwise written only by session_start. The idle +// reaper reclaims a row older than [session] persist_state_ttl_minutes (24 h by +// default) unless its proxy session is connected at that pass, and nothing is +// pruned at daemon start (#525). So a connected serve keeps its declarations +// however old, across restarts it reconnects through. The exemption does not +// cover a serve that is between connections when a pass runs, such as one +// reconnecting after its transport dropped, and there a conversation that +// declared once and then worked for a day would lose its row and be refused. +// +// Kept for that window. An ADMITTED state-changing call from an already-declared +// linkage refreshes the row, at most once per declarationRefreshEvery +// (min(TTL/4, 1 h)), so a working conversation's row is never old enough for +// such a pass to reclaim it. The cost is one UPDATE per linkage per interval. +// Unlike a pin's or a logical agent's, a declaration's updated_at is no other +// evidence (LogicalAgentIDsFor windows by the pin and logical-agent rows, +// never by it), so refreshing it skews nothing — the reason #525 gave for not +// refreshing every live row does not apply. +// +// UPDATE only, never insert: admission is not declaration. A call admitted +// under the one-conversation exemption must not become durable evidence that +// its linkage declared itself. +func (s *connSession) refreshDeclaration(ctx context.Context, toolName string) { + id := mcp.LogicalAgentFromCtx(ctx) + if id == "" || s.sessionState == nil || !slices.Contains(tools.StateChangingToolNames(), toolName) { + return + } + v := s.view() + if !v.session.PersistState || v.proxySessionID == "" { + return + } + linkage := linkageIDOf(id) + now := time.Now() + prev, due := s.logicalAgents.refreshDue(linkage, now, declarationRefreshEvery(v.session.PersistStateTTLMinutes)) + if !due { + return + } + if err := s.sessionState.TouchDeclaredLinkage(v.proxySessionID, linkage); err != nil { + s.logicalAgents.releaseRefresh(linkage, now, prev) + s.log().Debug("daemon: refreshing the declared linkage failed", "linkage", logicalAgentLabel(linkage), "err", err) + } +} + +// declarationRefreshEvery is how often an active conversation refreshes its +// durable declaration: a quarter of the TTL, capped at an hour, so a row is +// refreshed several times inside any TTL without a write per call. +func declarationRefreshEvery(ttlMinutes int) time.Duration { + every := time.Hour + if q := time.Duration(ttlMinutes) * time.Minute / 4; ttlMinutes > 0 && q < every { + every = q + } + return every +} + +// declareSessionStartCaller declares the per-call identity a SUCCESSFUL +// session_start ran under. linkExternalID already declares the session_id +// argument; this covers a client that stamps every call (a per-call _meta, or +// the hook's argument) and calls session_start without a session_id, which +// sharedIdentityRemedy sanctions. A failed session_start declares nothing — an +// agent whose re-pin was refused never attached. +func (s *connSession) declareSessionStartCaller(ctx context.Context, toolName string, isError bool) { + if toolName != "session_start" || isError { + return + } + if id := mcp.LogicalAgentFromCtx(ctx); id != "" { + s.declareLogicalAgent(id) + } +} + +// restoreDeclaredLinkages brings back, after a daemon restart, which +// conversations had declared themselves on this connection through +// session_start — the evidence refusal needs before it admits a stamped +// state-changing call on a shared connection (issue #513). +// +// Wired, unlike seedLogicalAgentsFromState, and the difference is the direction +// each can move the gate. That seed ADDS identities to seen, which ARMS the +// ceiling and so can refuse more — it locked out a client that cannot stamp. +// This only adds to declared, which can only ADMIT more: a declared linkage +// never causes a refusal, so restoring one cannot lock anybody out. Not +// restoring was the lockout: seen refills from the per-call stamps the hook +// keeps sending, the connection is shared again at the second identity, and +// every agent that declared before the restart — none of which will call +// session_start again unprompted — would have each state-changing call refused +// until it did. +// +// The sources are the declared_linkage rows and the identity record's own +// external linkage, both keyed on the proxy session ID. That ID is the serve +// process's own 122-bit secret (see inheritSessionID), so presenting it is +// proof of being the connection that made those declarations. The identity +// record is included because Prune never reclaims it, while declared_linkage +// ages out with the TTL like every other expendable row of a serve that is not +// connected when the reaper passes; logical_agent is +// deliberately NOT a source — it records every OBSERVED identity, so an +// invented id that made one admitted read would come back declared. +// +// Retried on a converged degraded recovery (retryRestoreIdentity): the record +// and any declaration written while the connection sat degraded are only in +// hand then. +// +// KNOWN LIMITS, both of which cost one refusal whose remedy (session_start) +// re-declares: +// - With [session] persist_state off nothing is saved, so nothing comes back. +// - declared_linkage ages out with persist_state_ttl_minutes, but only when an +// idle-reaper pass finds the serve disconnected (connected sessions are +// exempt, and nothing is pruned at daemon start, #525). refreshDeclaration +// keeps a WORKING conversation's row young (see its slack); one idle past +// the TTL at such a pass, that is not the identity record's linkage, comes +// back undeclared. +func (s *connSession) restoreDeclaredLinkages(proxySessionID string) { + if s.sessionState == nil || !s.view().session.PersistState || proxySessionID == "" { + return + } + linkages, err := s.sessionState.DeclaredLinkagesFor(proxySessionID) + if err != nil { + s.log().Debug("daemon: restoring declared linkages failed", "err", err) + } + if ext := s.view().persistedIdentity.ExternalID; ext != "" { + linkages = append(linkages, ext) + } + s.logicalAgents.restoreDeclared(linkages) +} diff --git a/internal/cli/conn_logical_agent_declared_test.go b/internal/cli/conn_logical_agent_declared_test.go new file mode 100644 index 00000000..a0f547ee --- /dev/null +++ b/internal/cli/conn_logical_agent_declared_test.go @@ -0,0 +1,318 @@ +package cli + +import ( + "context" + "encoding/json" + "slices" + "strings" + "testing" + "time" + + "github.com/plumbkit/plumb/internal/mcp" + "github.com/plumbkit/plumb/internal/tools" +) + +// Issue #513: on a shared connection, an identity nobody declared through +// session_start used to be admitted, handed a fresh shard seeded from the +// connection's root, and its relative write landed in the connection's +// workspace. The gate now asks whether the identity's LINKAGE was declared. +func TestRefuseUndeclaredIdentityOnSharedConnection(t *testing.T) { + var s connSession + s.recordLogicalAgentAttach("conv-a") // session_start(session_id: conv-a) + s.recordLogicalAgentAttach("conv-b") // session_start(session_id: conv-b) + + err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up") + if err == nil { + t.Fatal("a write under an identity no session_start declared was admitted on a shared connection") + } + // The refusal must name a remedy the caller can reach from inside a tool + // call: session_start is not state-changing, so it is never refused. + for _, want := range []string{"session_start", `"made-up"`, "no session_start on this connection has declared it"} { + if !strings.Contains(err.Error(), want) { + t.Errorf("refusal missing %q: %v", want, err) + } + } + if slices.Contains(tools.StateChangingToolNames(), "session_start") { + t.Fatal("session_start became state-changing; the remedy the refusal names is no longer reachable") + } + // Reads are never refused, declared or not. + if err := s.refuseSharedStateChange(context.Background(), "read_file", "made-up"); err != nil { + t.Errorf("a read under an undeclared identity must not be refused: %v", err) + } + + // Positive controls: the declared agents and a hook-stamped subagent of a + // declared conversation are admitted. Without them the refusal above could + // be the gate refusing every identified call. + for _, id := range []string{"conv-a", "conv-b", "conv-a/agent-1"} { + if err := s.refuseSharedStateChange(context.Background(), "write_file", id); err != nil { + t.Errorf("%s must be admitted: %v", id, err) + } + } + // A subagent whose CONVERSATION was never declared is not vouched for. + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up/agent-1"); err == nil { + t.Error("a subagent of an undeclared conversation was admitted") + } + + // The remedy works: once session_start runs under the identity, it is admitted. + s.declareSessionStartCaller(mcp.WithLogicalAgent(context.Background(), "made-up"), "session_start", false) + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err != nil { + t.Errorf("an identity declared by a successful session_start must be admitted: %v", err) + } +} + +// Only a SUCCESSFUL session_start declares: a failed one (its re-pin refused) +// never attached, and no other tool is a declaration. +func TestOnlyASuccessfulSessionStartDeclares(t *testing.T) { + var s connSession + s.recordLogicalAgentAttach("conv-a") + s.recordLogicalAgentAttach("conv-b") + ctx := mcp.WithLogicalAgent(context.Background(), "made-up") + + s.declareSessionStartCaller(ctx, "session_start", true) + s.declareSessionStartCaller(ctx, "read_file", false) + s.declareSessionStartCaller(ctx, "write_file", false) + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err == nil { + t.Fatal("a failed session_start or a non-session_start call declared the identity") + } + s.declareSessionStartCaller(ctx, "session_start", false) + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err != nil { + t.Fatalf("positive control: a successful session_start must declare: %v", err) + } +} + +// A single-agent connection is unaffected: whether or not the one identity it +// knows ever called session_start, and whatever identity a first call carries. +func TestUndeclaredIdentityOnSingleAgentConnectionIsAdmitted(t *testing.T) { + var none connSession + if err := none.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err != nil { + t.Errorf("the first identity a connection sees is the connection; it must not be refused: %v", err) + } + var one connSession + one.recordLogicalAgentCall("stamped-never-declared") + if err := one.refuseSharedStateChange(context.Background(), "write_file", "stamped-never-declared"); err != nil { + t.Errorf("a single-agent connection's own identity must not be refused: %v", err) + } + if err := one.refuseSharedStateChange(context.Background(), "write_file", ""); err != nil { + t.Errorf("a single-agent connection's anonymous call must not be refused: %v", err) + } +} + +// The refusal is judged against the caller COUNTED — the predicate shardFor +// routes on — so the second identity on a connection is refused from its first +// write rather than admitted once (onto a fresh shard of the connection's root) +// and refused only from its second. +func TestUndeclaredSecondIdentityIsRefusedFromItsFirstWrite(t *testing.T) { + var s connSession + s.recordLogicalAgentAttach("conv-a") + if !s.logicalAgents.sharedWith("made-up") { + t.Fatal("precondition: shardFor would give made-up a shard of its own") + } + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err == nil { + t.Fatal("the first write of an undeclared second identity was admitted") + } + if err := s.refuseSharedStateChange(context.Background(), "write_file", "conv-a/sub"); err != nil { + t.Errorf("positive control: a subagent of the declared conversation must be admitted: %v", err) + } + // A refusal commits nothing: the connection is still single-agent. + if got := s.logicalAgents.count(); got != 1 { + t.Errorf("the gate recorded an identity: %d observed, want 1", got) + } +} + +// A daemon restart must not lock out an agent that declared itself before it. +// The in-memory declared set is gone, while seen refills from the stamps the +// hook keeps sending, so without restoration the connection turns shared again +// and every declared agent's write is refused. An identity that was only ever +// OBSERVED (an invented id's admitted read) must not come back declared. +func TestDeclarationsSurviveADaemonRestart(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-restart" + + before := newPersistSession(t, store, ss, proxyID) + before.linkExternalID("conv-a") // session_start(session_id: conv-a) + before.declareLogicalAgent("conv-b") // session_start under conv-b's stamp + before.recordLogicalAgentCall("made-up") + before.close() + + after := newPersistSession(t, store, ss, proxyID) // onProxySession restores + after.recordLogicalAgentCall("conv-a") + after.recordLogicalAgentCall("conv-b") + for _, id := range []string{"conv-a", "conv-b", "conv-a/agent-1"} { + if err := after.refuseSharedStateChange(context.Background(), "write_file", id); err != nil { + t.Errorf("%s declared before the restart and was refused after it: %v", id, err) + } + } + if err := after.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err == nil { + t.Error("an identity only ever observed came back declared after the restart") + } +} + +// declared_linkage ages out with the TTL when the idle reaper runs while its +// serve is not connected; the identity record's own linkage is never pruned, so +// the conversation the session is linked to stays declared even when such a +// pass reclaimed its row. +func TestIdentityRecordLinkageSurvivesAPrunedDeclaration(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-pruned" + + before := newPersistSession(t, store, ss, proxyID) + before.linkExternalID("conv-a") + before.declareLogicalAgent("conv-b") + before.close() + // Every expendable row is older than this cutoff; session_names is kept. + if err := ss.Prune(time.Now().Add(time.Hour)); err != nil { + t.Fatalf("prune: %v", err) + } + + after := newPersistSession(t, store, ss, proxyID) + after.recordLogicalAgentCall("conv-a") + after.recordLogicalAgentCall("conv-b") + if err := after.refuseSharedStateChange(context.Background(), "write_file", "conv-a"); err != nil { + t.Errorf("the identity record's linkage must survive a prune: %v", err) + } + // The documented limit, and the control that the prune really removed the + // declared_linkage row: conv-b must re-declare. + if err := after.refuseSharedStateChange(context.Background(), "write_file", "conv-b"); err == nil { + t.Error("conv-b's pruned declaration came back; the prune control is vacuous") + } +} + +// Review of #535, item 1: the Claude Code CLI's everyday topology. A main +// thread that never called session_start and its own hook-stamped subagents +// share one linkage, so there is no other conversation to misroute into and +// nothing is refused. An identity with a different linkage breaks that. +func TestOneConversationNeedsNoDeclaration(t *testing.T) { + var s connSession + s.recordLogicalAgentCall("conv") + s.recordLogicalAgentCall("conv/sub") // a subagent's read, admitted and recorded + for _, id := range []string{"conv", "conv/sub", "conv/sub2"} { + if err := s.refuseSharedStateChange(context.Background(), "write_file", id); err != nil { + t.Errorf("one conversation: %s refused: %v", id, err) + } + } + if err := s.refuseSharedStateChange(context.Background(), "write_file", ""); err == nil { + t.Error("an anonymous write on a shared connection must still refuse") + } + // An invented identity with another linkage is refused on its first write... + if err := s.refuseSharedStateChange(context.Background(), "write_file", "made-up"); err == nil { + t.Fatal("an invented identity was admitted on a one-conversation connection") + } + // ...and once it has been observed (an admitted read), the connection holds + // two conversations, so the undeclared main thread must declare too. + s.recordLogicalAgentCall("made-up") + if err := s.refuseSharedStateChange(context.Background(), "write_file", "conv"); err == nil { + t.Error("two conversations are observed; an undeclared one must be refused") + } + s.declareSessionStartCaller(idCtx("conv"), "session_start", false) + if err := s.refuseSharedStateChange(context.Background(), "write_file", "conv/sub"); err != nil { + t.Errorf("control: after conv declared, its subagent must be admitted: %v", err) + } +} + +// Review of #535, item 4: the session_id the refusal names is the LINKAGE, so +// for a subagent it must not claim to be the identity being stamped. +func TestUndeclaredRefusalNamesTheConversationForASubagent(t *testing.T) { + var s connSession + s.recordLogicalAgentAttach("conv-a") + s.recordLogicalAgentAttach("conv-b") + sub := s.refuseSharedStateChange(context.Background(), "write_file", "made-up/agent-1") + if sub == nil { + t.Fatal("precondition: the subagent must be refused") + } + for _, want := range []string{`"made-up/agent-1"`, `session_id "made-up"`, "your conversation's id"} { + if !strings.Contains(sub.Error(), want) { + t.Errorf("subagent refusal missing %q: %v", want, sub) + } + } + if strings.Contains(sub.Error(), "the identity you are stamping") { + t.Errorf("subagent refusal calls the linkage the stamped identity: %v", sub) + } + plain := s.refuseSharedStateChange(context.Background(), "write_file", "made-up") + if plain == nil || !strings.Contains(plain.Error(), "the identity you are stamping") { + t.Errorf("control: a plain id's refusal names it as the stamped identity: %v", plain) + } + // Review N5: declaring alone leaves the caller on the connection's root, + // possibly another agent's checkout, so both refusals ask for workspace too. + for _, err := range []error{sub, plain} { + if err != nil && !strings.Contains(err.Error(), "and workspace set to the absolute path of the project you are working in") { + t.Errorf("the refusal does not ask for workspace: %v", err) + } + } +} + +// Review of #535, item 3: a declared_linkage row is written by session_start, +// and the idle reaper reclaims it once it is older than the TTL if its serve is +// not connected at that pass (one between connections, say; nothing is pruned at +// daemon start since #525). So a conversation that keeps working must keep its +// row young. An admitted state-changing call from a declared linkage refreshes +// it; a read does not; an undeclared linkage never gains a row (update, never +// insert). The Prune below, with no live sessions, is such a pass. +func TestAdmittedWritesKeepADeclarationYoung(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-refresh" + + before := newPersistSession(t, store, ss, proxyID) + before.linkExternalID("conv-a") + before.declareLogicalAgent("conv-b") + before.close() + if err := ss.BackdateLogicalAgents(proxyID, time.Now().Add(-48*time.Hour)); err != nil { + t.Fatalf("backdate: %v", err) + } + + after := newPersistSession(t, store, ss, proxyID) + for _, id := range []string{"conv-a", "conv-b", "made-up"} { + after.recordLogicalAgentCall(id) + } + after.onBeforeTool(idCtx("conv-a/sub"), "write_file", json.RawMessage(`{}`)) + after.onBeforeTool(idCtx("conv-b"), "read_file", json.RawMessage(`{}`)) + after.onBeforeTool(idCtx("made-up"), "write_file", json.RawMessage(`{}`)) + + if err := ss.Prune(time.Now().Add(-24 * time.Hour)); err != nil { + t.Fatalf("prune: %v", err) + } + got, err := ss.DeclaredLinkagesFor(proxyID) + if err != nil { + t.Fatalf("DeclaredLinkagesFor: %v", err) + } + if len(got) != 1 || got[0] != "conv-a" { + t.Errorf("declarations after prune = %v, want only conv-a (refreshed by its subagent's write); conv-b only read, made-up never declared", got) + } +} + +// Review of #535, item 6: a degraded identity restore converges later, on the +// bounded retry. Declarations written while the connection sat degraded (here, +// by the predecessor that was still detaching) must be restored then too. +func TestDeclarationsRestoreWhenADegradedRecoveryConverges(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-degraded" + first := newPersistSession(t, store, ss, proxyID) + + release := make(chan struct{}) + overlapping := newPersistSessionWithBackoff(t, store, ss, proxyID, func(int) time.Duration { + <-release + return 0 + }) + if overlapping.recovery() != recoveryDegraded { + close(release) + t.Fatal("precondition: the overlap did not degrade") + } + first.linkExternalID("conv-late") + first.close() + close(release) + + deadline := time.Now().Add(5 * time.Second) + for overlapping.recovery() != recoveryRestored { + if time.Now().After(deadline) { + t.Fatal("the degraded connection never converged") + } + time.Sleep(10 * time.Millisecond) + } + overlapping.recordLogicalAgentCall("conv-late") + overlapping.recordLogicalAgentCall("other") + if err := overlapping.refuseSharedStateChange(context.Background(), "write_file", "conv-late"); err != nil { + t.Errorf("a declaration made while degraded was not restored on convergence: %v", err) + } + if err := overlapping.refuseSharedStateChange(context.Background(), "write_file", "other"); err == nil { + t.Error("control: an undeclared identity must be refused, or the connection is not shared") + } +} diff --git a/internal/cli/conn_logical_agent_refresh_test.go b/internal/cli/conn_logical_agent_refresh_test.go new file mode 100644 index 00000000..2b0dc4cb --- /dev/null +++ b/internal/cli/conn_logical_agent_refresh_test.go @@ -0,0 +1,138 @@ +package cli + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/plumbkit/plumb/internal/config" + "github.com/plumbkit/plumb/internal/sessionstate" +) + +// The refresh is throttled: N admitted state-changing calls inside one +// interval write the durable row ONCE, not N times. The row itself is the +// counter: after the first call refreshed it, it is aged again, and only a +// further write could bring it back from the prune. +func TestDeclarationRefreshWritesOncePerInterval(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-throttle" + before := newPersistSession(t, store, ss, proxyID) + before.linkExternalID("conv-a") + before.close() + aged := time.Now().Add(-48 * time.Hour) + if err := ss.BackdateLogicalAgents(proxyID, aged); err != nil { + t.Fatalf("backdate: %v", err) + } + + after := newPersistSession(t, store, ss, proxyID) // restored, never refreshed + after.onBeforeTool(idCtx("conv-a"), "write_file", json.RawMessage(`{}`)) + if got, _ := ss.DeclaredLinkagesFor(proxyID); len(got) != 1 { + t.Fatalf("precondition: %v", got) + } + if err := ss.Prune(time.Now().Add(-24 * time.Hour)); err != nil { + t.Fatalf("prune: %v", err) + } + if got, _ := ss.DeclaredLinkagesFor(proxyID); len(got) != 1 { + t.Fatal("positive control: the first admitted write must refresh the row") + } + + if err := ss.BackdateLogicalAgents(proxyID, aged); err != nil { + t.Fatalf("backdate: %v", err) + } + for range 20 { + after.onBeforeTool(idCtx("conv-a"), "write_file", json.RawMessage(`{}`)) + } + if err := ss.Prune(time.Now().Add(-24 * time.Hour)); err != nil { + t.Fatalf("prune: %v", err) + } + if got, _ := ss.DeclaredLinkagesFor(proxyID); len(got) != 0 { + t.Errorf("20 writes inside one refresh interval touched the row again: %v", got) + } +} + +func TestRefreshDueClaimsOneSlotPerInterval(t *testing.T) { + var l logicalAgentState + l.declare("conv") + now := time.Now() + due := 0 + for range 50 { + if _, ok := l.refreshDue("conv", now, time.Hour); ok { + due++ + } + } + if due != 1 { + t.Fatalf("%d refreshes claimed inside one interval, want 1", due) + } + if _, ok := l.refreshDue("conv", now.Add(time.Hour), time.Hour); !ok { + t.Error("the next interval must be due again") + } + if _, ok := l.refreshDue("undeclared", now, time.Hour); ok { + t.Error("an undeclared linkage is never due") + } +} + +// A failed refresh hands its slot back, so the next call retries instead of +// waiting out the interval. +func TestFailedRefreshReleasesItsSlot(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-513-release" + s := newPersistSession(t, store, ss, proxyID) + s.logicalAgents.restoreDeclared([]string{"conv-a"}) // never refreshed: due + ss.Close() // every write now fails + + s.onBeforeTool(idCtx("conv-a"), "write_file", json.RawMessage(`{}`)) + if _, ok := s.logicalAgents.refreshDue("conv-a", time.Now(), time.Hour); !ok { + t.Error("a failed refresh kept its slot; the next call would not retry for an interval") + } +} + +// Round-2 review of #535: the LEGACY-HEAL branch of the bounded retry (a v3 +// record with a name and no session ID, whose name was held by a live session +// at reconnect) converges too, and must restore declarations the same way. +func TestDeclarationsRestoreWhenALegacyHealConverges(t *testing.T) { + t.Setenv("XDG_DATA_HOME", t.TempDir()) + store := config.NewStore(config.Defaults()) + ss := openStateStore(t) + // A live session holds the name; it persists nothing, so it reserves nothing. + holder := newConnSession(context.Background(), detectTestPool(), nil, store, nil, nil, newSharedBudgets()) + t.Cleanup(holder.close) + if _, err := holder.renameSession("legacy-stag"); err != nil { + t.Fatalf("holder rename: %v", err) + } + const proxyID = "proxy-513-legacy" + if err := ss.SaveIdentity(proxyID, sessionstate.Identity{Name: "legacy-stag"}); err != nil { + t.Fatalf("save legacy record: %v", err) + } + + release := make(chan struct{}) + s := newPersistSessionWithBackoff(t, store, ss, proxyID, func(int) time.Duration { + <-release + return 0 + }) + if s.recovery() != recoveryDegraded { + close(release) + t.Fatalf("precondition: the held legacy name must degrade the restore, got %q", s.recovery()) + } + if err := ss.RecordDeclaredLinkage(proxyID, "conv-late"); err != nil { + t.Fatalf("record declaration: %v", err) + } + holder.close() + close(release) + + deadline := time.Now().Add(5 * time.Second) + for s.recovery() != recoveryEstablished { + if time.Now().After(deadline) { + t.Fatalf("the legacy heal never converged (recovery %q)", s.recovery()) + } + time.Sleep(10 * time.Millisecond) + } + s.recordLogicalAgentCall("conv-late") + s.recordLogicalAgentCall("other") + if err := s.refuseSharedStateChange(context.Background(), "write_file", "conv-late"); err != nil { + t.Errorf("a declaration written while degraded was not restored by the legacy heal: %v", err) + } + if err := s.refuseSharedStateChange(context.Background(), "write_file", "other"); err == nil { + t.Error("control: an undeclared identity must be refused") + } +} diff --git a/internal/cli/conn_logical_agent_seed_test.go b/internal/cli/conn_logical_agent_seed_test.go index c6b29ab4..cc38f463 100644 --- a/internal/cli/conn_logical_agent_seed_test.go +++ b/internal/cli/conn_logical_agent_seed_test.go @@ -21,6 +21,9 @@ import ( func TestSeedArmsTheGateBeforeAnyAgentRedeclares(t *testing.T) { var l logicalAgentState l.seed([]string{"coordinator", "subagent"}) + // The coordinator's declaration is restored alongside (restoreDeclared); + // an undeclared identity would be refused, see issue #513. + l.restoreDeclared([]string{"coordinator"}) if !l.refuse("") { t.Error("a connection with two persisted identities must refuse an anonymous state-changing call immediately after a restart") @@ -93,6 +96,11 @@ func TestSeedFromStateArmsTheGate(t *testing.T) { t.Fatalf("persist subagent pin: %v", err) } + // The coordinator declared itself through session_start before the restart. + if err := ss.RecordDeclaredLinkage(proxyID, "coordinator"); err != nil { + t.Fatalf("persist coordinator declaration: %v", err) + } + // newPersistSession fires onProxySession, exactly as handleInitialize does. s := newPersistSession(t, store, ss, proxyID) s.seedLogicalAgentsFromState(proxyID) // deliberately unwired in production diff --git a/internal/cli/conn_logical_agent_stamp_test.go b/internal/cli/conn_logical_agent_stamp_test.go index ecfd87be..80a577bb 100644 --- a/internal/cli/conn_logical_agent_stamp_test.go +++ b/internal/cli/conn_logical_agent_stamp_test.go @@ -65,6 +65,11 @@ func TestStampChannelStateAgreesWithTheWriteGate(t *testing.T) { ctx := mcp.WithLogicalAgent(context.Background(), tc.callerID) st := s.stampChannelState(ctx) + // The note is read INSIDE session_start, and a successful + // session_start is itself the declaration of the identity it ran + // under (issue #513), so the gate the note predicts is the one after + // this call — exactly as the after-tool hook leaves it. + s.declareSessionStartCaller(ctx, "session_start", false) // "The note would warn that writes are being refused" must hold // exactly when the gate would in fact refuse one. noteWarnsRefused := st.Shared && !st.PerCallStamped diff --git a/internal/cli/conn_logical_agent_test.go b/internal/cli/conn_logical_agent_test.go index 4b2cd6e5..64387620 100644 --- a/internal/cli/conn_logical_agent_test.go +++ b/internal/cli/conn_logical_agent_test.go @@ -29,8 +29,18 @@ func TestLogicalAgentStateRefuse(t *testing.T) { t.Fatal("an explicit call ID must not refuse") } l.recordCall("B") // a second agent arrives per-call + // Issue #513: presenting an id is not declaring one. B only ever stamped a + // call, so on a shared connection it is refused until session_start + // declares it; A declared at attach and is admitted. + if !l.refuse("B") { + t.Fatal("an identity no session_start declared must refuse on a shared connection") + } + if l.refuse("A") { + t.Fatal("an identity declared at attach must not refuse") + } + l.declare("B") if l.refuse("B") { - t.Fatal("an explicit ID on a shared connection must not refuse") + t.Fatal("a declared ID on a shared connection must not refuse") } // PLAN-394: once the connection is shared the attach-time id is whichever // peer attached LAST, not the caller — attributing on its strength was the @@ -47,8 +57,9 @@ func TestLogicalAgentStateRefuseNoAttach(t *testing.T) { if !l.refuse("") { t.Fatal("an anonymous call on a shared, no-attach connection must refuse") } + l.declare("A") // session_start ran under A's per-call identity if l.refuse("A") { - t.Fatal("an explicit call ID on a shared connection must not refuse") + t.Fatal("a declared call ID on a shared connection must not refuse") } } @@ -56,6 +67,7 @@ func TestRefuseSharedStateChange(t *testing.T) { var s connSession s.recordLogicalAgentCall("agent-1") s.recordLogicalAgentCall("agent-2") // shared, no attach ID + s.declareLogicalAgent("agent-1") // agent-1's session_start succeeded if err := s.refuseSharedStateChange(context.Background(), "read_file", ""); err != nil { t.Fatalf("a read must never refuse: %v", err) diff --git a/internal/cli/conn_persist.go b/internal/cli/conn_persist.go index 6dcafcfe..0014db20 100644 --- a/internal/cli/conn_persist.go +++ b/internal/cli/conn_persist.go @@ -49,6 +49,7 @@ func (s *connSession) onProxySession(id string) { } s.mutate(func(v *sessionView) { v.proxySessionID = id }) s.restoreIdentity(id) + s.restoreDeclaredLinkages(id) // DISABLED — see seedLogicalAgentsFromState. Re-arming the ceiling from // durable state locks out every client that cannot stamp a per-call // identity, which is the client this whole card is about. diff --git a/internal/cli/conn_repin_scope.go b/internal/cli/conn_repin_scope.go index 79a3e5b2..25168da7 100644 --- a/internal/cli/conn_repin_scope.go +++ b/internal/cli/conn_repin_scope.go @@ -82,30 +82,49 @@ func (s *connSession) repinConnection(ctx context.Context, folder, langOverride // re-pin resolves against afterwards. // // It is decided by WHETHER the caller's shard follows the connection, not by -// comparing roots after the fact. A shard that never chose a root follows the -// connection, so it resolves against the root this move left the connection -// at — even if a peer's concurrent connection move has already dragged it on, -// which a root comparison mislabelled as "your own pin". A caller with no shard -// yet is seeded from the connection on its next call, so the same holds. Only -// a shard that chose its root (selfPinned) or restored its own persisted pin -// keeps a root of its own, and that root is reported. +// comparing roots after the fact. A shard that follows resolves against the +// root this move left the connection at — even if a peer's concurrent +// connection move has already dragged it on, which a root comparison +// mislabelled as "your own pin". Whether it follows is followsConnectionLocked, +// the same predicate followConnectionShards drags by, so the two agree on WHICH +// shards follow: a subagent on its conversation's chosen root was left there by +// the move but told it now worked in the connection's new root (review of #535 +// merged with #533). They can still differ where the move has nothing to drag +// from: on a connection's first pin a fresh shard stays at "" (#567). +// +// A shard that does not follow, or was restored from its own persisted pin +// (dragged only off the connection's previous root, so its current root is the +// truth either way), reports its own root. +// +// Two more things the caller's next call does are done here too, so the +// report names what workspaceFor will resolve. A refused declaration, the +// caller's own or its conversation's, resolves it against nothing. A caller +// with no shard yet gets one as its next call would, which seeds a subagent +// from its conversation's chosen root rather than the connection's. // // The shard is read after the move: the move and the shard are guarded by // different locks, and holding both would invert the documented order // (shardsMu before sh.mu, s.mu innermost). selfPinned only ever goes from false // to true, and a self-pinned shard's root changes only through its own agent's -// repinAgent, so the one window left is this agent's own concurrent -// session_start. +// repinAgent. Two transient, report-only windows remain: this agent's own +// concurrent session_start, and its conversation self-pinning between the +// parentChose read and the shard read (followParentShard can move the subagent +// in that gap). func (s *connSession) connScopeCallerRoot(id string, out repinOutcome) string { - s.shardsMu.Lock() - sh := s.shards[id] - s.shardsMu.Unlock() + ctx := mcp.WithLogicalAgent(context.Background(), id) + if _, _, pending := s.pendingDeclarationForCall(ctx); pending { + return "" + } + sh := s.shardFor(ctx) if sh == nil { return out.root } + s.shardsMu.Lock() + parentChose := s.parentChoseLocked(id) + s.shardsMu.Unlock() sh.mu.RLock() defer sh.mu.RUnlock() - if !sh.selfPinned && !sh.restored { + if followsConnectionLocked(sh, parentChose) && !sh.restored { return out.root } return sh.root diff --git a/internal/cli/conn_repin_scope_subagent_test.go b/internal/cli/conn_repin_scope_subagent_test.go new file mode 100644 index 00000000..8f5b1008 --- /dev/null +++ b/internal/cli/conn_repin_scope_subagent_test.go @@ -0,0 +1,155 @@ +package cli + +// conn_repin_scope_subagent_test.go — a connection-scoped re-pin's report must +// name the root a SUBAGENT caller really resolves against afterwards (#535 +// review of the merge with #533). The report decided "follows the connection" +// from the shard's own flags alone, while followConnectionShards also leaves a +// subagent on its conversation's chosen root; the report then told the +// subagent it worked in the connection's new root while its relative paths +// still landed in its conversation's worktree. + +import ( + "context" + "testing" + + "github.com/plumbkit/plumb/internal/sessionstate" +) + +// Every case ends with one connection-scoped move to `other` made BY a +// subagent, and asserts that what the move reports as the caller's next root is +// what workspaceFor then resolves for it. +func TestConnScopeReportMatchesWhereASubagentResolves(t *testing.T) { + cases := []struct { + name string + // setup arranges the connection and returns the expected next root, + // given the three roots newThreeRootConn made. + setup func(t *testing.T, s *connSession, rootMain, worktree, other string) string + // persisted builds the session over a durable store instead. + persisted bool + }{ + { + // The reviewer's reproduction: the conversation chose a worktree, + // the subagent sits on it, so the connection move must leave it there. + name: "existing shard on its conversation's chosen root", + setup: func(t *testing.T, s *connSession, _, worktree, _ string) string { + t.Helper() + repinOrFail(t, s, "conv-y", worktree) + if got := s.workspaceFor(idCtx("conv-y/sub")); got != worktree { + t.Fatalf("precondition: sub on its conversation's worktree, got %q", got) + } + return worktree + }, + }, + { + // No shard yet: the next call seeds it from the conversation's + // chosen root, not from the connection's new one. + name: "no shard yet, conversation chose a root", + setup: func(t *testing.T, s *connSession, _, worktree, _ string) string { + t.Helper() + repinOrFail(t, s, "conv-y", worktree) + return worktree + }, + }, + { + // Positive control: a subagent whose conversation never chose a + // root follows the connection, so the new root is the true answer. + name: "conversation never chose: the subagent follows", + setup: func(t *testing.T, s *connSession, rootMain, _, other string) string { + t.Helper() + if got := s.workspaceFor(idCtx("conv-y/sub")); got != rootMain { + t.Fatalf("precondition: sub on the connection root, got %q", got) + } + return other + }, + }, + { + // The conversation's refused declaration is inherited: the + // subagent resolves against nothing, and the report must say so. + name: "conversation's declaration refused", + setup: func(t *testing.T, s *connSession, rootMain, worktree, _ string) string { + t.Helper() + // An explicit (sticky) connection pin, so the seeded shard's + // unforced move to an unrelated root is refused and recorded. + if _, err := s.repinWorkspace(context.Background(), "file://"+rootMain, "", false, false); err != nil { + t.Fatalf("explicit connection pin: %v", err) + } + if _, err := s.repinWorkspace(idCtx("conv-y"), "file://"+worktree, "", false, false); err == nil { + t.Fatal("precondition: an unforced move of a seeded shard to an unrelated root is refused") + } + return "" + }, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + s, rootMain, worktree, other := newThreeRootConn(t) + want := tc.setup(t, s, rootMain, worktree, other) + assertConnScopeReportTrue(t, s, other, want) + }) + } +} + +// After a daemon restart the conversation is not in memory, so a new subagent +// shard is seeded from the conversation's PERSISTED pin (parentSeeded) and a +// connection move must neither drag it nor report that it did. +func TestConnScopeReportMatchesAParentSeededSubagent(t *testing.T) { + store, ss := newOriginStore(t) + const proxyID = "proxy-535-parent-seeded" + r := repinReportRoots(t, 3) + rootMain, worktree, other := r[0], r[1], r[2] + if err := ss.UpsertPinForAgent(proxyID, "conv-y", worktree, "", sessionstate.PinSourceSessionStart); err != nil { + t.Fatalf("persist the conversation's pin: %v", err) + } + s := newPersistSession(t, store, ss, proxyID) + newSharedConn(t, s, rootMain, "conv-x", "conv-y") + if got := s.workspaceFor(idCtx("conv-y/sub")); got != worktree { + t.Fatalf("precondition: sub seeded from its conversation's persisted pin, got %q", got) + } + assertConnScopeReportTrue(t, s, other, worktree) +} + +// The rendered line is what the agent reads, so check it end to end once: the +// subagent is told it stays on its conversation's worktree, labelled as not the +// connection's pin, and that is where its relative paths go. +func TestConnScopeReportRendersASubagentsOwnRoot(t *testing.T) { + s, _, worktree, other := newThreeRootConn(t) + repinOrFail(t, s, "conv-y", worktree) + ctxSub := idCtx("conv-y/sub") + out := runRepinReport(t, s, ctxSub, map[string]any{"workspace": other, "scope": "connection", "force": true}, "brief") + wantLines(t, out, "Next relative-path call resolves against: "+worktree+" (your own pin, not the connection's)\n") + if got := s.workspaceFor(ctxSub); got != worktree { + t.Errorf("sub resolves to %q, want its conversation's worktree %q", got, worktree) + } + if got := s.workspace(); got != other { + t.Errorf("control: the connection pin = %q, want it moved to %q", got, other) + } +} + +func repinOrFail(t *testing.T, s *connSession, id, root string) { + t.Helper() + if _, err := s.repinWorkspace(idCtx(id), "file://"+root, "", true, false); err != nil { + t.Fatalf("%s re-pin to %s: %v", id, root, err) + } +} + +// assertConnScopeReportTrue has conv-y/sub move the connection to other and +// checks the report's next root against both want and what workspaceFor +// resolves for the subagent afterwards. +func assertConnScopeReportTrue(t *testing.T, s *connSession, other, want string) { + t.Helper() + ctxSub := idCtx("conv-y/sub") + out, err := s.repinConnection(ctxSub, "file://"+other, "", true) + if err != nil { + t.Fatalf("connection-scoped move: %v", err) + } + if got := s.workspace(); got != other { + t.Fatalf("precondition: the connection moved to %q, got %q", other, got) + } + actual := s.workspaceFor(ctxSub) + if out.effective != actual { + t.Errorf("the report tells the subagent it resolves against %q, but its relative paths resolve against %q", out.effective, actual) + } + if actual != want { + t.Errorf("the subagent resolves against %q, want %q", actual, want) + } +} diff --git a/internal/cli/conn_restore.go b/internal/cli/conn_restore.go index 2b9d4754..1dbe64d7 100644 --- a/internal/cli/conn_restore.go +++ b/internal/cli/conn_restore.go @@ -180,6 +180,10 @@ func (s *connSession) retryRestoreIdentity(proxyID string) bool { if adoption == idResumed && named { s.repairBlankLinkage(rec) s.setRecovery(recoveryRestored) + // onProxySession restored declarations from what it could read THEN; + // the record (and any declaration written while this connection sat + // degraded) is only now in hand. Additive, so re-running is safe. + s.restoreDeclaredLinkages(proxyID) return true } if adoption == idAbsent && named { @@ -190,6 +194,7 @@ func (s *connSession) retryRestoreIdentity(proxyID string) bool { // connection stays degraded and the next attempt tries again. s.setRecovery(recoveryEstablished) if s.persistIdentity() { + s.restoreDeclaredLinkages(proxyID) return true } s.log().Warn("daemon: could not record the healed legacy identity on retry; staying degraded for the next attempt") diff --git a/internal/cli/conn_subsystems.go b/internal/cli/conn_subsystems.go index 965415fe..26c06407 100644 --- a/internal/cli/conn_subsystems.go +++ b/internal/cli/conn_subsystems.go @@ -437,6 +437,7 @@ func (s *connSession) onAfterTool(toolName string, args json.RawMessage, output, // derived INSIDE its Execute, which this hook never sees, so a session_start // row is attributed only when the call itself also carried an identity. func (s *connSession) afterToolFromCtx(ctx context.Context, toolName string, args json.RawMessage, output, errMsg string, dur time.Duration, isError bool, failure *toolerror.Error) { + s.declareSessionStartCaller(ctx, toolName, isError) s.onAfterTool(toolName, args, output, errMsg, dur, isError, failure, mcp.LogicalAgentFromCtx(ctx)) } diff --git a/internal/cli/daemon_sessionstate_test.go b/internal/cli/daemon_sessionstate_test.go index 5ec6712a..597647e1 100644 --- a/internal/cli/daemon_sessionstate_test.go +++ b/internal/cli/daemon_sessionstate_test.go @@ -56,8 +56,9 @@ func TestSweepLegacyWidePins_NilStore(t *testing.T) { } // seedAgedSession writes every kind of expendable row a long-lived serve -// accumulates — a connection-level read and pin, a per-agent read and pin, and a -// logical-agent declaration — then ages them all a month past the default TTL. +// accumulates — a connection-level read and pin, a per-agent read and pin, a +// logical-agent record and a conversation's declaration (#513) — then ages them +// all a month past the default TTL. func seedAgedSession(t *testing.T, ss *sessionstate.Store, proxyID, ws string) { t.Helper() file := filepath.Join(ws, "a.go") @@ -72,14 +73,17 @@ func seedAgedSession(t *testing.T, ss *sessionstate.Store, proxyID, ws string) { if err := ss.RecordLogicalAgent(proxyID, "agent-a"); err != nil { t.Fatal(err) } + if err := ss.RecordDeclaredLinkage(proxyID, "agent-a"); err != nil { + t.Fatal(err) + } if err := ss.BackdateSession(proxyID, time.Now().Add(-30*24*time.Hour)); err != nil { t.Fatal(err) } } // expendableRows counts what survives of seedAgedSession's rows, in its order: -// two reads, two pins, one logical-agent declaration. -func expendableRows(t *testing.T, ss *sessionstate.Store, proxyID, ws string) (reads, pins, agents int) { +// two reads, two pins, one logical agent, one declared conversation. +func expendableRows(t *testing.T, ss *sessionstate.Store, proxyID, ws string) (reads, pins, agents, declared int) { t.Helper() for _, agent := range []string{"", "agent-a"} { recs, err := ss.LoadReadsForAgent(proxyID, agent, ws) @@ -97,7 +101,11 @@ func expendableRows(t *testing.T, ss *sessionstate.Store, proxyID, ws string) (r if err != nil { t.Fatal(err) } - return reads, pins, len(ids) + linkages, err := ss.DeclaredLinkagesFor(proxyID) + if err != nil { + t.Fatal(err) + } + return reads, pins, len(ids), len(linkages) } // Issue #525: the daemon's start-up maintenance must not age out expendable @@ -112,10 +120,10 @@ func TestStartupMaintenance_KeepsAgedStateOfSessionsThatMayReconnect(t *testing. maintainSessionStateAtStart(ss) - if reads, pins, agents := expendableRows(t, ss, "proxyX", ws); reads != 2 || pins != 2 || agents != 1 { - t.Fatalf("after start-up maintenance: %d reads, %d pins, %d agent declarations; want 2, 2, 1 — "+ + if reads, pins, agents, declared := expendableRows(t, ss, "proxyX", ws); reads != 2 || pins != 2 || agents != 1 || declared != 1 { + t.Fatalf("after start-up maintenance: %d reads, %d pins, %d logical agents, %d declared conversations; want 2, 2, 1, 1 — "+ "start-up cannot know whether a serve is about to reconnect, so it must leave aged state to the reaper", - reads, pins, agents) + reads, pins, agents, declared) } } @@ -145,13 +153,13 @@ func TestReaper_PrunesAgedStateOfDeadSessionsOnly(t *testing.T) { close(ticks) <-done - if reads, pins, agents := expendableRows(t, ss, "live", ws); reads != 2 || pins != 2 || agents != 1 { - t.Errorf("connected session after a reaper pass: %d reads, %d pins, %d agent declarations; want 2, 2, 1", - reads, pins, agents) + if reads, pins, agents, declared := expendableRows(t, ss, "live", ws); reads != 2 || pins != 2 || agents != 1 || declared != 1 { + t.Errorf("connected session after a reaper pass: %d reads, %d pins, %d logical agents, %d declared conversations; want 2, 2, 1, 1", + reads, pins, agents, declared) } - if reads, pins, agents := expendableRows(t, ss, "dead", ws); reads != 0 || pins != 0 || agents != 0 { - t.Errorf("disconnected session after a reaper pass: %d reads, %d pins, %d agent declarations; want none — "+ - "a dead session's state must still be reclaimed", reads, pins, agents) + if reads, pins, agents, declared := expendableRows(t, ss, "dead", ws); reads != 0 || pins != 0 || agents != 0 || declared != 0 { + t.Errorf("disconnected session after a reaper pass: %d reads, %d pins, %d logical agents, %d declared conversations; want none — "+ + "a dead session's state must still be reclaimed", reads, pins, agents, declared) } for _, id := range []string{"live", "dead"} { if _, ok, err := ss.LoadIdentity(id); err != nil || !ok { diff --git a/internal/cli/desktop_connector_identity_integration_test.go b/internal/cli/desktop_connector_identity_integration_test.go index d5384bd9..0091cbc8 100644 --- a/internal/cli/desktop_connector_identity_integration_test.go +++ b/internal/cli/desktop_connector_identity_integration_test.go @@ -52,6 +52,9 @@ func newDesktopConn(t *testing.T, stampKey string) *desktopConn { srv.Register(tools.NewWriteFile(s.buildWriteDeps())) srv.Register(tools.NewReadFile(s.readTracker).WithReadsFor(s.readTrackerFor).WithWorkspace(s.workspaceFor)) srv.OnToolRefusal = s.refuseSharedStateChange + // As registerHooks wires it: a successful session_start declares the + // per-call identity it ran under (issue #513). + srv.OnAfterTool = s.afterToolFromCtx srv.OnBeforeTool = func(ctx context.Context, name string, args json.RawMessage, agent string) { s.recordLogicalAgentCall(agent) s.onBeforeTool(ctx, name, args) @@ -211,4 +214,123 @@ func TestDesktopConnector(t *testing.T) { } } }) + + // Issue #513. The stamp survives the host, but it names an identity no + // session_start on this connection declared — a model typing plumb_agent + // where the hook does not run. It used to be admitted onto a fresh shard + // seeded from the connection's root and land in the main checkout. It is + // refused, lands nowhere, and leaves no identity behind. + t.Run("InventedIdentityIsRefusedNotMisrouted", func(t *testing.T) { + c := newDesktopConn(t, mcp.ArgLogicalAgentDeclaredKey) + mainCheckout, worktree, text, isErr := incident(t, c) + if isErr { + t.Fatalf("precondition: the declared agent's write was refused: %s", text) + } + observed := c.s.logicalAgents.count() + + text, isErr = c.call(t, "my-session", "write_file", map[string]any{"file_path": "INVENTED.md", "content": "who\n"}) + if !isErr { + t.Fatalf("a write under an undeclared identity was admitted: %s", text) + } + if !strings.Contains(text, "session_start") { + t.Errorf("the refusal does not name session_start as the remedy: %s", text) + } + for _, dir := range []string{mainCheckout, worktree} { + if _, err := os.Stat(filepath.Join(dir, "INVENTED.md")); err == nil { + t.Errorf("the refused write landed in %s", dir) + } + } + if got := c.s.logicalAgents.count(); got != observed { + t.Errorf("the refused call registered an identity: %d observed, want %d", got, observed) + } + + // Controls. The declared agent still writes into its own worktree — + // the same relative write through the same host, so the refusal above + // is about the identity, not the call. + if text, isErr := c.call(t, "conv-y", "write_file", map[string]any{"file_path": "Y2.md", "content": "y\n"}); isErr { + t.Fatalf("declared agent's write refused: %s", text) + } + if _, err := os.Stat(filepath.Join(worktree, "Y2.md")); err != nil { + t.Fatalf("the declared agent's write did not land in its worktree: %v", err) + } + // A hook-stamped subagent of a declared conversation is admitted + // without ever calling session_start itself, and works where its + // conversation chose to: conv-y re-pinned itself to the worktree, so + // its subagent's relative write lands there, not in conv-x's checkout. + if text, isErr := c.call(t, "conv-y/sub", "write_file", map[string]any{"file_path": "SUB.md", "content": "sub\n"}); isErr { + t.Fatalf("a subagent of a declared conversation was refused: %s", text) + } + if _, err := os.Stat(filepath.Join(worktree, "SUB.md")); err != nil { + t.Fatalf("the subagent's write did not land in its conversation's worktree: %v", err) + } + if _, err := os.Stat(filepath.Join(mainCheckout, "SUB.md")); err == nil { + t.Fatal("the subagent's write landed in the other agent's checkout") + } + // The same root answers every implicit resolver, git's default + // repository included (the git tool resolves through workspaceFor). + if got := c.s.workspaceFor(mcp.WithLogicalAgent(context.Background(), "conv-y/sub")); got != worktree { + t.Errorf("the subagent's workspace (git's default repository) is %q, want %q", got, worktree) + } + // The count above is not vacuous: an admitted new identity IS recorded. + if got := c.s.logicalAgents.count(); got != observed+1 { + t.Errorf("an admitted subagent was not recorded: %d observed, want %d", got, observed+1) + } + + // The remedy is reachable: session_start under the identity declares it. + if text, isErr := c.call(t, "my-session", "session_start", map[string]any{"session_id": "my-session"}); isErr { + t.Fatalf("session_start, the named remedy, was refused: %s", text) + } + if text, isErr := c.call(t, "my-session", "write_file", map[string]any{"file_path": "INVENTED.md", "content": "who\n"}); isErr { + t.Fatalf("a write after declaring through session_start was refused: %s", text) + } + }) + + // A client that stamps every call and calls session_start WITHOUT a + // session_id — the _meta-only channel sharedIdentityRemedy sanctions. The + // declaration comes from the after-tool hook, not the session_id linker. + t.Run("MetaOnlySessionStartDeclares", func(t *testing.T) { + c := newDesktopConn(t, mcp.ArgLogicalAgentDeclaredKey) + mainCheckout, _, text, isErr := incident(t, c) + if isErr { + t.Fatalf("precondition: the declared agent's write was refused: %s", text) + } + write := `{"file_path":"META.md","content":"meta\n"}` + if text, isErr := c.callMeta(t, "meta-only", "write_file", write); !isErr { + t.Fatalf("an undeclared _meta identity's write was admitted: %s", text) + } + if text, isErr := c.callMeta(t, "meta-only", "session_start", `{}`); isErr { + t.Fatalf("session_start without a session_id failed: %s", text) + } + if text, isErr := c.callMeta(t, "meta-only", "write_file", write); isErr { + t.Fatalf("a write after a _meta-only session_start was refused: %s", text) + } + if _, err := os.Stat(filepath.Join(mainCheckout, "META.md")); err != nil { + t.Fatalf("the admitted write did not land in the root session_start reported: %v", err) + } + }) +} + +// callMeta serves a call whose identity rides _meta[dev.plumbkit/logical-agent] +// rather than an argument, with argsJSON passed through verbatim. +func (c *desktopConn) callMeta(t *testing.T, agent, name, argsJSON string) (string, bool) { + t.Helper() + c.id++ + out := c.serve(t, fmt.Sprintf(`{"jsonrpc":"2.0","id":%d,"method":"tools/call","params":{"name":%q,"arguments":%s,"_meta":{%q:%q}}}`, + c.id, name, argsJSON, mcp.MetaLogicalAgentKey, agent)) + var resp struct { + Result struct { + Content []struct { + Text string `json:"text"` + } `json:"content"` + IsError bool `json:"isError"` + } `json:"result"` + } + if err := json.Unmarshal(out, &resp); err != nil { + t.Fatalf("decode response %s: %v", out, err) + } + text := "" + if len(resp.Result.Content) > 0 { + text = resp.Result.Content[0].Text + } + return text, resp.Result.IsError } diff --git a/internal/cli/multiagent_pin_integration_test.go b/internal/cli/multiagent_pin_integration_test.go index 59fd92b2..da3ff4f4 100644 --- a/internal/cli/multiagent_pin_integration_test.go +++ b/internal/cli/multiagent_pin_integration_test.go @@ -466,9 +466,11 @@ func testAnonymousStateChangeStillRefused(t *testing.T) { t.Errorf("refusal missing %q: %s", want, err) } } - // An identified call on the same connection is not refused. + // An identified call on the same connection is not refused once its + // identity has been declared through session_start (issue #513). + m.s.declareSessionStartCaller(mcp.WithLogicalAgent(context.Background(), "agent-a"), "session_start", false) if err := m.s.refuseSharedStateChange(context.Background(), "write_file", "agent-a"); err != nil { - t.Errorf("an identified state-changing call must not be refused: %v", err) + t.Errorf("a declared, identified state-changing call must not be refused: %v", err) } // Reads are never refused: sharing read-only state is safe. if err := m.s.refuseSharedStateChange(context.Background(), "read_file", ""); err != nil { diff --git a/internal/sessionstate/agents.go b/internal/sessionstate/agents.go index f518c0c7..9e2822c0 100644 --- a/internal/sessionstate/agents.go +++ b/internal/sessionstate/agents.go @@ -57,6 +57,9 @@ func (s *Store) BackdateLogicalAgents(proxySessionID string, to time.Time) error if _, err := s.db.Exec(`UPDATE pinned_workspace SET updated_at=? WHERE proxy_session_id=?`, to.UnixMilli(), proxySessionID); err != nil { return fmt.Errorf("sessionstate: backdate pins: %w", err) } + if _, err := s.db.Exec(`UPDATE declared_linkage SET updated_at=? WHERE proxy_session_id=?`, to.UnixMilli(), proxySessionID); err != nil { + return fmt.Errorf("sessionstate: backdate declarations: %w", err) + } return nil } @@ -103,3 +106,75 @@ func (s *Store) LogicalAgentIDsFor(proxySessionID string, since time.Time) ([]st } return out, nil } + +// RecordDeclaredLinkage durably notes that the conversation linkage declared +// itself on this connection through session_start (issue #513). It is the +// evidence that lets a reconnecting daemon keep admitting that conversation's +// stamped writes instead of refusing them as undeclared. Refreshed on every +// declaration, so a conversation that re-orients keeps its row young against +// Prune. nil-safe; blank values are dropped. +func (s *Store) RecordDeclaredLinkage(proxySessionID, linkage string) error { + if s == nil || proxySessionID == "" || linkage == "" { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + _, err := s.db.Exec( + `INSERT INTO declared_linkage (proxy_session_id, linkage, updated_at) + VALUES (?, ?, ?) + ON CONFLICT(proxy_session_id, linkage) + DO UPDATE SET updated_at=excluded.updated_at`, + proxySessionID, linkage, time.Now().UnixMilli(), + ) + if err != nil { + return fmt.Errorf("sessionstate: record declared linkage: %w", err) + } + return nil +} + +// TouchDeclaredLinkage refreshes an EXISTING declaration's timestamp, so a +// conversation that keeps working stays ahead of Prune. It never inserts: an +// admitted call is not a declaration, and a row Prune already reclaimed stays +// gone until session_start declares again. nil-safe. +func (s *Store) TouchDeclaredLinkage(proxySessionID, linkage string) error { + if s == nil || proxySessionID == "" || linkage == "" { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + if _, err := s.db.Exec(`UPDATE declared_linkage SET updated_at=? WHERE proxy_session_id=? AND linkage=?`, + time.Now().UnixMilli(), proxySessionID, linkage); err != nil { + return fmt.Errorf("sessionstate: touch declared linkage: %w", err) + } + return nil +} + +// DeclaredLinkagesFor returns every linkage declared under proxySessionID that +// Prune has not yet reclaimed. nil-safe; an empty result means "no evidence". +func (s *Store) DeclaredLinkagesFor(proxySessionID string) ([]string, error) { + if s == nil || proxySessionID == "" { + return nil, nil + } + s.mu.Lock() + defer s.mu.Unlock() + rows, err := s.db.Query( + `SELECT linkage FROM declared_linkage WHERE proxy_session_id=? AND linkage<>''`, + proxySessionID, + ) + if err != nil { + return nil, fmt.Errorf("sessionstate: list declared linkages: %w", err) + } + defer rows.Close() + var out []string + for rows.Next() { + var l string + if err := rows.Scan(&l); err != nil { + return nil, fmt.Errorf("sessionstate: scan declared linkage: %w", err) + } + out = append(out, l) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("sessionstate: list declared linkages: %w", err) + } + return out, nil +} diff --git a/internal/sessionstate/db.go b/internal/sessionstate/db.go index 22e4f1d1..843f87f0 100644 --- a/internal/sessionstate/db.go +++ b/internal/sessionstate/db.go @@ -91,7 +91,10 @@ CREATE TABLE IF NOT EXISTS pinned_workspace ( // durable identity record carries its own authorised external linkage (so // recovery no longer depends on a prunable ended-session JSON file) and a // revision that orders name updates (PLAN-426) -const SchemaVersion = 8 +// 8 — logical_agent: every identity observed on a connection (PLAN-440) +// 9 — declared_linkage: the conversations that declared themselves through +// session_start, so a restart does not refuse them as undeclared (#513) +const SchemaVersion = 9 // PinSource records WHY a workspace was pinned. It is the discriminator that // lets a reconnecting connection tell a deliberate re-pin from a stale copy of @@ -416,22 +419,27 @@ func (s *Store) Prune(olderThan time.Time, live ...string) error { if _, err := s.db.Exec(`DELETE FROM logical_agent WHERE updated_at < ?`+keep, args...); err != nil { //nolint:gosec // G202: keep is a placeholder-only fragment, IDs are bound args return fmt.Errorf("sessionstate: prune logical agents: %w", err) } + if _, err := s.db.Exec(`DELETE FROM declared_linkage WHERE updated_at < ?`+keep, args...); err != nil { //nolint:gosec // G202: keep is a placeholder-only fragment, IDs are bound args + return fmt.Errorf("sessionstate: prune declared linkages: %w", err) + } // session_names is intentionally absent — see the doc comment. Do not add a // DELETE here without an explicit retirement signal to gate it on. return nil } // BackdateSession ages every EXPENDABLE row under a proxy session — its reads, -// its pins and its logical-agent declarations — for tests that need state older -// than the TTL without sleeping. The identity record is left alone, as Prune -// leaves it. +// its pins, the logical agents it observed and the conversations that declared +// themselves on it — for tests that need state older than the TTL without +// sleeping. The identity record is left alone, as Prune leaves it. The list +// must name every table Prune sweeps, or a test of the sweep cannot age, and so +// cannot cover, the table it leaves out. func (s *Store) BackdateSession(proxySessionID string, to time.Time) error { if s == nil || proxySessionID == "" { return nil } s.mu.Lock() defer s.mu.Unlock() - for _, tbl := range []string{"read_tracking", "pinned_workspace", "logical_agent"} { + for _, tbl := range []string{"read_tracking", "pinned_workspace", "logical_agent", "declared_linkage"} { if _, err := s.db.Exec(`UPDATE `+tbl+` SET updated_at=? WHERE proxy_session_id=?`, to.UnixMilli(), proxySessionID); err != nil { //nolint:gosec // G202: tbl is a constant from the list above return fmt.Errorf("sessionstate: backdate %s: %w", tbl, err) } diff --git a/internal/sessionstate/db_agent_test.go b/internal/sessionstate/db_agent_test.go index e1c1e036..cd1c9c91 100644 --- a/internal/sessionstate/db_agent_test.go +++ b/internal/sessionstate/db_agent_test.go @@ -54,3 +54,79 @@ func TestPerAgentPinIsolation(t *testing.T) { t.Fatalf("agent pin = %q/%q ok=%v err=%v, want /ws-a/zig", aRoot, aLang, aOK, err) } } + +// TestDeclaredLinkageRoundTripAndPrune pins the v9 declared_linkage record +// (issue #513): declarations are scoped by proxy session, blank values are +// dropped, a re-declaration refreshes rather than duplicates, and Prune reclaims +// an aged row unless its proxy session is live. +func TestDeclaredLinkageRoundTripAndPrune(t *testing.T) { + s := newTestStore(t) + for _, l := range []string{"conv-a", "conv-b", "conv-a", ""} { + if err := s.RecordDeclaredLinkage("proxyX", l); err != nil { + t.Fatalf("RecordDeclaredLinkage %q: %v", l, err) + } + } + if err := s.RecordDeclaredLinkage("proxyY", "conv-other"); err != nil { + t.Fatalf("RecordDeclaredLinkage proxyY: %v", err) + } + got, err := s.DeclaredLinkagesFor("proxyX") + if err != nil { + t.Fatalf("DeclaredLinkagesFor: %v", err) + } + if len(got) != 2 || !containsAll(got, "conv-a", "conv-b") { + t.Fatalf("declared linkages = %v, want exactly conv-a and conv-b", got) + } + + // proxyY is live, so only proxyX's rows are reclaimed. + if err := s.Prune(time.Now().Add(time.Hour), "proxyY"); err != nil { + t.Fatalf("Prune: %v", err) + } + if got, _ := s.DeclaredLinkagesFor("proxyX"); len(got) != 0 { + t.Errorf("aged declarations survived Prune: %v", got) + } + if got, _ := s.DeclaredLinkagesFor("proxyY"); len(got) != 1 { + t.Errorf("a live proxy's declaration was pruned: %v", got) + } +} + +func containsAll(have []string, want ...string) bool { + set := map[string]bool{} + for _, h := range have { + set[h] = true + } + for _, w := range want { + if !set[w] { + return false + } + } + return true +} + +// TestTouchDeclaredLinkageRefreshesButNeverInserts: an admitted call keeps an +// existing declaration young, but admission is not declaration, so a linkage +// with no row never gains one. +func TestTouchDeclaredLinkageRefreshesButNeverInserts(t *testing.T) { + s := newTestStore(t) + if err := s.RecordDeclaredLinkage("proxyX", "conv-a"); err != nil { + t.Fatalf("record: %v", err) + } + if err := s.BackdateLogicalAgents("proxyX", time.Now().Add(-48*time.Hour)); err != nil { + t.Fatalf("backdate: %v", err) + } + if err := s.TouchDeclaredLinkage("proxyX", "conv-a"); err != nil { + t.Fatalf("touch: %v", err) + } + if err := s.TouchDeclaredLinkage("proxyX", "never-declared"); err != nil { + t.Fatalf("touch undeclared: %v", err) + } + if err := s.Prune(time.Now().Add(-24 * time.Hour)); err != nil { + t.Fatalf("prune: %v", err) + } + got, err := s.DeclaredLinkagesFor("proxyX") + if err != nil { + t.Fatalf("DeclaredLinkagesFor: %v", err) + } + if len(got) != 1 || got[0] != "conv-a" { + t.Fatalf("after touch + prune = %v, want only the refreshed conv-a", got) + } +} diff --git a/internal/sessionstate/migrate.go b/internal/sessionstate/migrate.go index 23f2ee03..bd0c9243 100644 --- a/internal/sessionstate/migrate.go +++ b/internal/sessionstate/migrate.go @@ -28,6 +28,7 @@ var migrationSteps = map[int]func(*sql.Tx) error{ 6: migrateV6, 7: migrateV7, 8: migrateV8, + 9: migrateV9, } // runMigrationStep applies one step and advances user_version to that step's @@ -270,3 +271,26 @@ func migrateV8(tx *sql.Tx) error { return nil } + +// migrateV9 adds the durable record of which conversations DECLARED themselves +// on a connection through session_start (issue #513). +// +// logical_agent cannot answer that question, and not for want of a column: it +// records every identity OBSERVED on the connection, per-call stamps included, +// so an id a model invented and used on an admitted read is in it too. +// Restoring declarations from it after a daemon restart would re-admit exactly +// the undeclared identity the gate exists to refuse. Keyed on the linkage (the +// conversation half of the id), because that is what a declaration vouches for: +// a hook-stamped subagent `/` rides its parent's. +func migrateV9(tx *sql.Tx) error { + const addDeclared = `CREATE TABLE IF NOT EXISTS declared_linkage ( + proxy_session_id TEXT NOT NULL, + linkage TEXT NOT NULL, + updated_at INTEGER NOT NULL, + PRIMARY KEY (proxy_session_id, linkage) +) WITHOUT ROWID` + if _, err := tx.Exec(addDeclared); err != nil { + return fmt.Errorf("sessionstate: migrate v9 (declared_linkage): %w", err) + } + return nil +} diff --git a/internal/sessionstate/migrate_test.go b/internal/sessionstate/migrate_test.go index 0d779c92..ba1723a0 100644 --- a/internal/sessionstate/migrate_test.go +++ b/internal/sessionstate/migrate_test.go @@ -299,6 +299,14 @@ func TestMigrateV1ToV8_LogicalAgentRecordExistsAndWorks(t *testing.T) { if !tableExists(t, s, "logical_agent") { t.Fatal("logical_agent is missing after migrating a v1 database") } + // v9 (issue #513) on the same upgrade path: declarations restore after a + // restart only if every installed database gains the table. + if !tableExists(t, s, "declared_linkage") { + t.Fatal("declared_linkage is missing after migrating a v1 database") + } + if err := s.RecordDeclaredLinkage("legacy", "conv"); err != nil { + t.Fatalf("record a declaration on a migrated database: %v", err) + } // The legacy pin row survives and still counts as evidence, alongside a // declaration the new table records. diff --git a/internal/tools/session_start.go b/internal/tools/session_start.go index 610a83b0..2db8bd16 100644 --- a/internal/tools/session_start.go +++ b/internal/tools/session_start.go @@ -386,6 +386,13 @@ func (t *SessionStart) Execute(ctx context.Context, raw json.RawMessage) (string if err := t.applyPurpose(raw); err != nil { return "", err } + // Validate `detail` BEFORE resolveLinkage: linking is a commitment (it + // declares the caller's identity on a shared connection, issue #513), and a + // call that is about to fail on a malformed argument must commit nothing. + // The real resolution stays below, where the auto-brief signal exists. + if _, err := resolveDetail(raw, false); err != nil { + return "", err + } inheritedName, linked := t.resolveLinkage(raw) lang, lspKey := detectLanguageInfo(ws) // A forced/attached primary may have no root marker (e.g. swift pinned on an diff --git a/internal/tools/session_start_link_test.go b/internal/tools/session_start_link_test.go index fc2d3938..d0d94cde 100644 --- a/internal/tools/session_start_link_test.go +++ b/internal/tools/session_start_link_test.go @@ -68,3 +68,24 @@ func TestSessionStart_UnlinkedNotice(t *testing.T) { } }) } + +// A session_start that fails on a malformed `detail` must not link: linking +// declares the caller's identity on a shared connection (issue #513), and a +// failed call must commit nothing. The valid call is the positive control. +func TestSessionStart_InvalidDetailDoesNotLink(t *testing.T) { + var linked []string + tool := NewSessionStart(func(context.Context) string { return t.TempDir() }, nil, nil, nil, func() string { return "" }, nil). + WithExternalID(func(id string) string { linked = append(linked, id); return "" }) + if _, err := tool.Execute(context.Background(), json.RawMessage(`{"session_id":"abc-123","detail":"bogus"}`)); err == nil { + t.Fatal("an invalid detail must fail") + } + if len(linked) != 0 { + t.Fatalf("a failed session_start linked %v", linked) + } + if _, err := tool.Execute(context.Background(), json.RawMessage(`{"session_id":"abc-123","detail":"brief"}`)); err != nil { + t.Fatalf("control: a valid session_start failed: %v", err) + } + if len(linked) != 1 || linked[0] != "abc-123" { + t.Fatalf("control: a valid session_start must link once, got %v", linked) + } +}