返回 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 sessionID := agent.BranchID(st.SessionPath())
71 if id, canonical := strings.CutPrefix(st.SessionPath(), "session-id:"); canonical {
72 sessionID = id
73 }
74 rec, err := st.Enqueue(sessioninbox.EnqueueRequest{
75 Intent: intent,
76 Envelope: env,
77 Source: req.Source,
78 Idempotency: req.Idempotency,
79 SessionID: sessionID,
80 })
81 if err != nil {
82 if errors.Is(err, sessioninbox.ErrCapacityItems) || errors.Is(err, sessioninbox.ErrCapacityBytes) || errors.Is(err, sessioninbox.ErrItemTooLarge) {
83 sessioninbox.NoteCapacityReject()
84 } else {
85 sessioninbox.NoteTxFail()
86 }
87 return sessioninbox.InboxReceipt{}, err
88 }
89 if !rec.Idempotent && len(env.ReferenceErrors) > 0 {
90 reason := strings.Join(env.ReferenceErrors, "; ")
91 if stateErr := st.SetState(rec.ItemID, sessioninbox.StateBlocked, reason); stateErr != nil {
92 return sessioninbox.InboxReceipt{}, stateErr
93 }
94 if pauseErr := st.SetPaused(true); pauseErr != nil {
95 return sessioninbox.InboxReceipt{}, pauseErr
96 }
97 rec.Paused = true
98 }
99 sessioninbox.NoteEnqueue(int64(len(env.SubmitText)))
100 return rec, nil
101 }
102
102 lines GO