Files
orchestra/internal/orchestrator/orchestrator.go
T
2026-07-26 20:06:44 +04:00

127 lines
3.3 KiB
Go

// Package orchestrator connects router lease events to an opaque herdr
// session. It is deliberately small: scheduling remains in router and the
// adapter remains the only component that knows how to drive a harness.
package orchestrator
import (
"context"
"encoding/json"
"fmt"
"orchestra/internal/domain"
"orchestra/internal/herdr"
"orchestra/internal/store"
"os"
"os/exec"
"path/filepath"
"sync"
)
type Worktrees interface {
Create(context.Context, domain.Task) (string, error)
}
type Adapters interface {
Adapter(string) (herdr.Adapter, error)
}
// GitWorktrees creates one isolated checkout per task. The root is expected
// to be a clone containing the project's remote; callers may set a separate
// root per deployment.
type GitWorktrees struct {
Root string
Repo string
}
func (w GitWorktrees) Create(ctx context.Context, t domain.Task) (string, error) {
if w.Root == "" || w.Repo == "" {
return "", fmt.Errorf("worktree: root and repo required")
}
if err := os.MkdirAll(w.Root, 0755); err != nil {
return "", err
}
p := filepath.Join(w.Root, t.ID)
if _, err := os.Stat(p); err == nil {
return p, nil
}
branch := "orchestra/" + t.ID
cmd := exec.CommandContext(ctx, "git", "-C", w.Repo, "worktree", "add", "-b", branch, p, "HEAD")
if out, err := cmd.CombinedOutput(); err != nil {
return "", fmt.Errorf("%s: %w", string(out), err)
}
return p, nil
}
type AdapterFactory struct{ Herdrs map[string]herdr.Adapter }
func (f AdapterFactory) Adapter(id string) (herdr.Adapter, error) {
a, ok := f.Herdrs[id]
if !ok {
return nil, fmt.Errorf("adapter %q not registered", id)
}
return a, nil
}
type Coordinator struct {
Store *store.Store
Worktrees Worktrees
Adapters Adapters
mu sync.Mutex
sessions map[string]herdr.Session
}
func (c *Coordinator) Start(ctx context.Context, e domain.Event) error {
if e.Type != "TaskLeased" {
return nil
}
if c.Store == nil || c.Worktrees == nil || c.Adapters == nil {
return fmt.Errorf("orchestrator: dependencies required")
}
t, ok := c.Store.Task(e.TaskID)
if !ok {
return domain.ErrNotFound
}
var p struct {
HarnessID string `json:"harness_id"`
HandoffRef string `json:"handoff_ref"`
}
if err := json.Unmarshal(e.Payload, &p); err != nil || p.HarnessID == "" {
return fmt.Errorf("orchestrator: invalid lease")
}
w, err := c.Worktrees.Create(ctx, t)
if err != nil {
return c.block(t, "worktree: "+err.Error())
}
a, err := c.Adapters.Adapter(p.HarnessID)
if err != nil {
return c.block(t, "adapter: "+err.Error())
}
s, err := a.Lease(ctx, t.ID, w)
if err != nil {
return c.block(t, "lease: "+err.Error())
}
if p.HandoffRef != "" {
if err = a.Bootstrap(ctx, s, p.HandoffRef); err != nil {
_ = a.Kill(ctx, s)
return c.block(t, "bootstrap: "+err.Error())
}
}
c.mu.Lock()
if c.sessions == nil {
c.sessions = map[string]herdr.Session{}
}
c.sessions[t.ID] = s
c.mu.Unlock()
return nil
}
func (c *Coordinator) block(t domain.Task, reason string) error {
b, _ := json.Marshal(map[string]string{"blocker": reason})
return c.Store.Append(domain.Event{ID: domain.NewID(), Type: "TaskBlocked", TaskID: t.ID, Version: t.Version + 1, Payload: b})
}
func (c *Coordinator) Session(taskID string) (herdr.Session, bool) {
c.mu.Lock()
defer c.mu.Unlock()
s, ok := c.sessions[taskID]
return s, ok
}