返回 DeepSeek-Reasonix
api.go
根目录 / internal / session / api.go
1 package session
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "path/filepath"
9 "sort"
10 "strings"
11 "sync"
12 "time"
13 )
14
15 // VisitCommits streams the complete durable commit prefix from one canonical
16 // session directory. It is intended for exports and diagnostics that must not
17 // materialize the cumulative event log.
18 func VisitCommits(ctx context.Context, dir string, visit func(Commit) error) error {
19 if err := ctx.Err(); err != nil {
20 return err
21 }
22 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
23 if err != nil {
24 return err
25 }
26 file, err := os.Open(logPathForManifest(dir, manifest))
27 if os.IsNotExist(err) {
28 return nil
29 }
30 if err != nil {
31 return err
32 }
33 defer file.Close()
34 var visitErr error
35 adapter := func(_ int64, commit Commit) bool {
36 if visit != nil {
37 visitErr = visit(commit)
38 }
39 return visitErr == nil
40 }
41 if manifest.Codec == Codec {
42 err = scanV4CommitFile(ctx, file, 0, 1, contentStoreForSessionDir(dir), nil, adapter)
43 } else {
44 err = scanCommitFileCodecBoundaries(ctx, file, 0, 1, manifest.Codec, nil, func(start, _ int64, commit Commit) bool { return adapter(start, commit) })
45 }
46 return errors.Join(err, visitErr)
47 }
48
49 // AccessMode separates cold readers from the single leased writer. Read-only
50 // access never repairs, migrates, truncates, or advances writer generation.
51 type AccessMode string
52
53 const (
54 ReadOnly AccessMode = "read"
55 ReadWrite AccessMode = "write"
56 )
57
58 // SessionKind records what created a store; empty is an ordinary conversation.
59 type SessionKind string
60
61 // SessionKindHeadlessRun marks a one-shot `reasonix run`/`-p` store. It stays
62 // resumable, but it is not a conversation a sidebar should offer.
63 const SessionKindHeadlessRun SessionKind = "headless-run"
64
65 type CreateOptions struct {
66 SessionID string
67 CWD string
68 ParentSessionID string
69 Origin SessionOrigin
70 Kind SessionKind
71 }
72
73 type SessionInfo struct {
74 SessionID string
75 Ref SessionRef
76 Codec string
77 Title string
78 TitleSequence uint64
79 ModelRef string
80 ModelIdentity string
81 Turns int
82 CreatedAt time.Time
83 UpdatedAt time.Time
84 EventSequence uint64
85 ResultSequence uint64
86 Preview string
87 MetadataStatus string
88 CWD string
89 ParentSessionID string
90 Origin SessionOrigin
91 Kind SessionKind
92 Path string
93 Error string
94 }
95
96 type SessionPage struct {
97 Sessions []SessionInfo
98 NextCursor string
99 }
100
101 type EventPage struct {
102 Commits []Commit
103 Next uint64
104 Truncated bool
105 }
106
107 // eventPageReader is the paged durable read surface shared by the physical
108 // handle and a cold Session, so catalog rebuilds never depend on replaying a
109 // whole in-memory log.
110 type eventPageReader interface {
111 Read(context.Context, uint64, int) (EventPage, error)
112 }
113
114 // SessionHandle is the physical persistence contract. It reads and writes
115 // bytes for one session identity and owns the writer lease; it holds no
116 // projection, operation table, or accepted commit list.
117 type SessionHandle interface {
118 ID() string
119 Manifest() Manifest
120 Read(context.Context, uint64, int) (EventPage, error)
121 Append(context.Context, []Commit) error
122 Sync(context.Context) (DurableReceipt, error)
123 Close(context.Context) error
124 }
125
126 // WritableSessionHandle is the leased physical handle. Byte-level maintenance
127 // operations such as fork, export, and interrupted-turn recovery are Session
128 // operations because they read or extend the in-memory log; this contract only
129 // distinguishes a leased writer from a cold reader.
130 type WritableSessionHandle interface {
131 SessionHandle
132 SessionID() string
133 }
134
135 type SessionPersistence interface {
136 Create(CreateOptions) (*Session, error)
137 Open(sessionID string, mode AccessMode) (*Session, error)
138 Stat(context.Context, string) (SessionInfo, error)
139 List(context.Context, string, int) (SessionPage, error)
140 }
141
142 // FilesystemPersistence owns a versioned sessions-v4 root.
143 type FilesystemPersistence struct {
144 Root string
145 metadataReads sessionMetadataReads
146 }
147
148 func NewFilesystemPersistence(root string) *FilesystemPersistence {
149 return &FilesystemPersistence{Root: filepath.Clean(root)}
150 }
151
152 // RootForLegacyDir maps a host's legacy transcript catalog to the sibling
153 // final-format store. The mapping lives in the persistence package so Boot and
154 // controllers never derive a v3 identity from a transcript path.
155 func RootForLegacyDir(sessionDir string) string {
156 dir := filepath.Clean(strings.TrimSpace(sessionDir))
157 if dir == "." || dir == "" {
158 return ""
159 }
160 if filepath.Base(dir) == "sessions" {
161 return filepath.Join(filepath.Dir(dir), "sessions-v4")
162 }
163 return filepath.Join(dir, "sessions-v4")
164 }
165
166 func (p *FilesystemPersistence) Create(options CreateOptions) (*Session, error) {
167 id := strings.TrimSpace(options.SessionID)
168 if id == "" {
169 id = randomID()
170 }
171 if err := validateSessionID(id); err != nil {
172 return nil, err
173 }
174 dir, err := p.sessionDir(id, false)
175 if err != nil {
176 return nil, err
177 }
178 options.SessionID = id
179 if options.Kind != "" && options.Kind != SessionKindHeadlessRun {
180 return nil, fmt.Errorf("session: unsupported session kind %q", options.Kind)
181 }
182 header, err := headerForCreate(options)
183 if err != nil {
184 return nil, err
185 }
186 return createWithOptions(dir, id, OpenOptions{ExternalHistory: true}, header, options.Kind)
187 }
188
189 func (p *FilesystemPersistence) Open(sessionID string, mode AccessMode) (*Session, error) {
190 return p.OpenContext(context.Background(), sessionID, mode)
191 }
192
193 func (p *FilesystemPersistence) OpenContext(ctx context.Context, sessionID string, mode AccessMode) (*Session, error) {
194 if err := ctx.Err(); err != nil {
195 return nil, err
196 }
197 id := strings.TrimSpace(sessionID)
198 if err := validateSessionID(id); err != nil {
199 return nil, err
200 }
201 dir, err := p.sessionDir(id, true)
202 if err != nil {
203 return nil, err
204 }
205 if _, _, err := readSessionHeader(dir, id); err != nil {
206 return nil, err
207 }
208 if mode == ReadOnly {
209 return openReadSession(dir, id, filepath.Join(p.Root, ".query-cache", filepath.Base(id)))
210 }
211 if mode != ReadWrite {
212 return nil, fmt.Errorf("session: unsupported access mode %q", mode)
213 }
214 if _, err := os.Stat(dir); os.IsNotExist(err) {
215 return nil, fmt.Errorf("%w: %s", ErrSessionNotFound, id)
216 } else if err != nil {
217 return nil, err
218 }
219 return OpenWithOptions(dir, id, OpenOptions{ExternalHistory: true, Context: ctx})
220 }
221
222 func (p *FilesystemPersistence) Stat(ctx context.Context, sessionID string) (SessionInfo, error) {
223 return p.cachedSessionInfo(ctx, sessionID)
224 }
225
226 func (p *FilesystemPersistence) sessionDir(id string, mustExist bool) (string, error) {
227 if err := validateSessionID(id); err != nil {
228 return "", err
229 }
230 if !mustExist {
231 if err := os.MkdirAll(p.Root, 0o700); err != nil {
232 return "", err
233 }
234 }
235 root, err := os.OpenRoot(p.Root)
236 if os.IsNotExist(err) {
237 return "", fmt.Errorf("%w: %s", ErrSessionNotFound, id)
238 }
239 if err != nil {
240 return "", err
241 }
242 defer root.Close()
243 // Root.Lstat rejects traversal and follows the platform's reparse-point
244 // boundary rules. The single-segment validation above also keeps lock and
245 // cache names portable on Windows.
246 info, err := root.Lstat(id)
247 if os.IsNotExist(err) && !mustExist {
248 return filepath.Join(p.Root, id), nil
249 }
250 if os.IsNotExist(err) {
251 return "", fmt.Errorf("%w: %s", ErrSessionNotFound, id)
252 }
253 if err != nil {
254 return "", err
255 }
256 if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() {
257 return "", fmt.Errorf("session: session identity %q is not a confined directory", id)
258 }
259 // Resolve the physical path through the rooted handle. This keeps protocol
260 // input out of the path propagated through the storage stack.
261 child, err := root.OpenRoot(id)
262 if err != nil {
263 return "", err
264 }
265 dir := child.Name()
266 if err := child.Close(); err != nil {
267 return "", err
268 }
269 return dir, nil
270 }
271
272 func (p *FilesystemPersistence) List(ctx context.Context, cursor string, limit int) (SessionPage, error) {
273 if err := ctx.Err(); err != nil {
274 return SessionPage{}, err
275 }
276 if limit == 0 {
277 limit = 50
278 }
279 if limit < 1 || limit > 100 {
280 return SessionPage{}, fmt.Errorf("session: list limit must be 1..100")
281 }
282 entries, err := os.ReadDir(p.Root)
283 if os.IsNotExist(err) {
284 return SessionPage{Sessions: []SessionInfo{}}, nil
285 }
286 if err != nil {
287 return SessionPage{}, err
288 }
289 ids := make([]string, 0, len(entries))
290 for _, entry := range entries {
291 if err := ctx.Err(); err != nil {
292 return SessionPage{}, err
293 }
294 if entry.IsDir() && !strings.HasPrefix(entry.Name(), ".") && entry.Name() > cursor {
295 ids = append(ids, entry.Name())
296 }
297 }
298 sort.Strings(ids)
299 page := SessionPage{Sessions: []SessionInfo{}}
300 for _, id := range ids {
301 info, statErr := p.Stat(ctx, id)
302 if statErr != nil {
303 info = SessionInfo{SessionID: id, Path: filepath.Join(p.Root, id), Error: statErr.Error()}
304 }
305 if len(page.Sessions) == limit {
306 page.NextCursor = page.Sessions[len(page.Sessions)-1].SessionID
307 break
308 }
309 page.Sessions = append(page.Sessions, info)
310 }
311 return page, nil
312 }
313
314 func validateSessionID(id string) error {
315 id = strings.TrimSpace(id)
316 if len(id) == 0 || len(id) > 255 || strings.HasPrefix(id, ".") || strings.HasSuffix(id, ".") ||
317 !filepath.IsLocal(id) || id == "." || filepath.Base(id) != id || strings.ContainsAny(id, `/\\<>:"|?*`) {
318 return fmt.Errorf("session: invalid session id %q", id)
319 }
320 for _, char := range id {
321 if char < 0x20 {
322 return fmt.Errorf("session: invalid session id %q", id)
323 }
324 }
325 base := strings.ToUpper(strings.SplitN(id, ".", 2)[0])
326 if base == "CON" || base == "PRN" || base == "AUX" || base == "NUL" ||
327 (len(base) == 4 && (strings.HasPrefix(base, "COM") || strings.HasPrefix(base, "LPT")) && base[3] >= '1' && base[3] <= '9') {
328 return fmt.Errorf("session: reserved session id %q", id)
329 }
330 return nil
331 }
332
333 // ValidateSessionID applies the canonical session storage identity rules.
334 // Callers that accept compatibility routes must validate the extracted ID
335 // before deciding whether the input names a SessionRef or a legacy path.
336 func ValidateSessionID(id string) error {
337 return validateSessionID(id)
338 }
339
340 type readHandle struct {
341 id string
342 dir string
343 cacheDir string
344 manifest Manifest
345 mu sync.Mutex
346 closed bool
347 }
348
349 // openReadSession returns a cold Session backed only by the durable prefix. It
350 // performs no replay and takes no writer lease: paged reads remain the only way
351 // to consume it, which is what keeps catalog and history queries cheap.
352 func openReadSession(dir, id string, cacheDirs ...string) (*Session, error) {
353 handle, err := openReadHandle(dir, id, cacheDirs...)
354 if err != nil {
355 return nil, err
356 }
357 return newReadSession(handle), nil
358 }
359
360 func openReadHandle(dir, id string, cacheDirs ...string) (*readHandle, error) {
361 cacheDir := dir
362 if len(cacheDirs) > 0 && strings.TrimSpace(cacheDirs[0]) != "" {
363 cacheDir = cacheDirs[0]
364 }
365 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
366 if err != nil {
367 if os.IsNotExist(err) {
368 return nil, fmt.Errorf("%w: %s", ErrSessionNotFound, id)
369 }
370 return nil, err
371 }
372 if manifest.SessionID != id {
373 return nil, fmt.Errorf("session: manifest belongs to %q", manifest.SessionID)
374 }
375 return &readHandle{id: id, dir: dir, cacheDir: cacheDir, manifest: manifest}, nil
376 }
377
378 func (h *readHandle) ID() string { return h.id }
379
380 func (h *readHandle) Manifest() Manifest { return h.manifest }
381
382 // Dir reports the directory this cold reader opened. A fork from a session with
383 // no live runtime still has to locate the parent's owned files, and that must
384 // not require acquiring the writer lease the cold reader deliberately avoids.
385 func (h *readHandle) Dir() string {
386 if h == nil {
387 return ""
388 }
389 return h.dir
390 }
391
392 func (h *readHandle) Read(ctx context.Context, offset uint64, limit int) (EventPage, error) {
393 if h == nil {
394 return EventPage{}, os.ErrClosed
395 }
396 h.mu.Lock()
397 closed := h.closed
398 dir, cacheDir := h.dir, h.cacheDir
399 h.mu.Unlock()
400 if closed {
401 return EventPage{}, os.ErrClosed
402 }
403 return readCommitPageWithCache(ctx, dir, cacheDir, offset, limit)
404 }
405
406 func (h *readHandle) Append(context.Context, []Commit) error {
407 return ErrReadOnly
408 }
409
410 func (h *readHandle) Sync(context.Context) (DurableReceipt, error) {
411 if h == nil {
412 return DurableReceipt{}, os.ErrClosed
413 }
414 h.mu.Lock()
415 closed := h.closed
416 dir := h.dir
417 h.mu.Unlock()
418 if closed {
419 return DurableReceipt{}, os.ErrClosed
420 }
421 sequence, err := lastDurableSequence(dir)
422 return DurableReceipt{DurableSequence: sequence}, err
423 }
424
425 func (h *readHandle) Close(context.Context) error {
426 if h == nil {
427 return nil
428 }
429 h.mu.Lock()
430 h.closed = true
431 h.mu.Unlock()
432 return nil
433 }
434
435 func (s *Store) Read(ctx context.Context, offset uint64, limit int) (EventPage, error) {
436 if s == nil {
437 return EventPage{}, os.ErrClosed
438 }
439 s.mu.Lock()
440 closed := s.closed
441 dir := s.dir
442 s.mu.Unlock()
443 if closed {
444 return EventPage{}, os.ErrClosed
445 }
446 return readCommitPage(ctx, dir, offset, limit)
447 }
448
449 func readCommitPage(ctx context.Context, dir string, offset uint64, limit int) (EventPage, error) {
450 return readCommitPageWithCache(ctx, dir, dir, offset, limit)
451 }
452
453 func readCommitPageWithCache(ctx context.Context, dir, cacheDir string, offset uint64, limit int) (EventPage, error) {
454 if err := ctx.Err(); err != nil {
455 return EventPage{}, err
456 }
457 if limit == 0 {
458 limit = 100
459 }
460 if limit < 1 || limit > 1000 {
461 return EventPage{}, fmt.Errorf("session: read limit must be 1..1000 commits")
462 }
463 index, err := loadOrBuildSparseIndex(ctx, dir, cacheDir)
464 if err != nil {
465 return EventPage{}, err
466 }
467 page := EventPage{Commits: []Commit{}}
468 if index.LastSequence <= offset {
469 return page, nil
470 }
471 checkpoint := index.checkpoint(offset)
472 manifest, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
473 if err != nil {
474 return EventPage{}, err
475 }
476 file, err := os.Open(logPathForManifest(dir, manifest))
477 if err != nil {
478 return EventPage{}, err
479 }
480 defer file.Close()
481 visit := func(_ int64, commit Commit) bool {
482 if ctx.Err() != nil {
483 return false
484 }
485 if commit.LastSequence() <= offset {
486 return true
487 }
488 if len(page.Commits) == limit {
489 page.Truncated = true
490 return false
491 }
492 page.Commits = append(page.Commits, commit)
493 page.Next = commit.LastSequence()
494 return true
495 }
496 if manifest.Codec == Codec {
497 err = scanV4CommitFile(ctx, file, checkpoint.Offset, checkpoint.FirstSequence, contentStoreForSessionDir(dir), nil, visit)
498 } else {
499 err = scanCommitFileCodecBoundaries(ctx, file, checkpoint.Offset, checkpoint.FirstSequence, manifest.Codec, nil, func(start, _ int64, commit Commit) bool { return visit(start, commit) })
500 }
501 if err != nil {
502 return EventPage{}, err
503 }
504 if err := ctx.Err(); err != nil {
505 return EventPage{}, err
506 }
507 return page, nil
508 }
509
510 func lastDurableSequence(dir string) (uint64, error) {
511 return lastDurableSequenceWithCache(dir, dir)
512 }
513
514 func lastDurableSequenceWithCache(dir, cacheDir string) (uint64, error) {
515 index, err := loadOrBuildSparseIndex(context.Background(), dir, cacheDir)
516 return index.LastSequence, err
517 }
518
519 var _ SessionPersistence = (*FilesystemPersistence)(nil)
520 var _ SessionHandle = (*Store)(nil)
521 var _ SessionHandle = (*readHandle)(nil)
522
522 lines GO