返回 DeepSeek-Reasonix
inbox_run.go
根目录 / internal / control / inbox_run.go
1 package control
2
3 import (
4 "context"
5 "fmt"
6 "strings"
7
8 "reasonix/internal/attachment"
9 "reasonix/internal/sessioninbox"
10 )
11
12 // RunInboxTurn synchronously claims and executes one durable item. Bot and ACP
13 // use this path so their blocking response sink remains attached through every
14 // queued follow-up while Controller still owns durable state and ack semantics.
15 func (c *Controller) RunInboxTurn(ctx context.Context, id string) error {
16 st, err := c.ensureInbox()
17 if err != nil {
18 return err
19 }
20 meta, env, err := st.ReadItem(id)
21 if err != nil {
22 return err
23 }
24 if meta.State != sessioninbox.StateQueued {
25 return sessioninbox.ErrInvalidState
26 }
27 run, block, err := c.prepareInboxRunContext(ctx, env)
28 if err != nil {
29 return err
30 }
31 if block != "" {
32 if err := st.TransitionPrepared(id, sessioninbox.ContentVersion(meta), sessioninbox.StateBlocked, block, true); err != nil {
33 return err
34 }
35 return fmt.Errorf("%w: %s", sessioninbox.ErrInvalidState, block)
36 }
37 err = c.runSynchronousTurn(ctx, func() error {
38 c.inbox.admissionMu.Lock()
39 defer c.inbox.admissionMu.Unlock()
40 c.inbox.trackAdmission(id)
41 defer c.inbox.untrackAdmission(id)
42 if err := st.TransitionPrepared(id, sessioninbox.ContentVersion(meta), sessioninbox.StateRunning, "", true); err != nil {
43 return err
44 }
45 c.inbox.mu.Lock()
46 c.inbox.trackActive(id)
47 c.inbox.mu.Unlock()
48 return nil
49 }, run)
50 if err != nil {
51 return err
52 }
53 return c.waitForGoalTerminal(ctx)
54 }
55
56 func (c *Controller) prepareInboxRun(env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) {
57 return c.prepareInboxRunContext(c.attachmentContext(), env)
58 }
59
60 func (c *Controller) prepareInboxRunContext(ctx context.Context, env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) {
61 var sources []attachment.Source
62 for _, input := range env.ImageInputs {
63 if input.Attachment != nil {
64 sources = append(sources, attachment.Source{Existing: input.Attachment, DisplayName: input.Attachment.DisplayName})
65 }
66 }
67 if _, err := c.attachmentService().PrepareBatch(ctx, sources); err != nil {
68 return nil, ImageReferenceFailures(imageFailuresFromAttachment(err)).Error(), nil
69 }
70 submit, frozenImages, block, err := applyInboxReferences(env)
71 if err != nil || block != "" {
72 return nil, block, err
73 }
74 display := firstNonEmptyStr(env.DisplayText, submit)
75 raw := firstNonEmptyStr(env.RawText, submit)
76 requests := controlInvocationsFromInbox(env)
77 if len(requests) == 0 {
78 return func(ctx context.Context) error {
79 ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs})
80 return c.runGoalLoopWithFrozenImagesRawDisplay(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), submit, raw, display, frozenImages)
81 }, "", nil
82 }
83 prepared, err := c.prepareInvocationTurn(submit, requests)
84 if err != nil {
85 return nil, err.Error(), nil
86 }
87 return func(ctx context.Context) error {
88 ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs})
89 return c.runPreparedInvocationTurn(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), prepared, submit, raw, display, frozenImages)
90 }, "", nil
91 }
92
93 // submitPreparedInboxTurn starts an already-classified inbox envelope without
94 // interpreting slash commands, shell shortcuts, or @references a second time.
95 func (c *Controller) submitPreparedInboxTurn(itemID string, run func(context.Context) error) admissionResult {
96 return c.runGuardedInbox(run, func() {
97 c.inbox.mu.Lock()
98 c.inbox.trackActive(itemID)
99 c.inbox.mu.Unlock()
100 })
101 }
102
103 func sessionInboxInvocations(requests []InvocationRequest) []sessioninbox.StructuredInvocation {
104 if len(requests) == 0 {
105 return nil
106 }
107 out := make([]sessioninbox.StructuredInvocation, 0, len(requests))
108 for _, request := range requests {
109 out = append(out, sessioninbox.StructuredInvocation{Name: request.Name, Kind: request.Kind, Offset: request.Offset})
110 }
111 return out
112 }
113
114 func controlInvocationsFromInbox(env sessioninbox.PromptEnvelope) []InvocationRequest {
115 stored := env.Invocations
116 if len(stored) == 0 && env.Invocation != nil {
117 stored = []sessioninbox.StructuredInvocation{*env.Invocation}
118 }
119 out := make([]InvocationRequest, 0, len(stored))
120 for _, invocation := range stored {
121 out = append(out, InvocationRequest{Name: invocation.Name, Kind: invocation.Kind, Offset: invocation.Offset})
122 }
123 return out
124 }
125
125 lines GO