diff --git a/cmd/orchestra/main.go b/cmd/orchestra/main.go index 4da15d5..d082b0e 100644 --- a/cmd/orchestra/main.go +++ b/cmd/orchestra/main.go @@ -5,6 +5,7 @@ import ( "log" "net/http" "orchestra/internal/domain" + "orchestra/internal/registry" "orchestra/internal/store" "os" "strconv" @@ -22,6 +23,11 @@ func main() { if err != nil { log.Fatal(err) } + if config := os.Getenv("ORCHESTRA_CONFIG"); config != "" { + if _, err := registry.Load(config); err != nil { + log.Fatalf("load orchestra config: %v", err) + } + } mux := http.NewServeMux() mux.HandleFunc("/v1/tasks", func(w http.ResponseWriter, r *http.Request) { if r.Method == "GET" { diff --git a/internal/registry/registry.go b/internal/registry/registry.go new file mode 100644 index 0000000..4702c80 --- /dev/null +++ b/internal/registry/registry.go @@ -0,0 +1,159 @@ +// Package registry contains the static project, machine, and herdr topology. +package registry + +import ( + "encoding/json" + "errors" + "fmt" + "net" + "os" + "sort" + "strings" + "time" +) + +var ( + ErrUnknownProject = errors.New("unknown project") + ErrUnknownMachine = errors.New("unknown machine") + ErrNoAffinity = errors.New("project has no machine affinity") +) + +type Project struct { + ID string `json:"id"` + MachineAffinity []string `json:"machine_affinity"` +} +type Machine struct { + ID string `json:"id"` + Address string `json:"address"` +} +type Herdr struct { + ID string `json:"id"` + MachineID string `json:"machine_id"` + Address string `json:"address,omitempty"` + Capabilities []string `json:"capabilities"` + Concurrency int `json:"concurrency"` +} +type Config struct { + Projects []Project `json:"projects"` + Machines []Machine `json:"machines"` + Herdrs []Herdr `json:"herdrs"` +} +type Registry struct { + projects map[string]Project + machines map[string]Machine + herdrs map[string]Herdr +} + +func Load(path string) (Registry, error) { + b, err := os.ReadFile(path) + if err != nil { + return Registry{}, err + } + var c Config + if err = json.Unmarshal(b, &c); err != nil { + return Registry{}, fmt.Errorf("registry config: %w", err) + } + return New(c) +} + +func New(c Config) (Registry, error) { + r := Registry{map[string]Project{}, map[string]Machine{}, map[string]Herdr{}} + for _, p := range c.Projects { + if err := putID(r.projects, p.ID, "project"); err != nil { + return Registry{}, err + } + if len(p.MachineAffinity) == 0 { + return Registry{}, fmt.Errorf("project %q: %w", p.ID, ErrNoAffinity) + } + r.projects[p.ID] = p + } + for _, m := range c.Machines { + if err := putID(r.machines, m.ID, "machine"); err != nil { + return Registry{}, err + } + if strings.TrimSpace(m.Address) == "" { + return Registry{}, fmt.Errorf("machine %q: address required", m.ID) + } + r.machines[m.ID] = m + } + for _, h := range c.Herdrs { + if err := putID(r.herdrs, h.ID, "herdr"); err != nil { + return Registry{}, err + } + if _, ok := r.machines[h.MachineID]; !ok { + return Registry{}, fmt.Errorf("herdr %q: %w %q", h.ID, ErrUnknownMachine, h.MachineID) + } + if h.Concurrency < 0 { + return Registry{}, fmt.Errorf("herdr %q: negative concurrency", h.ID) + } + r.herdrs[h.ID] = h + } + for _, p := range r.projects { + for _, m := range p.MachineAffinity { + if _, ok := r.machines[m]; !ok { + return Registry{}, fmt.Errorf("project %q: %w %q", p.ID, ErrUnknownMachine, m) + } + } + } + return r, nil +} +func putID[T any](m map[string]T, id, kind string) error { + if strings.TrimSpace(id) == "" { + return fmt.Errorf("%s id required", kind) + } + if _, ok := m[id]; ok { + return fmt.Errorf("duplicate %s %q", kind, id) + } + return nil +} +func (r Registry) Project(id string) (Project, bool) { p, ok := r.projects[id]; return p, ok } +func (r Registry) Machine(id string) (Machine, bool) { m, ok := r.machines[id]; return m, ok } +func (r Registry) Herdr(id string) (Herdr, bool) { h, ok := r.herdrs[id]; return h, ok } +func (r Registry) Projects() []Project { return projects(r.projects) } +func projects(m map[string]Project) []Project { + out := make([]Project, 0, len(m)) + for _, v := range m { + out = append(out, v) + } + sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID }) + return out +} + +type Reachability interface { + Reachable(address string, timeout time.Duration) bool +} +type TCPReachability struct{} + +func (TCPReachability) Reachable(address string, timeout time.Duration) bool { + c, err := net.DialTimeout("tcp", address, timeout) + if err != nil { + return false + } + _ = c.Close() + return true +} +func (r Registry) Candidates(project string, check Reachability, timeout time.Duration) ([]Herdr, error) { + p, ok := r.projects[project] + if !ok { + return nil, ErrUnknownProject + } + allowed := map[string]bool{} + for _, m := range p.MachineAffinity { + allowed[m] = true + } + out := []Herdr{} + for _, h := range r.herdrs { + if !allowed[h.MachineID] { + continue + } + addr := h.Address + if addr == "" { + addr = r.machines[h.MachineID].Address + } + if check == nil || check.Reachable(addr, timeout) { + out = append(out, h) + } + } + sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID }) + return out, nil +} diff --git a/internal/registry/registry_test.go b/internal/registry/registry_test.go new file mode 100644 index 0000000..f5611b9 --- /dev/null +++ b/internal/registry/registry_test.go @@ -0,0 +1,36 @@ +package registry + +import ( + "errors" + "testing" + "time" +) + +type reach map[string]bool + +func (r reach) Reachable(a string, _ time.Duration) bool { return r[a] } + +func TestNewValidatesTopologyAndResolvesHardAffinity(t *testing.T) { + r, err := New(Config{ + Projects: []Project{{ID: "work", MachineAffinity: []string{"pc"}}}, + Machines: []Machine{{ID: "server", Address: "server:1"}, {ID: "pc", Address: "pc:1"}}, + Herdrs: []Herdr{{ID: "offline", MachineID: "pc", Address: "off:1"}, {ID: "online", MachineID: "pc", Address: "on:1"}, {ID: "wrong", MachineID: "server", Address: "server:1"}}, + }) + if err != nil { + t.Fatal(err) + } + got, err := r.Candidates("work", reach{"off:1": false, "on:1": true, "server:1": true}, time.Second) + if err != nil || len(got) != 1 || got[0].ID != "online" { + t.Fatalf("candidates=%v err=%v", got, err) + } + if _, err = r.Candidates("missing", nil, time.Second); !errors.Is(err, ErrUnknownProject) { + t.Fatalf("err=%v", err) + } +} + +func TestNewRejectsBrokenReferences(t *testing.T) { + _, err := New(Config{Projects: []Project{{ID: "p", MachineAffinity: []string{"missing"}}}}) + if !errors.Is(err, ErrUnknownMachine) { + t.Fatalf("err=%v", err) + } +} diff --git a/progress.md b/progress.md index d249c13..75904f9 100644 --- a/progress.md +++ b/progress.md @@ -17,12 +17,12 @@ This is the implementation-oriented breakdown of the specification. It is a proj - Done: Gitea reflection for terminal task state, keyed by the task's stable external issue number. - Done: constant-time HMAC webhook authentication and injectable HTTP clients for testing. -3. **Projects and machine registry** — **not started** - - Project configuration - - Machine registry - - Herdr registry - - Reachability checks - - Hard machine affinity resolution +3. **Projects and machine registry** — **complete** + - Done: typed JSON project, machine, and herdr configuration with duplicate/reference validation. + - Done: machine-bound herdr registry with per-herdr capabilities, endpoint override, and concurrency configuration. + - Done: injectable reachability checks plus TCP reachability implementation. + - Done: hard project machine-affinity resolution; candidates are restricted to configured, reachable herdrs on allowed machines. + - Done: optional `ORCHESTRA_CONFIG` startup validation. 4. **Router and leases** — **partial groundwork** - Done: manual lease/release/complete/block endpoints and lease-expiry release. @@ -87,6 +87,10 @@ This is the implementation-oriented breakdown of the specification. It is a proj - `POST /v1/tasks/{id}/complete` - `POST /v1/tasks/{id}/block` +## Item 3 status + +Item 3 (projects and machine registry) is implemented in `internal/registry`. Static JSON configuration is loaded and validated, projects resolve only to their explicitly configured machines, and candidate herdrs are filtered by registration and injected reachability. Set `ORCHESTRA_CONFIG` to validate a configuration file at server startup. + ## Item 2 status Item 2 (provider layer) is implemented. `internal/provider` now includes `JSONLWatcher`, `Gitea.Poll`, `Gitea.WebhookHandler`, `Gitea.IngestWebhook`, and `Gitea.ReflectTask`. Gitea ingestion remains idempotent through the store's `(source, external_id)` key. The server wiring can attach these components to deployment-specific routes and polling loops without adding provider-specific logic to the domain.