返回 DeepSeek-Reasonix
turn_orchestrator.go
根目录 / internal / control / turn_orchestrator.go
1 package control
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "time"
9
10 "reasonix/internal/agent"
11 "reasonix/internal/event"
12 "reasonix/internal/jobs"
13 "reasonix/internal/provider"
14 "reasonix/internal/skill"
15 "reasonix/internal/tool"
16 )
17
18 // turnOrchestrator owns foreground turn execution while Controller keeps the
19 // public ports, run-state guard, and session-scoped dependencies.
20 type turnOrchestrator struct {
21 c *Controller
22 }
23
24 type orchestratedTurn struct {
25 input string
26 raw string
27 imageRefs string
28 userImages []string
29 imageCandidates []string
30 imagesResolved bool
31 display string
32 editedOriginal string
33 synthetic bool
34 goalRound *goalRoundReservation
35 }
36
37 func newTurnOrchestrator(c *Controller) *turnOrchestrator {
38 return &turnOrchestrator{c: c}
39 }
40
41 func (o *turnOrchestrator) runTurnWithRawDisplay(ctx context.Context, input, raw, display string) error {
42 return o.runOrchestratedTurn(ctx, orchestratedTurn{input: input, raw: raw, display: display})
43 }
44
45 func (o *turnOrchestrator) runTurnWithImageRefsRawDisplay(ctx context.Context, input, raw, imageRefs, display string) error {
46 return o.runOrchestratedTurn(ctx, orchestratedTurn{input: input, raw: raw, imageRefs: imageRefs, display: display})
47 }
48
49 func (o *turnOrchestrator) runSyntheticTurnWithRawDisplay(ctx context.Context, input, raw, display string) error {
50 return o.runOrchestratedTurn(ctx, orchestratedTurn{input: input, raw: raw, display: display, synthetic: true})
51 }
52
53 func (o *turnOrchestrator) runComposedSyntheticTurn(ctx context.Context, text string) error {
54 c := o.c
55 ctx = agent.WithRawUserInput(ctx, text)
56 ctx = withTurnInputOrigin(ctx, true)
57 // ctx may carry the identity already spent on the turn's own user message.
58 if c.executor != nil {
59 ctx = agent.WithUserMessageIdentity(ctx, c.executor.Session(), agent.NewMessageID())
60 }
61 ctx = c.withTurnContext(ctx, false)
62 ctx = c.withPlannerTurnMetadata(ctx, text, true, c.messageCount())
63 return c.runModelTurn(ctx, c.ComposeSynthetic(text))
64 }
65
66 // runSubagentSkillGoalLoop executes a slash-invoked runAs=subagent skill as a
67 // real isolated child turn, then lets an active goal continue just as an inline
68 // skill turn did before.
69 func (o *turnOrchestrator) runSubagentSkillGoalLoop(ctx context.Context, sk skill.Skill, task, raw, display string, runner skill.SubagentRunner, planMode bool) error {
70 return o.runSubagentSkillTurnsGoalLoop(ctx, []skill.Skill{sk}, task, raw, display, runner, planMode)
71 }
72
73 func (o *turnOrchestrator) runSubagentSkillTurnsGoalLoop(ctx context.Context, skills []skill.Skill, task, raw, display string, runner skill.SubagentRunner, planMode bool, frozen ...[]string) error {
74 var userImages, imageCandidates []string
75 if len(frozen) > 0 {
76 imageCandidates = append([]string(nil), frozen[0]...)
77 if o.c.imageInputEnabled() {
78 userImages = append([]string(nil), imageCandidates...)
79 }
80 } else {
81 userImages, imageCandidates = o.c.resolveTurnImages(raw)
82 }
83 ctx = agent.WithSubagentImageCandidates(ctx, imageCandidates)
84 ctx = o.c.withPreparedTurnImages(ctx)
85 return o.runSubagentSkillTurns(ctx, skills, task, raw, display, runner, planMode, userImages, imageCandidates)
86 }
87
88 // runSubagentSkillTurns records the composed user task and distilled child
89 // answers only. Child reasoning and tool chatter stay out of the
90 // provider-visible parent context while their UI events nest under synthetic
91 // top-level run_skill cards.
92 func (o *turnOrchestrator) runSubagentSkillTurns(ctx context.Context, skills []skill.Skill, task, raw, display string, runner skill.SubagentRunner, planMode bool, images, imageCandidates []string) (err error) {
93 c := o.c
94 turnStartedAt := time.Now()
95 c.maybeSessionStart(ctx)
96 parentSession := c.parentSessionID()
97 ctx = agent.WithParentSession(ctx, parentSession)
98 ctx = jobs.WithSession(ctx, parentSession)
99 ctx = agent.WithUserImages(ctx, images)
100 ctx = agent.WithSubagentImageCandidates(ctx, imageCandidates)
101 ctx = agent.WithResponseLanguagePreference(ctx, c.responseLanguage)
102 ctx = agent.WithReasoningLanguagePreference(ctx, c.reasoningLanguage)
103 ctx = c.withTurnContext(ctx, true)
104
105 input := c.compose(task, raw, true)
106 startMessages := c.messageCount()
107 var marker agent.InFlightTurnMeta
108 defer func() { c.finishInFlightTurn(startMessages, marker) }()
109 defer c.recordDisplayForNewUser(startMessages, display)
110 // The checkpoint prompt labels the turn in the rewind picker (and is
111 // prefilled into the composer after a conversation rewind), so it must be
112 // the user's own text — never the composed provider input with its
113 // transient <response-language>/<reasoning-language>/memory/hook blocks.
114 c.beginCheckpoint(ctx, firstNonEmpty(raw, task))
115 if c.guardianSess != nil {
116 c.guardianSess.ResetTurn()
117 }
118 if c.hooks.Enabled() {
119 c.mu.Lock()
120 c.turn++
121 turn := c.turn
122 c.mu.Unlock()
123 if block, _ := c.hooks.PromptSubmit(ctx, input, turn); block {
124 return nil
125 }
126 defer func() { c.hooks.StopResult(context.Background(), lastAssistantText(c.History()), turn, err) }()
127 }
128
129 marker = c.markInFlightTurn(startMessages, true)
130 c.sink.Emit(event.Event{Kind: event.TurnStarted})
131 if c.executor == nil {
132 return fmt.Errorf("subagent slash invocation requires an active session")
133 }
134 message := persistedUserTurn(input, firstNonEmpty(raw, task), images, time.Now().UnixMilli())
135 if prepared, ok := ctx.Value(preparedImageReferencesContextKey{}).(preparedImageReferences); ok && len(prepared.inputs) > 0 {
136 message.Images = nil
137 message.ImageInputs = prepared.inputs
138 }
139 if _, err := c.executor.AppendTurnContextAndUserChecked(ctx, message); err != nil {
140 return err
141 }
142
143 for _, sk := range skills {
144 sk = c.skills.prepare(sk)
145 callID := fmt.Sprintf("slash-skill-%d", c.slashSkillSeq.Add(1))
146 args, _ := json.Marshal(map[string]string{"name": sk.Name, "arguments": task})
147 toolEvent := event.Tool{
148 ID: callID,
149 Name: "run_skill",
150 Args: string(args),
151 ReadOnly: sk.ReadOnly,
152 }
153 if c.skillProfile != nil {
154 toolEvent.Profile = c.skillProfile(sk)
155 }
156 if err := event.EmitChecked(c.sink, event.Event{Kind: event.ToolDispatch, Tool: toolEvent}); err != nil {
157 return fmt.Errorf("persist skill dispatch: %w", err)
158 }
159 runCtx := agent.WithToolCallContext(ctx, callID, c.sink, c, planMode)
160 runCtx = agent.WithSubagentDepth(runCtx, 0)
161 answer, err := runner(runCtx, sk, input, skill.SubagentRunOptions{HostInitiated: true})
162 if err != nil {
163 toolEvent.Err = err.Error()
164 c.sink.Emit(event.Event{Kind: event.ToolResult, Tool: toolEvent})
165 return err
166 }
167 answer = tool.GuardSubagentHostDecisionText(answer)
168 toolEvent.Output = answer
169 c.sink.Emit(event.Event{Kind: event.ToolResult, Tool: toolEvent})
170 workDurationMs := max(int64(1), time.Since(turnStartedAt).Milliseconds())
171 messageID := agent.NewMessageID()
172 assistant := provider.Message{ID: messageID, Role: provider.RoleAssistant, Content: answer, WorkDurationMs: workDurationMs}
173 if err := c.RecordSessionMessages(ctx, "orchestrated-assistant", []provider.Message{assistant}); err != nil {
174 return err
175 }
176 c.executor.Session().Add(assistant)
177 display := agent.DisplayAssistantText(answer)
178 c.sink.Emit(event.Event{Kind: event.Text, MessageID: messageID, Text: display})
179 c.sink.Emit(event.Event{Kind: event.Message, MessageID: messageID, Text: display})
180 }
181
182 return nil
183 }
184
185 func (o *turnOrchestrator) runOrchestratedTurn(ctx context.Context, turn orchestratedTurn) (err error) {
186 c := o.c
187 c.maybeSessionStart(ctx)
188 parentSession := c.parentSessionID()
189 ctx = agent.WithParentSession(ctx, parentSession)
190 ctx = jobs.WithSession(ctx, parentSession)
191 userImages, imageCandidates := c.imagesForOrchestratedTurn(ctx, turn)
192 ctx = agent.WithUserImages(ctx, userImages)
193 ctx = agent.WithSubagentImageCandidates(ctx, imageCandidates)
194 ctx = agent.WithRawUserInput(ctx, turn.raw)
195 ctx = c.withPreparedTurnImages(ctx)
196 ctx = withTurnInputOrigin(ctx, turn.synthetic)
197 userMessageID := agent.NewMessageID()
198 if _, turnID, active := c.currentTurnToken(); active {
199 if receipt, ok := c.submissionForTurn(turnID); ok {
200 userMessageID = receipt.MessageID
201 }
202 }
203 if c.executor != nil {
204 ctx = agent.WithUserMessageIdentity(ctx, c.executor.Session(), userMessageID)
205 }
206 var input string
207 if turn.goalRound != nil {
208 input = c.ComposeSynthetic(turn.input)
209 } else {
210 input = c.compose(turn.input, turn.raw, !turn.synthetic)
211 }
212 // input.receive: the composed text crosses the extension chain before it
213 // enters the session (checkpoint, hooks, and the model all see the final
214 // text). A block ruling aborts the turn with the redacted reason surfaced,
215 // mirroring the PromptSubmit hook's abort path; a required-class extension
216 // failure fails the turn.
217 input, blocked, interceptErr := c.interceptInputReceive(ctx, input)
218 if interceptErr != nil {
219 return interceptErr
220 }
221 if blocked {
222 return nil
223 }
224 startMessages := c.messageCount()
225 fallback := persistedUserTurn(input, turn.raw, userImages, time.Now().UnixMilli())
226 fallback.ID = userMessageID
227 c.noteTerminationBoundary(fallback, !turn.synthetic)
228 var marker agent.InFlightTurnMeta
229 defer func() { c.finishInFlightTurn(startMessages, marker) }()
230 defer c.recordDisplayForNewUser(startMessages, turn.display)
231 if turn.editedOriginal != "" {
232 defer c.markEditedForNewUser(startMessages, turn.editedOriginal)
233 }
234 // Open a checkpoint only for visible user turns before the user message is
235 // appended, so the recorded message boundary precedes it and pre-edit
236 // snapshots land here. Synthetic continuations stay attached to the visible
237 // turn that spawned them; otherwise hidden user-role messages would advance
238 // backend checkpoint turns without a matching frontend turn. The label is
239 // the user's own text (raw, falling back to the expanded input) — the
240 // composed provider input carries transient prefab blocks that must never
241 // surface in the rewind picker or be prefilled into the composer.
242 if !turn.synthetic {
243 c.beginCheckpoint(ctx, firstNonEmpty(turn.raw, turn.input))
244 }
245 if c.guardianSess != nil {
246 c.guardianSess.ResetTurn()
247 }
248 // UserPromptSubmit / Stop hooks bracket the whole turn (incl. the plan
249 // research + approved-execution sub-turns below): a gating UserPromptSubmit
250 // aborts before any model call; Stop fires once when the turn returns.
251 if c.hooks.Enabled() {
252 c.mu.Lock()
253 c.turn++
254 turn := c.turn
255 c.mu.Unlock()
256 if block, _ := c.hooks.PromptSubmit(ctx, input, turn); block {
257 return nil // the hook's notify callback already surfaced the reason
258 }
259 defer func() { c.hooks.StopResult(context.Background(), lastAssistantText(c.History()), turn, err) }()
260 }
261 marker = c.markInFlightTurn(startMessages, !turn.synthetic)
262 ctx = c.withTurnContext(ctx, !turn.synthetic)
263 if turn.goalRound != nil {
264 if authority, ok := c.goalAuthorityForRound(turn.goalRound); ok {
265 ctx = tool.WithGoalLifecycle(ctx, c, authority)
266 }
267 } else if !turn.synthetic {
268 if authority, ok := c.directHumanGoalAuthority(); ok {
269 ctx = tool.WithGoalLifecycle(ctx, c, authority)
270 }
271 }
272 ctx = c.withPlannerTurnMetadata(ctx, turn.raw, turn.synthetic, startMessages)
273 modelInput := input
274 if !turn.synthetic {
275 modelInput = c.withCapabilityRoute(ctx, input, turn.raw)
276 }
277 modelInput, ctx, err = c.prepareVisionTurn(ctx, modelInput, imageCandidates)
278 if err != nil {
279 return err
280 }
281 err = c.runModelTurn(ctx, modelInput)
282 if err != nil {
283 fallback := persistedUserTurn(input, turn.raw, userImages, time.Now().UnixMilli())
284 fallback.ID = userMessageID
285 // When the user explicitly cancels, keep the real prompt and any fully
286 // paired tool work. Partial reasoning/output remains durable for display
287 // but is marked local-only, and a bounded recovery summary is folded into
288 // the next real user turn (#5499, #6680).
289 if errors.Is(err, context.Canceled) && c.CancelRequested() {
290 if turn.synthetic {
291 c.stripInterruptedSyntheticTurnMessagesAfter(startMessages)
292 } else {
293 c.stripCancelledVisibleTurnMessagesAfterWithFallback(startMessages, fallback)
294 }
295 } else if !turn.synthetic && c.hasInterruptedDisplayAfter(startMessages, fallback) {
296 // Provider/API failures use the same safe recovery path as an explicit
297 // stop once the agent has recorded a partial stream. Completed tool
298 // pairs survive; unsafe stream fragments stay local-only.
299 c.stripCancelledVisibleTurnMessagesAfterWithFallback(startMessages, fallback)
300 }
301 return err
302 }
303 return o.executeApprovedPlan(ctx)
304 }
305
306 func (o *turnOrchestrator) executeApprovedPlan(ctx context.Context) error {
307 c := o.c
308 c.mu.Lock()
309 plan := c.sessionSettings.planMode
310 c.mu.Unlock()
311 if !plan {
312 return nil
313 }
314 proposal := lastAssistantText(c.History())
315 if proposal == "" {
316 return nil // no substantive proposal to gate
317 }
318 // The plan is already visible as the assistant's answer, so the request
319 // carries no subject — it's purely the gate.
320 allow, _, err := c.requestApproval(ctx, planApprovalTool, "", nil)
321 if err != nil {
322 return err
323 }
324 if !allow {
325 // The host decides whether denial means "revise and keep planning" or
326 // "exit without executing" by leaving plan mode on or switching it off.
327 return nil
328 }
329 c.SetPlanMode(false)
330 execStart := c.sessionMessageCount()
331 // The plan is the go-ahead: don't re-prompt for each write of the approved
332 // work. Auto-approve writers for the duration of this execution turn only; a
333 // later turn (even "continue") falls back to the normal per-tool approval.
334 c.approval.setPlanAutoApprove(true)
335 defer c.approval.setPlanAutoApprove(false)
336 err = func() error {
337 marker := c.markInFlightTurn(execStart, false)
338 defer c.finishInFlightTurn(execStart, marker)
339 return o.runComposedSyntheticTurn(ctx, planApprovedMessage)
340 }()
341 if err != nil {
342 if errors.Is(err, context.Canceled) && c.CancelRequested() {
343 c.stripInterruptedSyntheticTurnMessagesAfter(execStart)
344 }
345 return err
346 }
347 return nil
348 }
349
350 func (o *turnOrchestrator) runGoalLoopWithRawDisplay(ctx context.Context, input, raw, display string) error {
351 return o.runGoalLoopWithImageRefsRawDisplay(ctx, input, raw, "", display)
352 }
353
354 func (o *turnOrchestrator) runGoalLoopWithImageRefsRawDisplay(ctx context.Context, input, raw, imageRefs, display string) error {
355 turn := o.c.prepareOrchestratedTurnImages(orchestratedTurn{input: input, raw: raw, imageRefs: imageRefs, display: display})
356 return o.runGoalLoopWithPreparedTurn(ctx, turn)
357 }
358
359 func (o *turnOrchestrator) runGoalLoopWithFrozenImagesRawDisplay(ctx context.Context, input, raw, display string, images []string) error {
360 turn := orchestratedTurn{
361 input: input,
362 raw: raw,
363 display: display,
364 imageCandidates: append([]string(nil), images...),
365 imagesResolved: true,
366 }
367 if o.c.imageInputEnabled() {
368 turn.userImages = append([]string(nil), images...)
369 }
370 return o.runGoalLoopWithPreparedTurn(ctx, turn)
371 }
372
373 func (o *turnOrchestrator) runGoalLoopWithPreparedTurn(ctx context.Context, turn orchestratedTurn) error {
374 // Every accepted input is exactly one top-level turn. Automatic Goal work is
375 // owned exclusively by goalRoundDriver after the runtime becomes idle.
376 ctx = agent.WithSubagentImageCandidates(ctx, turn.imageCandidates)
377 return o.runOrchestratedTurn(ctx, turn)
378 }
379
380 func (o *turnOrchestrator) runEditedGoalLoopWithRawDisplay(ctx context.Context, input, raw, display, original string) error {
381 return o.runEditedGoalLoopWithImageRefsRawDisplay(ctx, input, raw, "", display, original)
382 }
383
384 func (o *turnOrchestrator) runEditedGoalLoopWithImageRefsRawDisplay(ctx context.Context, input, raw, imageRefs, display, original string) error {
385 turn := o.c.prepareOrchestratedTurnImages(orchestratedTurn{
386 input: input, raw: raw, imageRefs: imageRefs, display: display, editedOriginal: original,
387 })
388 ctx = agent.WithSubagentImageCandidates(ctx, turn.imageCandidates)
389 return o.runOrchestratedTurn(ctx, turn)
390 }
391
392 func (o *turnOrchestrator) runEditedGoalLoopWithFrozenImagesRawDisplay(ctx context.Context, input, raw, display, original string, images []string) error {
393 turn := orchestratedTurn{
394 input: input, raw: raw, display: display, editedOriginal: original,
395 imageCandidates: append([]string(nil), images...), imagesResolved: true,
396 }
397 if o.c.imageInputEnabled() {
398 turn.userImages = append([]string(nil), images...)
399 }
400 ctx = agent.WithSubagentImageCandidates(ctx, turn.imageCandidates)
401 return o.runOrchestratedTurn(ctx, turn)
402 }
403
403 lines GO