b57894b183
Introduces the browser-facing surface and the worker-side protocol that backs it: - internal/ui: joined read model plus per-task lifecycle and approval controls, kept separate from the raw endpoints workers and harnesses depend on. - internal/webui + web/: Vite/React app, build output embedded via go:embed and served as an SPA fallback. - federation: per-(worker, task) captures with a monotonic revision that advances only when pane text actually changes, and a command queue restricted to grant_approval / deny_approval, each bound to the capture revision the operator acted on. - orchestra-worker: publishes captures and executes commands only after re-reading the pane and confirming the revision still matches. Sends keystrokes only for a visible y/n prompt or OpenCode's fully labelled selector, and refuses to deny through that selector rather than guess at unobservable navigation. This is the ownership boundary AUDIT.md's B14 and B17 call for: approval becomes an explicit, revision-bound operation executed by the worker that owns the pane, instead of a side effect of prompting over a coordinator-driven remote socket. Also ignores the web build inputs and outputs. node_modules ships vendored Go packages, so go build and go test walk into it if it is merely untracked; both node_modules and .node_modules are excluded. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01535A3Y8RtkAi8wYuWhtkEd
260 lines
9.7 KiB
Go
260 lines
9.7 KiB
Go
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"net"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"orchestra/internal/domain"
|
|
"orchestra/internal/federation"
|
|
"orchestra/internal/herdr"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
func TestWorkerReregistersAfterCoordinatorForgetsIt(t *testing.T) {
|
|
registered := false
|
|
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
if r.Method != http.MethodPost || r.URL.Path != "/v1/federation/workers" {
|
|
t.Fatalf("unexpected %s %s", r.Method, r.URL.Path)
|
|
}
|
|
registered = true
|
|
w.WriteHeader(http.StatusCreated)
|
|
}))
|
|
defer s.Close()
|
|
w := &worker{api: federation.Client{BaseURL: s.URL, WorkerID: "h", Token: "t"}, registration: federation.Worker{ID: "h", Capacity: 1}}
|
|
if !w.reRegisterAfterCoordinatorRestart(context.Background(), errors.New("federation: 401 Unauthorized: unknown worker")) {
|
|
t.Fatal("worker did not identify a coordinator restart")
|
|
}
|
|
if !registered {
|
|
t.Fatal("worker did not register")
|
|
}
|
|
}
|
|
|
|
func TestInitialReplayDoesNotResurrectReleasedLease(t *testing.T) {
|
|
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
switch r.URL.Path {
|
|
case "/v1/federation/events":
|
|
_, _ = w.Write([]byte(`{"cursor":0,"events":[{"seq":1,"id":"c","type":"TaskCreated","task_id":"t","version":1,"payload":{"source":"s","external_id":"x","project":"p"},"surface":"system"},{"seq":2,"id":"l","type":"TaskLeased","task_id":"t","version":2,"payload":{"harness_id":"h"},"surface":"system"},{"seq":3,"id":"r","type":"TaskReleased","task_id":"t","version":3,"payload":{"reason":"expired"},"surface":"system"}]}`))
|
|
case "/v1/federation/events/ack":
|
|
w.WriteHeader(http.StatusNoContent)
|
|
default:
|
|
t.Fatalf("unexpected %s", r.URL.Path)
|
|
}
|
|
}))
|
|
defer s.Close()
|
|
w := &worker{api: federation.Client{BaseURL: s.URL, WorkerID: "h", Token: "t"}, harnessID: "h", tasks: map[string]domain.Task{}, sessions: map[string]herdr.Session{}, leases: map[string]lease{}, statePath: t.TempDir() + "/state.json", hard: .75}
|
|
if err := w.once(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(w.leases) != 0 || len(w.sessions) != 0 {
|
|
t.Fatalf("replayed lease survived: leases=%v sessions=%v", w.leases, w.sessions)
|
|
}
|
|
if w.cursor != 3 {
|
|
t.Fatalf("cursor=%d want 3", w.cursor)
|
|
}
|
|
}
|
|
|
|
func TestWorkerHydratesLeaseMissingTaskCache(t *testing.T) {
|
|
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
switch r.URL.Path {
|
|
case "/v1/federation/events":
|
|
_, _ = w.Write([]byte(`{"cursor":2,"events":[{"seq":2,"id":"l","type":"TaskLeased","task_id":"t","version":2,"payload":{"harness_id":"h"},"surface":"system"}]}`))
|
|
case "/v1/tasks":
|
|
_, _ = w.Write([]byte(`[{"id":"t","source":"s","external_id":"x","project":"p","state":"leased"}]`))
|
|
case "/v1/federation/events/ack":
|
|
w.WriteHeader(http.StatusNoContent)
|
|
default:
|
|
t.Fatalf("unexpected %s", r.URL.Path)
|
|
}
|
|
}))
|
|
defer s.Close()
|
|
w := &worker{api: federation.Client{BaseURL: s.URL, WorkerID: "h", Token: "t"}, harnessID: "h", cursor: 1, tasks: map[string]domain.Task{}, sessions: map[string]herdr.Session{"t": {}}, leases: map[string]lease{}, statePath: t.TempDir() + "/state.json", hard: .75}
|
|
if err := w.once(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, ok := w.tasks["t"]; !ok || w.cursor != 2 {
|
|
t.Fatalf("task was not hydrated: tasks=%v cursor=%d", w.tasks, w.cursor)
|
|
}
|
|
}
|
|
|
|
func TestWorkerReconcilesLeaseAfterEmptyEventReplay(t *testing.T) {
|
|
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
switch r.URL.Path {
|
|
case "/v1/federation/events":
|
|
_, _ = w.Write([]byte(`{"cursor":42,"events":[]}`))
|
|
case "/v1/tasks":
|
|
_, _ = w.Write([]byte(`[{"id":"t","source":"s","external_id":"x","project":"p","state":"leased","lease":{"harness_id":"h"}}]`))
|
|
case "/v1/federation/events/ack":
|
|
w.WriteHeader(http.StatusNoContent)
|
|
default:
|
|
t.Fatalf("unexpected %s", r.URL.Path)
|
|
}
|
|
}))
|
|
defer s.Close()
|
|
w := &worker{api: federation.Client{BaseURL: s.URL, WorkerID: "h", Token: "t"}, harnessID: "h", cursor: 42, tasks: map[string]domain.Task{}, sessions: map[string]herdr.Session{"t": {}}, leases: map[string]lease{}, statePath: t.TempDir() + "/state.json", hard: .75}
|
|
if err := w.once(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, ok := w.leases["t"]; !ok {
|
|
t.Fatal("empty event replay did not recover current lease")
|
|
}
|
|
}
|
|
|
|
func TestWorkerStartsRouterIssuedLeaseInLocalGitWorktree(t *testing.T) {
|
|
remote := filepath.Join(t.TempDir(), "remote.git")
|
|
if out, err := exec.Command("git", "init", "--bare", remote).CombinedOutput(); err != nil {
|
|
t.Fatalf("remote: %v %s", err, out)
|
|
}
|
|
seed := t.TempDir()
|
|
for _, a := range [][]string{{"init", seed}, {"-C", seed, "config", "user.email", "t@t"}, {"-C", seed, "config", "user.name", "t"}, {"-C", seed, "commit", "--allow-empty", "-m", "init"}, {"-C", seed, "remote", "add", "origin", remote}, {"-C", seed, "push", "-u", "origin", "HEAD:master"}} {
|
|
if out, err := exec.Command("git", a...).CombinedOutput(); err != nil {
|
|
t.Fatalf("git %v: %v %s", a, err, out)
|
|
}
|
|
}
|
|
repo := filepath.Join(t.TempDir(), "repo")
|
|
if out, err := exec.Command("git", "clone", remote, repo).CombinedOutput(); err != nil {
|
|
t.Fatal(string(out))
|
|
}
|
|
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer ln.Close()
|
|
go func() {
|
|
for {
|
|
c, e := ln.Accept()
|
|
if e != nil {
|
|
return
|
|
}
|
|
go func() {
|
|
defer c.Close()
|
|
var r herdr.Request
|
|
if json.NewDecoder(bufio.NewReader(c)).Decode(&r) != nil {
|
|
return
|
|
}
|
|
var result string
|
|
switch r.Method {
|
|
case "worktree.create":
|
|
result = `{"path":"` + filepath.Join(filepath.Dir(repo), "worktrees", "task") + `","root_pane":{"pane_id":"p1"}}`
|
|
case "agent.start":
|
|
result = `{}`
|
|
case "pane.get":
|
|
result = `{"pane":{"agent":"opencode","agent_status":"idle"}}`
|
|
case "pane.read":
|
|
result = `{"read":{"text":""}}`
|
|
case "agent.prompt":
|
|
result = `{}`
|
|
default:
|
|
result = `{}`
|
|
}
|
|
_ = json.NewEncoder(c).Encode(herdr.Response{ID: r.ID, Result: json.RawMessage(result)})
|
|
}()
|
|
}
|
|
}()
|
|
root := filepath.Join(filepath.Dir(repo), "worktrees")
|
|
w := &worker{herdr: &herdr.Client{Path: ln.Addr().String()}, repo: repo, root: root, remote: "origin", harness: "opencode", tasks: map[string]domain.Task{}, sessions: map[string]herdr.Session{}, leases: map[string]lease{}, statePath: filepath.Join(t.TempDir(), "state.json"), hard: .75}
|
|
// New() is needed for its pane map; override the address for the fake.
|
|
w.herdr = herdr.New(ln.Addr().String())
|
|
task := domain.Task{ID: "task", Source: "s", ExternalID: "x", Project: "p", Title: "test", Description: "do work"}
|
|
if err := w.start(context.Background(), task, ""); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, ok := w.sessions[task.ID]; !ok {
|
|
t.Fatal("session not persisted")
|
|
}
|
|
if _, err := os.Stat(filepath.Join(root, "task", "TASK.md")); err != nil {
|
|
t.Fatalf("TASK.md: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestWorkerApprovalCommandIsRevisionBoundAndAcknowledged(t *testing.T) {
|
|
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer ln.Close()
|
|
sent := make(chan string, 1)
|
|
go func() {
|
|
for {
|
|
c, err := ln.Accept()
|
|
if err != nil {
|
|
return
|
|
}
|
|
go func() {
|
|
defer c.Close()
|
|
var request herdr.Request
|
|
if json.NewDecoder(c).Decode(&request) != nil {
|
|
return
|
|
}
|
|
result := `{}`
|
|
switch request.Method {
|
|
case "pane.read":
|
|
result = `{"read":{"text":"Permission required\n$ git status\nProceed? [y/n]"}}`
|
|
case "pane.send_text":
|
|
var p struct {
|
|
Text string `json:"text"`
|
|
}
|
|
_ = json.Unmarshal(mustJSON(request.Params), &p)
|
|
sent <- p.Text
|
|
}
|
|
_ = json.NewEncoder(c).Encode(herdr.Response{ID: request.ID, Result: json.RawMessage(result)})
|
|
}()
|
|
}
|
|
}()
|
|
resolved := make(chan map[string]string, 1)
|
|
api := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) {
|
|
switch {
|
|
case r.Method == http.MethodGet && r.URL.Path == "/v1/federation/commands":
|
|
_ = json.NewEncoder(rw).Encode([]federation.Command{{ID: "c1", TaskID: "task", Kind: "grant_approval", PaneID: "pane", CaptureRevision: 7}})
|
|
case r.Method == http.MethodPost && r.URL.Path == "/v1/federation/workers/h/captures":
|
|
_ = json.NewEncoder(rw).Encode(federation.Capture{TaskID: "task", PaneID: "pane", Revision: 7})
|
|
case r.Method == http.MethodPost && r.URL.Path == "/v1/federation/commands/c1":
|
|
var body map[string]string
|
|
_ = json.NewDecoder(r.Body).Decode(&body)
|
|
resolved <- body
|
|
rw.WriteHeader(http.StatusNoContent)
|
|
default:
|
|
t.Errorf("unexpected %s %s", r.Method, r.URL.Path)
|
|
rw.WriteHeader(http.StatusNotFound)
|
|
}
|
|
}))
|
|
defer api.Close()
|
|
w := &worker{api: federation.Client{BaseURL: api.URL, WorkerID: "h", Token: "t"}, herdr: herdr.New(ln.Addr().String()), harness: "opencode", sessions: map[string]herdr.Session{"task": {PaneID: "pane"}}}
|
|
w.runCommands(context.Background())
|
|
select {
|
|
case got := <-sent:
|
|
if got != "y\n" {
|
|
t.Fatalf("approval input=%q", got)
|
|
}
|
|
case <-time.After(time.Second):
|
|
t.Fatal("worker did not send approval")
|
|
}
|
|
select {
|
|
case got := <-resolved:
|
|
if got["status"] != "acknowledged" {
|
|
t.Fatalf("resolution=%v", got)
|
|
}
|
|
case <-time.After(time.Second):
|
|
t.Fatal("worker did not resolve command")
|
|
}
|
|
}
|
|
|
|
func TestApprovalResponseOpenCodeAllowOnce(t *testing.T) {
|
|
text := "Permission required\nAllow once Allow always Reject\n⇆ select enter confirm"
|
|
if got, ok := approvalResponse(text, "grant_approval"); !ok || got != "\n" {
|
|
t.Fatalf("grant response = %q, %v", got, ok)
|
|
}
|
|
if got, ok := approvalResponse(text, "deny_approval"); ok || got != "" {
|
|
t.Fatalf("deny response = %q, %v; reject must not guess selector navigation", got, ok)
|
|
}
|
|
}
|
|
|
|
func mustJSON(v any) []byte { b, _ := json.Marshal(v); return b }
|