Files
orchestra/cmd/orchestra-worker/main.go
T
kami edbe98fc5e Confirm every write Orchestra sends, and count none of them as progress
F20. Only the launch confirmed its submit. A decision notice at a turn
boundary, and /clear or @HANDOFF.md during a context reset, were
fire-and-forget through the same transport that loses an Enter often enough
that the launch needed three resubmits. A lost Enter on the context-reset path
is the worst of them: it strands the session mid-rollover and nothing retries
it. LaunchConfirmer is therefore InputConfirmer, ConfirmLaunch is ConfirmInput,
and sendPrompt and sendLine both go through it.

Orchestra does not try to guarantee delivery of input it did not originate.
But it must never read that input as work, which is the F16 half. Burn-in run
3 stalled with an unexplained "go ahead and implement it" in the editor, and
the renewal check hashed the whole capture, so those keystrokes read as
progress and the lease kept renewing around an idle agent. PaneProgress drops
input lines from the capture, which the -J join makes exact: a wrapped input
block is one line beginning with the prompt marker.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-27 14:44:04 +04:00

1689 lines
59 KiB
Go

// orchestra-worker consumes router-issued leases for one or more local
// execution backends. Each declared harness is a separate federation identity
// with its own token, cursor, backend, and state file, because the coordinator
// authorizes a lease call by comparing the URL's worker id against the lease's
// harness id. Homesrv remains the scheduler and CAS authority. This process
// owns only local Git and pane operations.
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"orchestra/internal/agentctx"
"orchestra/internal/buildinfo"
"orchestra/internal/continuity"
"orchestra/internal/domain"
"orchestra/internal/federation"
"orchestra/internal/herdr"
"orchestra/internal/orchestrator"
"orchestra/internal/review"
"orchestra/internal/workphase"
"os"
"os/exec"
"os/signal"
"path/filepath"
"reflect"
"sort"
"strconv"
"strings"
"syscall"
"time"
)
type worker struct {
api federation.Client
backend herdr.Backend
// herdr is retained as a test/backward-compatibility alias. Production
// workers set backend; executionBackend keeps older state-machine tests
// from needing protocol-irrelevant rewrites.
herdr *herdr.Client
harnessID, harness, repo, root, remote string
projects map[string]projectConfig
cursor uint64
tasks map[string]domain.Task
sessions map[string]herdr.Session
leases map[string]lease
releases map[string]releaseTransaction
quarantined map[string]bool
statePath string
hard float64
registration federation.Worker
lastError string
lastErrorAt time.Time
soft float64
window int64
}
func (w *worker) executionBackend() herdr.Backend {
if w.backend != nil {
return w.backend
}
if w.herdr != nil {
return w.herdr
}
return nil
}
func (w *worker) recordError(err error) {
if err == nil {
return
}
w.lastError = err.Error()
w.lastErrorAt = time.Now().UTC()
}
func (w *worker) health(ctx context.Context) federation.WorkerHealth {
h := federation.WorkerHealth{HerdrStatus: "unknown"}
if backend := w.executionBackend(); backend != nil {
h.Backend = backend.Kind()
}
for taskID, session := range w.sessions {
// Workers currently advertise capacity one. Pick deterministically so a
// recovered legacy state with more sessions remains intelligible.
if h.ActiveTask == "" || taskID < h.ActiveTask {
h.ActiveTask, h.ActivePane = taskID, session.PaneID
}
}
if backend := w.executionBackend(); backend != nil {
checkCtx, cancel := context.WithTimeout(ctx, 3*time.Second)
err := backend.Check(checkCtx)
cancel()
h.CheckedAt = time.Now().UTC()
if err == nil {
h.HerdrStatus = "reachable"
} else {
h.HerdrStatus = "unreachable"
w.recordError(fmt.Errorf("local %s backend: %w", backend.Kind(), err))
}
}
h.LastError, h.ErrorAt = w.lastError, w.lastErrorAt
return h
}
type lease struct {
Epoch string `json:"epoch"`
HandoffRef string `json:"handoff_ref,omitempty"`
TransactionID string `json:"transaction_id,omitempty"`
AnchorSHA string `json:"anchor_sha,omitempty"`
PickupAcknowledged bool `json:"pickup_acknowledged,omitempty"`
Version int `json:"version"`
Until time.Time `json:"until"`
UsageBaseline float64 `json:"usage_baseline,omitempty"`
// ProgressSHA hashes the pane capture taken at the last renewal. Renewal
// requires the pane to have changed since then, or the agent to be busy.
ProgressSHA string `json:"progress_sha,omitempty"`
}
type releaseTransaction struct {
ID string `json:"id"`
LeaseVersion int `json:"lease_version"`
Ref string `json:"handoff_ref,omitempty"`
AnchorSHA string `json:"anchor_sha,omitempty"`
Phase string `json:"phase"` // prepared, anchor_pushed, event_committed, pickup_validated, predecessor_retired
AgentReleased bool `json:"agent_released,omitempty"`
LastError string `json:"last_error,omitempty"`
UpdatedAt time.Time `json:"updated_at"`
}
type projectConfig struct {
Repo string `json:"repo"`
Root string `json:"worktree_root"`
Remote string `json:"remote"`
QualityGate string `json:"quality_gate,omitempty"`
SafeOperations []string `json:"safe_operations,omitempty"`
}
type completionEvidence struct {
TaskID string `json:"task_id"`
Project string `json:"project"`
Worker string `json:"worker"`
Harness string `json:"harness"`
PaneID string `json:"pane_id"`
BaseSHA string `json:"base_sha"`
ResultSHA string `json:"result_sha"`
Branch string `json:"branch"`
Remote string `json:"remote"`
QualityGate string `json:"quality_gate,omitempty"`
GateExit int `json:"gate_exit"`
CompletedAt time.Time `json:"completed_at"`
}
type workerState struct {
Cursor uint64 `json:"cursor"`
Sessions map[string]herdr.Session `json:"sessions"`
Tasks map[string]domain.Task `json:"tasks"`
Leases map[string]lease `json:"leases"`
Releases map[string]releaseTransaction `json:"releases"`
Quarantined map[string]bool `json:"quarantined,omitempty"`
}
func (w *worker) load() error {
b, e := os.ReadFile(w.statePath)
if e == nil {
var s workerState
if err := json.Unmarshal(b, &s); err != nil {
return fmt.Errorf("corrupt worker state %s: %w", w.statePath, err)
}
w.cursor = s.Cursor
w.sessions = s.Sessions
w.tasks = s.Tasks
w.leases = s.Leases
w.releases = s.Releases
w.quarantined = s.Quarantined
} else if !errors.Is(e, os.ErrNotExist) {
return fmt.Errorf("read worker state %s: %w", w.statePath, e)
}
if w.sessions == nil {
w.sessions = map[string]herdr.Session{}
}
if w.tasks == nil {
w.tasks = map[string]domain.Task{}
}
if w.leases == nil {
w.leases = map[string]lease{}
}
if w.releases == nil {
w.releases = map[string]releaseTransaction{}
}
if w.quarantined == nil {
w.quarantined = map[string]bool{}
}
return nil
}
func (w *worker) save() error {
b, e := json.Marshal(workerState{Cursor: w.cursor, Sessions: w.sessions, Tasks: w.tasks, Leases: w.leases, Releases: w.releases, Quarantined: w.quarantined})
if e != nil {
return e
}
if e := os.MkdirAll(filepath.Dir(w.statePath), 0700); e != nil {
return e
}
tmp := w.statePath + ".tmp"
f, e := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0600)
if e != nil {
return e
}
if _, e = f.Write(b); e == nil {
e = f.Sync()
}
if closeErr := f.Close(); e == nil {
e = closeErr
}
if e != nil {
_ = os.Remove(tmp)
return e
}
if e = os.Rename(tmp, w.statePath); e != nil {
return e
}
dir, e := os.Open(filepath.Dir(w.statePath))
if e != nil {
return e
}
defer dir.Close()
return dir.Sync()
}
// quarantine stops a pane before its lost lease mapping can be forgotten.
// A failed close remains durable and is retried; it is never treated as a
// harmless cleanup error while the old harness could still be working.
func (w *worker) quarantine(ctx context.Context, taskID string, s herdr.Session) {
if w.executionBackend() == nil {
w.quarantined[taskID] = true
return
}
if err := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).Kill(ctx, s); err != nil {
w.quarantined[taskID] = true
w.recordError(fmt.Errorf("quarantine %s: %w", taskID, err))
return
}
delete(w.sessions, taskID)
delete(w.quarantined, taskID)
}
func (w *worker) retryQuarantines(ctx context.Context) {
for taskID := range w.quarantined {
if s, ok := w.sessions[taskID]; ok {
w.quarantine(ctx, taskID, s)
}
}
}
type artifactCAS struct{ api federation.Client }
func (c artifactCAS) PutArtifact(b []byte) (string, error) {
return c.api.PutArtifact(context.Background(), b)
}
func (c artifactCAS) Artifact(ref string) ([]byte, error) {
return c.api.Artifact(context.Background(), ref)
}
func created(e domain.Event) (domain.Task, bool) {
if e.Type != "TaskCreated" {
return domain.Task{}, false
}
var p struct {
Source string `json:"source"`
ExternalID string `json:"external_id"`
Project string `json:"project"`
Capability []string `json:"capability"`
Title string `json:"title"`
Description string `json:"description"`
Acceptance []string `json:"acceptance"`
QualityGate string `json:"quality_gate"`
}
if json.Unmarshal(e.Payload, &p) != nil || p.Source == "" || p.ExternalID == "" || p.Project == "" {
return domain.Task{}, false
}
return domain.Task{ID: e.TaskID, Source: p.Source, ExternalID: p.ExternalID, Project: p.Project, Capability: p.Capability, Title: p.Title, Description: p.Description, Acceptance: p.Acceptance, QualityGate: p.QualityGate}, true
}
func (w *worker) project(t domain.Task) (projectConfig, error) {
if w.projects != nil {
if p, ok := w.projects[t.Project]; ok && p.Repo != "" && p.Root != "" && p.Remote != "" {
return p, nil
}
return projectConfig{}, fmt.Errorf("project %q is not configured on worker", t.Project)
}
return projectConfig{Repo: w.repo, Root: w.root, Remote: w.remote}, nil
}
func (w *worker) syncBase(ctx context.Context, p projectConfig) error {
if out, err := exec.CommandContext(ctx, "git", "-C", p.Repo, "fetch", p.Remote, "--prune").CombinedOutput(); err != nil {
return fmt.Errorf("fetch base checkout: %s: %w", out, err)
}
branch, err := exec.CommandContext(ctx, "git", "-C", p.Repo, "symbolic-ref", "--quiet", "--short", "HEAD").Output()
if err != nil {
return fmt.Errorf("identify base branch: %w", err)
}
branchName := strings.TrimSpace(string(branch))
if out, err := exec.CommandContext(ctx, "git", "-C", p.Repo, "merge", "--ff-only", p.Remote+"/"+branchName).CombinedOutput(); err != nil {
return fmt.Errorf("fast-forward base checkout: %s: %w", out, err)
}
return nil
}
func (w *worker) start(ctx context.Context, t domain.Task, ref string) error {
var wt string
var h continuity.Handoff
var err error
p, err := w.project(t)
if err != nil {
return err
}
// Synchronize the local base before any worktree operation. A worker never
// treats a coordinator-side path as truth; the Git remote is the only
// cross-machine transport.
if err := w.syncBase(ctx, p); err != nil {
return err
}
if ref != "" {
b, err := w.api.Artifact(ctx, ref)
if err != nil {
return err
}
h, err = continuity.Decode(b)
if err != nil {
return err
}
if out, err := exec.CommandContext(ctx, "git", "-C", p.Repo, "fetch", p.Remote, "--prune").CombinedOutput(); err != nil {
return fmt.Errorf("fetch pickup anchor: %s: %w", out, err)
}
wt = filepath.Join(p.Root, t.ID)
if _, err := os.Stat(wt); os.IsNotExist(err) {
if out, err := exec.CommandContext(ctx, "git", "-C", p.Repo, "worktree", "add", "-b", "orchestra/"+t.ID, wt, h.Anchor.GitSHA).CombinedOutput(); err != nil {
return fmt.Errorf("create pickup worktree: %s: %w", out, err)
}
}
if err = continuity.ValidatePickup(wt, h, taskHash(t)); err != nil {
return err
}
} else {
wt, err = (orchestrator.GitWorktrees{Repo: p.Repo, Root: p.Root}).Create(ctx, t)
if err != nil {
return err
}
}
backend := w.executionBackend()
if backend == nil {
return fmt.Errorf("execution backend is not configured")
}
if _, err = backend.Worktree(ctx, p.Repo, wt, "orchestra/"+t.ID); err != nil {
return err
}
s, err := backend.StartAgent(ctx, wt, wt, "orchestra/"+t.ID, w.harness, t.ID)
if err != nil {
return err
}
s.TaskFileSHA = taskHash(t)
if w.harness == "claude" {
s.ContextHandoffSHA, _ = fileSHA256(filepath.Join(wt, "HANDOFF.md"))
}
// One renderer, on both machines. The worker fetches the reduced
// authority as data and renders it with agentctx, so a decision the human
// recorded before this session existed is visible from its first turn.
intent, err := w.api.Intent(ctx, t.ID)
if err != nil {
return fmt.Errorf("effective intent: %w", err)
}
in := agentctx.Input{
Task: t, Intent: intent, Phase: t.WorkPhase, DecisionRequest: t.DecisionRequest,
Git: agentctx.GitState{Worktree: wt, Branch: "orchestra/" + t.ID},
RepoRules: agentctx.DiscoverRepoRules(wt),
}
if sha, shaErr := herdr.HeadSHA(wt); shaErr == nil {
in.Git.HeadSHA = sha
}
if ref != "" {
in.Handoff = &h
}
if len(p.SafeOperations) > 0 {
in.Policy = []string{
"Permitted without an operator grant, inside this worktree only: " + strings.Join(p.SafeOperations, ", ") + ".",
"Network access, secrets, destructive actions, and paths outside this worktree require an explicit operator approval.",
}
}
if t.ResearchRef != "" {
b, artErr := w.api.Artifact(ctx, t.ResearchRef)
if artErr != nil {
return fmt.Errorf("research artifact: %w", artErr)
}
r, decErr := workphase.DecodeResearch(b)
if decErr != nil {
return fmt.Errorf("research artifact: %w", decErr)
}
in.Research = &r
}
if t.PlanRef != "" {
b, artErr := w.api.Artifact(ctx, t.PlanRef)
if artErr != nil {
return fmt.Errorf("plan artifact: %w", artErr)
}
pl, decErr := workphase.DecodePlan(b)
if decErr != nil {
return fmt.Errorf("plan artifact: %w", decErr)
}
in.Plan = &pl
}
if t.Review != nil {
b, artErr := w.api.Artifact(ctx, t.Review.ArtifactRef)
if artErr != nil {
return fmt.Errorf("review artifact: %w", artErr)
}
r, decErr := review.Decode(b)
if decErr != nil {
return fmt.Errorf("review artifact: %w", decErr)
}
in.Review = &r
}
built, err := agentctx.Build(in)
if err != nil {
return fmt.Errorf("build context: %w", err)
}
prompt := built.System + "\n\n" + built.Task
// The transport decides what is submitted, never what the agent receives:
// the file holds the exact bytes agentctx rendered either way.
submitted, transport := prompt, herdr.LaunchInline
if lt, ok := backend.(herdr.LaunchTransporter); ok {
transport = lt.LaunchTransport(w.harness)
}
if transport == herdr.LaunchFileRef {
submitted = herdr.LaunchReference
}
if writeErr := herdr.WriteLaunchContext(s.Worktree, prompt); writeErr != nil {
// Under LaunchFileRef the file is the instruction, so a failed write
// is a failed launch rather than lost evidence.
if transport == herdr.LaunchFileRef {
return fmt.Errorf("launch context %s: %w", t.ID, writeErr)
}
w.recordError(fmt.Errorf("launch context %s: %w", t.ID, writeErr))
}
// The launch instruction carried these, so the first turn boundary must
// not re-announce them as news.
for _, d := range intent.Decisions {
s.DeliveredDecisions = append(s.DeliveredDecisions, d.ID)
}
w.sessions[t.ID] = s
if err := w.save(); err != nil {
return err
}
// A prompt response can be lost after the backend accepted it. Persist the
// session first so the worker can reconcile/release it after restart.
if err := backend.Prompt(ctx, s.PaneID, submitted, 0); err != nil {
return err
}
// Acknowledging a launch means the harness accepted the instruction, not
// that the adapter call returned nil. Without this the worker reported a
// started agent while the prompt sat unsubmitted in the input editor.
if c, ok := backend.(herdr.InputConfirmer); ok {
evidence, confirmErr := c.ConfirmInput(ctx, s, submitted)
if confirmErr != nil {
// An unsubmitted prompt leaves a live pane that nothing owns, and
// a retained session would make the retry skip this task
// entirely. Reclaim both so the released lease can be re-leased.
if killErr := backend.Kill(ctx, s); killErr != nil {
w.recordError(fmt.Errorf("kill unlaunched pane %s: %w", t.ID, killErr))
}
delete(w.sessions, t.ID)
if saveErr := w.save(); saveErr != nil {
w.recordError(fmt.Errorf("save after failed launch %s: %w", t.ID, saveErr))
}
return fmt.Errorf("launch %s: %w", t.ID, confirmErr)
}
log.Printf("launch %s confirmed: %s", t.ID, evidence)
}
if l, ok := w.leases[t.ID]; ok {
if err := w.api.Start(ctx, t.ID, l.Epoch, l.Version, w.sessionEvidence(ctx, t.ID, s)); err != nil {
return fmt.Errorf("ack start: %w", err)
}
l.Version++
w.leases[t.ID] = l
if err := w.save(); err != nil {
return err
}
}
if ref != "" {
return w.ackPickup(ctx, t.ID, s)
}
return nil
}
func classifyLaunchError(err error, sessionStarted bool) string {
// Positive evidence that the harness never accepted the prompt is not
// uncertainty. Release the lease so the existing retry path can take it.
if errors.Is(err, herdr.ErrPromptNotSubmitted) {
return "prompt_not_submitted"
}
if sessionStarted {
// A prompt response can be lost after herdr accepted it. Never reclaim
// that pane just because its acknowledgement was uncertain.
return "launch_uncertain"
}
text := strings.ToLower(err.Error())
if strings.Contains(text, "handoff") || strings.Contains(text, "pickup") || strings.Contains(text, "task.md") {
return "invalid_handoff"
}
return "launch_transient"
}
func taskHash(t domain.Task) string { b := continuity.RenderTaskFile(t); return domain.Hash(b) }
func fileSHA256(path string) (string, error) {
b, err := os.ReadFile(path)
if err != nil {
return "", err
}
return domain.Hash(b), nil
}
func (w *worker) releaseReady(ctx context.Context) {
for id, s := range w.sessions {
if w.quarantined[id] {
continue
}
if l := w.leases[id]; l.HandoffRef != "" && !l.PickupAcknowledged {
if err := w.ackPickup(ctx, id, s); err != nil {
w.recordError(err)
continue
}
}
if _, err := os.Stat(filepath.Join(s.Worktree, ".orchestra", "done")); err == nil {
// A done marker is an intent, not enough on its own: do not race a
// still-running native harness into committing half-written work.
status, statusErr := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).AgentStatus(ctx, s)
if statusErr != nil {
w.recordError(fmt.Errorf("completion identity %s: %w", id, statusErr))
continue
}
if herdr.IsBusy(status) {
w.recordError(fmt.Errorf("completion %s deferred: agent status %s", id, status))
continue
}
evidence, err := w.finalize(ctx, id, s)
if err != nil {
w.recordError(fmt.Errorf("complete %s: %w", id, err))
log.Printf("complete %s: %v", id, err)
continue
}
report, _ := json.Marshal(evidence)
ref, err := w.api.PutArtifact(ctx, report)
if err != nil {
w.recordError(fmt.Errorf("upload completion %s: %w", id, err))
log.Printf("upload completion %s: %v", id, err)
continue
}
if err = w.api.Complete(ctx, id, ref, evidence.ResultSHA, evidence.Branch, evidence.Remote, w.leases[id].Epoch, w.leases[id].Version, w.usageReceipt(s, w.leases[id]), w.sessionEvidence(ctx, id, s)); err != nil {
w.recordError(fmt.Errorf("complete %s: %w", id, err))
log.Printf("complete %s: %v", id, err)
continue
}
// Completion is durable before closing the exact pane. If close
// fails, retain the session mapping for a later explicit cleanup.
a := herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}
if err := a.Kill(ctx, s); err != nil {
w.recordError(fmt.Errorf("close completed pane %s: %w", id, err))
log.Printf("close completed pane %s: %v", id, err)
continue
}
_ = os.Remove(filepath.Join(s.Worktree, ".orchestra", "done"))
_ = os.Remove(filepath.Join(s.Worktree, ".orchestra"))
delete(w.sessions, id)
delete(w.leases, id)
_ = w.save()
continue
}
if _, err := os.Stat(filepath.Join(s.Worktree, herdr.HandoffReportFile)); err == nil || w.releases[id].ID != "" {
w.advanceRelease(ctx, id, s)
continue
}
w.rotationTick(ctx, id, s)
}
}
func (w *worker) adapter(s herdr.Session, remote string) herdr.CLIAdapter {
a := herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness, Window: w.window, CAS: artifactCAS{w.api}, Remote: remote}
switch w.harness {
case "claude":
a.Usage = herdr.ClaudeUsage
case "codex":
a.Usage = herdr.CodexUsage
case "opencode":
a.Usage = herdr.OpenCodeUsage
}
return a
}
// rotationTick is the checkout-owner state machine. Occupancy, tool activity,
// and pane status are all read from the persisted harness session identity;
// any unknown source is recorded and never treated as zero usage.
func (w *worker) rotationTick(ctx context.Context, id string, s herdr.Session) {
// Claude Code owns its context threshold through the installed
// context-handoff hook. A changed HANDOFF.md means that hook has landed a
// durable local continuation. Resume in the same process with Claude's
// native context reset instead of manufacturing Orchestra's cross-worker
// release artifact. Codex and OpenCode continue through the existing
// occupancy/release state machine below.
if w.harness == "claude" {
if err := w.advanceClaudeContextReset(ctx, id, s); err != nil {
w.recordError(fmt.Errorf("Claude context reset %s: %w", id, err))
}
return
}
t, ok := w.tasks[id]
if !ok {
w.recordError(fmt.Errorf("rotation %s: task cache missing", id))
return
}
p, err := w.project(t)
if err != nil {
w.recordError(err)
return
}
a := w.adapter(s, p.Remote)
resolved, err := a.ResolveSessionIdentity(s)
if err != nil {
w.recordError(fmt.Errorf("rotation %s occupancy degraded: %w", id, err))
return
}
// DeepEqual, not !=: Session carries a slice since decisions are tracked
// per session, so it is no longer comparable with ==.
if !reflect.DeepEqual(resolved, s) {
w.sessions[id] = resolved
s = resolved
_ = w.save()
}
d := (orchestrator.RotationStateMachine{Soft: w.soft, Hard: w.hard}).Evaluate(ctx, a, s)
if d.ActivityDegraded != nil {
w.recordError(fmt.Errorf("rotation %s activity degraded: %w", id, d.ActivityDegraded))
}
if d.Degraded != nil {
w.recordError(fmt.Errorf("rotation %s degraded: %w", id, d.Degraded))
if d.Action == orchestrator.TurnContinue || d.Action == "" {
return
}
}
if d.Action == orchestrator.TurnContinue || d.Action == "" {
w.federatedTurn(ctx, id, a, orchestrator.TurnContinue)
return
}
if d.Action == orchestrator.TurnRefuse || s.HandoffRequested {
return
}
if d.Reason == "milestone" || d.Reason == "thrash" {
if err := a.RequestHandoffReason(ctx, s, d.Reason, d.DeadEnds); err != nil {
w.recordError(fmt.Errorf("rotation %s %s prompt: %w", id, d.Reason, err))
return
}
} else if err := a.RequestHandoff(ctx, s); err != nil {
w.recordError(fmt.Errorf("rotation %s threshold prompt: %w", id, err))
return
}
s.HandoffRequested, s.HandoffReason = true, d.Reason
w.sessions[id] = s
_ = w.save()
}
func (w *worker) sendLine(ctx context.Context, s herdr.Session, line string) error {
backend := w.executionBackend()
if backend == nil {
return fmt.Errorf("execution backend is not configured")
}
if err := backend.SendText(ctx, s, line); err != nil {
return fmt.Errorf("send %q: %w", line, err)
}
if err := backend.SendKeys(ctx, s, []string{"ENTER"}); err != nil {
return fmt.Errorf("submit %q: %w", line, err)
}
// A lost Enter here leaves the session mid-rollover with /clear sitting in
// the editor, which is worse than a lost launch: nothing retries it.
if err := w.confirmInput(ctx, s, line); err != nil {
return fmt.Errorf("submit %q: %w", line, err)
}
return nil
}
func (w *worker) advanceClaudeContextReset(ctx context.Context, id string, s herdr.Session) error {
backend := w.executionBackend()
if backend == nil {
return fmt.Errorf("execution backend is not configured")
}
if s.ContextResetSHA == "" {
sha, err := fileSHA256(filepath.Join(s.Worktree, "HANDOFF.md"))
if os.IsNotExist(err) {
return nil
}
if err != nil {
return fmt.Errorf("read HANDOFF.md: %w", err)
}
if sha == s.ContextHandoffSHA {
return nil
}
status, err := backend.AgentStatus(ctx, s)
if err != nil {
return fmt.Errorf("confirm Claude stopped after handoff: %w", err)
}
if status != "idle" {
return nil
}
s.ContextResetSHA = sha
s.ContextResetPhase = "clear"
w.sessions[id] = s
if err := w.save(); err != nil {
return err
}
}
switch s.ContextResetPhase {
case "clear":
if err := w.sendLine(ctx, s, "/clear"); err != nil {
return err
}
s.ContextResetPhase = "handoff"
w.sessions[id] = s
if err := w.save(); err != nil {
return err
}
// /clear redraws Claude's input UI asynchronously. Give it a small,
// bounded interval before submitting the new-session file mention.
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(500 * time.Millisecond):
}
fallthrough
case "handoff":
if err := w.sendLine(ctx, s, "@HANDOFF.md"); err != nil {
return err
}
s.ContextHandoffSHA = s.ContextResetSHA
s.ContextResetSHA = ""
s.ContextResetPhase = ""
// Claude normally opens a fresh transcript for /clear. Force the next
// observation to discover it instead of retaining the exhausted path.
s.SessionFile = ""
w.sessions[id] = s
return w.save()
default:
return fmt.Errorf("unknown persisted context-reset phase %q", s.ContextResetPhase)
}
}
func (w *worker) advanceRelease(ctx context.Context, id string, s herdr.Session) {
t, ok := w.tasks[id]
if !ok {
w.recordError(fmt.Errorf("release %s: task cache missing", id))
return
}
p, err := w.project(t)
if err != nil {
w.recordError(fmt.Errorf("release %s: %w", id, err))
return
}
if w.releases == nil {
w.releases = map[string]releaseTransaction{}
}
tx := w.releases[id]
if tx.ID == "" {
l, ok := w.leases[id]
if !ok {
w.recordError(fmt.Errorf("release %s: lease missing", id))
return
}
tx = releaseTransaction{ID: domain.NewID(), LeaseVersion: l.Version, Phase: "prepared", UpdatedAt: time.Now().UTC()}
w.releases[id] = tx
_ = w.save()
}
a := w.adapter(s, p.Remote)
if tx.Phase == "prepared" {
prepared, err := a.PrepareRelease(ctx, s)
if err != nil {
tx.LastError, tx.UpdatedAt = err.Error(), time.Now().UTC()
w.releases[id] = tx
_ = w.save()
w.recordError(fmt.Errorf("release %s prepare: %w", id, err))
return
}
tx.Ref, tx.AnchorSHA, tx.Phase, tx.LastError, tx.UpdatedAt = prepared.Ref, prepared.AnchorSHA, "anchor_pushed", "", time.Now().UTC()
w.releases[id] = tx
_ = w.save()
}
if tx.Phase == "anchor_pushed" {
l := w.leases[id]
if err := w.api.Release(ctx, id, tx.Ref, tx.AnchorSHA, tx.ID, l.Epoch, tx.LeaseVersion, w.sessionEvidence(ctx, id, s)); err != nil {
tx.LastError, tx.UpdatedAt = err.Error(), time.Now().UTC()
w.releases[id] = tx
_ = w.save()
w.recordError(fmt.Errorf("release %s commit: %w", id, err))
return
}
tx.Phase, tx.LastError, tx.UpdatedAt = "event_committed", "", time.Now().UTC()
w.releases[id] = tx
_ = w.save()
}
if tx.Phase == "event_committed" && !tx.AgentReleased {
if err := a.ReleaseAgent(ctx, s); err != nil {
tx.LastError, tx.UpdatedAt = err.Error(), time.Now().UTC()
w.releases[id] = tx
_ = w.save()
w.recordError(fmt.Errorf("release %s release agent: %w", id, err))
return
}
tx.AgentReleased, tx.LastError, tx.UpdatedAt = true, "", time.Now().UTC()
w.releases[id] = tx
_ = w.save()
}
if tx.Phase == "pickup_validated" {
if err := a.Kill(ctx, s); err != nil {
tx.LastError, tx.UpdatedAt = err.Error(), time.Now().UTC()
w.releases[id] = tx
_ = w.save()
w.recordError(fmt.Errorf("release %s retire predecessor: %w", id, err))
return
}
_ = os.Remove(filepath.Join(s.Worktree, herdr.HandoffReportFile))
tx.Phase, tx.UpdatedAt = "predecessor_retired", time.Now().UTC()
w.releases[id] = tx
_ = w.save()
delete(w.sessions, id)
delete(w.releases, id)
_ = w.save()
}
}
// ackPickup retries the successor acknowledgement from persisted lease state.
// It is safe after a lost response: the coordinator recognizes the exact
// transaction/lease epoch as an idempotent pickup.
func (w *worker) ackPickup(ctx context.Context, id string, s herdr.Session) error {
l, ok := w.leases[id]
if !ok || l.HandoffRef == "" || l.TransactionID == "" || l.AnchorSHA == "" {
return fmt.Errorf("pickup %s: lease is missing its release transaction", id)
}
if l.PickupAcknowledged {
return nil
}
if err := w.api.Pickup(ctx, id, l.HandoffRef, l.AnchorSHA, l.TransactionID, l.Epoch, l.Version, w.sessionEvidence(ctx, id, s)); err != nil {
return fmt.Errorf("pickup %s acknowledgement: %w", id, err)
}
l.PickupAcknowledged = true
l.Version++ // TaskPickupValidated increments the task version.
w.leases[id] = l
return w.save()
}
func (w *worker) sessionEvidence(ctx context.Context, taskID string, s herdr.Session) domain.SessionEvidence {
e := domain.SessionEvidence{PaneID: s.PaneID, HarnessID: w.harnessID, PaneState: "open", Source: "worker", CheckedAt: time.Now().UTC()}
text, err := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).PaneCapture(ctx, s, "recent")
if err != nil {
e.PaneState = "unreachable"
return e
}
if capture, err := w.api.PublishCapture(ctx, federation.Capture{TaskID: taskID, PaneID: s.PaneID, Text: text}); err == nil {
e.CapturedAt = capture.At
}
return e
}
func (w *worker) usageReceipt(s herdr.Session, l lease) map[string]any {
if s.SessionFile == "" && !(w.harness == "opencode" && s.SessionID != "") {
return map[string]any{"harness_id": w.harnessID, "consumed": 0, "known": false, "error": "native usage identity missing"}
}
var usage herdr.Usage
var err error
switch w.harness {
case "claude":
usage, err = herdr.ClaudeUsage(s.SessionFile)
case "codex":
usage, err = herdr.CodexUsage(s.SessionFile)
case "opencode":
usage, err = herdr.OpenCodeSessionUsage(s.SessionID)
}
if err != nil {
return map[string]any{"harness_id": w.harnessID, "consumed": 0, "known": false, "error": err.Error()}
}
delta := float64(usage.Numerator()) - l.UsageBaseline
if delta < 0 {
delta = 0
}
return map[string]any{"harness_id": w.harnessID, "input_tokens": usage.Input, "cache_read_tokens": usage.CacheRead, "cache_write_tokens": usage.CacheWrite, "consumed": delta, "lease_usage_delta": delta, "known": true}
}
func git(ctx context.Context, dir string, args ...string) ([]byte, error) {
return exec.CommandContext(ctx, "git", append([]string{"-C", dir}, args...)...).CombinedOutput()
}
// finalize performs only mechanical delivery work. It never asks the harness
// to narrate Git state, gates, or a report; those are generated from the
// worker-owned checkout and then verified against the configured remote.
func (w *worker) finalize(ctx context.Context, id string, s herdr.Session) (completionEvidence, error) {
t, ok := w.tasks[id]
if !ok {
return completionEvidence{}, fmt.Errorf("task cache missing")
}
p, err := w.project(t)
if err != nil {
return completionEvidence{}, err
}
if s.TaskFileSHA != "" {
if err := continuity.VerifyTaskFile(s.Worktree, s.TaskFileSHA); err != nil {
return completionEvidence{}, fmt.Errorf("verify immutable TASK.md: %w", err)
}
}
base, err := git(ctx, s.Worktree, "rev-parse", "HEAD")
if err != nil {
return completionEvidence{}, fmt.Errorf("base sha: %s: %w", base, err)
}
e := completionEvidence{TaskID: id, Project: t.Project, Worker: w.harnessID, Harness: w.harness, PaneID: s.PaneID, BaseSHA: strings.TrimSpace(string(base)), Remote: p.Remote, QualityGate: t.QualityGate, CompletedAt: time.Now().UTC()}
gateCommand := t.QualityGate
if gateCommand == "" {
gateCommand = p.QualityGate
}
e.QualityGate = gateCommand
if gateCommand != "" {
gate := exec.CommandContext(ctx, "sh", "-c", gateCommand)
gate.Dir = s.Worktree
if out, err := gate.CombinedOutput(); err != nil {
e.GateExit = 1
return completionEvidence{}, fmt.Errorf("quality gate %q: %s: %w", gateCommand, out, err)
}
}
if _, err := git(ctx, s.Worktree, "diff", "--quiet", "--", "TASK.md"); err != nil {
return completionEvidence{}, errors.New("TASK.md was modified")
}
if out, err := git(ctx, s.Worktree, "add", "-A", "--", ".", ":!.orchestra/done"); err != nil {
return completionEvidence{}, fmt.Errorf("stage result: %s: %w", out, err)
}
if _, err := git(ctx, s.Worktree, "diff", "--cached", "--quiet"); err != nil {
if out, err := git(ctx, s.Worktree, "commit", "-m", "orchestra: complete "+id); err != nil {
return completionEvidence{}, fmt.Errorf("commit result: %s: %w", out, err)
}
}
branch, err := git(ctx, s.Worktree, "branch", "--show-current")
if err != nil || strings.TrimSpace(string(branch)) == "" {
return completionEvidence{}, fmt.Errorf("result branch: %s: %w", branch, err)
}
e.Branch = strings.TrimSpace(string(branch))
sha, err := git(ctx, s.Worktree, "rev-parse", "HEAD")
if err != nil {
return completionEvidence{}, fmt.Errorf("result sha: %s: %w", sha, err)
}
e.ResultSHA = strings.TrimSpace(string(sha))
if out, err := git(ctx, s.Worktree, "push", p.Remote, "HEAD:refs/heads/"+e.Branch); err != nil {
return completionEvidence{}, fmt.Errorf("push result: %s: %w", out, err)
}
remote, err := git(ctx, s.Worktree, "ls-remote", p.Remote, "refs/heads/"+e.Branch)
if err != nil || !strings.HasPrefix(string(remote), e.ResultSHA+"\t") {
return completionEvidence{}, fmt.Errorf("verify pushed sha: got %q: %w", strings.TrimSpace(string(remote)), err)
}
return e, nil
}
func (w *worker) renewLeases(ctx context.Context) {
if w.executionBackend() == nil {
return
}
now := time.Now()
for taskID, l := range w.leases {
s, ok := w.sessions[taskID]
if !ok || l.Version == 0 || l.Until.After(now.Add(10*time.Minute)) {
continue
}
adapter := herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}
// Input lines are excluded: keystrokes arriving at a pane, from
// Orchestra or from anyone else, are not the agent doing work.
text, err := w.paneProgress(ctx, adapter, s)
if err != nil {
w.recordError(fmt.Errorf("validate lease %s: %w", taskID, err))
continue
}
// A live pane is not progress. Renewing on pane existence alone let a
// pane that opened and never started hold its lease forever, which is
// what orphaned the July task once the launch itself had failed.
status, err := adapter.AgentStatus(ctx, s)
if err != nil {
w.recordError(fmt.Errorf("validate lease %s: %w", taskID, err))
continue
}
progress := domain.Hash([]byte(text))
switch {
case herdr.IsBusy(status):
case progress != l.ProgressSHA && l.ProgressSHA != "":
case l.ProgressSHA == "":
// First renewal has no baseline to compare against. Record one and
// allow this renewal; the next one must show real movement.
default:
w.recordError(fmt.Errorf("lease %s not renewed: agent status %s and pane unchanged since the last renewal", taskID, status))
continue
}
if err := w.api.Renew(ctx, taskID, l.Epoch, l.Version, int((30 * time.Minute).Seconds())); err != nil {
w.recordError(fmt.Errorf("renew lease %s: %w", taskID, err))
log.Printf("renew lease %s: %v", taskID, err)
} else {
// RenewLease appends one event. Retain that epoch locally until its
// replay arrives so a release transaction uses the same version.
l.Version++
l.Until = now.Add(30 * time.Minute)
l.ProgressSHA = progress
w.leases[taskID] = l
_ = w.save()
}
}
}
// publishCaptures makes remote panes observable without allowing the
// coordinator to touch their unix herdr socket.
func (w *worker) publishCaptures(ctx context.Context) {
for taskID, session := range w.sessions {
text, err := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).PaneCapture(ctx, session, "recent")
if err != nil {
continue
}
if _, err := w.api.PublishCapture(ctx, federation.Capture{TaskID: taskID, PaneID: session.PaneID, Text: text}); err != nil {
log.Printf("publish capture %s: %v", taskID, err)
}
}
}
type approvalInput struct {
Text string
Keys []string
}
func approvalResponse(text, kind string) (approvalInput, bool) {
low := strings.ToLower(text)
// Never invent a keystroke. y/n prompts label both decisions directly.
if strings.Contains(low, "[y/n]") || strings.Contains(low, "(y/n)") {
if kind == "grant_approval" {
return approvalInput{Text: "y\n"}, true
}
return approvalInput{Text: "n\n"}, true
}
// OpenCode's explicit selector states "Allow once Allow always Reject"
// and "enter confirm". Send a real ENTER key, not a newline through
// pane.send_text: OpenCode's selector does not treat the latter as input.
// Enter is consequently a bounded one-time grant;
// rejection would require unobservable selector navigation, so refuse it.
if kind == "grant_approval" && strings.Contains(low, "allow once") && strings.Contains(low, "allow always") && strings.Contains(low, "reject") && strings.Contains(low, "enter confirm") {
return approvalInput{Keys: []string{"ENTER"}}, true
}
return approvalInput{}, false
}
func (w *worker) runCommands(ctx context.Context) {
commands, err := w.api.Commands(ctx)
if err != nil {
log.Printf("poll controls: %v", err)
return
}
for _, command := range commands {
session, ok := w.sessions[command.TaskID]
if !ok || session.PaneID != command.PaneID {
_ = w.api.ResolveCommand(ctx, command.ID, "stale", "session or pane changed")
continue
}
text, err := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).PaneCapture(ctx, session, "recent")
if err != nil {
_ = w.api.ResolveCommand(ctx, command.ID, "rejected", "capture unavailable: "+err.Error())
continue
}
capture, err := w.api.PublishCapture(ctx, federation.Capture{TaskID: command.TaskID, PaneID: session.PaneID, Text: text})
if err != nil {
_ = w.api.ResolveCommand(ctx, command.ID, "rejected", "cannot publish capture: "+err.Error())
continue
}
if capture.Revision != command.CaptureRevision {
_ = w.api.ResolveCommand(ctx, command.ID, "stale", "capture revision changed")
continue
}
input, ok := approvalResponse(text, command.Kind)
if !ok {
_ = w.api.ResolveCommand(ctx, command.ID, "rejected", "prompt does not expose an executable approval control")
continue
}
backend := w.executionBackend()
var inputErr error
if len(input.Keys) > 0 {
inputErr = backend.SendKeys(ctx, session, input.Keys)
} else {
inputErr = backend.SendText(ctx, session, input.Text)
}
if inputErr != nil {
_ = w.api.ResolveCommand(ctx, command.ID, "rejected", backend.Kind()+" backend did not acknowledge input: "+inputErr.Error())
continue
}
if err := w.api.ResolveCommand(ctx, command.ID, "acknowledged", ""); err != nil {
log.Printf("ack command %s: %v", command.ID, err)
}
}
}
func (w *worker) once(ctx context.Context) error {
es, _, err := w.api.Events(ctx, w.cursor)
if err != nil {
return err
}
for _, e := range es {
if t, ok := created(e); ok {
w.tasks[t.ID] = t
}
if e.Type == "TaskLeased" {
var p struct {
HarnessID string `json:"harness_id"`
Epoch string `json:"lease_epoch"`
HandoffRef string `json:"handoff_ref"`
TransactionID string `json:"transaction_id"`
AnchorSHA string `json:"anchor_sha"`
}
if json.Unmarshal(e.Payload, &p) == nil && p.HarnessID == w.harnessID {
var until struct {
UntilNS int64 `json:"until_ns"`
}
_ = json.Unmarshal(e.Payload, &until)
w.leases[e.TaskID] = lease{Epoch: p.Epoch, HandoffRef: p.HandoffRef, TransactionID: p.TransactionID, AnchorSHA: p.AnchorSHA, Version: e.Version, Until: time.Unix(0, until.UntilNS)}
}
}
if e.Type == "TaskLeaseRenewed" {
var p struct {
HarnessID string `json:"harness_id"`
Epoch string `json:"lease_epoch"`
UntilNS int64 `json:"until_ns"`
}
if json.Unmarshal(e.Payload, &p) == nil && p.HarnessID == w.harnessID {
l := w.leases[e.TaskID]
l.Version, l.Epoch, l.Until = e.Version, p.Epoch, time.Unix(0, p.UntilNS)
w.leases[e.TaskID] = l
}
}
if e.Type == "TaskLaunchAcknowledged" {
if l, ok := w.leases[e.TaskID]; ok {
l.Version = e.Version
w.leases[e.TaskID] = l
}
}
if e.Type == "TaskNeedsAttention" {
// The diagnostic event increments the aggregate version but leaves
// ownership intact. Keep our locally persisted expected version in
// sync so a late, otherwise valid completion is not self-staled.
if l, ok := w.leases[e.TaskID]; ok {
l.Version = e.Version
w.leases[e.TaskID] = l
}
}
if e.Type == "TaskCompleted" {
delete(w.leases, e.TaskID)
if session, active := w.sessions[e.TaskID]; active {
if w.executionBackend() == nil {
delete(w.sessions, e.TaskID)
} else if err := (herdr.CLIAdapter{Backend: w.executionBackend(), Harness: w.harness}).Kill(ctx, session); err != nil {
w.quarantined[e.TaskID] = true
w.recordError(fmt.Errorf("close completed pane %s: %w", e.TaskID, err))
continue
} else {
_ = os.Remove(filepath.Join(session.Worktree, ".orchestra", "done"))
_ = os.Remove(filepath.Join(session.Worktree, ".orchestra"))
delete(w.sessions, e.TaskID)
}
}
}
if e.Type == "TaskPickupValidated" {
var p struct {
TransactionID string `json:"transaction_id"`
}
if json.Unmarshal(e.Payload, &p) == nil {
if tx := w.releases[e.TaskID]; tx.ID != "" && tx.ID == p.TransactionID {
tx.Phase, tx.UpdatedAt = "pickup_validated", time.Now().UTC()
w.releases[e.TaskID] = tx
}
}
if l, ok := w.leases[e.TaskID]; ok {
l.Version = e.Version
w.leases[e.TaskID] = l
}
}
if e.Type == "TaskReleased" || e.Type == "TaskFailed" || e.Type == "TaskBlocked" {
if e.Type == "TaskReleased" {
var p struct {
TransactionID string `json:"transaction_id"`
}
if json.Unmarshal(e.Payload, &p) == nil {
if tx := w.releases[e.TaskID]; tx.ID != "" && tx.ID == p.TransactionID && tx.Phase == "anchor_pushed" {
tx.Phase, tx.LastError, tx.UpdatedAt = "event_committed", "", time.Now().UTC()
w.releases[e.TaskID] = tx
}
}
}
delete(w.leases, e.TaskID)
// A releasing predecessor remains intentionally recoverable until
// TaskPickupValidated for its transaction. Do not erase its pane
// mapping merely because our own release event was replayed.
if _, releasing := w.releases[e.TaskID]; !releasing {
if session, active := w.sessions[e.TaskID]; active {
w.quarantine(ctx, e.TaskID, session)
}
}
}
if e.Seq > w.cursor {
w.cursor = e.Seq
}
}
// A coordinator restart can restore its task snapshot without retaining
// the in-memory event tail. In that state a worker with a persisted cursor
// receives an empty page even though a lease is currently assigned to it.
// Reconcile the authoritative projection before treating an empty page as
// "nothing to do"; otherwise the lease remains invisible until expiry.
if len(es) == 0 {
if err := w.reconcileLeases(ctx); err != nil {
return err
}
}
// State is only a cache. If a lease survived but its TaskCreated event is
// older than the worker's cursor (or the event has been compacted), hydrate
// the authoritative task projection before deciding whether to start.
for taskID := range w.leases {
if _, ok := w.tasks[taskID]; !ok {
tasks, err := w.api.Tasks(ctx)
if err != nil {
return fmt.Errorf("hydrate leased task %s: %w", taskID, err)
}
for _, task := range tasks {
w.tasks[task.ID] = task
}
break
}
}
// Project the whole batch before launching. This prevents a new worker
// from resurrecting every historical lease during its initial replay.
for taskID, l := range w.leases {
if _, started := w.sessions[taskID]; started {
continue
}
if t, ok := w.tasks[taskID]; ok {
if err := w.start(ctx, t, l.HandoffRef); err != nil {
log.Printf("lease %s: %v", t.ID, err)
_, started := w.sessions[taskID]
class := classifyLaunchError(err, started)
var evidence domain.SessionEvidence
if session, ok := w.sessions[taskID]; ok {
evidence = w.sessionEvidence(ctx, taskID, session)
}
if nackErr := w.api.NackStart(ctx, taskID, l.Epoch, l.Version, class, err.Error(), evidence); nackErr != nil {
w.recordError(fmt.Errorf("nack launch %s: %w", taskID, nackErr))
continue
}
if class != "launch_uncertain" {
delete(w.leases, taskID)
}
}
}
}
// Unit/replay-only workers intentionally have no execution backend. A
// production worker always does, and only then participates in the live
// capture/control protocol.
if w.executionBackend() != nil {
w.publishCaptures(ctx)
w.runCommands(ctx)
w.renewLeases(ctx)
}
w.retryQuarantines(ctx)
w.releaseReady(ctx)
if err := w.save(); err != nil {
return err
}
return w.api.Ack(ctx, w.cursor)
}
func (w *worker) reconcileLeases(ctx context.Context) error {
tasks, err := w.api.Tasks(ctx)
if err != nil {
return fmt.Errorf("reconcile leased tasks: %w", err)
}
active := make(map[string]lease)
for _, task := range tasks {
w.tasks[task.ID] = task
if (task.State == domain.StateLeased || task.State == domain.StateNeedsAttention) && task.Lease != nil && task.Lease.HarnessID == w.harnessID {
active[task.ID] = lease{Epoch: task.Lease.Epoch, HandoffRef: task.HandoffRef, TransactionID: task.ReleaseTransaction, AnchorSHA: task.ReleaseAnchor, Version: task.Version, Until: task.Lease.Until}
}
}
for taskID := range w.leases {
if _, ok := active[taskID]; !ok {
delete(w.leases, taskID)
}
}
for taskID, l := range active {
// The coordinator is authoritative for the lease, not for what this
// worker has observed under it. Rebuilding the struct wholesale wiped
// every worker-local field each tick, which silently disabled the
// renewal progress check and reset the usage baseline. Carry them
// across, but only while the epoch is the same lease.
if prev, ok := w.leases[taskID]; ok && prev.Epoch == l.Epoch {
l.ProgressSHA = prev.ProgressSHA
l.UsageBaseline = prev.UsageBaseline
l.PickupAcknowledged = prev.PickupAcknowledged
}
w.leases[taskID] = l
}
return nil
}
// reRegisterAfterCoordinatorRestart restores the coordinator's in-memory
// worker registry. A worker must survive a server restart without operator
// intervention; its persisted cursor and sessions remain valid.
func (w *worker) reRegisterAfterCoordinatorRestart(ctx context.Context, cause error) bool {
if cause == nil || !strings.Contains(cause.Error(), "401 Unauthorized: unknown worker") {
return false
}
if err := w.api.Register(ctx, w.registration); err != nil {
log.Printf("re-register: %v", err)
return false
}
log.Printf("re-registered after coordinator restart")
return true
}
func required(k string) string {
v := os.Getenv(k)
if v == "" {
log.Fatalf("%s required", k)
}
return v
}
// harnessSpec is one harness identity this process serves. Each spec becomes a
// separate federation identity because the coordinator authorizes lease calls by
// comparing the URL's worker id against the lease's harness id
// (cmd/orchestra/main.go). One process may therefore hold several identities,
// but it may never present one identity for several harnesses.
type harnessSpec struct {
ID string `json:"id"`
Harness string `json:"harness"`
Token string `json:"token,omitempty"`
Backend string `json:"backend,omitempty"`
// Herdr is the JSON-RPC address for Backend "herdr".
Herdr string `json:"herdr,omitempty"`
// TmuxSocket and Command configure Backend "tmux".
TmuxSocket string `json:"tmux_socket,omitempty"`
Command string `json:"command,omitempty"`
State string `json:"state,omitempty"`
Address string `json:"address,omitempty"`
}
// tokenEnvKey maps a harness id onto a per-identity token variable, so a
// multi-harness deployment keeps its tokens in the protected environment file
// rather than in the harness config file.
func tokenEnvKey(id string) string {
var b strings.Builder
b.WriteString("ORCHESTRA_WORKER_TOKEN_")
for _, r := range strings.ToUpper(id) {
if (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') {
b.WriteRune(r)
continue
}
b.WriteRune('_')
}
return b.String()
}
// harnessSpecs reads the multi-harness declaration, falling back to the legacy
// single-harness environment so an existing deployment upgrades unchanged.
func harnessSpecs() []harnessSpec {
path := os.Getenv("ORCHESTRA_WORKER_HARNESS_CONFIG_FILE")
if path == "" {
return []harnessSpec{{
ID: required("ORCHESTRA_WORKER_HERDR_ID"),
Harness: required("ORCHESTRA_WORKER_HARNESS"),
Token: required("ORCHESTRA_WORKER_TOKEN"),
Backend: os.Getenv("ORCHESTRA_WORKER_BACKEND"),
Herdr: os.Getenv("ORCHESTRA_WORKER_HERDR"),
TmuxSocket: os.Getenv("ORCHESTRA_WORKER_TMUX_SOCKET"),
Command: os.Getenv("ORCHESTRA_WORKER_HARNESS_COMMAND"),
State: os.Getenv("ORCHESTRA_WORKER_STATE"),
Address: os.Getenv("ORCHESTRA_WORKER_ADDRESS"),
}}
}
b, err := os.ReadFile(path)
if err != nil {
log.Fatalf("read ORCHESTRA_WORKER_HARNESS_CONFIG_FILE: %v", err)
}
var specs []harnessSpec
if err := json.Unmarshal(b, &specs); err != nil {
log.Fatalf("parse ORCHESTRA_WORKER_HARNESS_CONFIG_FILE: %v", err)
}
if len(specs) == 0 {
log.Fatal("ORCHESTRA_WORKER_HARNESS_CONFIG_FILE declares no harnesses")
}
seen := map[string]bool{}
for i, spec := range specs {
if spec.ID == "" || spec.Harness == "" {
log.Fatalf("harness %d requires id and harness", i)
}
if seen[spec.ID] {
log.Fatalf("harness id %q is declared twice", spec.ID)
}
seen[spec.ID] = true
if spec.Token == "" {
key := tokenEnvKey(spec.ID)
if spec.Token = os.Getenv(key); spec.Token == "" {
log.Fatalf("harness %s has no token: set %s or its config-file token", spec.ID, key)
}
specs[i] = spec
}
}
return specs
}
// statePathFor gives every identity its own state file. Sharing one across
// harnesses would let a persisted session and its lease cross backends, handing
// a tmux backend herdr pane ids it cannot act on. Only the single-harness form
// keeps the historical default path, so an existing deployment recovers its
// sessions after the upgrade instead of orphaning them.
func statePathFor(spec harnessSpec, root, stateDir string, single bool) string {
switch {
case spec.State != "":
return spec.State
case stateDir != "":
return filepath.Join(stateDir, "state-"+spec.ID+".json")
case single:
return filepath.Join(root, ".orchestra-worker-state.json")
default:
return filepath.Join(root, ".orchestra-worker-state-"+spec.ID+".json")
}
}
// backendFor builds the machine-local execution backend for one harness.
func backendFor(spec harnessSpec) herdr.Backend {
name := strings.ToLower(strings.TrimSpace(spec.Backend))
if name == "" {
name = "herdr"
}
switch name {
case "herdr":
address := spec.Herdr
if address == "" {
log.Fatalf("harness %s: backend herdr requires an address", spec.ID)
}
return herdr.New(address)
case "tmux":
if spec.Harness != "claude" {
log.Fatalf("harness %s: backend tmux currently supports only harness claude, got %q", spec.ID, spec.Harness)
}
return herdr.NewTmuxBackend(spec.TmuxSocket, spec.Command)
default:
log.Fatalf("harness %s: unsupported backend %q (want herdr or tmux)", spec.ID, name)
return nil
}
}
func main() {
workerID := required("ORCHESTRA_WORKER_ID")
hard := .75
if v, err := strconv.ParseFloat(os.Getenv("ORCHESTRA_OCCUPANCY_HARD"), 64); err == nil && v > 0 && v < 1 {
hard = v
}
projects := map[string]projectConfig{}
if path := os.Getenv("ORCHESTRA_WORKER_PROJECT_CONFIG_FILE"); path != "" {
b, err := os.ReadFile(path)
if err != nil {
log.Fatalf("read ORCHESTRA_WORKER_PROJECT_CONFIG_FILE: %v", err)
}
if err := json.Unmarshal(b, &projects); err != nil {
log.Fatalf("parse ORCHESTRA_WORKER_PROJECT_CONFIG_FILE: %v", err)
}
} else {
for _, project := range strings.Split(os.Getenv("ORCHESTRA_WORKER_PROJECTS"), ",") {
if project = strings.TrimSpace(project); project != "" {
projects[project] = projectConfig{Repo: required("ORCHESTRA_REPO"), Root: required("ORCHESTRA_WORKTREE_ROOT"), Remote: required("ORCHESTRA_GIT_REMOTE")}
}
}
}
if len(projects) == 0 {
// Existing deployments may be upgraded before their protected systemd
// environment is amended. Stay observable and fail closed in that
// interval: an empty declaration makes this worker ineligible for all
// new leases instead of turning a configuration rollout into a crash
// loop or treating its legacy global checkout as every project.
log.Printf("no ORCHESTRA_WORKER_PROJECT_CONFIG_FILE/ORCHESTRA_WORKER_PROJECTS; registering with no supported projects")
}
for project, config := range projects {
if config.Repo == "" || config.Root == "" || config.Remote == "" {
log.Fatalf("project %q requires repo, worktree_root, and remote", project)
}
}
supported := make([]string, 0, len(projects))
for project := range projects {
supported = append(supported, project)
}
sort.Strings(supported)
repo, root, remote := required("ORCHESTRA_REPO"), required("ORCHESTRA_WORKTREE_ROOT"), required("ORCHESTRA_GIT_REMOTE")
if len(supported) > 0 {
first := projects[supported[0]]
repo, root, remote = first.Repo, first.Root, first.Remote
}
soft, window := .55, int64(200000)
if v, err := strconv.ParseFloat(os.Getenv("ORCHESTRA_OCCUPANCY_SOFT"), 64); err == nil && v > 0 && v < hard {
soft = v
}
if v, err := strconv.ParseInt(os.Getenv("ORCHESTRA_CONTEXT_WINDOW"), 10, 64); err == nil && v > 0 {
window = v
}
url, admit := required("ORCHESTRA_URL"), os.Getenv("ORCHESTRA_FEDERATION_ADMIT_TOKEN")
specs := harnessSpecs()
if len(specs) == 1 && specs[0].ID != workerID {
// A single-harness deployment keeps the historical invariant: its lease
// owner and its process identity are the same name.
log.Fatal("ORCHESTRA_WORKER_ID must equal ORCHESTRA_WORKER_HERDR_ID so leases and offline recovery have one owner")
}
stateDir := os.Getenv("ORCHESTRA_WORKER_STATE_DIR")
workers := make([]*worker, 0, len(specs))
for _, spec := range specs {
w := &worker{
api: federation.Client{BaseURL: url, WorkerID: spec.ID, Token: spec.Token, AdmitToken: admit},
harnessID: spec.ID,
harness: spec.Harness,
repo: repo,
root: root,
remote: remote,
projects: projects,
tasks: map[string]domain.Task{},
sessions: map[string]herdr.Session{},
leases: map[string]lease{},
releases: map[string]releaseTransaction{},
quarantined: map[string]bool{},
statePath: spec.State,
hard: hard,
soft: soft,
window: window,
// Capacity stays one per identity because a herdr's declared
// concurrency is one. Serving N harnesses gives the process N slots.
registration: federation.Worker{ID: spec.ID, Address: spec.Address, Capacity: 1, SupportedProjects: supported, Build: buildinfo.Current()},
}
if w.registration.Address == "" {
w.registration.Address = os.Getenv("ORCHESTRA_WORKER_ADDRESS")
}
w.statePath = statePathFor(spec, root, stateDir, len(specs) == 1)
if err := w.load(); err != nil {
log.Fatal(err)
}
w.backend = backendFor(spec)
if err := w.api.Register(context.Background(), w.registration); err != nil {
log.Fatalf("register %s: %v", spec.ID, err)
}
// Same reason as the coordinator's: the worker half of the pair must be
// checkable from journalctl, not only from the coordinator's
// credentialed worker list.
b := buildinfo.Current()
log.Printf("orchestra-worker revision %s built %s dirty %s", b.Revision, b.Time, b.Dirty)
log.Printf("serving harness %s (%s) on %s backend, state %s", spec.ID, spec.Harness, w.backend.Kind(), w.statePath)
workers = append(workers, w)
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
// Identities are served sequentially. Their checkouts and Git remotes are
// shared, so concurrent ticks would race two fetch/worktree operations on
// one repository for no useful latency gain at this fan-out.
for _, w := range workers {
if err := w.api.Heartbeat(ctx, w.health(ctx)); err != nil {
w.recordError(fmt.Errorf("heartbeat: %w", err))
log.Printf("[%s] heartbeat: %v", w.harnessID, err)
if w.reRegisterAfterCoordinatorRestart(ctx, err) {
continue
}
}
if err := w.once(ctx); err != nil {
w.recordError(fmt.Errorf("poll: %w", err))
log.Printf("[%s] poll: %v", w.harnessID, err)
w.reRegisterAfterCoordinatorRestart(ctx, err)
}
}
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}
// federatedTurn is the worker half of a turn boundary. The coordinator
// reconciles human input and answers with the decisions this session has not
// been shown; the worker delivers them into its own pane.
//
// Nothing is preempted. The boundary is confirmed against the live pane
// first, so a correction never lands mid tool call.
func (w *worker) federatedTurn(ctx context.Context, id string, a herdr.Adapter, verdict string) {
l, ok := w.leases[id]
if !ok {
return
}
s, ok := w.sessions[id]
if !ok {
return
}
boundary, ok := a.(herdr.TurnBoundary)
if !ok {
return
}
at, err := boundary.AtTurnBoundary(ctx, s)
if err != nil {
w.recordError(fmt.Errorf("turn boundary %s: %w", id, err))
return
}
if !at {
return
}
answer, err := w.api.Turn(ctx, id, l.Epoch, verdict, s.DeliveredDecisions)
if err != nil {
// Observable, not fatal. A coordinator that cannot be reached does not
// make this session's current intent any more stale than it already is.
w.recordError(fmt.Errorf("federated turn %s: %w", id, err))
return
}
if answer.Verdict == orchestrator.TurnPrepareHandoff && !s.HandoffRequested {
// The coordinator has lost the ability to refresh this task's intent.
// Ask for a handoff; the release loop takes over as soon as the agent
// writes the report, exactly as it does for a local session.
requester, ok := a.(herdr.ReasonedHandoffRequester)
if !ok {
w.recordError(fmt.Errorf("reconcile failure handoff %s: adapter cannot state a reason", id))
return
}
if err := requester.RequestHandoffReason(ctx, s, "reconcile_failure", nil); err != nil {
w.recordError(fmt.Errorf("reconcile failure handoff %s: %w", id, err))
return
}
s.HandoffRequested, s.HandoffReason = true, "reconcile_failure"
w.sessions[id] = s
_ = w.save()
return
}
if len(answer.Decisions) == 0 {
return
}
if err := w.sendPrompt(ctx, s, agentctx.DecisionNotice(answer.Decisions)); err != nil {
// Not recorded as delivered, so the next boundary retries.
w.recordError(fmt.Errorf("deliver decisions %s: %w", id, err))
return
}
for _, d := range answer.Decisions {
s.DeliveredDecisions = append(s.DeliveredDecisions, d.ID)
}
w.sessions[id] = s
_ = w.save()
}
// sendPrompt delivers Orchestra-originated input and confirms the harness took
// it. A phase continuation or a decision notice whose Enter is lost strands the
// session exactly as a lost launch does.
func (w *worker) sendPrompt(ctx context.Context, s herdr.Session, text string) error {
backend := w.executionBackend()
if backend == nil {
return fmt.Errorf("execution backend is not configured")
}
if err := backend.Prompt(ctx, s.PaneID, text, time.Minute); err != nil {
return err
}
return w.confirmInput(ctx, s, text)
}
// confirmInput is the one place Orchestra proves a write landed. A backend
// whose own protocol acknowledges input does not implement InputConfirmer and
// needs no second opinion.
func (w *worker) confirmInput(ctx context.Context, s herdr.Session, text string) error {
c, ok := w.executionBackend().(herdr.InputConfirmer)
if !ok {
return nil
}
evidence, err := c.ConfirmInput(ctx, s, text)
if err != nil {
return err
}
log.Printf("input to %s confirmed: %s", s.PaneID, evidence)
return nil
}
// paneProgress reports pane content with input lines removed where the backend
// can separate them, and falls back to the raw capture where it cannot.
func (w *worker) paneProgress(ctx context.Context, adapter herdr.CLIAdapter, s herdr.Session) (string, error) {
if p, ok := w.executionBackend().(herdr.PaneProgress); ok {
return p.PaneProgress(ctx, s)
}
return adapter.PaneCapture(ctx, s, "recent")
}