返回 DeepSeek-Reasonix
transcript_api.go
根目录 / internal / serve / transcript_api.go
1 package serve
2
3 import (
4 "encoding/base64"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "net/http"
9 "strings"
10
11 "reasonix/internal/agent"
12 "reasonix/internal/control"
13 "reasonix/internal/session"
14 "reasonix/internal/sessioncontent"
15 "reasonix/internal/transcript"
16 )
17
18 func (s *Server) registerTranscriptRoutes(mux *http.ServeMux) {
19 mux.HandleFunc("GET /session-export/snapshot", s.sessionExportSnapshot)
20 mux.HandleFunc("POST /session-export/document", s.sessionExportDocument)
21 mux.HandleFunc("POST /session-export/validate", s.sessionExportValidate)
22 mux.HandleFunc("POST /session-export/diagnostic", s.sessionExportDiagnostic)
23 mux.HandleFunc("GET /transcript/follow", s.transcriptFollow)
24 mux.HandleFunc("GET /transcript/snapshot", s.transcriptSnapshot)
25 mux.HandleFunc("GET /transcript/page", s.transcriptSnapshot)
26 mux.HandleFunc("GET /transcript/content", s.transcriptContent)
27 mux.HandleFunc("GET /transcript/outline", s.transcriptOutline)
28 mux.HandleFunc("GET /transcript/replay", s.transcriptReplay)
29 mux.HandleFunc("GET /session/open", s.sessionOpen)
30 mux.HandleFunc("GET /session-history/page", s.sessionHistoryPage)
31 mux.HandleFunc("GET /session-history/search", s.sessionHistorySearch)
32 mux.HandleFunc("GET /session-history/locate", s.sessionHistoryLocate)
33 mux.HandleFunc("GET /session-history/content", s.sessionHistoryContent)
34 mux.HandleFunc("GET /session-history/window", s.sessionHistoryWindow)
35 mux.HandleFunc("GET /session-history/outline", s.sessionHistoryOutline)
36 mux.HandleFunc("GET /session-message-field", s.sessionMessageField)
37 }
38
39 // sessionHistoryWindow serves history-window-v1: a bounded window around an
40 // anchor in either direction of a fixed durable snapshot.
41 func (s *Server) sessionHistoryWindow(w http.ResponseWriter, r *http.Request) {
42 s.bindMu.Lock()
43 defer s.bindMu.Unlock()
44 query, ref, ok := s.canonicalSessionQuery(w, r)
45 if !ok {
46 return
47 }
48 req := session.HistoryWindowRequest{Anchor: r.URL.Query().Get("anchor"), MessageID: r.URL.Query().Get("messageId"), Cursor: r.URL.Query().Get("cursor"), Direction: r.URL.Query().Get("direction")}
49 req.Generation = r.URL.Query().Get("generation")
50 if raw := r.URL.Query().Get("snapshotSequence"); raw != "" {
51 var cut uint64
52 if _, err := fmt.Sscan(raw, &cut); err != nil {
53 http.Error(w, "invalid snapshot sequence", http.StatusBadRequest)
54 return
55 }
56 req.SnapshotSequence = &cut
57 }
58 if raw := r.URL.Query().Get("turn"); raw != "" {
59 if _, err := fmt.Sscan(raw, &req.Turn); err != nil {
60 http.Error(w, "invalid history window turn", http.StatusBadRequest)
61 return
62 }
63 }
64 if raw := r.URL.Query().Get("limit"); raw != "" {
65 if _, err := fmt.Sscan(raw, &req.Limit); err != nil {
66 http.Error(w, "invalid history window limit", http.StatusBadRequest)
67 return
68 }
69 }
70 page, err := query.ReadHistoryWindow(r.Context(), ref, req)
71 if err != nil {
72 http.Error(w, err.Error(), http.StatusConflict)
73 return
74 }
75 w.Header().Set("Cache-Control", "no-store")
76 w.Header().Set("Content-Type", "application/json")
77 _ = json.NewEncoder(w).Encode(page)
78 }
79
80 // sessionMessageField serves one bounded, UTF-8 aligned fragment of one
81 // top-level message field.
82 func (s *Server) sessionMessageField(w http.ResponseWriter, r *http.Request) {
83 s.bindMu.Lock()
84 defer s.bindMu.Unlock()
85 query, ref, ok := s.canonicalSessionQuery(w, r)
86 if !ok {
87 return
88 }
89 q := r.URL.Query()
90 var version int
91 var offset, length int64
92 if raw := q.Get("version"); raw != "" {
93 if _, err := fmt.Sscan(raw, &version); err != nil {
94 http.Error(w, "invalid message field version", http.StatusBadRequest)
95 return
96 }
97 }
98 if raw := q.Get("offset"); raw != "" {
99 if _, err := fmt.Sscan(raw, &offset); err != nil {
100 http.Error(w, "invalid message field offset", http.StatusBadRequest)
101 return
102 }
103 }
104 if raw := q.Get("length"); raw != "" {
105 if _, err := fmt.Sscan(raw, &length); err != nil {
106 http.Error(w, "invalid message field length", http.StatusBadRequest)
107 return
108 }
109 }
110 page, err := query.ReadMessageField(r.Context(), ref, q.Get("messageId"), version, q.Get("field"), offset, length)
111 if err != nil {
112 http.Error(w, err.Error(), http.StatusConflict)
113 return
114 }
115 w.Header().Set("Cache-Control", "no-store")
116 w.Header().Set("Content-Type", "application/json")
117 _ = json.NewEncoder(w).Encode(page)
118 }
119
120 func (s *Server) sessionOpen(w http.ResponseWriter, r *http.Request) {
121 s.bindMu.Lock()
122 defer s.bindMu.Unlock()
123 query, ref, ok := s.canonicalSessionQuery(w, r)
124 if !ok {
125 return
126 }
127 view, err := query.OpenSession(r.Context(), ref)
128 if err != nil {
129 http.Error(w, err.Error(), http.StatusConflict)
130 return
131 }
132 w.Header().Set("Cache-Control", "no-store")
133 w.Header().Set("Content-Type", "application/json")
134 _ = json.NewEncoder(w).Encode(view)
135 }
136
137 type sessionHistoryContentRequest struct {
138 Ref sessioncontent.Ref `json:"ref"`
139 Offset int64 `json:"offset"`
140 Length int64 `json:"length"`
141 }
142
143 type sessionHistoryContentResponse struct {
144 Data string `json:"data"`
145 NextOffset int64 `json:"nextOffset"`
146 Done bool `json:"done"`
147 }
148
149 func (s *Server) canonicalSessionQuery(w http.ResponseWriter, r *http.Request) (*session.Query, session.SessionRef, bool) {
150 identity, ok := s.ctl().(control.IdentityLifecycle)
151 if !ok || !identity.UsesExclusiveSession() {
152 http.Error(w, "canonical session history is unavailable", http.StatusNotImplemented)
153 return nil, session.SessionRef{}, false
154 }
155 service := identity.SessionService()
156 if service == nil || service.Query() == nil {
157 http.Error(w, "canonical session identity is unavailable", http.StatusConflict)
158 return nil, session.SessionRef{}, false
159 }
160 ref, bound := identity.SessionRef()
161 requested := strings.TrimSpace(r.URL.Query().Get("sessionId"))
162 if requested == "" || (bound && requested == ref.SessionID) {
163 if !bound {
164 http.Error(w, "canonical session identity is unavailable", http.StatusConflict)
165 return nil, session.SessionRef{}, false
166 }
167 return service.Query(), ref, true
168 }
169 // A remote tab renders its persisted first page before POST /resume
170 // activates the session, and a spectator reads history the foreground no
171 // longer owns: both are cold reads the store answers without a runtime.
172 requested = strings.TrimPrefix(requested, remoteSessionIDQueryPrefix)
173 cold := session.SessionRef{HostID: service.HostID(), SessionID: requested}
174 if _, err := service.Query().Stat(r.Context(), cold); err != nil {
175 http.Error(w, "session history is not bound to this runtime", http.StatusConflict)
176 return nil, session.SessionRef{}, false
177 }
178 return service.Query(), cold, true
179 }
180
181 func (s *Server) sessionHistoryPage(w http.ResponseWriter, r *http.Request) {
182 s.bindMu.Lock()
183 defer s.bindMu.Unlock()
184 query, ref, ok := s.canonicalSessionQuery(w, r)
185 if !ok {
186 return
187 }
188 limit := 0
189 if raw := r.URL.Query().Get("limit"); raw != "" {
190 if _, err := fmt.Sscan(raw, &limit); err != nil {
191 http.Error(w, "invalid history limit", http.StatusBadRequest)
192 return
193 }
194 }
195 page, err := query.HistoryPage(r.Context(), ref, r.URL.Query().Get("cursor"), limit)
196 if err != nil {
197 http.Error(w, err.Error(), http.StatusConflict)
198 return
199 }
200 w.Header().Set("Cache-Control", "no-store")
201 w.Header().Set("Content-Type", "application/json")
202 _ = json.NewEncoder(w).Encode(page)
203 }
204
205 func (s *Server) sessionHistoryContent(w http.ResponseWriter, r *http.Request) {
206 var request sessionHistoryContentRequest
207 if !transcriptRequest(w, r, &request) {
208 return
209 }
210 if request.Length <= 0 || request.Length > 1<<20 {
211 http.Error(w, "invalid content range", http.StatusBadRequest)
212 return
213 }
214 s.bindMu.Lock()
215 defer s.bindMu.Unlock()
216 query, ref, ok := s.canonicalSessionQuery(w, r)
217 if !ok {
218 return
219 }
220 data, err := query.ReadContent(r.Context(), ref, request.Ref, request.Offset, request.Length)
221 if err != nil {
222 http.Error(w, err.Error(), http.StatusConflict)
223 return
224 }
225 next := request.Offset + int64(len(data))
226 w.Header().Set("Cache-Control", "no-store")
227 w.Header().Set("Content-Type", "application/json")
228 _ = json.NewEncoder(w).Encode(sessionHistoryContentResponse{Data: base64.StdEncoding.EncodeToString(data), NextOffset: next, Done: next == request.Ref.Bytes})
229 }
230
231 func (s *Server) sessionHistorySearch(w http.ResponseWriter, r *http.Request) {
232 s.bindMu.Lock()
233 defer s.bindMu.Unlock()
234 query, ref, ok := s.canonicalSessionQuery(w, r)
235 if !ok {
236 return
237 }
238 textQuery := r.URL.Query().Get("q")
239 if len(textQuery) > 4096 {
240 http.Error(w, "history search query is too large", http.StatusBadRequest)
241 return
242 }
243 limit := 0
244 if raw := r.URL.Query().Get("limit"); raw != "" {
245 if _, err := fmt.Sscan(raw, &limit); err != nil {
246 http.Error(w, "invalid history search limit", http.StatusBadRequest)
247 return
248 }
249 }
250 page, err := query.SearchHistory(r.Context(), ref, textQuery, r.URL.Query().Get("cursor"), limit)
251 if err != nil {
252 http.Error(w, err.Error(), http.StatusConflict)
253 return
254 }
255 w.Header().Set("Cache-Control", "no-store")
256 w.Header().Set("Content-Type", "application/json")
257 _ = json.NewEncoder(w).Encode(page)
258 }
259
260 func (s *Server) sessionHistoryLocate(w http.ResponseWriter, r *http.Request) {
261 s.bindMu.Lock()
262 defer s.bindMu.Unlock()
263 query, ref, ok := s.canonicalSessionQuery(w, r)
264 if !ok {
265 return
266 }
267 var snapshot uint64
268 if raw := r.URL.Query().Get("snapshot"); raw != "" {
269 if _, err := fmt.Sscan(raw, &snapshot); err != nil {
270 http.Error(w, "invalid history snapshot", http.StatusBadRequest)
271 return
272 }
273 }
274 location, err := query.LocateMessage(r.Context(), ref, r.URL.Query().Get("messageId"), snapshot)
275 if err != nil {
276 http.Error(w, err.Error(), http.StatusConflict)
277 return
278 }
279 w.Header().Set("Cache-Control", "no-store")
280 w.Header().Set("Content-Type", "application/json")
281 _ = json.NewEncoder(w).Encode(location)
282 }
283
284 // errTranscriptCapabilityMissing lets a read decline an optional capability
285 // without colliding with a genuine read failure, which must stay a conflict.
286 var errTranscriptCapabilityMissing = errors.New("transcript capability is missing")
287
288 // transcriptRead binds each read to the controller that owns the referenced
289 // session: the foreground when the reference is absent or matches it, else
290 // the detached session holding that identity or path. A file mirror cannot
291 // claim a live event cursor and explicitly declines this protocol.
292 func (s *Server) transcriptRead(w http.ResponseWriter, r *http.Request, read func(control.TranscriptProjectionAPI) (any, error)) {
293 s.transcriptBoundRead(w, r, func(ctrl control.SessionAPI) (any, error) {
294 api, ok := ctrl.(control.TranscriptProjectionAPI)
295 if !ok {
296 return nil, errTranscriptCapabilityMissing
297 }
298 return read(api)
299 })
300 }
301
302 // transcriptBoundRead resolves the selected controller, enforces the session
303 // binding every transcript read shares, and encodes one JSON response. An
304 // unimplemented optional capability is reported as not implemented rather than
305 // silently answered with an empty page.
306 func (s *Server) transcriptBoundRead(w http.ResponseWriter, r *http.Request, read func(control.SessionAPI) (any, error)) {
307 s.bindMu.Lock()
308 raw := strings.TrimSpace(r.URL.Query().Get("session"))
309 if raw != "" && !strings.HasPrefix(raw, remoteSessionIDQueryPrefix) {
310 if resolved, err := s.resolveSessionPath(raw); err == nil && s.sessionMirrored(agent.CanonicalSessionPath(resolved)) {
311 s.bindMu.Unlock()
312 http.Error(w, "transcript projection is unavailable", http.StatusNotImplemented)
313 return
314 }
315 }
316 ctrl := s.resolveReadControllerLocked(raw)
317 if ctrl == nil {
318 s.bindMu.Unlock()
319 http.Error(w, "transcript session is not bound to this runtime", http.StatusConflict)
320 return
321 }
322 if ctrl == s.ctl() {
323 if path := agent.CanonicalSessionPath(ctrl.SessionPath()); path != "" && s.sessionMirrored(path) {
324 s.bindMu.Unlock()
325 http.Error(w, "transcript projection is unavailable", http.StatusNotImplemented)
326 return
327 }
328 }
329 s.bindMu.Unlock()
330 value, err := read(ctrl)
331 s.bindMu.Lock()
332 // A detached or identity-routed read cannot be verified by a foreground
333 // path comparison; re-resolving the same reference and comparing the
334 // controller covers every routing case.
335 current := s.resolveReadControllerLocked(raw) == ctrl
336 s.bindMu.Unlock()
337 if !current {
338 http.Error(w, "transcript runtime changed during read", http.StatusConflict)
339 return
340 }
341 if errors.Is(err, errTranscriptCapabilityMissing) {
342 http.Error(w, "transcript projection is unavailable", http.StatusNotImplemented)
343 return
344 }
345 if err != nil {
346 http.Error(w, err.Error(), http.StatusConflict)
347 return
348 }
349 w.Header().Set("Cache-Control", "no-store")
350 w.Header().Set("Content-Type", "application/json")
351 _ = json.NewEncoder(w).Encode(value)
352 }
353
354 // resolveReadControllerLocked returns the controller a transcript read
355 // targets. An empty reference selects the foreground; an identity or path
356 // reference selects the foreground when it matches, otherwise the detached
357 // session holding it. bindMu must be held; detachedMu nests inside it.
358 func (s *Server) resolveReadControllerLocked(raw string) control.SessionAPI {
359 foreground := s.ctl()
360 if raw == "" {
361 return foreground
362 }
363 if id, ok := strings.CutPrefix(raw, remoteSessionIDQueryPrefix); ok {
364 if controllerBoundToIdentity(foreground, id) {
365 return foreground
366 }
367 s.detachedMu.Lock()
368 defer s.detachedMu.Unlock()
369 for _, detached := range s.detached {
370 if controllerBoundToIdentity(detached.ctrl, id) {
371 return detached.ctrl
372 }
373 }
374 return nil
375 }
376 path := agent.CanonicalSessionPath(raw)
377 if resolved, err := s.resolveSessionPath(raw); err == nil {
378 path = agent.CanonicalSessionPath(resolved)
379 }
380 if agent.CanonicalSessionPath(foreground.SessionPath()) == path {
381 return foreground
382 }
383 s.detachedMu.Lock()
384 defer s.detachedMu.Unlock()
385 for _, detached := range s.detached {
386 if agent.CanonicalSessionPath(detached.ctrl.SessionPath()) == path {
387 return detached.ctrl
388 }
389 }
390 return nil
391 }
392
393 // remoteSessionIDQueryPrefix marks a session reference as an identity ID
394 // rather than a legacy transcript path; it matches the desktop's routing
395 // prefix for exclusive identity sessions.
396 const remoteSessionIDQueryPrefix = "session-id:"
397
398 // controllerBoundToIdentity reports whether ctrl currently runs the exclusive
399 // identity session the caller referenced.
400 func controllerBoundToIdentity(ctrl control.SessionAPI, id string) bool {
401 ref, ok := ctrl.(interface {
402 SessionRef() (session.SessionRef, bool)
403 })
404 if !ok {
405 return false
406 }
407 bound, has := ref.SessionRef()
408 return has && bound.SessionID == id
409 }
410
411 func (s *Server) transcriptFollow(w http.ResponseWriter, r *http.Request) {
412 var req transcript.FollowRequest
413 if !transcriptRequest(w, r, &req) {
414 return
415 }
416 s.transcriptBoundRead(w, r, func(ctrl control.SessionAPI) (any, error) {
417 api, ok := ctrl.(control.TranscriptFollowAPI)
418 if !ok {
419 return nil, errTranscriptCapabilityMissing
420 }
421 return api.TranscriptFollow(r.Context(), req)
422 })
423 }
424
425 func transcriptRequest(w http.ResponseWriter, r *http.Request, dst any) bool {
426 encoded := r.URL.Query().Get("request")
427 if encoded == "" {
428 return true
429 }
430 if len(encoded) > 8192 || json.Unmarshal([]byte(encoded), dst) != nil {
431 http.Error(w, "invalid transcript request", http.StatusBadRequest)
432 return false
433 }
434 return true
435 }
436
437 func (s *Server) transcriptSnapshot(w http.ResponseWriter, r *http.Request) {
438 var req transcript.PageRequest
439 if !transcriptRequest(w, r, &req) {
440 return
441 }
442 s.transcriptRead(w, r, func(api control.TranscriptProjectionAPI) (any, error) { return api.TranscriptSnapshot(req) })
443 }
444
445 func (s *Server) transcriptContent(w http.ResponseWriter, r *http.Request) {
446 var req transcript.ContentRequest
447 if !transcriptRequest(w, r, &req) {
448 return
449 }
450 s.transcriptRead(w, r, func(api control.TranscriptProjectionAPI) (any, error) { return api.TranscriptContent(req) })
451 }
452
453 func (s *Server) transcriptOutline(w http.ResponseWriter, r *http.Request) {
454 var req transcript.OutlineRequest
455 if !transcriptRequest(w, r, &req) {
456 return
457 }
458 s.transcriptBoundRead(w, r, func(ctrl control.SessionAPI) (any, error) {
459 api, ok := ctrl.(control.TranscriptOutlineAPI)
460 if !ok {
461 return nil, errTranscriptCapabilityMissing
462 }
463 return api.TranscriptOutline(req)
464 })
465 }
466
467 func (s *Server) transcriptReplay(w http.ResponseWriter, r *http.Request) {
468 var req control.TranscriptReplayRequest
469 if !transcriptRequest(w, r, &req) {
470 return
471 }
472 s.transcriptRead(w, r, func(api control.TranscriptProjectionAPI) (any, error) { return api.TranscriptReplay(req) })
473 }
474
474 lines GO