返回 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 _ = st.SetState(id, sessioninbox.StateBlocked, block)
33 _ = st.SetPaused(true)
34 return fmt.Errorf("%w: %s", sessioninbox.ErrInvalidState, block)
35 }
36 err = c.runSynchronousTurn(ctx, func() error {
37 c.inbox.admissionMu.Lock()
38 defer c.inbox.admissionMu.Unlock()
39 c.inbox.trackAdmission(id)
40 defer c.inbox.untrackAdmission(id)
41 if err := st.ClaimItem(id); err != nil {
42 return err
43 }
44 c.inbox.mu.Lock()
45 c.inbox.trackActive(id)
46 c.inbox.mu.Unlock()
47 return nil
48 }, run)
49 if err != nil {
50 return err
51 }
52 return c.waitForGoalTerminal(ctx)
53 }
54
55 func (c *Controller) prepareInboxRun(env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) {
56 return c.prepareInboxRunContext(c.attachmentContext(), env)
57 }
58
59 func (c *Controller) prepareInboxRunContext(ctx context.Context, env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) {
60 var sources []attachment.Source
61 for _, input := range env.ImageInputs {
62 if input.Attachment != nil {
63 sources = append(sources, attachment.Source{Existing: input.Attachment, DisplayName: input.Attachment.DisplayName})
64 }
65 }
66 if _, err := c.attachmentService().PrepareBatch(ctx, sources); err != nil {
67 return nil, ImageReferenceFailures(imageFailuresFromAttachment(err)).Error(), nil
68 }
69 submit, frozenImages, block, err := applyInboxReferences(env)
70 if err != nil || block != "" {
71 return nil, block, err
72 }
73 display := firstNonEmptyStr(env.DisplayText, submit)
74 raw := firstNonEmptyStr(env.RawText, submit)
75 requests := controlInvocationsFromInbox(env)
76 if len(requests) == 0 {
77 return func(ctx context.Context) error {
78 ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs})
79 return c.runGoalLoopWithFrozenImagesRawDisplay(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), submit, raw, display, frozenImages)
80 }, "", nil
81 }
82 prepared, err := c.prepareInvocationTurn(submit, requests)
83 if err != nil {
84 return nil, err.Error(), nil
85 }
86 return func(ctx context.Context) error {
87 ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs})
88 return c.runPreparedInvocationTurn(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), prepared, submit, raw, display, frozenImages)
89 }, "", nil
90 }
91
92 // submitPreparedInboxTurn starts an already-classified inbox envelope without
93 // interpreting slash commands, shell shortcuts, or @references a second time.
94 func (c *Controller) submitPreparedInboxTurn(itemID string, run func(context.Context) error) admissionResult {
95 return c.runGuardedInbox(run, func() {
96 c.inbox.mu.Lock()
97 c.inbox.trackActive(itemID)
98 c.inbox.mu.Unlock()
99 })
100 }
101
102 func sessionInboxInvocations(requests []InvocationRequest) []sessioninbox.StructuredInvocation {
103 if len(requests) == 0 {
104 return nil
105 }
106 out := make([]sessioninbox.StructuredInvocation, 0, len(requests))
107 for _, request := range requests {
108 out = append(out, sessioninbox.StructuredInvocation{Name: request.Name, Kind: request.Kind, Offset: request.Offset})
109 }
110 return out
111 }
112
113 func controlInvocationsFromInbox(env sessioninbox.PromptEnvelope) []InvocationRequest {
114 stored := env.Invocations
115 if len(stored) == 0 && env.Invocation != nil {
116 stored = []sessioninbox.StructuredInvocation{*env.Invocation}
117 }
118 out := make([]InvocationRequest, 0, len(stored))
119 for _, invocation := range stored {
120 out = append(out, InvocationRequest{Name: invocation.Name, Kind: invocation.Kind, Offset: invocation.Offset})
121 }
122 return out
123 }
124
124 lines GO