返回 DeepSeek-Reasonix
inbox_submission.go
根目录 / internal / control / inbox_submission.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "maps"
7 "reasonix/internal/agent"
8 "reasonix/internal/sessioninbox"
9 "strings"
10 )
11
12 // EnqueueInbox durably queues an instruction. Only returns a receipt after
13 // blob+manifest commit. Does not auto-start a turn (call TrySubmit / dispatcher).
14 func (c *Controller) EnqueueInbox(req InboxRequest) (sessioninbox.InboxReceipt, error) {
15 return c.EnqueueInboxContext(c.attachmentContext(), req)
16 }
17
18 func (c *Controller) EnqueueInboxContext(ctx context.Context, req InboxRequest) (sessioninbox.InboxReceipt, error) {
19 ctx, cancel := c.NewAttachmentOperationContext(ctx)
20 defer cancel()
21 c.inbox.prepareMu.Lock()
22 defer c.inbox.prepareMu.Unlock()
23 st, err := c.ensureInbox()
24 if err != nil {
25 return sessioninbox.InboxReceipt{}, err
26 }
27 if req.ExpectedSessionPath != "" && st.SessionPath() != req.ExpectedSessionPath {
28 return sessioninbox.InboxReceipt{}, ErrInboxSessionChanged
29 }
30 submit := strings.TrimSpace(firstNonEmptyStr(req.Submit, req.Raw))
31 if submit == "" && len(req.Invocations) == 0 {
32 submit = strings.TrimSpace(req.Display)
33 }
34 if submit == "" && len(req.Invocations) == 0 {
35 return sessioninbox.InboxReceipt{}, sessioninbox.ErrEmpty
36 }
37 display := firstNonEmptyStr(req.Display, submit)
38 raw := firstNonEmptyStr(req.Raw, submit)
39 env := sessioninbox.PromptEnvelope{
40 DisplayText: display,
41 RawText: raw,
42 SubmitText: submit,
43 Format: req.Format,
44 Source: req.Source,
45 Idempotency: req.Idempotency,
46 ExplicitRefs: append([]string(nil), req.FreezeRefs...),
47 Invocations: sessionInboxInvocations(req.Invocations),
48 Extra: maps.Clone(req.Extra),
49 }
50 for _, item := range req.Attachments {
51 env.AttachmentIdentities = append(env.AttachmentIdentities, item.ClientAttachmentID)
52 }
53 if len(req.Attachments) > 0 {
54 env.FingerprintVersion = 1
55 env.RequestFingerprint = inboxAttachmentFingerprint(req, env)
56 }
57 if receipt, found, err := st.LookupEnvelopeReceipt(req.Idempotency, env); found || err != nil {
58 return receipt, err
59 }
60 if err := c.freezeInboxEnvelopeReferences(ctx, &env, submit, req.FreezeRefs, req.Attachments...); err != nil {
61 return sessioninbox.InboxReceipt{}, err
62 }
63 if err := ctx.Err(); err != nil {
64 return sessioninbox.InboxReceipt{}, err
65 }
66 intent := req.Intent
67 if intent != sessioninbox.IntentSteer {
68 intent = sessioninbox.IntentFollowup
69 }
70 rec, err := st.Enqueue(sessioninbox.EnqueueRequest{
71 Intent: intent,
72 Envelope: env,
73 Source: req.Source,
74 Idempotency: req.Idempotency,
75 SessionID: agent.BranchID(st.SessionPath()),
76 })
77 if err != nil {
78 if errors.Is(err, sessioninbox.ErrCapacityItems) || errors.Is(err, sessioninbox.ErrCapacityBytes) || errors.Is(err, sessioninbox.ErrItemTooLarge) {
79 sessioninbox.NoteCapacityReject()
80 } else {
81 sessioninbox.NoteTxFail()
82 }
83 return sessioninbox.InboxReceipt{}, err
84 }
85 if !rec.Idempotent && len(env.ReferenceErrors) > 0 {
86 reason := strings.Join(env.ReferenceErrors, "; ")
87 if stateErr := st.SetState(rec.ItemID, sessioninbox.StateBlocked, reason); stateErr != nil {
88 return sessioninbox.InboxReceipt{}, stateErr
89 }
90 if pauseErr := st.SetPaused(true); pauseErr != nil {
91 return sessioninbox.InboxReceipt{}, pauseErr
92 }
93 rec.Paused = true
94 }
95 sessioninbox.NoteEnqueue(int64(len(env.SubmitText)))
96 return rec, nil
97 }
98
98 lines GO