返回 DeepSeek-Reasonix
fleet.go
根目录 / internal / agent / fleet.go
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(&params); 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
457 lines GO