返回 DeepSeek-Reasonix
inbox.go
根目录 / internal / serve / inbox.go
1 package serve
2
3 import (
4 "encoding/json"
5 "errors"
6 "net/http"
7 "strings"
8
9 "reasonix/internal/control"
10 "reasonix/internal/sessioninbox"
11 )
12
13 func (s *Server) registerInboxRoutes(mux *http.ServeMux) {
14 mux.HandleFunc("POST /inbox/queue", s.foregroundMutation(s.inboxQueueCommand))
15 mux.HandleFunc("GET /inbox", s.inboxList)
16 mux.HandleFunc("GET /inbox/receipt", s.inboxReceipt)
17 mux.HandleFunc("POST /inbox/items", s.foregroundMutation(s.inboxEnqueue))
18 mux.HandleFunc("GET /inbox/items/{id}", s.inboxGet)
19 mux.HandleFunc("PATCH /inbox/items/{id}", s.foregroundMutation(s.inboxUpdate))
20 mux.HandleFunc("DELETE /inbox/items/{id}", s.foregroundMutation(s.inboxDelete))
21 mux.HandleFunc("POST /inbox/move", s.foregroundMutation(s.inboxMove))
22 mux.HandleFunc("POST /inbox/pause", s.foregroundMutation(s.inboxPause))
23 mux.HandleFunc("POST /inbox/resume", s.foregroundMutation(s.inboxResume))
24 mux.HandleFunc("POST /inbox/items/{id}/retry", s.foregroundMutation(s.inboxRetry))
25 mux.HandleFunc("POST /inbox/items/{id}/refresh", s.foregroundMutation(s.inboxRefresh))
26 }
27
28 func (s *Server) inboxQueueCommand(w http.ResponseWriter, r *http.Request) {
29 var body struct {
30 SessionPath string `json:"sessionPath"`
31 Request control.InboxQueueRequest `json:"request"`
32 }
33 if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, sessioninbox.DefaultMaxItemBytes+4096)).Decode(&body); err != nil || body.SessionPath == "" {
34 http.Error(w, "missing queue target", http.StatusBadRequest)
35 return
36 }
37 api, ok := s.inboxAPI().(interface {
38 InboxQueue(string, control.InboxQueueRequest) (control.InboxQueueResult, error)
39 })
40 if !ok {
41 http.Error(w, "unsupported", http.StatusNotImplemented)
42 return
43 }
44 result, err := api.InboxQueue(body.SessionPath, body.Request)
45 if err != nil {
46 writeInboxError(w, err)
47 return
48 }
49 writeJSON(w, result)
50 }
51
52 func (s *Server) inboxAPI() control.SessionAPI {
53 return s.ctl()
54 }
55
56 func writeInboxError(w http.ResponseWriter, err error) {
57 switch {
58 case errors.Is(err, sessioninbox.ErrItemTooLarge):
59 http.Error(w, err.Error(), http.StatusRequestEntityTooLarge) // 413
60 case errors.Is(err, sessioninbox.ErrCapacityItems), errors.Is(err, sessioninbox.ErrCapacityBytes),
61 errors.Is(err, sessioninbox.ErrInvalidState), errors.Is(err, sessioninbox.ErrPaused),
62 errors.Is(err, sessioninbox.ErrNotFound), errors.Is(err, sessioninbox.ErrIdempotencyConflict):
63 http.Error(w, err.Error(), http.StatusConflict) // 409
64 case errors.Is(err, sessioninbox.ErrEmpty):
65 http.Error(w, err.Error(), http.StatusBadRequest)
66 default:
67 http.Error(w, err.Error(), http.StatusInternalServerError)
68 }
69 }
70
71 func (s *Server) inboxList(w http.ResponseWriter, r *http.Request) {
72 s.bindMu.Lock()
73 defer s.bindMu.Unlock()
74 if !s.validateInboxReadSessionLocked(w, r) {
75 return
76 }
77 snap := s.inboxAPI().InboxSnapshot()
78 w.Header().Set("Content-Type", "application/json")
79 _ = json.NewEncoder(w).Encode(snap)
80 }
81
82 // validateInboxReadSessionLocked keeps legacy unscoped reads compatible while
83 // fencing modern Desktop reads against a concurrent foreground replacement.
84 func (s *Server) validateInboxReadSessionLocked(w http.ResponseWriter, r *http.Request) bool {
85 if !s.validateExpectedSessionLocked(w, r) {
86 return false
87 }
88 if err := s.expectedSessionPathErrorLocked(r.URL.Query().Get("session")); err != nil {
89 http.Error(w, err.Error(), http.StatusConflict)
90 return false
91 }
92 return true
93 }
94
95 func (s *Server) inboxEnqueue(w http.ResponseWriter, r *http.Request) {
96 var body struct {
97 Input string `json:"input"`
98 Display string `json:"display"`
99 Invocations []control.InvocationRequest `json:"invocations"`
100 Intent string `json:"intent"`
101 IdempotencyKey string `json:"idempotencyKey"`
102 }
103 if err := json.NewDecoder(r.Body).Decode(&body); err != nil || strings.TrimSpace(body.Input) == "" {
104 http.Error(w, "missing input", http.StatusBadRequest)
105 return
106 }
107 intent := sessioninbox.IntentFollowup
108 if strings.EqualFold(body.Intent, "steer") {
109 intent = sessioninbox.IntentSteer
110 }
111 api := s.inboxAPI()
112 if ensurer, ok := any(api).(interface{ EnsureSessionPath() }); ok {
113 ensurer.EnsureSessionPath()
114 }
115 req := control.InboxRequest{
116 Intent: intent,
117 Display: body.Display,
118 Raw: body.Input,
119 Submit: body.Input,
120 Source: "http",
121 Idempotency: body.IdempotencyKey,
122 Invocations: body.Invocations,
123 }
124 if req.Display == "" {
125 req.Display = body.Input
126 }
127 var rec sessioninbox.InboxReceipt
128 var err error
129 if intent == sessioninbox.IntentSteer {
130 rec, err = api.TryEnqueueAndSteer(req)
131 } else {
132 rec, err = api.TryEnqueueFollowup(req)
133 }
134 if err != nil {
135 writeInboxError(w, err)
136 return
137 }
138 w.Header().Set("Content-Type", "application/json")
139 w.WriteHeader(http.StatusAccepted)
140 _ = json.NewEncoder(w).Encode(rec)
141 }
142
143 func (s *Server) inboxReceipt(w http.ResponseWriter, r *http.Request) {
144 s.bindMu.Lock()
145 defer s.bindMu.Unlock()
146 if !s.validateInboxReadSessionLocked(w, r) {
147 return
148 }
149 ctrl := s.ctl()
150 reader, ok := ctrl.(interface {
151 LookupInboxReceipt(string) (sessioninbox.InboxReceipt, bool, error)
152 })
153 if !ok {
154 http.NotFound(w, r)
155 return
156 }
157 receipt, found, err := reader.LookupInboxReceipt(r.URL.Query().Get("key"))
158 if err != nil {
159 writeInboxError(w, err)
160 return
161 }
162 if !found {
163 http.NotFound(w, r)
164 return
165 }
166 writeJSON(w, receipt)
167 }
168
169 func (s *Server) inboxGet(w http.ResponseWriter, r *http.Request) {
170 id := r.PathValue("id")
171 meta, env, err := s.inboxAPI().ReadInboxItem(id)
172 if err != nil {
173 writeInboxError(w, err)
174 return
175 }
176 w.Header().Set("Content-Type", "application/json")
177 _ = json.NewEncoder(w).Encode(map[string]any{"meta": meta, "envelope": env})
178 }
179
180 func (s *Server) inboxUpdate(w http.ResponseWriter, r *http.Request) {
181 id := r.PathValue("id")
182 var body struct {
183 Input string `json:"input"`
184 }
185 if err := json.NewDecoder(r.Body).Decode(&body); err != nil || strings.TrimSpace(body.Input) == "" {
186 http.Error(w, "missing input", http.StatusBadRequest)
187 return
188 }
189 meta, err := s.inboxAPI().UpdateInboxItem(id, body.Input, body.Input, body.Input)
190 if err != nil {
191 writeInboxError(w, err)
192 return
193 }
194 w.Header().Set("Content-Type", "application/json")
195 _ = json.NewEncoder(w).Encode(meta)
196 }
197
198 func (s *Server) inboxDelete(w http.ResponseWriter, r *http.Request) {
199 id := r.PathValue("id")
200 if err := s.inboxAPI().DeleteInboxItem(id); err != nil {
201 writeInboxError(w, err)
202 return
203 }
204 w.WriteHeader(http.StatusNoContent)
205 }
206
207 func (s *Server) inboxMove(w http.ResponseWriter, r *http.Request) {
208 var body struct {
209 ID string `json:"id"`
210 ToIndex int `json:"toIndex"`
211 }
212 if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.ID == "" {
213 http.Error(w, "missing id", http.StatusBadRequest)
214 return
215 }
216 if err := s.inboxAPI().MoveInboxItem(body.ID, body.ToIndex); err != nil {
217 writeInboxError(w, err)
218 return
219 }
220 w.WriteHeader(http.StatusNoContent)
221 }
222
223 func (s *Server) inboxPause(w http.ResponseWriter, r *http.Request) {
224 _ = r
225 if err := s.inboxAPI().SetInboxPaused(true); err != nil {
226 writeInboxError(w, err)
227 return
228 }
229 w.WriteHeader(http.StatusNoContent)
230 }
231
232 func (s *Server) inboxResume(w http.ResponseWriter, r *http.Request) {
233 _ = r
234 if err := s.inboxAPI().SetInboxPaused(false); err != nil {
235 writeInboxError(w, err)
236 return
237 }
238 w.WriteHeader(http.StatusNoContent)
239 }
240
241 func (s *Server) inboxRetry(w http.ResponseWriter, r *http.Request) {
242 id := r.PathValue("id")
243 if err := s.inboxAPI().RetryInboxItem(id); err != nil {
244 writeInboxError(w, err)
245 return
246 }
247 w.WriteHeader(http.StatusNoContent)
248 }
249
250 func (s *Server) inboxRefresh(w http.ResponseWriter, r *http.Request) {
251 id := r.PathValue("id")
252 if err := s.inboxAPI().RefreshInboxReferences(id); err != nil {
253 writeInboxError(w, err)
254 return
255 }
256 w.WriteHeader(http.StatusNoContent)
257 }
258
258 lines GO