| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "context" |
| 6 | "encoding/json" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | "io" |
| 10 | "strings" |
| 11 | "sync" |
| 12 | "time" |
| 13 | |
| 14 | "reasonix/internal/event" |
| 15 | "reasonix/internal/evidence" |
| 16 | "reasonix/internal/jobs" |
| 17 | ) |
| 18 | |
| 19 | const ( |
| 20 | fleetMinTasks = 2 |
| 21 | fleetMaxTasks = 64 |
| 22 | ) |
| 23 | |
| 24 | // FleetTool dispatches multiple profile-aware sub-agent tasks in parallel |
| 25 | // under the session scheduler. Write tasks must predeclare non-overlapping |
| 26 | // write_paths; preflight failure starts nothing. |
| 27 | type FleetTool struct { |
| 28 | taskTool *TaskTool |
| 29 | } |
| 30 | |
| 31 | // NewFleetTool creates a fleet dispatcher that reuses TaskTool infrastructure. |
| 32 | func NewFleetTool(taskTool *TaskTool) *FleetTool { |
| 33 | return &FleetTool{taskTool: taskTool} |
| 34 | } |
| 35 | |
| 36 | func (*FleetTool) Name() string { return "fleet" } |
| 37 | |
| 38 | func (*FleetTool) Description() string { |
| 39 | return "Dispatch 2–64 sub-agent tasks in parallel and return bounded previews plus stable Subagent references for full-result retrieval from completed persisted children with read_subagent_result. Each item may select a profile, model, effort, tools, write_paths, or read_only. Multiple writers must declare non-overlapping write_paths; omitted write_paths claim the whole workspace, so two or more writers without paths fail preflight before any task starts. Independent failure is the default: one failure does not cancel others. Background mode returns a fleet job id collectable with wait." |
| 40 | } |
| 41 | |
| 42 | func (*FleetTool) Schema() json.RawMessage { |
| 43 | return json.RawMessage(`{ |
| 44 | "type":"object", |
| 45 | "properties":{ |
| 46 | "tasks":{ |
| 47 | "type":"array", |
| 48 | "description":"Array of 2–64 sub-tasks to run under the session scheduler.", |
| 49 | "minItems":2, |
| 50 | "maxItems":64, |
| 51 | "items":{ |
| 52 | "type":"object", |
| 53 | "properties":{ |
| 54 | "prompt":{"type":"string","description":"Task prompt for the sub-agent."}, |
| 55 | "description":{"type":"string","description":"Optional short label shown in the job list."}, |
| 56 | "profile":{"type":"string","description":"Optional runAs=subagent profile name."}, |
| 57 | "write_paths":{"type":"array","items":{"type":"string"},"description":"Write targets for this item. Parallel writers must declare non-overlapping paths. Omitting write_paths claims the whole workspace; multiple whole-workspace claims (or any path overlap) fail preflight and start nothing."}, |
| 58 | "read_only":{"type":"boolean","description":"Force the read-only registry even if the profile is writable."}, |
| 59 | "tools":{"type":"array","items":{"type":"string"},"description":"Optional tool whitelist (intersected with profile allowed-tools)."}, |
| 60 | "max_steps":{"type":"integer","description":"Optional max tool-call rounds.","minimum":1}, |
| 61 | "model":{"type":"string","description":"Optional model override."}, |
| 62 | "effort":{"type":"string","description":"Optional reasoning effort override."} |
| 63 | }, |
| 64 | "required":["prompt"] |
| 65 | } |
| 66 | }, |
| 67 | "run_in_background":{"type":"boolean","description":"Run the whole fleet asynchronously and return a job id collectable with wait. Items queue for concurrency/write slots inside the job."} |
| 68 | }, |
| 69 | "required":["tasks"] |
| 70 | }`) |
| 71 | } |
| 72 | |
| 73 | func (*FleetTool) ReadOnly() bool { return false } |
| 74 | |
| 75 | type fleetTaskItem struct { |
| 76 | Prompt string `json:"prompt"` |
| 77 | Description string `json:"description"` |
| 78 | Profile string `json:"profile"` |
| 79 | WritePaths []string `json:"write_paths"` |
| 80 | ReadOnly bool `json:"read_only"` |
| 81 | Tools []string `json:"tools"` |
| 82 | MaxSteps int `json:"max_steps"` |
| 83 | Model string `json:"model"` |
| 84 | Effort string `json:"effort"` |
| 85 | } |
| 86 | |
| 87 | type fleetItemStatus string |
| 88 | |
| 89 | const ( |
| 90 | fleetItemPending fleetItemStatus = "pending" |
| 91 | fleetItemCompleted fleetItemStatus = "completed" |
| 92 | fleetItemFailed fleetItemStatus = "failed" |
| 93 | fleetItemCancelled fleetItemStatus = "cancelled" |
| 94 | fleetItemSkipped fleetItemStatus = "skipped" |
| 95 | ) |
| 96 | |
| 97 | type fleetItemResult struct { |
| 98 | index int |
| 99 | status fleetItemStatus |
| 100 | profile string |
| 101 | output string |
| 102 | err error |
| 103 | ref string |
| 104 | } |
| 105 | |
| 106 | // fleetGroupTerminalPhase classifies a fleet group's single terminal status: |
| 107 | // cancellation/deadline wins, then any failed child, then any error |
| 108 | // (including validation failures), then completed. |
| 109 | func fleetGroupTerminalPhase(ctx context.Context, err error, results []fleetItemResult) subagentProgressPhase { |
| 110 | if ctx.Err() != nil { |
| 111 | return subagentPhaseCancelled |
| 112 | } |
| 113 | for _, r := range results { |
| 114 | if r.status == fleetItemFailed { |
| 115 | return subagentPhaseFailed |
| 116 | } |
| 117 | } |
| 118 | if err != nil { |
| 119 | return subagentPhaseFailed |
| 120 | } |
| 121 | return subagentPhaseCompleted |
| 122 | } |
| 123 | |
| 124 | func (f *FleetTool) Execute(ctx context.Context, args json.RawMessage) (result string, err error) { |
| 125 | if f == nil || f.taskTool == nil { |
| 126 | return "", fmt.Errorf("fleet is not configured") |
| 127 | } |
| 128 | // Group lifecycle: the group card's terminal is an explicit event from |
| 129 | // the tool (running once children start, exactly one terminal at the |
| 130 | // end) so frontends never infer group completion from the children they |
| 131 | // happen to have observed. Validation failures emit a failed terminal; |
| 132 | // once runFleet starts it owns the lifecycle (the background job runs |
| 133 | // runFleet inside the job, after this function has returned). |
| 134 | groupParentID, groupSink, _, ok := CallContext(ctx) |
| 135 | if !ok || groupSink == nil { |
| 136 | groupParentID = "fleet" |
| 137 | groupSink = event.Discard |
| 138 | } |
| 139 | // The merger emits already-namespaced group/child IDs, so it must use the |
| 140 | // raw call sink. A nested subSink would prefix the group ID a second time |
| 141 | // (group/group), leaving the frontend unable to match its lifecycle card. |
| 142 | merger := newSubagentProgressMerger(realProgressClock{}, groupSink, groupParentID) |
| 143 | lifecycleHandoff := false |
| 144 | mergerCloseHandoff := false |
| 145 | defer func() { |
| 146 | if !mergerCloseHandoff { |
| 147 | merger.Close() |
| 148 | } |
| 149 | }() |
| 150 | defer func() { |
| 151 | if lifecycleHandoff { |
| 152 | return |
| 153 | } |
| 154 | merger.directStatus(groupParentID, fleetGroupTerminalPhase(ctx, err, nil)) |
| 155 | }() |
| 156 | ctx = withSubagentProgressMerger(ctx, merger) |
| 157 | |
| 158 | var params struct { |
| 159 | Tasks []fleetTaskItem `json:"tasks"` |
| 160 | RunInBackground bool `json:"run_in_background"` |
| 161 | } |
| 162 | dec := json.NewDecoder(bytes.NewReader(args)) |
| 163 | dec.DisallowUnknownFields() |
| 164 | if err := dec.Decode(¶ms); err != nil { |
| 165 | return "", fmt.Errorf("invalid args: %w", err) |
| 166 | } |
| 167 | if n := len(params.Tasks); n < fleetMinTasks || n > fleetMaxTasks { |
| 168 | return "", fmt.Errorf("fleet requires between %d and %d tasks (got %d)", fleetMinTasks, fleetMaxTasks, n) |
| 169 | } |
| 170 | |
| 171 | specs := make([]ProfileExecSpec, len(params.Tasks)) |
| 172 | // Keep one claim slot per original task so preflight errors report the |
| 173 | // caller-visible task numbers even when read-only items are interleaved. |
| 174 | claims := make([]WritePathSet, len(params.Tasks)) |
| 175 | for i, item := range params.Tasks { |
| 176 | if strings.TrimSpace(item.Prompt) == "" { |
| 177 | return "", fmt.Errorf("task %d: prompt is required", i+1) |
| 178 | } |
| 179 | // Fleet writers without write_paths claim the whole workspace so the |
| 180 | // preflight can detect multi-writer collisions before anything starts. |
| 181 | forceBackgroundClaim := !item.ReadOnly |
| 182 | spec, err := f.taskTool.buildTaskSpec(ctx, item.Prompt, item.Description, item.Profile, item.WritePaths, item.Tools, item.MaxSteps, item.Model, item.Effort, "", "", false, item.ReadOnly) |
| 183 | if err != nil { |
| 184 | return "", fmt.Errorf("task %d: %w", i+1, err) |
| 185 | } |
| 186 | if forceBackgroundClaim && !spec.ReadOnly && spec.WritePaths.Empty() { |
| 187 | whole, werr := WholeWorkspaceWriteClaim(f.taskTool.workspaceRoot) |
| 188 | if werr != nil { |
| 189 | return "", fmt.Errorf("task %d: %w", i+1, werr) |
| 190 | } |
| 191 | spec.WritePaths = whole |
| 192 | } |
| 193 | spec.Nested = SubagentDepth(ctx) > 0 |
| 194 | spec.RunInBackground = false // fleet owns backgrounding |
| 195 | if spec.Description == "" { |
| 196 | spec.Description = fmt.Sprintf("fleet-%d", i+1) |
| 197 | } |
| 198 | specs[i] = spec |
| 199 | if !spec.ReadOnly { |
| 200 | claims[i] = spec.WritePaths |
| 201 | } |
| 202 | } |
| 203 | if err := ValidateNonOverlappingWriteClaims(claims); err != nil { |
| 204 | return "", fmt.Errorf("fleet preflight: %w", err) |
| 205 | } |
| 206 | |
| 207 | if params.RunInBackground { |
| 208 | for i := range specs { |
| 209 | specs[i].BackgroundWriter = !specs[i].ReadOnly |
| 210 | } |
| 211 | jm, ok := jobs.FromContext(ctx) |
| 212 | if !ok { |
| 213 | return "", fmt.Errorf("background execution is not available in this context") |
| 214 | } |
| 215 | parentID := groupParentID |
| 216 | parentSession := ParentSession(ctx) |
| 217 | label := fmt.Sprintf("fleet(%d)", len(specs)) |
| 218 | backgroundEvidence := evidence.NewLedger() |
| 219 | writerID := fmt.Sprintf("background-fleet:%s:%d", parentID, time.Now().UnixNano()) |
| 220 | writerRegistered := false |
| 221 | observer := f.taskTool.mutationObserver |
| 222 | if observer != nil { |
| 223 | hasWriter := false |
| 224 | for i := range specs { |
| 225 | if specs[i].BackgroundWriter { |
| 226 | hasWriter = true |
| 227 | break |
| 228 | } |
| 229 | } |
| 230 | if hasWriter { |
| 231 | if err := observer.RegisterWriter(writerID, "background_fleet", observer.OwnershipTurn()); err != nil { |
| 232 | return "", err |
| 233 | } |
| 234 | writerRegistered = true |
| 235 | } |
| 236 | } |
| 237 | job := jm.StartForSession(jobs.SessionFromContext(ctx), "fleet", label, func(jobCtx context.Context, _ io.Writer) (string, error) { |
| 238 | // Execute returns as soon as the job is registered, so the job owns |
| 239 | // the handed-off merger until every child preview and terminal has |
| 240 | // flushed. Closing it in Execute would strand child cards at running. |
| 241 | defer merger.Close() |
| 242 | if writerRegistered { |
| 243 | defer observer.UnregisterWriter(writerID) |
| 244 | } |
| 245 | jobCtx = WithParentSession(jobCtx, parentSession) |
| 246 | jobCtx = evidence.WithLedger(jobCtx, backgroundEvidence) |
| 247 | defer func() { jobs.PublishEvidence(jobCtx, backgroundEvidence.Summary()) }() |
| 248 | // The job shares the Execute-level merger so the group lifecycle |
| 249 | // events and the child previews ride the same pacing budget. |
| 250 | jobCtx = withSubagentProgressMerger(jobCtx, merger) |
| 251 | return f.runFleet(jobCtx, groupSink, specs, parentID) |
| 252 | }) |
| 253 | // runFleet (inside the job) owns the terminal and merger close from |
| 254 | // here on. Foreground runFleet hands off only the terminal; Execute |
| 255 | // still closes the merger after the synchronous call returns. |
| 256 | lifecycleHandoff = true |
| 257 | mergerCloseHandoff = true |
| 258 | return fmt.Sprintf("Started background fleet %q (%s). Collect results with wait; you will be notified when it finishes.", job.ID, label), nil |
| 259 | } |
| 260 | |
| 261 | lifecycleHandoff = true |
| 262 | return f.runFleet(ctx, groupSink, specs, groupParentID) |
| 263 | } |
| 264 | |
| 265 | func (f *FleetTool) runFleet(ctx context.Context, sink event.Sink, specs []ProfileExecSpec, groupParentID string) (result string, err error) { |
| 266 | if sink == nil { |
| 267 | sink = event.Discard |
| 268 | } |
| 269 | // Child IDs are namespaced exactly once under the group call. Background |
| 270 | // jobs no longer carry the original call context, so groupParentID is the |
| 271 | // authoritative identity there; direct callers fall back to CallContext. |
| 272 | parentID := strings.TrimSpace(groupParentID) |
| 273 | if parentID == "" { |
| 274 | var ok bool |
| 275 | parentID, _, _, ok = CallContext(ctx) |
| 276 | if !ok || parentID == "" { |
| 277 | parentID = "fleet" |
| 278 | } |
| 279 | } |
| 280 | groupParentID = parentID |
| 281 | // The Execute-level merger (or a fallback for direct callers) paces the |
| 282 | // group; runFleet owns the lifecycle once it starts: running up front |
| 283 | // and exactly one terminal after every child settles. |
| 284 | merger := subagentProgressMergerFromContext(ctx) |
| 285 | ownsMerger := false |
| 286 | if merger == nil { |
| 287 | merger = newSubagentProgressMerger(realProgressClock{}, sink, groupParentID) |
| 288 | ownsMerger = true |
| 289 | ctx = withSubagentProgressMerger(ctx, merger) |
| 290 | } |
| 291 | if ownsMerger { |
| 292 | defer merger.Close() |
| 293 | } |
| 294 | merger.directStatus(groupParentID, subagentPhaseRunning) |
| 295 | var results []fleetItemResult |
| 296 | defer func() { |
| 297 | merger.directStatus(groupParentID, fleetGroupTerminalPhase(ctx, err, results)) |
| 298 | }() |
| 299 | |
| 300 | n := len(specs) |
| 301 | results = make([]fleetItemResult, n) |
| 302 | for i := range results { |
| 303 | results[i] = fleetItemResult{index: i, status: fleetItemPending, profile: specs[i].Profile} |
| 304 | } |
| 305 | |
| 306 | var wg sync.WaitGroup |
| 307 | doneCh := make(chan fleetItemResult, n) |
| 308 | |
| 309 | startOne := func(idx int) { |
| 310 | spec := specs[idx] |
| 311 | label := spec.Description |
| 312 | subID := fmt.Sprintf("%s/fleet-%d", parentID, idx+1) |
| 313 | dispatchArgs, _ := json.Marshal(map[string]any{ |
| 314 | "prompt": spec.Prompt, |
| 315 | "description": label, |
| 316 | "profile": spec.Profile, |
| 317 | }) |
| 318 | sink.Emit(event.Event{ |
| 319 | Kind: event.ToolDispatch, |
| 320 | Tool: event.Tool{ |
| 321 | ID: subID, ParentID: parentID, Name: "task", |
| 322 | Args: string(dispatchArgs), ReadOnly: spec.ReadOnly, |
| 323 | }, |
| 324 | }) |
| 325 | |
| 326 | wg.Add(1) |
| 327 | go func() { |
| 328 | defer wg.Done() |
| 329 | // Each fleet item runs as its own task-shaped execution so |
| 330 | // transcripts, evidence, and scheduler claims stay independent. |
| 331 | itemCtx := withCallContext(ctx, subID, subSinkFor(subID, sink), nil, false) |
| 332 | out, err := f.taskTool.RunProfileSpec(itemCtx, spec) |
| 333 | answer, ref := splitSubagentRunResult(out) |
| 334 | res := fleetItemResult{index: idx, profile: spec.Profile, output: answer, ref: ref, err: err} |
| 335 | if err == nil { |
| 336 | res.status = fleetItemCompleted |
| 337 | sink.Emit(event.Event{ |
| 338 | Kind: event.ToolResult, |
| 339 | Tool: event.Tool{ID: subID, ParentID: parentID, Name: "task", Output: out}, |
| 340 | }) |
| 341 | } else { |
| 342 | if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { |
| 343 | res.status = fleetItemCancelled |
| 344 | } else { |
| 345 | res.status = fleetItemFailed |
| 346 | } |
| 347 | sink.Emit(event.Event{ |
| 348 | Kind: event.ToolResult, |
| 349 | Tool: event.Tool{ID: subID, ParentID: parentID, Name: "task", Err: err.Error()}, |
| 350 | }) |
| 351 | } |
| 352 | doneCh <- res |
| 353 | }() |
| 354 | } |
| 355 | |
| 356 | // Mark only genuinely unstarted items skipped. Started items always publish |
| 357 | // a terminal result, including after cancellation, so partial writer work is |
| 358 | // never misreported as a task that did not run. |
| 359 | started := 0 |
| 360 | for i := range specs { |
| 361 | if ctx.Err() != nil { |
| 362 | for j := i; j < n; j++ { |
| 363 | results[j].status = fleetItemSkipped |
| 364 | results[j].err = ctx.Err() |
| 365 | } |
| 366 | break |
| 367 | } |
| 368 | startOne(i) |
| 369 | started++ |
| 370 | } |
| 371 | |
| 372 | completed := 0 |
| 373 | cancelled := false |
| 374 | for completed < started && !cancelled { |
| 375 | select { |
| 376 | case r := <-doneCh: |
| 377 | results[r.index] = r |
| 378 | completed++ |
| 379 | case <-ctx.Done(): |
| 380 | cancelled = true |
| 381 | } |
| 382 | } |
| 383 | // doneCh is buffered for every item, so workers can always publish their |
| 384 | // terminal result while this goroutine waits. Once they stop, drain exactly |
| 385 | // the outstanding started items and preserve their real completed/cancelled |
| 386 | // status instead of replacing it with skipped. |
| 387 | wg.Wait() |
| 388 | for completed < started { |
| 389 | r := <-doneCh |
| 390 | results[r.index] = r |
| 391 | completed++ |
| 392 | } |
| 393 | for _, r := range results { |
| 394 | if r.status == fleetItemCancelled || r.status == fleetItemSkipped { |
| 395 | cancelled = true |
| 396 | break |
| 397 | } |
| 398 | } |
| 399 | if cancelled { |
| 400 | err := ctx.Err() |
| 401 | if err == nil { |
| 402 | err = context.Canceled |
| 403 | } |
| 404 | return formatFleetAggregate(results, true), err |
| 405 | } |
| 406 | return formatFleetAggregate(results, false), nil |
| 407 | } |
| 408 | |
| 409 | func formatFleetAggregate(results []fleetItemResult, cancelled bool) string { |
| 410 | n := len(results) |
| 411 | var prefix string |
| 412 | if cancelled { |
| 413 | completed := 0 |
| 414 | for _, r := range results { |
| 415 | if r.status == fleetItemCompleted { |
| 416 | completed++ |
| 417 | } |
| 418 | } |
| 419 | prefix = fmt.Sprintf("Cancelled fleet after completing %d of %d tasks:\n", completed, n) |
| 420 | } else { |
| 421 | prefix = fmt.Sprintf("Completed fleet of %d tasks:\n", n) |
| 422 | } |
| 423 | items := make([]subagentAggregateItem, 0, n) |
| 424 | for i, r := range results { |
| 425 | header := fmt.Sprintf("── task-%d", i+1) |
| 426 | if r.profile != "" { |
| 427 | header += " profile=" + boundedInline(r.profile, 80) |
| 428 | } |
| 429 | header += " ──\n" |
| 430 | item := subagentAggregateItem{header: header, ref: r.ref} |
| 431 | switch r.status { |
| 432 | case fleetItemCompleted: |
| 433 | item.status = "status: completed\n" |
| 434 | item.answer = strings.TrimSpace(r.output) |
| 435 | case fleetItemFailed: |
| 436 | item.status = "status: failed\n" |
| 437 | if r.err != nil { |
| 438 | item.detail = fmt.Sprintf("[FAILED] %s\n", boundedInline(r.err.Error(), 256)) |
| 439 | } |
| 440 | case fleetItemCancelled: |
| 441 | item.status = "status: cancelled\n" |
| 442 | if r.err != nil { |
| 443 | item.detail = fmt.Sprintf("[CANCELLED] %s\n", boundedInline(r.err.Error(), 256)) |
| 444 | } |
| 445 | case fleetItemSkipped: |
| 446 | item.status = "status: skipped\n" |
| 447 | if r.err != nil { |
| 448 | item.detail = fmt.Sprintf("[SKIPPED] %s\n", boundedInline(r.err.Error(), 256)) |
| 449 | } |
| 450 | default: |
| 451 | item.status = "status: pending\n" |
| 452 | } |
| 453 | items = append(items, item) |
| 454 | } |
| 455 | return formatBoundedSubagentAggregate(prefix, items) |
| 456 | } |
| 457 |