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