返回 DeepSeek-Reasonix
inbox_queue.go
根目录 / internal / control / inbox_queue.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "slices"
8 "strings"
9
10 "reasonix/internal/sessioninbox"
11 )
12
13 // InboxQueueRequest is the additive, session-fenced queue editing protocol.
14 type InboxQueueRequest struct {
15 Kind string `json:"kind"`
16 ItemID string `json:"itemId,omitempty"`
17 Text string `json:"text,omitempty"`
18 ContentVersion string `json:"contentVersion,omitempty"`
19 BeforeItemID *string `json:"beforeItemId"`
20 QueueRevision int64 `json:"queueRevision"`
21 Paused bool `json:"paused,omitempty"`
22 TurnID string `json:"turnId,omitempty"`
23 Display string `json:"display,omitempty"`
24 IdempotencyKey string `json:"idempotencyKey,omitempty"`
25 }
26
27 type InboxQueueEdit struct {
28 ID string `json:"id"`
29 Text string `json:"text"`
30 ContentVersion string `json:"contentVersion"`
31 References []string `json:"references"`
32 }
33
34 type InboxQueueResult struct {
35 Outcome string `json:"outcome"`
36 Reason string `json:"reason,omitempty"`
37 Snapshot sessioninbox.InboxSnapshot `json:"snapshot"`
38 Edit *InboxQueueEdit `json:"edit,omitempty"`
39 Receipt *sessioninbox.InboxReceipt `json:"receipt,omitempty"`
40 }
41
42 // InboxQueue pins one store for the whole operation. Never re-resolve the
43 // controller's current session after reference preparation or an async write.
44 func (c *Controller) InboxQueue(path string, req InboxQueueRequest) (InboxQueueResult, error) {
45 st, err := c.ensureInbox()
46 if err != nil {
47 return InboxQueueResult{}, err
48 }
49 if path == "" || st.SessionPath() != path {
50 return InboxQueueResult{Outcome: "unavailable", Reason: "session_changed"}, nil
51 }
52 result := InboxQueueResult{Outcome: "applied"}
53 switch req.Kind {
54 case "snapshot":
55 result.Outcome = "unchanged"
56 case "read":
57 result.Edit, err = readInboxQueueEdit(st, req.ItemID)
58 case "edit":
59 err = c.editInboxQueue(st, req)
60 case "move":
61 err = st.MoveItemBefore(req.ItemID, req.BeforeItemID, req.QueueRevision)
62 case "delete":
63 err = st.DeleteItem(req.ItemID)
64 case "pause":
65 err = st.SetPaused(req.Paused)
66 case "retry":
67 err = st.RetryItem(req.ItemID)
68 case "steer":
69 if req.TurnID == "" {
70 return result, fmt.Errorf("turnId is required")
71 }
72 _, err = c.trySteerInboxItemForSession(req.ItemID, req.TurnID, path)
73 case "enqueue_steer":
74 if req.TurnID == "" || req.IdempotencyKey == "" {
75 return result, fmt.Errorf("turnId and idempotencyKey are required")
76 }
77 var receipt sessioninbox.InboxReceipt
78 receipt, err = c.TryEnqueueAndSteerForTurn(req.TurnID, InboxRequest{ExpectedSessionPath: path, Display: req.Display, Raw: req.Text, Submit: req.Text, Idempotency: req.IdempotencyKey, Source: "desktop"})
79 if err == nil {
80 result.Receipt = &receipt
81 }
82 default:
83 return result, fmt.Errorf("unknown inbox queue operation %q", req.Kind)
84 }
85 result.Snapshot = st.Snapshot()
86 if err != nil {
87 result.Outcome, result.Reason = inboxQueueFailure(err)
88 if result.Reason == "" {
89 return result, err
90 }
91 return result, nil
92 }
93 if req.Kind != "read" && req.Kind != "snapshot" && !result.Snapshot.Paused {
94 c.maybeDispatchInbox()
95 }
96 return result, nil
97 }
98
99 func inboxQueueFailure(err error) (string, string) {
100 for _, item := range []struct {
101 err error
102 reason string
103 }{
104 {sessioninbox.ErrContentChanged, "content_changed"},
105 {sessioninbox.ErrOrderChanged, "order_changed"},
106 {sessioninbox.ErrAnchorMissing, "anchor_missing"},
107 } {
108 if errors.Is(err, item.err) {
109 return "conflict", item.reason
110 }
111 }
112 for _, item := range []struct {
113 err error
114 reason string
115 }{
116 {sessioninbox.ErrNotFound, "item_missing"},
117 {sessioninbox.ErrInvalidState, "item_not_pending"},
118 {sessioninbox.ErrSchemaReadonly, "read_only"},
119 {errInboxStructuredEdit, "structured_input"},
120 {sessioninbox.ErrClosed, "session_changed"},
121 {ErrInboxSessionChanged, "session_changed"},
122 } {
123 if errors.Is(err, item.err) {
124 return "unavailable", item.reason
125 }
126 }
127 return "", ""
128 }
129
130 var errInboxStructuredEdit = errors.New("structured invocation cannot be edited as plain text")
131
132 func readInboxQueueEdit(st *sessioninbox.Store, id string) (*InboxQueueEdit, error) {
133 meta, env, err := st.ReadItem(id)
134 if err != nil {
135 return nil, err
136 }
137 if meta.State != sessioninbox.StateQueued && meta.State != sessioninbox.StateBlocked && meta.State != sessioninbox.StateUncertain {
138 return nil, sessioninbox.ErrInvalidState
139 }
140 if env.Invocation != nil || len(env.Invocations) > 0 {
141 return nil, errInboxStructuredEdit
142 }
143 refs := append([]string{}, env.Attachments...)
144 refs = append(refs, env.ExplicitRefs...)
145 for _, ref := range env.Refs {
146 refs = append(refs, firstNonEmptyStr(ref.DisplayPath, ref.Path))
147 }
148 for _, input := range env.ImageInputs {
149 if input.Attachment != nil {
150 refs = append(refs, input.Attachment.DisplayName)
151 }
152 }
153 return &InboxQueueEdit{ID: id, Text: env.SubmitText, ContentVersion: sessioninbox.ContentVersion(meta), References: refs}, nil
154 }
155
156 func (c *Controller) editInboxQueue(st *sessioninbox.Store, req InboxQueueRequest) error {
157 if strings.TrimSpace(req.Text) == "" {
158 return sessioninbox.ErrEmpty
159 }
160 meta, env, err := st.ReadItem(req.ItemID)
161 if err != nil {
162 return err
163 }
164 if req.ContentVersion == "" || sessioninbox.ContentVersion(meta) != req.ContentVersion {
165 return sessioninbox.ErrContentChanged
166 }
167 if env.Invocation != nil || len(env.Invocations) > 0 {
168 return errInboxStructuredEdit
169 }
170 // Keep the original immutable reference bytes when only prose changes.
171 // A changed @ reference is resolved by the existing authorization path.
172 if !slices.Equal(parseRefTokens(env.SubmitText), parseRefTokens(req.Text)) {
173 if err := c.freezeInboxEnvelopeReferences(context.Background(), &env, req.Text, env.ExplicitRefs); err != nil {
174 return err
175 }
176 }
177 env.DisplayText, env.RawText, env.SubmitText = req.Text, req.Text, req.Text
178 _, err = st.UpdateItemIfVersion(req.ItemID, env, req.ContentVersion)
179 return err
180 }
181
181 lines GO