diff --git a/README.md b/README.md index 6e9531a..ac3d29b 100644 --- a/README.md +++ b/README.md @@ -113,7 +113,9 @@ the arrows keep working inside `claude`. | `↵` | open / attach | | `n` | new session | | `a` | add project | +| `e` | rename project | | `x` | close session | +| `c` | connect this session to one on another project | | `t` | theme picker | | `?` | help | | `q` | quit | @@ -226,12 +228,33 @@ you would want to watch it. The reviewer is handed the diff, given no tools, and run in a mode that answers but cannot act. `analysis` collects the answer along with what it cost. -The sidebar shows `⊙ n` for claims held, `✉ n` for messages waiting, and -`⚗ n · $x.xx` for reviews running and what they have cost. That last one is -the only thing in Deck that spends money with nobody watching it: the review -has no pane, and the session that asked for it has moved on. The figure stays -after the last review finishes, so a total is not lost the moment it stops -moving. +The sidebar shows `⊙ n` for claims held, `✉ n` for messages waiting, and `⚗` for +reviews. That last one is the only thing in Deck that spends money with nobody +watching it: the review has no pane, and the session that asked for it has +moved on. While a review runs the badge shows how many are in flight and the +tokens they have produced; once they land it shows the total spent, and that +figure stays so it is not lost the moment it stops moving. It is tokens during +and dollars after because the CLI reports a cost only when a turn ends — a +dollar figure while the review ran would sit at zero and read as free. + +## Connecting sessions across projects + +Everything above is scoped to one project, which is where sessions share a +repository. Press `c` to connect the selected session to one on **another** +project — an API changing in one repository while its consumer changes in +another is the case that scope cannot express. + +A connection is a pair, and it widens rather than narrows. The two sessions +appear in each other's `sessions`, can read each other's `work`, can `analyse` +it, can `message` each other, and see each other's notes. Connect A to B and A +to C and A sees both, while B and C stay invisible to each other. Pressing `c` +on a session you are already connected to disconnects it. Connections are +saved, so they survive restarting Deck. + +Claims are the deliberate exception: they stay inside one project. A claim is a +path relative to a repository root, so two sessions in different repositories +touching `internal/api/client.go` are not in each other's way, and saying they +were would make the whole mechanism worth ignoring. ## Status diff --git a/docs/architecture.md b/docs/architecture.md index 91b758d..f77d594 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -271,7 +271,8 @@ Sessions are isolated by construction — separate worktrees, separate claude transcripts — so two agents will happily refactor the same interface in parallel and find out at merge time. `internal/coord` is the one channel through that isolation: an in-process MCP server exposing `sessions`, `claim`, -`release`, `note`, and `notes`. +`release`, `note`, `notes`, `message`, `inbox`, `work`, `analyse` and +`analysis`. It is modelled on cathode's `approvals.go` (hand-rolled JSON-RPC over Streamable HTTP, honouring the client's `Accept` header for SSE framing) with @@ -373,6 +374,44 @@ with anything else lying around in the tree. later, for the reason in (1): by the time anyone asks, the branch it came from has moved. +### Connections widen a project, never narrow one (`coord.sees`) + +Everything above scopes to a project, because that is where sessions share a +repository. A **connection** is the one exception, and it goes outward: it +joins two sessions on *different* projects so each can see the other's work. +The case it exists for is an API changing in one repository while its consumer +changes in another, which project scope excludes by construction. + +It is a pair — `store.Connection{A, B}` — rather than a named set. Sessions on +one project already see each other, so the motivating case is exactly two +sessions, and a set would need a name, a member editor and a rule for the last +member leaving. Connecting A to B and A to C lets A see both without making B +and C visible to each other. The pair lives on `State`, not on `Session`, so +there is no two-way link to keep consistent. + +Three things follow, and the first is the one a later reader is most likely to +"fix": + +1. **Claiming does not go through `sees`.** A claim is a repo-relative path, so + two sessions in different repositories both claiming + `internal/api/client.go` would be told they collide over a file they do not + share. One false conflict is enough to teach an agent to ignore the + mechanism. A connection widens what a session can read; it must never widen + the soft lock. +2. **The shared log widens on read only.** A note still goes to the writer's + own project log. A reader gets that merged with what connected sessions + wrote in theirs, filtered to those sessions — a connection joins two + sessions, so handing over the far project's whole log would publish the + notes of every session there. +3. **The coordinator is given the whole set, never a delta.** `SetConnections` + replaces. The store owns the document and the coordinator holds a copy of + it; a copy maintained by deltas drifts the first time an update is missed, + and nothing observes the drift. + +A result that crosses the boundary says so. `work` and the `sessions` rows +carry a `project` only when the far session is on another one, so a row without +it is on yours and its paths resolve against the tree you are looking at. + ### A spawned review (`coord.Analyse`, `agent.RunClaude`) `analyse` starts a **separate agent** to review a sibling's work. Four @@ -406,6 +445,22 @@ on whether its context was read from cache or written to it — so the running total is what makes a pattern visible. It is dropped with the session, like claims and the inbox, because a review belongs to the session that paid for it. +**A run in flight reports tokens, not dollars.** The CLI is read as +line-delimited events (`--output-format stream-json` with +`--include-partial-messages`), and `RunClaude` takes a callback that fires on +every event carrying usage. Two facts about that stream decide what the UI can +show. Only the **output** count moves: the first event of a turn already +carries the final input, cache-read and cache-write figures — one measured run +knew 15,888 read and 7,954 written before generating a word. And **cost appears +only in the result event**, so a running badge showing dollars would sit at +zero for the whole review and read as free. The sidebar therefore shows tokens +while a review runs and the total once it lands. + +Reading stdout to the end before waiting on the process is what makes the +paragraph above about failed turns actually true. It used to be `cmd.Output`, +which turns any non-zero exit into a Go error and takes the result envelope +down with it — the exact case the design says must keep its accounting. + Jobs are bounded like the inbox and the log. Dropping the oldest is safe: `Spend` is a running total kept separately, so a discarded record costs the reader an old answer and never the bill. diff --git a/docs/backlog.md b/docs/backlog.md index 1b097c3..5566d1e 100644 --- a/docs/backlog.md +++ b/docs/backlog.md @@ -143,46 +143,98 @@ currently awkward without a mouse. It lived under "Not built yet" until the deferral reason expired, and was listed in both places for a while — a deferred item and a planned one are different claims. -## 10. Connections between sessions in a project - -Sessions on one project can already read each other's work (`work`) and have it -reviewed by a spawned agent (`analyse`). Both are scoped to the project, which -is the same scope `Siblings`, `notes` and `message` use. - -A **connection** would be a smaller grouping inside that: a set of sessions -that share context and can review each other, persisted as -`Connections []Connection` on `store.State` rather than as a field on -`Session`, so there is no two-way link to keep consistent. `Load` already -back-fills missing fields, so an older state file stays readable. - -It was deliberately deferred rather than built first. A connection gates two -things — the shared log and the analysis — and until the analysis existed a -`Connection` type would have been stored, rendered and read by nothing. Now -that both exist, the question is answerable from use rather than from -prediction: **is project scope actually too coarse?** With three or four -sessions on one repository it is not, and a grouping inside it would add a -concept without removing a problem. - -The case that project scope genuinely cannot express is a connection *across* -projects — an API changing in one repository while its consumer is updated in -another. `Siblings` excludes that by construction. If connections are built, -that is the motivating case, and it inverts the framing: a connection is not a -narrowing of the project, it is an escape from it. - -One question decides the shape and should be answered before any code: does a -connection **narrow the shared notes log**, or only gate the analysis? -Narrowing changes behaviour every sibling relies on today; gating is additive. - -## 11. Live token counts while a review runs - -The sidebar shows elapsed time while a spawned analysis is in flight and the -exact cost once it lands, because the figures arrive only in the final result -envelope. Showing them as they accumulate needs `--output-format stream-json` -with `--include-partial-messages`, which is a parser rather than a field read. - -Worth doing only if a review ever runs long enough that watching the number -move tells you something a spinner does not. The measured runs so far finish in -a few seconds. +## ~~10. Connections between sessions~~ — done 2026-08-28 + +Shipped, and the entry that stood here answered its own question. It asked +whether a connection should **narrow** a project, and concluded that the case +project scope cannot express is the opposite one: an API changing in one +repository while its consumer changes in another. So a connection is an escape +from the project, never a subdivision of it, and the shared log question it +left open resolves the same way — nothing a sibling relies on today changed. + +Four decisions are worth keeping. + +**A pair, not a named set.** `store.Connection{A, B}`. Sessions on one project +already see each other, so the motivating case is exactly two sessions. A set +would need a name, a member editor and a rule for the last member leaving, and +none of that is asked for by the case that justified the feature. Connecting A +to B and A to C lets A see both without making B and C visible to each other. + +**Claims deliberately do not widen.** Every other scoping test became `sees`; +`Claim` kept `ProjectID`. A claim is a repo-relative path, so two sessions in +different repositories both claiming `internal/api/client.go` would be reported +as colliding over a file they do not share, and an agent that meets one false +conflict stops trusting the mechanism. This is the exception most likely to be +tidied away by a later reader, which is why it has a test named for it. + +**The read widens, the write does not.** A note still goes to the writer's own +project log. A reader gets that log merged with what connected sessions wrote +in theirs, filtered to those sessions — a connection joins two sessions, so +handing over the far project's whole log would publish the notes of every +session there. + +**The coordinator is told the whole set, never a delta.** `SetConnections` +replaces. The store owns the document and the coordinator holds a copy; a copy +updated by deltas is free to drift the first time an update is missed, and the +drift is invisible. + +Still open, and deliberately not built: connecting sessions whose agents are +not running. The registry is live, so a connection to a stopped session is +recorded and does nothing until it starts. That matches `work`, which cannot +read an exited session either, and it is the same underlying question — what a +session means after its agent stops. + +## ~~11. Live token counts while a review runs~~ — done 2026-08-28 + +The run now uses `--output-format stream-json --verbose +--include-partial-messages` and `agent.RunClaude` takes a callback that fires +on every event carrying usage. + +What a captured stream showed, and what the design follows from: **only the +output count moves.** The first event of a turn already carries the final +input, cache-read and cache-write figures — in one measured run, 15,888 read +and 7,954 written were known before a single word was generated. And **cost is +not in any event but the last**, so a run in flight can report what it is using +and not what it will cost. The sidebar therefore shows tokens while a review +runs and dollars once it lands, rather than a dollar figure that would sit at +zero for the whole run and read as free. + +One format, not two. The non-live path could have kept `--output-format json`, +but a second parser is a second place for claude's field names to drift, and +the totals are the thing least affordable to get quietly wrong. + +It also closed a gap the old code's own doc comment denied. `RunClaude` +promised to return an error "only when there is no accounting at all", while +`cmd.Output` turned any non-zero exit into an error and discarded the result +envelope with it. Reading stdout to the end before waiting means a result that +arrived is returned whatever the process does afterwards. + +## 13. A shell as a session + +Not every session wants an agent. Reaching a server, running a migration, or +watching a log is work that belongs beside the agents rather than in a separate +terminal. + +Almost all of it already works: `agent.Start` runs whatever `store.Session.Agent` +names, `coordArgs` gives no coordination flags to a program it does not know, +and `willResume` refuses `--continue` to anything but claude. `deck -agent +/bin/zsh` is a working shell session today. What is missing is the menu entry, +and three decisions around it. + +1. **Which command.** `$SHELL`, falling back to `/bin/sh`. It holds the login + shell on both macOS and Linux and Deck inherits it from the terminal it was + started in. The authoritative record is per-platform and needs a subprocess + to read — `getent passwd` on Linux, Open Directory on macOS, where + `/etc/passwd` holds only system accounts — to reproduce a value already in + hand. +2. **`-agent-args` must not reach it.** `agentArgsFor` prepends them + unconditionally. That is harmless while `-agent` and `-agent-args` are set + together, and stops being harmless once a shell is on the menu. +3. **The choice must not stick.** Submitting the form writes the agent to + settings as the next session's default, which is wrong for a one-off. + +Numbered 13 rather than reusing 12: that number already names the public-tree +notice, and a closed entry should not change meaning. ## ~~12. A public-repo notice a cloner will meet~~ — done 2026-08-28 diff --git a/internal/agent/claudestream.go b/internal/agent/claudestream.go new file mode 100644 index 0000000..5ef9b37 --- /dev/null +++ b/internal/agent/claudestream.go @@ -0,0 +1,131 @@ +package agent + +// Reading what `claude -p --output-format stream-json` emits: one JSON object +// per line, and the totals in the result event at the end. + +import ( + "bufio" + "encoding/json" + "fmt" + "io" + "strings" + "time" +) + +// maxEventLine bounds one line of the stream. +// +// The result event carries the model's whole answer, so the longest line is as +// long as a review. Four megabytes is far above any answer measured and far +// below a size that would matter; a line past it is reported rather than +// truncated, because half a JSON object parses as nothing and would otherwise +// read as "the run produced no result". +const maxEventLine = 4 << 20 + +// claudeEvent is one line of the stream. +// +// Usage arrives in four different places depending on the event, which is why +// this carries the shape rather than a flat set of fields: the result event +// holds it at the top level, an assistant event under message, and a +// stream_event under event or under event.message. +type claudeEvent struct { + Type string `json:"type"` + + // The result event's own fields. Absent everywhere else. + Subtype string `json:"subtype"` + IsError bool `json:"is_error"` + Result string `json:"result"` + DurationMS int `json:"duration_ms"` + CostUSD float64 `json:"total_cost_usd"` + Usage *Tokens `json:"usage"` + + Message *usageHolder `json:"message"` + Event *struct { + Type string `json:"type"` + Usage *Tokens `json:"usage"` + Message *usageHolder `json:"message"` + } `json:"event"` +} + +type usageHolder struct { + Usage *Tokens `json:"usage"` +} + +// usage is the accounting this event carries, or nil. +func (e claudeEvent) usage() *Tokens { + switch { + case e.Usage != nil: + return e.Usage + case e.Message != nil && e.Message.Usage != nil: + return e.Message.Usage + case e.Event == nil: + return nil + case e.Event.Usage != nil: + return e.Event.Usage + case e.Event.Message != nil && e.Event.Message.Usage != nil: + return e.Event.Message.Usage + } + return nil +} + +// readClaudeStream consumes the stream and returns what the run reported, +// calling onUsage with each accounting it passes. +// +// Only Output moves once the run has started: the first event of a turn +// already carries the final input, cache-read and cache-write counts, which is +// why a live display fills in almost at once and then creeps. Cost is not in +// any of them — it appears only in the result — so a run in flight can report +// how much it is using and not what it will cost. +// +// A line that does not parse is skipped rather than fatal. The stream is +// several event kinds wide and gains more over time, and refusing a run +// because one line was unfamiliar would throw away the accounting on the next. +func readClaudeStream(r io.Reader, onUsage func(Tokens)) (ClaudeRun, error) { + sc := bufio.NewScanner(r) + sc.Buffer(make([]byte, 0, 64<<10), maxEventLine) + + var run ClaudeRun + var got bool + for sc.Scan() { + line := sc.Bytes() + if len(line) == 0 { + continue + } + var e claudeEvent + if json.Unmarshal(line, &e) != nil { + continue + } + if u := e.usage(); u != nil && onUsage != nil { + onUsage(*u) + } + if e.Type != "result" { + continue + } + run = ClaudeRun{ + Text: e.Result, + CostUSD: e.CostUSD, + Took: time.Duration(e.DurationMS) * time.Millisecond, + } + if e.Usage != nil { + run.Tokens = *e.Usage + } + if e.IsError { + run.Text = "" + run.Failure = strings.TrimSpace(e.Subtype + ": " + e.Result) + } + got = true + } + // The read error is subordinate to the result, not the other way round. got + // is set only by a fully parsed result event, so once one has arrived the + // accounting is in hand and a later oversized or unreadable line cannot + // take it back. Returning the error here regardless would lose the cost, + // the tokens and the answer of a run that had already reported all three — + // the exact case ClaudeRun's shape exists to prevent. + if err := sc.Err(); err != nil && !got { + return ClaudeRun{}, fmt.Errorf("read claude output: %w", err) + } + if !got { + // No result event means no accounting: whatever it spent, we cannot say. + return ClaudeRun{}, fmt.Errorf("claude produced no result event") + } + return run, nil +} diff --git a/internal/agent/headless.go b/internal/agent/headless.go index 9a2d1db..5a3d584 100644 --- a/internal/agent/headless.go +++ b/internal/agent/headless.go @@ -6,8 +6,6 @@ package agent import ( "context" - "encoding/json" - "errors" "fmt" "os/exec" "strings" @@ -47,24 +45,28 @@ type ClaudeRun struct { // Failed reports whether the turn ran without producing an answer. func (r ClaudeRun) Failed() bool { return r.Failure != "" } -// resultEnvelope is the one event in the stream that carries the totals. The -// field names are claude's; another agent would need its own parser, which is -// why this function is named for the one it understands. -type resultEnvelope struct { - Type string `json:"type"` - Subtype string `json:"subtype"` - IsError bool `json:"is_error"` - Result string `json:"result"` - DurationMS int `json:"duration_ms"` - CostUSD float64 `json:"total_cost_usd"` - Usage Tokens `json:"usage"` +// streamArgs put the CLI in line-delimited mode. +// +// stream-json rather than json because it is the only format that reports +// usage before the turn ends, and one format is kept rather than two: a second +// parser for the non-live path would be a second place for the field names to +// drift. --verbose is required alongside it, and --include-partial-messages is +// what makes the output count move during a turn rather than once at the end. +var streamArgs = []string{ + "-p", "--output-format", "stream-json", "--verbose", "--include-partial-messages", } // RunClaude runs one non-interactive turn in dir and reports what it cost, // with the answer when there is one. // +// onUsage, when not nil, is called with each accounting the run reports as it +// goes, on the reading goroutine. See readClaudeStream for what actually moves. +// // A turn that ends in a refusal or an API error comes back as a ClaudeRun with -// Failure set and its cost intact, not as an error. See ClaudeRun. +// Failure set and its cost intact, not as an error. That now holds for a +// non-zero exit status too: the result event is the accounting, so if one +// arrived it is returned whatever the process did afterwards. Reading stdout to +// the end before waiting is what makes that possible. // // The prompt goes over stdin rather than as an argument: with stdin empty // claude reports "input must be provided" and ignores a positional prompt. @@ -72,61 +74,44 @@ type resultEnvelope struct { // The environment is scrubbed exactly as an interactive session's is, so a // spawned run cannot inherit credentials or the child-session marker from // whatever started Deck. -func RunClaude(ctx context.Context, dir, prompt string, args ...string) (ClaudeRun, error) { - full := append([]string{"-p", "--output-format", "json"}, args...) - cmd := exec.CommandContext(ctx, "claude", full...) +func RunClaude(ctx context.Context, dir, prompt string, onUsage func(Tokens), args ...string) (ClaudeRun, error) { + cmd := exec.CommandContext(ctx, "claude", append(append([]string{}, streamArgs...), args...)...) cmd.Dir = dir cmd.Env = ScrubbedEnv() cmd.Stdin = strings.NewReader(prompt) - out, err := cmd.Output() + // Captured by hand because StdoutPipe rules out cmd.Output, which is what + // used to collect it. The CLI reports its own diagnosis here, and a bare + // "exit status 1" tells the caller nothing about whether it was auth, a bad + // flag or a refusal. + var errOut strings.Builder + cmd.Stderr = &errOut + + out, err := cmd.StdoutPipe() if err != nil { - // The CLI reports its own diagnosis on stderr; a bare "exit status 1" - // tells the caller nothing about whether it was auth, a bad flag or a - // refusal. - var ee *exec.ExitError - if errors.As(err, &ee) && len(ee.Stderr) > 0 { - return ClaudeRun{}, fmt.Errorf("claude: %s", strings.TrimSpace(string(ee.Stderr))) - } return ClaudeRun{}, fmt.Errorf("claude: %w", err) } - return parseClaudeJSON(out) -} + if err := cmd.Start(); err != nil { + return ClaudeRun{}, fmt.Errorf("claude: %w", err) + } + run, readErr := readClaudeStream(out, onUsage) + waitErr := cmd.Wait() -// parseClaudeJSON pulls the totals out of a --output-format json stream. -// -// The stream is an array of events, not a single object: an init event, the -// assistant turns, and one result event carrying the totals. Reading the last -// element would work today and break the moment anything is appended after -// it, so this selects by type. -func parseClaudeJSON(out []byte) (ClaudeRun, error) { - var events []resultEnvelope - if err := json.Unmarshal(out, &events); err != nil { - // A single object rather than an array is also valid JSON output; try - // it before giving up, so a change of shape degrades to one parse - // failure rather than to a wrong answer. - var one resultEnvelope - if json.Unmarshal(out, &one) != nil { - return ClaudeRun{}, fmt.Errorf("claude produced no readable result: %w", err) + if readErr != nil { + // Nothing to report but the failure, so the exit status and whatever + // the CLI said about it are worth more than the parse error. + if waitErr != nil { + return ClaudeRun{}, claudeFailure(waitErr, errOut.String()) } - events = []resultEnvelope{one} + return ClaudeRun{}, readErr } - for _, e := range events { - if e.Type != "result" { - continue - } - run := ClaudeRun{ - Text: e.Result, - CostUSD: e.CostUSD, - Tokens: e.Usage, - Took: time.Duration(e.DurationMS) * time.Millisecond, - } - if e.IsError { - run.Text = "" - run.Failure = strings.TrimSpace(e.Subtype + ": " + e.Result) - } - return run, nil + return run, nil +} + +// claudeFailure turns a failed run into the most informative error available. +func claudeFailure(err error, stderr string) error { + if s := strings.TrimSpace(stderr); s != "" { + return fmt.Errorf("claude: %s", s) } - // No result event means no accounting: whatever it spent, we cannot say. - return ClaudeRun{}, fmt.Errorf("claude produced no result event") + return fmt.Errorf("claude: %w", err) } diff --git a/internal/agent/headless_test.go b/internal/agent/headless_test.go index 1fcf2da..df257ff 100644 --- a/internal/agent/headless_test.go +++ b/internal/agent/headless_test.go @@ -1,32 +1,37 @@ package agent import ( + "bufio" "strings" "testing" "time" ) // A captured stream, trimmed to the events that matter. The shape is the one -// `claude -p --output-format json` actually emits: an array, with the totals -// in a result event rather than in the last element by position. -const captured = `[ - {"type":"system","subtype":"init","session_id":"abc"}, - {"type":"assistant","message":{"role":"assistant"}}, - {"type":"result","subtype":"success","is_error":false,"result":"ZEPHYR_QUOTA_GUARD", - "duration_ms":1624,"num_turns":1,"total_cost_usd":0.23709, - "usage":{"input_tokens":2,"output_tokens":4, - "cache_read_input_tokens":0,"cache_creation_input_tokens":23698}}, - {"type":"trailing_event_added_later"} -]` +// `claude -p --output-format stream-json --include-partial-messages` actually +// emits: one JSON object per line, usage in four different places, and the +// totals in a result event rather than in the last line by position. +const captured = `{"type":"system","subtype":"init","session_id":"abc"} +{"type":"stream_event","event":{"type":"message_start","message":{"usage":{"input_tokens":2,"cache_creation_input_tokens":23698,"cache_read_input_tokens":0,"output_tokens":1}}}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"input_tokens":2,"cache_creation_input_tokens":23698,"cache_read_input_tokens":0,"output_tokens":4}}} +{"type":"assistant","message":{"usage":{"input_tokens":2,"cache_creation_input_tokens":23698,"cache_read_input_tokens":0,"output_tokens":4}}} +{"type":"result","subtype":"success","is_error":false,"result":"ZEPHYR_QUOTA_GUARD","duration_ms":1624,"num_turns":1,"total_cost_usd":0.23709,"usage":{"input_tokens":2,"output_tokens":4,"cache_read_input_tokens":0,"cache_creation_input_tokens":23698}} +{"type":"trailing_event_added_later"}` -// TestParseClaudeJSONSelectsByType is why the parser does not read the last -// element. A stream with anything appended after the result would otherwise -// report zero cost and no answer. -func TestParseClaudeJSONSelectsByType(t *testing.T) { - run, err := parseClaudeJSON([]byte(captured)) +func parse(t *testing.T, stream string) ClaudeRun { + t.Helper() + run, err := readClaudeStream(strings.NewReader(stream), nil) if err != nil { t.Fatal(err) } + return run +} + +// TestReadClaudeStreamSelectsByType is why the parser does not read the last +// line. A stream with anything appended after the result would otherwise +// report zero cost and no answer. +func TestReadClaudeStreamSelectsByType(t *testing.T) { + run := parse(t, captured) if run.Text != "ZEPHYR_QUOTA_GUARD" { t.Errorf("text = %q", run.Text) } @@ -38,17 +43,44 @@ func TestParseClaudeJSONSelectsByType(t *testing.T) { } } -// TestParseClaudeJSONSplitsCacheTokens covers the split that explains a bill. +// TestReadClaudeStreamSplitsCacheTokens covers the split that explains a bill. // Reads and writes of the same context price differently, so collapsing them // into one number would hide the largest influence on a short run's cost. -func TestParseClaudeJSONSplitsCacheTokens(t *testing.T) { - run, err := parseClaudeJSON([]byte(captured)) +func TestReadClaudeStreamSplitsCacheTokens(t *testing.T) { + want := Tokens{Input: 2, Output: 4, CacheRead: 0, CacheWrite: 23698} + if got := parse(t, captured).Tokens; got != want { + t.Errorf("tokens = %+v, want %+v", got, want) + } +} + +// TestUsageIsReportedBeforeTheResult is what a live figure depends on. The +// callback has to fire on events other than the result — otherwise the number +// arrives at the same moment as the answer and there is nothing to watch. +// +// It also pins where usage is read from. Each of the three events below holds +// it in a different place, and a parser that understood only the result would +// pass every other test in this file. +func TestUsageIsReportedBeforeTheResult(t *testing.T) { + var seen []Tokens + run, err := readClaudeStream(strings.NewReader(captured), func(tk Tokens) { + seen = append(seen, tk) + }) if err != nil { t.Fatal(err) } - want := Tokens{Input: 2, Output: 4, CacheRead: 0, CacheWrite: 23698} - if run.Tokens != want { - t.Errorf("tokens = %+v, want %+v", run.Tokens, want) + if len(seen) < 3 { + t.Fatalf("usage reported %d times, want one per event that carries it: %+v", len(seen), seen) + } + // message_start, then message_delta: the output count is the only figure + // that moves, which is why the badge shows it and not the input. + if seen[0].Output != 1 || seen[1].Output != 4 { + t.Errorf("output did not climb across events: %+v", seen) + } + if seen[0].CacheWrite != 23698 { + t.Errorf("the first event did not carry the full cache cost: %+v", seen[0]) + } + if run.Tokens.Output != 4 { + t.Errorf("final tokens = %+v", run.Tokens) } } @@ -57,10 +89,10 @@ func TestParseClaudeJSONSplitsCacheTokens(t *testing.T) { // would tell every caller to discard the value — putting the bill out of reach // in the one case where it is surprising. func TestFailedTurnKeepsItsCost(t *testing.T) { - stream := `[{"type":"result","subtype":"error_max_turns","is_error":true, - "result":"ran out of turns","total_cost_usd":0.4, - "usage":{"input_tokens":7,"output_tokens":0}}]` - run, err := parseClaudeJSON([]byte(stream)) + stream := `{"type":"result","subtype":"error_max_turns","is_error":true,` + + `"result":"ran out of turns","total_cost_usd":0.4,` + + `"usage":{"input_tokens":7,"output_tokens":0}}` + run, err := readClaudeStream(strings.NewReader(stream), nil) if err != nil { t.Fatalf("a turn that ran was reported as unusable: %v", err) } @@ -83,14 +115,75 @@ func TestFailedTurnKeepsItsCost(t *testing.T) { } } -// TestParseClaudeJSONReportsAMissingResult covers a stream that ends without +// TestReadClaudeStreamReportsAMissingResult covers a stream that ends without // totals — a crash mid-run. Returning a zero-cost success would under-report // the bill and hand the caller an empty answer as if it were real. -func TestParseClaudeJSONReportsAMissingResult(t *testing.T) { - if _, err := parseClaudeJSON([]byte(`[{"type":"system","subtype":"init"}]`)); err == nil { +func TestReadClaudeStreamReportsAMissingResult(t *testing.T) { + if _, err := readClaudeStream(strings.NewReader(`{"type":"system","subtype":"init"}`), nil); err == nil { t.Error("a stream with no result event parsed as a success") } - if _, err := parseClaudeJSON([]byte(`not json at all`)); err == nil { + if _, err := readClaudeStream(strings.NewReader("not json at all"), nil); err == nil { t.Error("unparseable output was accepted") } } + +// failAfter reads a string and then fails, standing in for a stream that dies +// after the result event: an oversized trailing line, or a pipe that breaks +// once the process is killed. +type failAfter struct { + rest string + err error +} + +func (f *failAfter) Read(p []byte) (int, error) { + if f.rest == "" { + return 0, f.err + } + n := copy(p, f.rest) + f.rest = f.rest[n:] + return n, nil +} + +// TestAReadFailureAfterTheResultKeepsTheAccounting is the property ClaudeRun's +// whole shape exists for, at the one boundary that used to break it. The +// result event *is* the accounting, so once it has been parsed a later read +// error cannot take back the cost, the tokens and the answer — otherwise a run +// that spent money reports JobFailed with nothing added to the session total. +func TestAReadFailureAfterTheResultKeepsTheAccounting(t *testing.T) { + r := &failAfter{ + rest: `{"type":"result","subtype":"success","result":"ZEPHYR_QUOTA_GUARD",` + + `"total_cost_usd":0.42,"usage":{"output_tokens":9}}` + "\n", + err: bufio.ErrTooLong, + } + run, err := readClaudeStream(r, nil) + if err != nil { + t.Fatalf("a read failure after the result discarded it: %v", err) + } + if run.CostUSD != 0.42 || run.Text != "ZEPHYR_QUOTA_GUARD" || run.Tokens.Output != 9 { + t.Errorf("run = %+v, want the totals the result reported", run) + } +} + +// And the other side of it: a read failure with no result is still a failure, +// because there is genuinely nothing to report. +func TestAReadFailureWithNoResultIsStillAFailure(t *testing.T) { + r := &failAfter{rest: `{"type":"system","subtype":"init"}` + "\n", err: bufio.ErrTooLong} + if _, err := readClaudeStream(r, nil); err == nil { + t.Error("a broken stream with no accounting was reported as a success") + } +} + +// TestAnUnreadableLineDoesNotLoseTheRun. The stream is several event kinds +// wide and gains more over time. Refusing the whole run because one line was +// unfamiliar would throw away the accounting carried by the next. +func TestAnUnreadableLineDoesNotLoseTheRun(t *testing.T) { + stream := "{ this is not json\n" + + `{"type":"result","subtype":"success","result":"ok","total_cost_usd":0.01}` + run, err := readClaudeStream(strings.NewReader(stream), nil) + if err != nil { + t.Fatalf("one bad line lost a run that reported its totals: %v", err) + } + if run.Text != "ok" || run.CostUSD != 0.01 { + t.Errorf("run = %+v", run) + } +} diff --git a/internal/coord/analyse.go b/internal/coord/analyse.go index 1ef91b6..f5c132a 100644 --- a/internal/coord/analyse.go +++ b/internal/coord/analyse.go @@ -13,6 +13,8 @@ import ( "fmt" "strings" "time" + + "github.com/tripledownab/deck/internal/agent" ) // jobTimeout bounds a spawned run. Long enough for a review of a substantial @@ -21,10 +23,10 @@ const jobTimeout = 10 * time.Minute // reviewPrompt is what the spawned agent is given. It receives the diff in the // prompt and no tools at all, so everything it needs to answer has to be here. -const reviewPrompt = `You are reviewing work done by another agent on this project. +const reviewPrompt = `You are reviewing work done by another agent. You cannot run anything or read any file: the change is reproduced in full below. -Session: %s — %s +Session: %s — %s%s Measured from: %s Files changed: @@ -48,15 +50,23 @@ const defaultQuestion = "What is wrong, risky or incomplete in this change? " + // caller. It runs in plan mode, which answers a question but refuses to act, // so a review cannot become an edit. func (c *Coordinator) Analyse(sessionID, target, question string) (*Job, error) { - found, w, err := c.workOf(sessionID, target) + found, w, elsewhere, err := c.workOf(sessionID, target) if err != nil { return nil, err } if strings.TrimSpace(question) == "" { question = defaultQuestion } + // The reviewer runs in the target's own directory, so for a connected + // session that is a different repository from the one the caller is in. + // Saying so is not decoration: without it the reviewer reads the change as + // belonging to whatever project it happens to have been asked about. + var from string + if elsewhere { + from = "\nProject: " + found.Project + } prompt := fmt.Sprintf(reviewPrompt, - found.Name, found.Title, w.Base, w.Stat, w.Patch, question) + found.Name, found.Title, from, w.Base, w.Stat, w.Patch, question) job := &Job{ ID: newJobID(), From: sessionID, Subject: found.Name, @@ -92,7 +102,15 @@ func (c *Coordinator) runAnalysis(job *Job, dir, prompt string) { ctx, cancel := context.WithTimeout(c.life, jobTimeout) defer cancel() - run, err := c.spawn(ctx, dir, prompt) + // Tokens land on the job as they accumulate, so a caller polling Analysis + // sees the run using context rather than only a spinner. Cost is not here: + // the CLI reports it once, in the result, so a job in flight can say what + // it is using and not what it will cost. + run, err := c.spawn(ctx, dir, prompt, func(t agent.Tokens) { + c.mu.Lock() + job.Tokens = t + c.mu.Unlock() + }) c.mu.Lock() defer c.mu.Unlock() diff --git a/internal/coord/claims.go b/internal/coord/claims.go index 126d618..dd82cad 100644 --- a/internal/coord/claims.go +++ b/internal/coord/claims.go @@ -22,8 +22,9 @@ type Claim struct { At time.Time `json:"at"` } -// Siblings returns the other live sessions on the same project, with the -// claims each of them holds. +// Siblings returns the other live sessions this one may see, with the claims +// each of them holds: everyone on the same project, plus anyone joined by a +// connection. func (c *Coordinator) Siblings(sessionID string) []map[string]any { c.mu.Lock() defer c.mu.Unlock() @@ -34,7 +35,7 @@ func (c *Coordinator) Siblings(sessionID string) []map[string]any { } out := make([]map[string]any, 0, len(c.sessions)) for id, s := range c.sessions { - if id == sessionID || s.ProjectID != me.ProjectID { + if id == sessionID || !c.sees(me, s) { continue } held := make([]string, 0, len(c.claims[id])) @@ -42,12 +43,19 @@ func (c *Coordinator) Siblings(sessionID string) []map[string]any { held = append(held, cl.Path) } sort.Strings(held) - out = append(out, map[string]any{ + row := map[string]any{ "session": s.Name, "title": s.Title, "branch": s.Branch, "claims": held, - }) + } + // Named only when it differs. A row without a project is one on your + // own, so the paths and the branch read the way they always did; a row + // with one is from another repository, where they do not. + if s.ProjectID != me.ProjectID { + row["project"] = s.Project + } + out = append(out, row) } sort.Slice(out, func(i, j int) bool { return out[i]["session"].(string) < out[j]["session"].(string) @@ -69,6 +77,11 @@ func (c *Coordinator) Claim(sessionID string, paths []string, reason string) (gr return nil, nil } + // Project-scoped on purpose, where Siblings above is not. A path here is + // repo-relative, so two sessions in different repositories both claiming + // internal/api/client.go would be reported as colliding over a file they do + // not share. A connection widens what a session can read; it must not widen + // a lock whose names only mean something inside one repository. held := map[string]Claim{} for id, list := range c.claims { if id == sessionID || c.sessions[id].ProjectID != me.ProjectID { diff --git a/internal/coord/connections.go b/internal/coord/connections.go new file mode 100644 index 0000000..4c9a3c3 --- /dev/null +++ b/internal/coord/connections.go @@ -0,0 +1,122 @@ +package coord + +// Connections: which sessions may see each other across the project boundary. +// +// Everything else in this package scopes to a project, because that is where +// sessions share a repository. A connection is the one exception, and it goes +// the other way: it does not narrow a project, it lets a session out of one. +// The case it exists for is an API changing in one repository while its +// consumer changes in another. + +import "sort" + +// Connection joins two sessions. It mirrors the pair the store persists, the +// way Session mirrors the one the store records: this package holds the live +// picture and does not read the document. +type Connection struct { + A string + B string +} + +// SetConnections replaces the picture of who is linked to whom. +// +// Replaces rather than adds, so the caller can hand over the whole set after +// any change without working out what moved. It is indexed here rather than +// searched per call because the predicate runs inside the loop of every tool +// that lists sessions. +func (c *Coordinator) SetConnections(links []Connection) { + c.mu.Lock() + defer c.mu.Unlock() + peers := make(map[string]map[string]bool, len(links)*2) + for _, l := range links { + if l.A == "" || l.B == "" || l.A == l.B { + continue + } + if peers[l.A] == nil { + peers[l.A] = map[string]bool{} + } + if peers[l.B] == nil { + peers[l.B] = map[string]bool{} + } + peers[l.A][l.B] = true + peers[l.B][l.A] = true + } + c.peers = peers +} + +// sees reports whether one session may see another: same project, or joined by +// a connection. Caller holds the lock. +// +// Claiming deliberately does not go through this. A claim is a repo-relative +// path, so two sessions in different repositories both claiming +// internal/api/client.go would collide on a name they do not share. A +// connection widens what a session can read, never the soft lock. +func (c *Coordinator) sees(me, other Session) bool { + return other.ProjectID == me.ProjectID || c.peers[me.ID][other.ID] +} + +// connectedNotes returns what the given sessions wrote in their own projects' +// logs. +// +// Filtered to those sessions by name, not merged whole. A connection joins two +// sessions, so handing over the other project's entire log would publish the +// notes of every session there, including ones nobody connected to. Each +// project's file is read once however many connected sessions live in it. +// +// A log that will not read is reported, not skipped. Notes already returns the +// error when the caller's own log fails, and one os.ReadFile failure must not +// have two answers depending on which project the file belongs to. The +// expected case — a project that has no notes yet — is not a failure at all: +// readNotes turns a missing file into an empty log, so what reaches the error +// here is a real fault. Silence would tell an agent that a connected session +// had written nothing, which is the one conclusion a connection exists to stop +// it drawing. +// +// Called without the lock: it reads files, and the caller has already taken +// the copy of the registry it needs. +func (c *Coordinator) connectedNotes(peers []Session) ([]Note, error) { + if len(peers) == 0 { + return nil, nil + } + byProject := map[string][]string{} + for _, p := range peers { + byProject[p.ProjectID] = append(byProject[p.ProjectID], p.Name) + } + var out []Note + for projectID, names := range byProject { + wanted := make(map[string]bool, len(names)) + for _, n := range names { + wanted[n] = true + } + notes, err := readNotes(c.notesPath(projectID)) + if err != nil { + return nil, err + } + for _, n := range notes { + if wanted[n.Session] { + out = append(out, n) + } + } + } + return out, nil +} + +// connectedElsewhere lists the live sessions a session is joined to that are +// not on its own project, newest picture first read under the lock. +// +// A connected session whose agent has exited is not here, because Unregister +// drops it from the registry. That matches Work, which reads the live registry +// for the same reason: what a finished session left behind is a question about +// what a session means after its agent stops, not something a reader of this +// list should answer by accident. +func (c *Coordinator) connectedElsewhere(me Session) []Session { + var out []Session + for id, s := range c.sessions { + if id == me.ID || s.ProjectID == me.ProjectID || !c.peers[me.ID][id] { + continue + } + out = append(out, s) + } + sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name }) + return out +} diff --git a/internal/coord/connections_test.go b/internal/coord/connections_test.go new file mode 100644 index 0000000..7d7c192 --- /dev/null +++ b/internal/coord/connections_test.go @@ -0,0 +1,265 @@ +package coord + +import ( + "os" + "path/filepath" + "strings" + "testing" +) + +// connectedPair registers one session on each of two projects and links them. +// The pair is the case connections exist for: an API changing in one +// repository while its consumer changes in another. +func connectedPair(t *testing.T, c *Coordinator, theirDir, head string) { + t.Helper() + c.Register(Session{ID: "me", ProjectID: "p1", Project: "gateway", + Name: "scheming-hawk-jhgk", Dir: t.TempDir()}) + c.Register(Session{ID: "them", ProjectID: "p2", Project: "billing-service", + Name: "wily-crane-bbbb", Title: "split the auth middleware", + Dir: theirDir, Branch: "session/wily-crane-bbbb", Isolated: true, BaseRef: head}) + c.SetConnections([]Connection{{A: "me", B: "them"}}) +} + +// TestAConnectionReachesAcrossProjects is the whole point of the feature. +// TestWorkStaysInsideTheProject is its other half: the same two sessions +// without a connection cannot see each other at all. +func TestAConnectionReachesAcrossProjects(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + + if err := os.WriteFile(filepath.Join(dir, "auth.go"), + []byte("package gateway\n\nfunc Auth() {}\n"), 0o644); err != nil { + t.Fatal(err) + } + + out, err := c.Work("me", "wily-crane-bbbb") + if err != nil { + t.Fatalf("a connected session is unreachable: %v", err) + } + if !strings.Contains(out["patch"].(string), "func Auth()") { + t.Errorf("patch missing the change: %v", out["patch"]) + } + // Named because it crosses. Without it the reader resolves the paths in + // that patch against its own tree, where they mean something else or + // nothing. + if got := out["project"]; got != "billing-service" { + t.Errorf("project = %v, want billing-service", got) + } +} + +// TestWorkOnTheSameProjectDoesNotNameAProject is the other side of that rule: +// a result with no project is one from your own, which is what makes the +// field's presence meaningful. +func TestWorkOnTheSameProjectDoesNotNameAProject(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + c.Register(Session{ID: "me", ProjectID: "p1", Project: "gateway", + Name: "scheming-hawk-jhgk", Dir: t.TempDir()}) + c.Register(Session{ID: "them", ProjectID: "p1", Project: "gateway", + Name: "wily-crane-bbbb", Dir: dir, Isolated: true, BaseRef: head}) + + out, err := c.Work("me", "wily-crane-bbbb") + if err != nil { + t.Fatal(err) + } + if _, named := out["project"]; named { + t.Errorf("a same-project result names a project: %v", out["project"]) + } +} + +// TestAConnectionDoesNotWidenClaims is the deliberate exception, and the one +// most likely to be "fixed" by a later reader who sees claims scoping to the +// project while everything beside it scopes to what a session can see. +// +// A claim is a repo-relative path. Two sessions in different repositories both +// working on internal/api/client.go are not in each other's way, and reporting +// a conflict would teach agents to ignore the mechanism. +func TestAConnectionDoesNotWidenClaims(t *testing.T) { + c := start(t) + c.Register(Session{ID: "me", ProjectID: "p1", Name: "scheming-hawk-jhgk", Dir: "/wt/mine"}) + c.Register(Session{ID: "them", ProjectID: "p2", Name: "wily-crane-bbbb", Dir: "/wt/theirs"}) + c.SetConnections([]Connection{{A: "me", B: "them"}}) + + if granted, _ := c.Claim("me", []string{"internal/api/client.go"}, "rewrite"); len(granted) != 1 { + t.Fatalf("granted = %v, want the path", granted) + } + granted, conflicts := c.Claim("them", []string{"internal/api/client.go"}, "consume") + if len(conflicts) != 0 { + t.Errorf("a connection made two repositories collide on one path: %v", conflicts) + } + if len(granted) != 1 { + t.Errorf("granted = %v, want the path", granted) + } +} + +// TestSessionsNamesTheProjectOfAConnectedRow. A connected row's branch and +// claims belong to another repository, and read wrongly without saying so. +func TestSessionsNamesTheProjectOfAConnectedRow(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + + rows := c.Siblings("me") + if len(rows) != 1 { + t.Fatalf("rows = %v, want the connected session", rows) + } + if rows[0]["session"] != "wily-crane-bbbb" { + t.Fatalf("wrong row: %v", rows[0]) + } + if rows[0]["project"] != "billing-service" { + t.Errorf("project = %v, want billing-service", rows[0]["project"]) + } +} + +// TestMessagesReachAConnectedSession keeps the mailbox on the same rule as the +// listing. A session named by the sessions tool that cannot be written to would +// be worse than not listing it. +func TestMessagesReachAConnectedSession(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + + sent, err := c.Send("me", "wily-crane-bbbb", "the response shape changed") + if err != nil { + t.Fatalf("send: %v", err) + } + if len(sent) != 1 || sent[0] != "wily-crane-bbbb" { + t.Fatalf("sent = %v", sent) + } + if got := c.Collect("them"); len(got) != 1 || got[0].From != "scheming-hawk-jhgk" { + t.Errorf("inbox = %v", got) + } +} + +// TestABroadcastReachesAConnectedSession guards the rule Send's doc comment +// states: an empty target reaches everyone the sessions tool lists, which now +// includes a connected session on another project. +// +// Without this the rule is unguarded, and worse than unguarded. +// TestBroadcastReachesEverySibling asserts that an unconnected session on +// another project receives nothing, so the suite reads as "a broadcast never +// leaves the project". Reverting Send to a project test would leave every test +// green while the documented behaviour was gone. +func TestABroadcastReachesAConnectedSession(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + // A second session on the connected project, joined to nobody. It is the + // half that proves the broadcast follows the connection rather than simply + // going everywhere. + c.Register(Session{ID: "other", ProjectID: "p2", Project: "billing-service", + Name: "quiet-vole-cccc", Dir: t.TempDir()}) + + sent, err := c.Send("me", "", "regenerating the schema") + if err != nil { + t.Fatalf("broadcast: %v", err) + } + if len(sent) != 1 || sent[0] != "wily-crane-bbbb" { + t.Fatalf("delivered to %v, want only the connected session", sent) + } + if c.Unread("them") != 1 { + t.Error("the connected session did not receive the broadcast") + } + if c.Unread("other") != 0 { + t.Error("an unconnected session on the connected project received it") + } +} + +// TestNotesFromAConnectedSessionAreVisible covers the read half of the shared +// log, and TestNotesOfAStrangerOnThatProjectStayHidden covers the filter that +// keeps it from being the other project's whole log. +func TestNotesFromAConnectedSessionAreVisible(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + + if err := c.AppendNote("me", "starting on the client"); err != nil { + t.Fatal(err) + } + if err := c.AppendNote("them", "the response shape changed"); err != nil { + t.Fatal(err) + } + + notes, err := c.Notes("me") + if err != nil { + t.Fatal(err) + } + var texts []string + for _, n := range notes { + texts = append(texts, n.Text) + } + if len(texts) != 2 { + t.Fatalf("notes = %v, want mine and the connected one", texts) + } + // Oldest first, merged across two files rather than concatenated. + if texts[0] != "starting on the client" || texts[1] != "the response shape changed" { + t.Errorf("notes out of order: %v", texts) + } +} + +func TestNotesOfAStrangerOnThatProjectStayHidden(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + // A third session on the connected project that nobody linked to. + c.Register(Session{ID: "other", ProjectID: "p2", Project: "billing-service", + Name: "quiet-vole-cccc", Dir: t.TempDir()}) + + if err := c.AppendNote("other", "unrelated refactor in billing"); err != nil { + t.Fatal(err) + } + notes, err := c.Notes("me") + if err != nil { + t.Fatal(err) + } + for _, n := range notes { + if n.Session == "quiet-vole-cccc" { + t.Fatalf("a connection published the whole project log: %v", notes) + } + } +} + +// TestAnUnreadableConnectedLogIsReported keeps one os.ReadFile failure to one +// answer. Notes already returns the error when the caller's own log will not +// read, and skipping the far one would tell an agent that a connected session +// had written nothing — the single conclusion a connection exists to stop it +// drawing. +// +// The log is made unreadable by putting a directory where the file goes, which +// fails for the same reason whatever user the tests run as. A missing file is +// not this case: readNotes turns that into an empty log on purpose. +func TestAnUnreadableConnectedLogIsReported(t *testing.T) { + dir, head := worktreeSession(t) + notes := t.TempDir() + c, err := Start(notes) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = c.Close() }) + connectedPair(t, c, dir, head) + + if err := os.MkdirAll(filepath.Join(notes, "p2.jsonl"), 0o755); err != nil { + t.Fatal(err) + } + if _, err := c.Notes("me"); err == nil { + t.Error("an unreadable connected log was skipped rather than reported") + } +} + +// TestSetConnectionsReplacesRatherThanAdds. The coordinator holds a copy of +// state the store owns, so the only safe update is the whole set: a stale link +// that survives a replace is one the user disconnected and the agents kept. +func TestSetConnectionsReplacesRatherThanAdds(t *testing.T) { + dir, head := worktreeSession(t) + c := start(t) + connectedPair(t, c, dir, head) + + if rows := c.Siblings("me"); len(rows) != 1 { + t.Fatalf("the fixture is not connected: %v", rows) + } + c.SetConnections(nil) + if rows := c.Siblings("me"); len(rows) != 0 { + t.Errorf("a disconnected session is still visible: %v", rows) + } +} diff --git a/internal/coord/coord.go b/internal/coord/coord.go index 2095bc6..61d596d 100644 --- a/internal/coord/coord.go +++ b/internal/coord/coord.go @@ -28,6 +28,12 @@ type Session struct { Branch string Dir string // worktree or project directory + // Project is the project's display name. Carried so a row that crosses the + // project boundary can say which repository it is from: a connected + // session's claims and branch name are read wrongly without it, and the id + // means nothing to an agent. + Project string + // Isolated and BaseRef are what the work tool needs. A session sharing the // project directory has no branch of its own, so its changes cannot be // told apart from anyone else's, and BaseRef is the commit its worktree @@ -47,6 +53,12 @@ type Coordinator struct { // re-read the file on every append. noteLines map[string]int + // peers is who each session may see beyond its own project, indexed both + // ways from the pairs the store holds. Set wholesale by SetConnections and + // never touched by Unregister: a connection outlives the agent process, and + // the session is still in the store when its agent exits. + peers map[string]map[string]bool + // status holds what each session's hooks last reported. Separate from the // sessions map because a report can arrive before or after a session is // registered, and losing one to a race would leave a stale dot. @@ -77,7 +89,11 @@ type Coordinator struct { // Reviewer performs one spawned analysis. Deck uses claude; a caller may // substitute another, and the tests do so the suite neither spends money nor // needs the CLI installed. -type Reviewer func(ctx context.Context, dir, prompt string) (agent.ClaudeRun, error) +// +// onUsage is called with the run's accounting as it accumulates, so a job in +// flight can report what it is using. A reviewer that cannot report until it +// finishes simply never calls it. +type Reviewer func(ctx context.Context, dir, prompt string, onUsage func(agent.Tokens)) (agent.ClaudeRun, error) // Option configures a Coordinator at startup. type Option func(*Coordinator) @@ -100,10 +116,11 @@ func Start(notesDir string, opts ...Option) (*Coordinator, error) { claims: map[string][]Claim{}, inbox: map[string][]Message{}, noteLines: map[string]int{}, + peers: map[string]map[string]bool{}, jobs: map[string]*Job{}, spend: map[string]float64{}, - spawn: func(ctx context.Context, dir, prompt string) (agent.ClaudeRun, error) { - return agent.RunClaude(ctx, dir, prompt, "--permission-mode", "plan") + spawn: func(ctx context.Context, dir, prompt string, onUsage func(agent.Tokens)) (agent.ClaudeRun, error) { + return agent.RunClaude(ctx, dir, prompt, onUsage, "--permission-mode", "plan") }, status: newStatusBoard(), notesDir: notesDir, @@ -120,7 +137,6 @@ func Start(notesDir string, opts ...Option) (*Coordinator, error) { return c, nil } -// Close stops the MCP endpoint. // Close stops the endpoint and everything the coordinator started. // // Cancelling first: a spawned analysis holds no lock and needs none to stop, @@ -136,45 +152,6 @@ func (c *Coordinator) Close() error { return c.server.close() } -// Register adds or updates a live session. -func (c *Coordinator) Register(s Session) { - c.mu.Lock() - defer c.mu.Unlock() - c.sessions[s.ID] = s -} - -// Registered lists the session ids the coordinator currently knows about, so -// the caller can reconcile them against the processes that are actually alive. -func (c *Coordinator) Registered() []string { - c.mu.Lock() - defer c.mu.Unlock() - ids := make([]string, 0, len(c.sessions)) - for id := range c.sessions { - ids = append(ids, id) - } - return ids -} - -// Unregister drops a session and every claim it held. -// -// Releasing on exit is why claims are in memory and not on disk: a claim held -// by a process that is gone is worse than no claim at all, because the next -// agent believes someone is working there. -func (c *Coordinator) Unregister(id string) { - c.mu.Lock() - defer c.mu.Unlock() - delete(c.sessions, id) - delete(c.claims, id) - delete(c.inbox, id) - delete(c.spend, id) - for jid, j := range c.jobs { - if j.From == id { - delete(c.jobs, jid) - } - } - c.status.clear(id) -} - // MCPConfigJSON is the inline --mcp-config that points one session's agent at // its own endpoint. The session id is in the URL path, so the server knows who // is calling without the agent having to prove it. diff --git a/internal/coord/jobs.go b/internal/coord/jobs.go index ceee1d9..b6ca392 100644 --- a/internal/coord/jobs.go +++ b/internal/coord/jobs.go @@ -123,19 +123,33 @@ func newJobID() string { return fmt.Sprintf("a%d", jobSeq.n) } -// Analyses is what a session's spawned reviews cost and how many are still -// running, for the sidebar badge. +// Badge is what the sidebar needs to draw a session's analyses. // -// Two numbers rather than the records themselves: the sidebar has room for a -// glyph and a figure, and returning fifty full reviews so the caller can count +// Figures rather than the records themselves: the sidebar has room for a glyph +// and a number or two, and returning fifty full reviews so the caller can count // them would hand a renderer the whole answer text on every frame. -func (c *Coordinator) Analyses(sessionID string) (running int, spent float64) { +type Badge struct { + // Running is how many analyses are in flight, and Output the tokens they + // have produced between them. Output is the only figure that moves during a + // run — a turn's input and cache counts are final at its first event, and + // its cost is not reported until its last. + Running int + Output int + + // Spent is the session's running total, which only completed runs add to. + Spent float64 +} + +// Analyses is a session's badge. +func (c *Coordinator) Analyses(sessionID string) Badge { c.mu.Lock() defer c.mu.Unlock() + b := Badge{Spent: c.spend[sessionID]} for _, j := range c.jobs { if j.From == sessionID && j.State == JobRunning { - running++ + b.Running++ + b.Output += j.Tokens.Output } } - return running, c.spend[sessionID] + return b } diff --git a/internal/coord/jobs_test.go b/internal/coord/jobs_test.go index dcd8813..91a6b90 100644 --- a/internal/coord/jobs_test.go +++ b/internal/coord/jobs_test.go @@ -36,7 +36,7 @@ func analysableWatching(t *testing.T, run agent.ClaudeRun, spawnErr error, seen t.Fatal(err) } var mu sync.Mutex - c := startWith(t, func(_ context.Context, _, prompt string) (agent.ClaudeRun, error) { + c := startWith(t, func(_ context.Context, _, prompt string, _ func(agent.Tokens)) (agent.ClaudeRun, error) { if seen != nil { mu.Lock() *seen = append(*seen, prompt) @@ -88,7 +88,7 @@ func TestAnalyseReturnsBeforeItFinishes(t *testing.T) { // run open is what makes "still running" a fact rather than a gamble. release := make(chan struct{}) defer close(release) - c := analysableWith(t, func(_ context.Context, _, _ string) (agent.ClaudeRun, error) { + c := analysableWith(t, func(_ context.Context, _, _ string, _ func(agent.Tokens)) (agent.ClaudeRun, error) { <-release return agent.ClaudeRun{}, nil }) @@ -439,7 +439,7 @@ func TestAGeneralReviewGetsTheDefaultQuestion(t *testing.T) { // show it on — the opposite of the visibility the cost reporting exists for. func TestCloseStopsARunningAnalysis(t *testing.T) { cancelled := make(chan bool, 1) - c := analysableWith(t, func(ctx context.Context, _, _ string) (agent.ClaudeRun, error) { + c := analysableWith(t, func(ctx context.Context, _, _ string, _ func(agent.Tokens)) (agent.ClaudeRun, error) { select { case <-ctx.Done(): cancelled <- true @@ -469,6 +469,58 @@ func TestCloseStopsARunningAnalysis(t *testing.T) { } } +// TestARunningAnalysisReportsWhatItIsUsing is the whole value of streaming the +// run rather than waiting for its envelope. Without it a review in flight can +// report only that it exists, and the figures land at the same moment as the +// answer — by which point nobody needs to watch them. +// +// The reviewer here reports usage and then blocks, which is what lets the +// assertions run against a job that is genuinely still going. +func TestARunningAnalysisReportsWhatItIsUsing(t *testing.T) { + reported := make(chan struct{}) + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + + c := analysableWith(t, func(_ context.Context, _, _ string, onUsage func(agent.Tokens)) (agent.ClaudeRun, error) { + onUsage(agent.Tokens{Input: 2, CacheWrite: 7954, Output: 120}) + close(reported) + <-release + return agent.ClaudeRun{Text: "looks sound", CostUSD: 0.05}, nil + }) + + job, err := c.Analyse("me", "wily-crane-bbbb", "") + if err != nil { + t.Fatal(err) + } + select { + case <-reported: + case <-time.After(3 * time.Second): + t.Fatal("the reviewer never reported usage") + } + + got, err := c.Analysis("me", job.ID) + if err != nil { + t.Fatal(err) + } + if got.State != JobRunning { + t.Fatalf("state = %v, want running", got.State) + } + if got.Tokens.Output != 120 || got.Tokens.CacheWrite != 7954 { + t.Errorf("tokens = %+v, want what the run reported mid-flight", got.Tokens) + } + + // And the same figure reaches the sidebar, which is the surface the user + // actually watches. Cost is deliberately absent: the CLI reports it once, + // at the end, so a running badge that showed dollars would show zero. + b := c.Analyses("me") + if b.Running != 1 || b.Output != 120 { + t.Errorf("badge = %+v, want one running analysis at 120 output tokens", b) + } + if b.Spent != 0 { + t.Errorf("spent = %v before the run finished; that figure cannot be known yet", b.Spent) + } +} + // startWith is start with a substituted reviewer. func startWith(t *testing.T, r Reviewer) *Coordinator { t.Helper() diff --git a/internal/coord/messages.go b/internal/coord/messages.go index 4d6f87d..e500deb 100644 --- a/internal/coord/messages.go +++ b/internal/coord/messages.go @@ -20,8 +20,14 @@ type Message struct { // should not grow without limit; the oldest go first. const maxInbox = 50 -// Send queues a message for one sibling by session name, or for every sibling -// on the project when to is empty. It returns the names it reached. +// Send queues a message for one session by name, or for every session the +// sender can see when to is empty. It returns the names it reached. +// +// "Can see" is the same rule the sessions tool lists by, so a broadcast reaches +// exactly the sessions that tool named and a connected session is not a +// surprise recipient. Two rules — one for listing, a narrower one for +// broadcasting — would mean a session you were told about that you cannot +// reach. // // Delivery is a mailbox the recipient collects, not a write into its terminal. // Typing into a live pane would corrupt the input of an agent mid-turn, and @@ -45,7 +51,7 @@ func (c *Coordinator) Send(fromID, to, text string) ([]string, error) { msg := Message{At: time.Now(), From: me.Name, Text: text} var sent []string for id, s := range c.sessions { - if id == fromID || s.ProjectID != me.ProjectID { + if id == fromID || !c.sees(me, s) { continue } if to != "" && s.Name != to { @@ -60,7 +66,7 @@ func (c *Coordinator) Send(fromID, to, text string) ([]string, error) { } sort.Strings(sent) if to != "" && len(sent) == 0 { - return nil, fmt.Errorf("no live session named %q on this project", to) + return nil, fmt.Errorf("no live session named %q that this session can see", to) } return sent, nil } diff --git a/internal/coord/notes.go b/internal/coord/notes.go index 5b4ecc4..ccee118 100644 --- a/internal/coord/notes.go +++ b/internal/coord/notes.go @@ -7,6 +7,7 @@ import ( "fmt" "os" "path/filepath" + "sort" "strings" "time" ) @@ -127,25 +128,59 @@ func countLines(path string) int { // maxNotes caps what a reader gets back. A shared log is re-read by every // agent that asks, and an uncapped one quietly becomes the most expensive file // in the project. +// +// It bounds the answer, not each source. A connected session's notes are +// merged in by time, so a busy far session can push a reader's own project +// notes past the cap: connecting can leave you seeing less of your own log, +// not more. That is the intended trade — the newest fifty a session can see is +// what it asked for — but it is the kind of thing nobody expects from a +// feature described as widening, so it is written down here rather than found. const maxNotes = 50 -// Notes returns the most recent notes for the session's project, oldest first. +// Notes returns the most recent notes the session can read, oldest first: its +// own project's log, plus what any connected session wrote in theirs. +// +// A note is still written to one log — the writer's own project's — so there is +// one write path and a project's log means what it always meant. Only the read +// widens. See connectedNotes for what a connection contributes. func (c *Coordinator) Notes(sessionID string) ([]Note, error) { c.mu.Lock() me, ok := c.sessions[sessionID] + var peers []Session + if ok { + peers = c.connectedElsewhere(me) + } c.mu.Unlock() if !ok { return nil, fmt.Errorf("unknown session") } - data, err := os.ReadFile(c.notesPath(me.ProjectID)) + out, err := readNotes(c.notesPath(me.ProjectID)) + if err != nil { + return nil, err + } + shared, err := c.connectedNotes(peers) + if err != nil { + return nil, err + } + out = append(out, shared...) + sort.Slice(out, func(i, j int) bool { return out[i].At.Before(out[j].At) }) + if len(out) > maxNotes { + out = out[len(out)-maxNotes:] + } + return out, nil +} + +// readNotes parses one log file. A missing file is an empty log, which is +// every project before its first note. +func readNotes(path string) ([]Note, error) { + data, err := os.ReadFile(path) if os.IsNotExist(err) { return nil, nil } if err != nil { return nil, fmt.Errorf("read notes: %w", err) } - var out []Note for _, line := range strings.Split(string(data), "\n") { if strings.TrimSpace(line) == "" { @@ -157,9 +192,6 @@ func (c *Coordinator) Notes(sessionID string) ([]Note, error) { } out = append(out, n) } - if len(out) > maxNotes { - out = out[len(out)-maxNotes:] - } return out, nil } diff --git a/internal/coord/registry.go b/internal/coord/registry.go new file mode 100644 index 0000000..a099a49 --- /dev/null +++ b/internal/coord/registry.go @@ -0,0 +1,48 @@ +package coord + +// Sessions arriving and leaving: what the coordinator knows about a live agent +// and what it forgets when that agent goes. + +// Register adds or updates a live session. +func (c *Coordinator) Register(s Session) { + c.mu.Lock() + defer c.mu.Unlock() + c.sessions[s.ID] = s +} + +// Registered lists the session ids the coordinator currently knows about, so +// the caller can reconcile them against the processes that are actually alive. +func (c *Coordinator) Registered() []string { + c.mu.Lock() + defer c.mu.Unlock() + ids := make([]string, 0, len(c.sessions)) + for id := range c.sessions { + ids = append(ids, id) + } + return ids +} + +// Unregister drops a session and every claim it held. +// +// Releasing on exit is why claims are in memory and not on disk: a claim held +// by a process that is gone is worse than no claim at all, because the next +// agent believes someone is working there. +// +// Connections are not dropped here, and that is the difference between them +// and everything else in this list. A claim, an inbox and a spend belong to a +// running process; a connection belongs to the session, survives its agent +// exiting, and is removed only when the store removes the session itself. +func (c *Coordinator) Unregister(id string) { + c.mu.Lock() + defer c.mu.Unlock() + delete(c.sessions, id) + delete(c.claims, id) + delete(c.inbox, id) + delete(c.spend, id) + for jid, j := range c.jobs { + if j.From == id { + delete(c.jobs, jid) + } + } + c.status.clear(id) +} diff --git a/internal/coord/work.go b/internal/coord/work.go index d3da70d..92b0de8 100644 --- a/internal/coord/work.go +++ b/internal/coord/work.go @@ -26,7 +26,7 @@ import ( // publishes and what message already accepts, so an agent has only one kind of // handle to learn. func (c *Coordinator) Work(sessionID, target string) (map[string]any, error) { - found, w, err := c.workOf(sessionID, target) + found, w, elsewhere, err := c.workOf(sessionID, target) if err != nil { return nil, err } @@ -38,29 +38,40 @@ func (c *Coordinator) Work(sessionID, target string) (map[string]any, error) { "summary": w.Stat, "patch": w.Patch, } + // Named only when it differs. A result without a project is from your own, + // so the paths in the patch resolve against the tree you are looking at; a + // result with one is from another repository, where they do not. + if elsewhere { + out["project"] = found.Project + } if w.Truncated { out["truncated"] = true } return out, nil } -// workOf resolves a sibling by name and reads its changes. Both the work tool +// workOf resolves a session by name and reads its changes. Both the work tool // and a spawned analysis go through it, so the scoping and the two refusals // are stated once rather than in each caller. -func (c *Coordinator) workOf(sessionID, target string) (*Session, gitx.Work, error) { +// +// elsewhere reports that the target is on another project, reached through a +// connection. Both callers need it and neither can derive it: they hold the +// target but not the caller's own session, and re-reading that would mean +// taking the lock a second time for a fact this one already knows. +func (c *Coordinator) workOf(sessionID, target string) (found *Session, w gitx.Work, elsewhere bool, err error) { c.mu.Lock() me, ok := c.sessions[sessionID] if !ok { c.mu.Unlock() - return nil, gitx.Work{}, fmt.Errorf("unknown session") + return nil, gitx.Work{}, false, fmt.Errorf("unknown session") } - var found *Session for id, s := range c.sessions { - if id == sessionID || s.ProjectID != me.ProjectID || s.Name != target { + if id == sessionID || s.Name != target || !c.sees(me, s) { continue } copied := s found = &copied + elsewhere = s.ProjectID != me.ProjectID break } // Unlocked by hand rather than deferred, unlike every other method here: @@ -69,19 +80,19 @@ func (c *Coordinator) workOf(sessionID, target string) (*Session, gitx.Work, err c.mu.Unlock() if found == nil { - return nil, gitx.Work{}, fmt.Errorf("no live session named %q on this project", target) + return nil, gitx.Work{}, false, fmt.Errorf("no live session named %q that this session can see", target) } // Refused rather than answered approximately. A shared project directory // holds everyone's edits at once, so a diff of it would credit this // session with work it may not have done. if !found.Isolated { - return nil, gitx.Work{}, fmt.Errorf("%s runs in the project directory, so its changes "+ + return nil, gitx.Work{}, false, fmt.Errorf("%s runs in the project directory, so its changes "+ "cannot be told apart from any other work there", found.Name) } - w, err := gitx.Diff(found.Dir, found.BaseRef) + w, err = gitx.Diff(found.Dir, found.BaseRef) if err != nil { - return nil, gitx.Work{}, fmt.Errorf("read %s: %w", found.Name, err) + return nil, gitx.Work{}, false, fmt.Errorf("read %s: %w", found.Name, err) } - return found, w, nil + return found, w, elsewhere, nil } diff --git a/internal/store/connections.go b/internal/store/connections.go new file mode 100644 index 0000000..9699240 --- /dev/null +++ b/internal/store/connections.go @@ -0,0 +1,92 @@ +package store + +// Connections between sessions: which pairs may see each other across the +// project boundary that otherwise separates them. + +import "sort" + +// Connection joins two sessions so each can see the other's work. +// +// A pair rather than a named set. Sessions on one project already see each +// other, so the case project scope cannot express is one repository's API +// changing while its consumer changes in another — and that case is two +// sessions. A set would need a name, a member editor, and a rule for what +// happens when the last member leaves, none of which that case asks for. +// Connecting A to B and A to C therefore lets A see both without making B and +// C visible to each other. +// +// It lives on State rather than as a field on Session so there is no two-way +// link to keep consistent: the pair is one record, and dropping it drops the +// whole relationship. +type Connection struct { + A string `json:"a"` + B string `json:"b"` +} + +// joins reports whether the pair names these two sessions, in either order. +// Order is not meaningful — A asked and B was picked, which is history, not a +// direction — so every comparison goes through here rather than testing the +// fields and forgetting the second case. +func (c Connection) joins(x, y string) bool { + return (c.A == x && c.B == y) || (c.A == y && c.B == x) +} + +// Connect links two sessions and reports whether the link was new. Asking +// twice is not an error: the caller is expressing a state, not an increment. +func (s *State) Connect(a, b string) bool { + if a == "" || b == "" || a == b { + return false + } + for _, c := range s.Connections { + if c.joins(a, b) { + return false + } + } + s.Connections = append(s.Connections, Connection{A: a, B: b}) + return true +} + +// Disconnect drops the link between two sessions and reports whether there was +// one. +func (s *State) Disconnect(a, b string) bool { + for i, c := range s.Connections { + if c.joins(a, b) { + s.Connections = append(s.Connections[:i], s.Connections[i+1:]...) + return true + } + } + return false +} + +// disconnectAll drops every link a session holds. +// +// Called from RemoveSession rather than exported for the caller to remember. +// A pair naming a session that no longer exists is not dangerous — ids are +// eight random bytes, so nothing inherits one — but it is a record that +// accumulates and that every reader has to skip. +func (s *State) disconnectAll(id string) { + out := s.Connections[:0] + for _, c := range s.Connections { + if c.A != id && c.B != id { + out = append(out, c) + } + } + s.Connections = out +} + +// ConnectedTo lists the sessions linked to id, sorted. Sorted because the +// order links were made is not information anyone wants, and an unstable order +// makes a rendered list move under the cursor. +func (s *State) ConnectedTo(id string) []string { + var out []string + for _, c := range s.Connections { + switch id { + case c.A: + out = append(out, c.B) + case c.B: + out = append(out, c.A) + } + } + sort.Strings(out) + return out +} diff --git a/internal/store/connections_test.go b/internal/store/connections_test.go new file mode 100644 index 0000000..6a2f6d4 --- /dev/null +++ b/internal/store/connections_test.go @@ -0,0 +1,81 @@ +package store + +import ( + "slices" + "testing" +) + +func TestConnectingTwiceIsOneLink(t *testing.T) { + s := &State{} + if !s.Connect("a", "b") { + t.Fatal("the first connection was refused") + } + // The same pair the other way round is the same relationship. Recording it + // twice would give Disconnect a link to remove and one to leave behind. + if s.Connect("b", "a") { + t.Error("connecting the same pair in reverse made a second link") + } + if len(s.Connections) != 1 { + t.Errorf("connections = %v, want one", s.Connections) + } +} + +func TestConnectedToReadsBothWays(t *testing.T) { + s := &State{} + s.Connect("a", "b") + s.Connect("c", "a") + + if got := s.ConnectedTo("a"); !slices.Equal(got, []string{"b", "c"}) { + t.Errorf("a is connected to %v, want [b c]", got) + } + // The link is not transitive. b and c are each joined to a and to nothing + // else, which is what makes a pair a pair rather than a group. + if got := s.ConnectedTo("b"); !slices.Equal(got, []string{"a"}) { + t.Errorf("b is connected to %v, want [a]", got) + } +} + +func TestASessionCannotConnectToItself(t *testing.T) { + s := &State{} + if s.Connect("a", "a") { + t.Error("a session connected to itself") + } + if len(s.Connections) != 0 { + t.Errorf("connections = %v, want none", s.Connections) + } +} + +// TestRemovingASessionDropsItsConnections is why disconnectAll is called from +// RemoveSession rather than left to the caller. A pair naming a session that +// is gone is a record every reader has to skip, and the caller that forgets is +// the one nobody reviews again. +func TestRemovingASessionDropsItsConnections(t *testing.T) { + s := &State{} + a := s.AddSession(Session{ProjectID: "p1", Name: "swift-otter-aaaa"}) + b := s.AddSession(Session{ProjectID: "p2", Name: "wily-crane-bbbb"}) + s.Connect(a.ID, b.ID) + + s.RemoveSession(a.ID) + + if len(s.Connections) != 0 { + t.Errorf("connections = %v, want none after the session went", s.Connections) + } + if got := s.ConnectedTo(b.ID); len(got) != 0 { + t.Errorf("the surviving session is still connected to %v", got) + } +} + +func TestDisconnectReportsWhetherThereWasALink(t *testing.T) { + s := &State{} + s.Connect("a", "b") + + if !s.Disconnect("b", "a") { + t.Error("disconnecting in reverse order found nothing") + } + if s.Disconnect("a", "b") { + t.Error("disconnecting twice reported a second removal") + } + if len(s.Connections) != 0 { + t.Errorf("connections = %v, want none", s.Connections) + } +} diff --git a/internal/store/query.go b/internal/store/query.go index 524ae46..3dd6855 100644 --- a/internal/store/query.go +++ b/internal/store/query.go @@ -104,9 +104,9 @@ func (s *State) AddSession(sess Session) *Session { return &s.Sessions[len(s.Sessions)-1] } -// RemoveSession drops a session from the store. It does not touch the -// worktree on disk; the caller decides that, because removing a worktree can -// destroy uncommitted work. +// RemoveSession drops a session from the store, and every connection it held. +// It does not touch the worktree on disk; the caller decides that, because +// removing a worktree can destroy uncommitted work. func (s *State) RemoveSession(id string) { out := s.Sessions[:0] for _, sess := range s.Sessions { @@ -115,4 +115,5 @@ func (s *State) RemoveSession(id string) { } } s.Sessions = out + s.disconnectAll(id) } diff --git a/internal/store/store.go b/internal/store/store.go index 58d631a..eb5ee77 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -17,8 +17,9 @@ import ( // State is the whole persisted document. type State struct { - Projects []Project `json:"projects"` - Sessions []Session `json:"sessions"` + Projects []Project `json:"projects"` + Sessions []Session `json:"sessions"` + Connections []Connection `json:"connections,omitempty"` } // Load reads the state file. A missing file is an empty store, which is the diff --git a/internal/ui/connections.go b/internal/ui/connections.go new file mode 100644 index 0000000..d7af16f --- /dev/null +++ b/internal/ui/connections.go @@ -0,0 +1,98 @@ +package ui + +// Connecting a session to one on another project: choosing the far end, and +// keeping the coordinator's picture in step with the store's. + +import ( + "fmt" + + tea "github.com/charmbracelet/bubbletea" + + "github.com/tripledownab/deck/internal/coord" + "github.com/tripledownab/deck/internal/store" +) + +// syncConnections hands the coordinator the whole set of links. +// +// The whole set every time, not the one that changed. The store owns the +// document and the coordinator owns the live picture; a coordinator that +// applied deltas would be a second copy of the same state, free to drift the +// moment one update is missed. Sending everything makes a missed update +// impossible to observe. +func (m Model) syncConnections() { + if m.coord == nil { + return + } + links := make([]coord.Connection, 0, len(m.state.Connections)) + for _, c := range m.state.Connections { + links = append(links, coord.Connection{A: c.A, B: c.B}) + } + m.coord.SetConnections(links) +} + +// openConnectPicker lists the sessions the selected one could be connected to. +// +// Only sessions on other projects are offered. Sessions on the same project +// already see each other, so listing one would offer a link that changes +// nothing — and a menu entry that does nothing is how a user learns to distrust +// the menu. +// +// Rows for sessions already connected read "connected", and choosing one +// disconnects. One key does both directions, which is what keeps this off a +// screen of its own. +func (m Model) openConnectPicker(sess *store.Session) (tea.Model, tea.Cmd) { + if sess == nil { + return m, nil + } + connected := map[string]bool{} + for _, id := range m.state.ConnectedTo(sess.ID) { + connected[id] = true + } + + var rows []pickerRow + for _, other := range m.state.Sessions { + if other.ProjectID == sess.ProjectID { + continue + } + desc := projectLabel(m.state, other.ProjectID) + if other.Branch != "" { + desc += " · " + other.Branch + } + if connected[other.ID] { + desc += " · connected" + } + rows = append(rows, pickerRow{id: other.ID, label: sessionLabel(other), desc: desc}) + } + if len(rows) == 0 { + m.notice = "no sessions on another project to connect to" + return m, nil + } + m.connectFrom = sess.ID + m.picker = newPicker(pickConnect, "Connect "+sessionLabel(*sess)+" to", rows, "") + return m, nil +} + +// connectPicked links or unlinks the chosen session and persists the result. +func (m Model) connectPicked(otherID string) (tea.Model, tea.Cmd) { + sess := m.state.Session(m.connectFrom) + other := m.state.Session(otherID) + if sess == nil || other == nil { + return m, nil + } + if m.state.Disconnect(sess.ID, other.ID) { + m.notice = "disconnected " + sessionLabel(*other) + } else { + m.state.Connect(sess.ID, other.ID) + m.notice = "connected to " + sessionLabel(*other) + + " on " + projectLabel(m.state, other.ProjectID) + } + // Saved before the coordinator is told. A link the agents can act on but + // the store never recorded would vanish at the next restart, and the user + // would have watched it work. + if err := m.state.Save(); err != nil { + m.fault = fmt.Errorf("save connection: %w", err) + return m, nil + } + m.syncConnections() + return m, nil +} diff --git a/internal/ui/format.go b/internal/ui/format.go new file mode 100644 index 0000000..cb7bc2b --- /dev/null +++ b/internal/ui/format.go @@ -0,0 +1,58 @@ +package ui + +// Turning model values into the short strings a row can hold: what a project +// and a session are called, and how a figure is abbreviated to fit beside them. + +import ( + "fmt" + "path/filepath" + + "github.com/tripledownab/deck/internal/store" +) + +// compactCount renders a token count for a sidebar badge. +// +// Thousands are abbreviated because the badge shares a status line that +// truncates from the right, and the digits past the first two are noise on a +// figure that changes several times a second. 950 stays 950; 1500 becomes +// 1.5k; 23800 becomes 24k, since a tenth of a thousand is below what anyone +// reads at that size. +func compactCount(n int) string { + switch { + case n < 1000: + return fmt.Sprintf("%d", n) + case n < 10000: + return fmt.Sprintf("%.1fk", float64(n)/1000) + default: + return fmt.Sprintf("%.0fk", float64(n)/1000) + } +} + +// sessionLabel is what a session is called on screen. +// +// The title if it has one, otherwise the generated name. A session is created +// with a title and Deck never clears it, so the fallback is for a session +// recorded before titles existed and for one titled with whitespace. +func sessionLabel(sess store.Session) string { + if sess.Title != "" { + return sess.Title + } + return sess.Name +} + +// projectLabel is what a project is called on screen. +// +// Name is filled in from the directory when the field is left empty, both when +// registering and when renaming, so the fallback here is not the usual path. It +// stands for a hand-edited state file: a blank row in a list of projects tells +// the user nothing about which one it is. +func projectLabel(state *store.State, projectID string) string { + p := state.Project(projectID) + if p == nil { + return "unknown project" + } + if p.Name != "" { + return p.Name + } + return filepath.Base(p.Path) +} diff --git a/internal/ui/format_test.go b/internal/ui/format_test.go new file mode 100644 index 0000000..5fc8f88 --- /dev/null +++ b/internal/ui/format_test.go @@ -0,0 +1,61 @@ +package ui + +import ( + "testing" + + "github.com/tripledownab/deck/internal/store" +) + +// TestCompactCountSwitchesUnitAtTheRightPlace pins both boundaries. The badge +// is the only report of what a spawned review is using, and a figure that +// changes unit one short of where it should reads as a review a thousand times +// larger or smaller than it is. +func TestCompactCountSwitchesUnitAtTheRightPlace(t *testing.T) { + for _, tc := range []struct { + n int + want string + }{ + {0, "0"}, + {999, "999"}, // last value printed in full + {1000, "1.0k"}, // first abbreviated one + {1500, "1.5k"}, + {9999, "10.0k"}, // last with a tenth + {10000, "10k"}, // first without one + {23800, "24k"}, + } { + if got := compactCount(tc.n); got != tc.want { + t.Errorf("compactCount(%d) = %q, want %q", tc.n, got, tc.want) + } + } +} + +// TestSessionLabelFallsBackToTheGeneratedName covers the rule the sidebar used +// to state inline. A session recorded before titles existed has none, and a +// blank row tells the reader nothing about which session it is. +func TestSessionLabelFallsBackToTheGeneratedName(t *testing.T) { + if got := sessionLabel(store.Session{Title: "rate limiting", Name: "swift-otter-aaaa"}); got != "rate limiting" { + t.Errorf("label = %q, want the title", got) + } + if got := sessionLabel(store.Session{Name: "swift-otter-aaaa"}); got != "swift-otter-aaaa" { + t.Errorf("label = %q, want the generated name", got) + } +} + +// TestProjectLabelFallsBackToTheDirectory. Name is filled in from the +// directory when the field is left empty, so this stands for a hand-edited +// state file rather than the usual path. +func TestProjectLabelFallsBackToTheDirectory(t *testing.T) { + st := &store.State{} + named := st.AddProject(store.Project{Name: "api-gateway", Path: "/src/gw"}) + blank := st.AddProject(store.Project{Path: "/src/billing-service"}) + + if got := projectLabel(st, named.ID); got != "api-gateway" { + t.Errorf("label = %q, want the name", got) + } + if got := projectLabel(st, blank.ID); got != "billing-service" { + t.Errorf("label = %q, want the directory", got) + } + if got := projectLabel(st, "gone"); got == "" { + t.Error("an unknown project rendered as an empty label") + } +} diff --git a/internal/ui/keyroutes.go b/internal/ui/keyroutes.go index bcf13de..dd694cd 100644 --- a/internal/ui/keyroutes.go +++ b/internal/ui/keyroutes.go @@ -30,6 +30,8 @@ func (m Model) dashboardKey(msg tea.KeyMsg) (tea.Model, tea.Cmd) { return m.openBrowser() case key.Matches(msg, m.keys.Rename): return m.openEditProjectForm() + case key.Matches(msg, m.keys.Connect): + return m.openConnectPicker(m.dashboardSession()) case key.Matches(msg, m.keys.Theme): return m.openThemePicker() case key.Matches(msg, m.keys.Delete): @@ -56,6 +58,8 @@ func (m Model) sessionKey(msg tea.KeyMsg) (tea.Model, tea.Cmd) { return m.openNewSessionForm() case key.Matches(msg, m.keys.Delete): m.stopCurrent() + case key.Matches(msg, m.keys.Connect): + return m.openConnectPicker(m.currentSession()) case key.Matches(msg, m.keys.Theme): return m.openThemePicker() case msg.String() == "d": diff --git a/internal/ui/keys.go b/internal/ui/keys.go index 4729fc9..722f576 100644 --- a/internal/ui/keys.go +++ b/internal/ui/keys.go @@ -32,6 +32,7 @@ type keyMap struct { AddProject key.Binding Rename key.Binding Delete key.Binding + Connect key.Binding Theme key.Binding // Command mode, reached through the prefix. @@ -64,6 +65,7 @@ func defaultKeys() keyMap { AddProject: key.NewBinding(key.WithKeys("a"), key.WithHelp("a", "add project")), Rename: key.NewBinding(key.WithKeys("e"), key.WithHelp("e", "rename project")), Delete: key.NewBinding(key.WithKeys("x"), key.WithHelp("x", "close session")), + Connect: key.NewBinding(key.WithKeys("c"), key.WithHelp("c", "connect to a session elsewhere")), Theme: key.NewBinding(key.WithKeys("t"), key.WithHelp("t", "theme")), Dashboard: key.NewBinding(key.WithKeys("d"), key.WithHelp("^g d", "dashboard")), @@ -105,7 +107,7 @@ func (k keyMap) helpGroups() []helpGroup { return []helpGroup{ {"CHROME", []key.Binding{ k.Up, k.Down, k.Left, k.Right, k.SwitchCol, k.Enter, - k.NewSession, k.AddProject, k.Rename, k.Delete, k.Theme, k.Help, k.Quit, + k.NewSession, k.AddProject, k.Rename, k.Delete, k.Connect, k.Theme, k.Help, k.Quit, }}, {"COMMAND — press " + PrefixKey + " first", []key.Binding{ k.Dashboard, k.Sessions, k.NextSess, k.JumpSess, k.NewSessCmd, k.StopSessCmd, diff --git a/internal/ui/model.go b/internal/ui/model.go index c608c28..ad76f9e 100644 --- a/internal/ui/model.go +++ b/internal/ui/model.go @@ -55,6 +55,15 @@ type Model struct { picker *picker showHelp bool + // connectFrom is the session the open connect picker is linking from. + // + // Captured when the picker opens rather than resolved on the commit key, + // for the same reason picker.restore is: the modal outlives the keystroke + // that opened it, and the dashboard and the session view resolve "the + // selected session" by different rules. Re-reading it would connect + // whichever session the other screen happens to be pointing at. + connectFrom string + // notice is a one-line message in the footer; fault is a failure the user // must see, and it outranks the notice. // @@ -92,8 +101,13 @@ func (m Model) WithSettings(s store.Settings) Model { } // WithCoordinator attaches the cross-session coordination server. +// +// The persisted connections go over immediately. A link survives the app, so +// the coordinator starting without them would leave every restored session +// blind to its far end until something happened to re-send the set. func (m Model) WithCoordinator(c *coord.Coordinator) Model { m.coord = c + m.syncConnections() return m } diff --git a/internal/ui/picker.go b/internal/ui/picker.go index a2506bc..a2bdb8e 100644 --- a/internal/ui/picker.go +++ b/internal/ui/picker.go @@ -17,6 +17,7 @@ type pickerKind int const ( pickTheme pickerKind = iota pickProject + pickConnect ) // picker is a modal list: a title, rows, a cursor. diff --git a/internal/ui/pickerctl.go b/internal/ui/pickerctl.go index e28d261..18aaed6 100644 --- a/internal/ui/pickerctl.go +++ b/internal/ui/pickerctl.go @@ -41,6 +41,21 @@ func (m Model) openFieldPicker() (tea.Model, tea.Cmd) { func (m Model) pickerKey(msg tea.KeyMsg) (tea.Model, tea.Cmd) { commit, cancel := m.picker.update(msg) + // Connecting previews nothing: a link is a change to the store, and + // applying one as the cursor moves would write a record for every session + // scrolled past. So it acts on commit only, and cancel just closes. + if m.picker.kind == pickConnect { + switch { + case cancel: + m.picker = nil + case commit: + id := m.picker.selected() + m.picker = nil + return m.connectPicked(id) + } + return m, nil + } + // A field picker floats over the open form and only writes back into it. // Nothing is previewed and nothing is persisted, so cancel is simply // closing it. diff --git a/internal/ui/runner.go b/internal/ui/runner.go index d7c77f2..4613278 100644 --- a/internal/ui/runner.go +++ b/internal/ui/runner.go @@ -44,7 +44,7 @@ func (m Model) attach() (tea.Model, tea.Cmd) { // process starts keeps the registry to sessions that actually exist. if m.coord != nil { m.coord.Register(coord.Session{ - ID: sess.ID, ProjectID: sess.ProjectID, + ID: sess.ID, ProjectID: sess.ProjectID, Project: projectLabel(m.state, sess.ProjectID), Name: sess.Name, Title: sess.Title, Branch: sess.Branch, Dir: sess.Dir, Isolated: sess.Isolated, BaseRef: sess.BaseRef, }) diff --git a/internal/ui/sessions.go b/internal/ui/sessions.go index ec43981..f613a41 100644 --- a/internal/ui/sessions.go +++ b/internal/ui/sessions.go @@ -116,6 +116,29 @@ func (m Model) openFromDashboard() (tea.Model, tea.Cmd) { return m.attach() } +// dashboardSession is the session the dashboard cursor is on, or nil. +// +// The dashboard resolves a session from the highlighted project and the list +// index, which is a different rule from the session view's cursor. +// +// openFromDashboard deliberately does not use this. It *clamps* listIx into +// range and writes it back, where this *guards* and answers nil, and the two +// disagree exactly when the index is past the end: opening lands on the last +// session, while closing and connecting do nothing. That is the right split — +// opening a session the cursor is near is helpful, closing or connecting one +// the user cannot see is not — so do not fold them together. +func (m Model) dashboardSession() *store.Session { + p := m.currentProject() + if p == nil { + return nil + } + sessions := m.state.SessionsFor(p.ID) + if len(sessions) == 0 || m.listIx >= len(sessions) { + return nil + } + return m.state.Session(sessions[m.listIx].ID) +} + // closeSelectedFromDashboard stops a session's agent and forgets it. // // The worktree is left on disk on purpose. It may hold uncommitted work, and @@ -123,24 +146,25 @@ func (m Model) openFromDashboard() (tea.Model, tea.Cmd) { // gets to make silently. The notice says where it went. func (m *Model) closeSelectedFromDashboard() { p := m.currentProject() - if p == nil { - return - } - sessions := m.state.SessionsFor(p.ID) - if len(sessions) == 0 || m.listIx >= len(sessions) { + selected := m.dashboardSession() + if p == nil || selected == nil { return } - sess := sessions[m.listIx] + sess := *selected if r, ok := m.runners[sess.ID]; ok { r.Stop() delete(m.runners, sess.ID) } m.releaseCoord(sess.ID) + // RemoveSession drops the links this session held, so the coordinator has + // to be told: releaseCoord frees claims and the inbox, but peers is set + // wholesale and outlives an agent exiting on purpose. m.state.RemoveSession(sess.ID) if err := m.state.Save(); err != nil { m.fault = err return } + m.syncConnections() m.rebuildRows() m.listIx = clamp(m.listIx, 0, max(len(m.state.SessionsFor(p.ID))-1, 0)) if sess.Isolated { diff --git a/internal/ui/sidebar.go b/internal/ui/sidebar.go index c0e0d68..8b8076f 100644 --- a/internal/ui/sidebar.go +++ b/internal/ui/sidebar.go @@ -87,10 +87,7 @@ func (m Model) sessionCard(sess *store.Session, active bool, width, nth int) []s num, numW = s.Accent.Render(strconv.Itoa(nth))+" ", 2 } - title := sess.Title - if title == "" { - title = sess.Name - } + title := sessionLabel(*sess) ref := sess.Branch if ref == "" { @@ -127,12 +124,23 @@ func (m Model) sessionCard(sess *store.Session, active bool, width, nth int) []s // Coloured by the same rule as the two above. A run in flight is // spending now, which is the accent's job; a settled total is a fact // about the past, like a claim. - if n, spent := m.coord.Analyses(sess.ID); n > 0 || spent > 0 { - if n > 0 { - label += s.Accent.Render(fmt.Sprintf(" ⚗ %d · $%.2f", n, spent)) - } else { - label += s.Faint.Render(fmt.Sprintf(" ⚗ $%.2f", spent)) - } + // + // A running badge shows tokens, not dollars. The CLI reports cost once, + // in the result event, so the only figure that can move while a review + // runs is the output it has produced — and a dollar figure that sat + // still for the whole run would read as a review that was not costing + // anything. + // The total is kept alongside the live count rather than replaced by it. + // A session that has spent two dollars over five reviews and is running + // a sixth should not read as though it had spent nothing. + switch b := m.coord.Analyses(sess.ID); { + case b.Running > 0 && b.Spent > 0: + label += s.Accent.Render(fmt.Sprintf(" ⚗ %d · %s · $%.2f", + b.Running, compactCount(b.Output), b.Spent)) + case b.Running > 0: + label += s.Accent.Render(fmt.Sprintf(" ⚗ %d · %s", b.Running, compactCount(b.Output))) + case b.Spent > 0: + label += s.Faint.Render(fmt.Sprintf(" ⚗ $%.2f", b.Spent)) } } diff --git a/internal/ui/status_test.go b/internal/ui/status_test.go index 45b8af0..02dbaee 100644 --- a/internal/ui/status_test.go +++ b/internal/ui/status_test.go @@ -189,7 +189,12 @@ func TestSidebarSurvivesALoadedStatusLine(t *testing.T) { func TestSidebarShowsWhatAnalysesCost(t *testing.T) { release := make(chan struct{}) c, err := coord.Start(t.TempDir(), coord.WithReviewer( - func(ctx context.Context, _, _ string) (agent.ClaudeRun, error) { + func(ctx context.Context, _, _ string, onUsage func(agent.Tokens)) (agent.ClaudeRun, error) { + // Reported before blocking, so the assertions below run against a + // review that is genuinely mid-flight. A stub that never reports + // leaves the live figure at zero and the badge's token path + // unexercised, which is how it went untested the first time. + onUsage(agent.Tokens{Input: 2, CacheWrite: 7954, Output: 1500}) <-release return agent.ClaudeRun{Text: "fine", CostUSD: 0.37}, nil })) @@ -224,7 +229,10 @@ func TestSidebarShowsWhatAnalysesCost(t *testing.T) { if _, err := c.Analyse(asker.ID, worker.Name, ""); err != nil { t.Fatal(err) } - waitForBadge(t, m, "⚗ 1") + // Tokens while it runs: the CLI reports a cost only when a turn ends, so a + // dollar figure here would sit at zero and read as a review costing + // nothing. + waitForBadge(t, m, "⚗ 1 · 1.5k") close(release) waitForBadge(t, m, "$0.37")