// Package router assigns queued tasks to registered, reachable herdrs. package router import ( "encoding/json" "errors" "orchestra/internal/authz" "orchestra/internal/domain" "orchestra/internal/registry" "orchestra/internal/store" "sort" "strings" "time" ) type Availability interface{ Available(h registry.Herdr) bool } type AlwaysAvailable struct{} func (AlwaysAvailable) Available(registry.Herdr) bool { return true } // QuotaWindowLimits are the two independent caps the spec (§7.2) requires: // the subscription pool's 5-hour rolling window and its weekly window. They // are tracked and evaluated separately — a harness deep into its 5h window // but fine on the week, or vice versa, must still be excluded. type QuotaWindowLimits struct { FiveHour float64 Weekly float64 } const ( fiveHourWindow = 5 * time.Hour weeklyWindow = 7 * 24 * time.Hour // quotaConservativeFraction is the degrade-safe default from spec §7.2/§9 // item 1: since no quota pool is authoritative, treat 80% reported as // full rather than trusting the exact number. quotaConservativeFraction = 0.8 ) // QuotaAvailability applies the conservative 80% rule independently to the // 5-hour rolling window and the weekly window, per harness. Receipts are // additive across rotations; a cumulative session total must never replace // earlier rotations' receipts (spec §5.2.1) — summing native per-report // `consumed` deltas is what keeps this correct across rotation. type QuotaAvailability struct { Store *store.Store Limits map[string]QuotaWindowLimits Now func() time.Time } func (q QuotaAvailability) sumSince(harnessID string, since time.Time) float64 { var consumed float64 for _, e := range q.Store.Events(0) { if e.Type != "QuotaReported" || e.At.Before(since) { continue } var p struct { HarnessID string `json:"harness_id"` Consumed float64 `json:"consumed"` } if json.Unmarshal(e.Payload, &p) == nil && p.HarnessID == harnessID && p.Consumed >= 0 { consumed += p.Consumed } } return consumed } func (q QuotaAvailability) Available(h registry.Herdr) bool { if q.Store == nil { return false } limits, bounded := q.Limits[h.ID] if !bounded || (limits.FiveHour <= 0 && limits.Weekly <= 0) { return true } now := time.Now() if q.Now != nil { now = q.Now() } if limits.FiveHour > 0 && q.sumSince(h.ID, now.Add(-fiveHourWindow)) >= limits.FiveHour*quotaConservativeFraction { return false } if limits.Weekly > 0 && q.sumSince(h.ID, now.Add(-weeklyWindow)) >= limits.Weekly*quotaConservativeFraction { return false } return true } type RetryPolicy struct { MaxAttempts int Backoff time.Duration } type Router struct { Store *store.Store Registry registry.Registry Reachability registry.Reachability Availability Availability Timeout time.Duration Retry RetryPolicy Now func() time.Time backoff map[string]time.Time attempts map[string]int OnLease func(domain.Event) error } func (r *Router) init() { if r.Availability == nil { r.Availability = AlwaysAvailable{} } if r.Now == nil { r.Now = time.Now } if r.backoff == nil { r.backoff = map[string]time.Time{} } if r.attempts == nil { r.attempts = map[string]int{} } } // HandleEvent evaluates the sink after creation and after a lease is freed. func (r *Router) HandleEvent(e domain.Event) ([]domain.Event, error) { r.init() if e.Type != "TaskCreated" && e.Type != "TaskReleased" { return nil, nil } if e.Type == "TaskReleased" { // Rotation *is* TaskReleased (spec §5.3: "rotation = intra-task // lease transfer") — a task healthy enough to rotate repeatedly // must not be killed by the retry limit meant for genuine failures // (expiry, crash). Only a release without a valid handoff_ref // (expiry/crash) counts against MaxAttempts. var p struct { HandoffRef string `json:"handoff_ref"` } isRotation := json.Unmarshal(e.Payload, &p) == nil && p.HandoffRef != "" if !isRotation { r.attempts[e.TaskID]++ if r.Retry.Backoff > 0 { r.backoff[e.TaskID] = r.Now().Add(r.Retry.Backoff) } } } return r.AssignPending() } func (r *Router) AssignPending() ([]domain.Event, error) { r.init() if r.Store == nil { return nil, errors.New("router: store required") } var queued []domain.Task for _, t := range r.Store.Tasks() { if t.State == domain.StateQueued && !r.Now().Before(r.backoff[t.ID]) { queued = append(queued, t) } } sort.SliceStable(queued, func(i, j int) bool { return importance(queued[i], r.Now()).Before(importance(queued[j], r.Now())) }) var out []domain.Event for _, t := range queued { if r.Retry.MaxAttempts > 0 && r.attempts[t.ID] >= r.Retry.MaxAttempts { e, err := r.fail(t) if err != nil { return out, err } out = append(out, e) continue } cs, err := r.Registry.Candidates(t.Project, r.Reachability, r.Timeout) if err != nil { continue } for _, h := range cs { if !matches(t.Capability, h.Capabilities) || !r.Availability.Available(h) || occupied(r.Store, h.ID, h.Concurrency) { continue } e, err := r.Store.Lease(t.ID, h.ID, 30*time.Minute) if err != nil { continue } // attempts is the failure counter checked against MaxAttempts // above; it advances only on a non-rotation TaskReleased (see // HandleEvent), not here, so a task that leases and rotates // repeatedly is not double-counted toward the retry limit. out = append(out, e) if r.OnLease != nil { if err := r.OnLease(e); err != nil { return out, err } } break } } return out, nil } func matches(need, have []string) bool { set := map[string]bool{} for _, x := range have { set[strings.ToLower(x)] = true } for _, x := range need { if !set[strings.ToLower(x)] { return false } } return true } func occupied(s *store.Store, id string, limit int) bool { if limit <= 0 { return false } n := 0 for _, t := range s.Tasks() { if t.State == domain.StateLeased && t.Lease != nil && t.Lease.HarnessID == id { n++ } } return n >= limit } func importance(t domain.Task, now time.Time) time.Time { if t.Due != nil { return t.Due.Add(-time.Duration(t.InherentPriority) * time.Hour) } return now.Add(-time.Duration(t.InherentPriority) * time.Hour) } func (r *Router) fail(t domain.Task) (domain.Event, error) { b, _ := json.Marshal(map[string]any{"reason": "retry_limit", "attempts": r.attempts[t.ID]}) e := domain.Event{ID: domain.NewID(), Type: "TaskFailed", TaskID: t.ID, Version: t.Version + 1, Payload: b, Surface: string(authz.System)} return e, r.Store.Append(e) }