返回 DeepSeek-Reasonix
store.go
1 package sessioninbox
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "path/filepath"
9 "strings"
10 "sync"
11 "time"
12
13 filelock "reasonix/internal/identitylock"
14 "reasonix/internal/store"
15 )
16
17 const (
18 manifestName = "manifest.json"
19 blobsDirName = "blobs"
20 quarantineName = "quarantine"
21 blobSuffix = ".json"
22 diskLockName = "transaction.lock"
23 diskLockWait = 5 * time.Second
24 maxManifestBytes = 8 << 20
25 )
26
27 // Store is the transactional durable inbox for one session locator.
28 // Disk I/O runs under store.mu only; callers must not hold Controller locks.
29 type Store struct {
30 mu sync.Mutex
31 dir string
32 session string // logical locator: transcript path or canonical session route
33 runID string
34 limits Limits
35 man *manifest
36 readonly bool
37 closed bool
38 // listeners receive revision bumps after durable commits (non-blocking).
39 listeners []func(InboxSnapshot)
40 }
41
42 // Open binds a Store to the session's inbox directory. The directory is created
43 // for its transaction lock; body blobs remain lazy. Cross-process recovery
44 // marks uncertain items and pauses.
45 func Open(sessionPath string, limits Limits) (*Store, error) {
46 return OpenAt(sessionPath, store.SessionInboxDir(sessionPath), limits)
47 }
48
49 // OpenAt separates the logical session locator from its physical inbox directory.
50 // Canonical runtimes have an immutable identity, not a transcript file path.
51 func OpenAt(sessionPath, dir string, limits Limits) (*Store, error) {
52 sessionPath = strings.TrimSpace(sessionPath)
53 if sessionPath == "" {
54 return nil, fmt.Errorf("sessioninbox: empty session path")
55 }
56 if strings.TrimSpace(dir) == "" {
57 return nil, fmt.Errorf("sessioninbox: empty inbox directory")
58 }
59 s := &Store{
60 dir: dir,
61 session: sessionPath,
62 runID: ProcessRunID(),
63 limits: limits.withDefaults(),
64 man: emptyManifest(ProcessRunID()),
65 }
66 if err := s.loadOrInit(); err != nil {
67 return nil, err
68 }
69 return s, nil
70 }
71
72 // Dir returns the on-disk inbox directory.
73 func (s *Store) Dir() string {
74 if s == nil {
75 return ""
76 }
77 s.mu.Lock()
78 defer s.mu.Unlock()
79 return s.dir
80 }
81
82 // SessionPath returns the bound logical session locator.
83 func (s *Store) SessionPath() string {
84 if s == nil {
85 return ""
86 }
87 s.mu.Lock()
88 defer s.mu.Unlock()
89 return s.session
90 }
91
92 // Rebind moves the store to a new session path without copying future work
93 // (used after rename migration that already relocated the directory).
94 func (s *Store) Rebind(sessionPath string) error {
95 if s == nil {
96 return ErrClosed
97 }
98 sessionPath = strings.TrimSpace(sessionPath)
99 if sessionPath == "" {
100 return fmt.Errorf("sessioninbox: empty session path")
101 }
102 s.mu.Lock()
103 defer s.mu.Unlock()
104 if s.closed {
105 return ErrClosed
106 }
107 s.session = sessionPath
108 s.dir = store.SessionInboxDir(sessionPath)
109 release, err := s.beginDiskTransactionLocked()
110 if err != nil {
111 return err
112 }
113 release()
114 return nil
115 }
116
117 // Close seals the store. Further mutations fail with ErrClosed.
118 func (s *Store) Close() {
119 if s == nil {
120 return
121 }
122 s.mu.Lock()
123 s.closed = true
124 s.mu.Unlock()
125 }
126
127 // OnChange registers a non-blocking snapshot listener.
128 func (s *Store) OnChange(fn func(InboxSnapshot)) {
129 if s == nil || fn == nil {
130 return
131 }
132 s.mu.Lock()
133 s.listeners = append(s.listeners, fn)
134 s.mu.Unlock()
135 }
136
137 func (s *Store) loadOrInit() error {
138 s.mu.Lock()
139 defer s.mu.Unlock()
140 release, err := s.beginDiskTransactionLocked()
141 if err != nil {
142 return err
143 }
144 release()
145 return nil
146 }
147
148 // beginDiskTransactionLocked serializes every manifest read/modify/write with
149 // other Store instances and processes, then refreshes the in-memory snapshot.
150 // The caller must hold s.mu and call the returned release function.
151 func (s *Store) beginDiskTransactionLocked() (func(), error) {
152 if err := ensurePrivateDir(s.dir); err != nil {
153 return nil, fmt.Errorf("sessioninbox: create inbox directory: %w", err)
154 }
155 ctx, cancel := context.WithTimeout(context.Background(), diskLockWait)
156 defer cancel()
157 release, err := filelock.Acquire(ctx, filepath.Join(s.dir, diskLockName))
158 if err != nil {
159 return nil, fmt.Errorf("sessioninbox: acquire disk lock: %w", err)
160 }
161 if err := s.loadOrInitLocked(); err != nil {
162 release()
163 return nil, err
164 }
165 return release, nil
166 }
167
168 func (s *Store) loadOrInitLocked() error {
169 path := filepath.Join(s.dir, manifestName)
170 data, err := readRegularFile(path, maxManifestBytes)
171 if errors.Is(err, os.ErrNotExist) {
172 s.man = emptyManifest(s.runID)
173 s.readonly = false
174 return nil
175 }
176 if err != nil {
177 return fmt.Errorf("sessioninbox: read manifest: %w", err)
178 }
179 man, err := decodeManifest(data)
180 if err != nil {
181 // Corrupt manifest → quarantine, salvage orphan blobs as uncertain
182 // items, pause for user inspection. Never present "0 recovered".
183 _ = s.quarantineFileLocked(path, "manifest-corrupt")
184 salvaged := s.salvageOrphanBlobsLocked()
185 s.man = emptyManifest(s.runID)
186 s.man.Paused = true
187 s.man.Recovered = true
188 s.man.RecoveredN = len(salvaged)
189 s.man.Items = salvaged
190 return s.commitManifestLocked(s.man)
191 }
192 if man.SchemaVersion > SchemaVersion {
193 s.man = man
194 s.readonly = true
195 s.man.Paused = true
196 return nil
197 }
198 migrated := man.SchemaVersion < SchemaVersion
199 if migrated {
200 for key, id := range man.Idempotency {
201 if man.IdempotencyHashes[key] != "" {
202 continue
203 }
204 meta, ok := man.item(id)
205 if !ok {
206 return fmt.Errorf("sessioninbox: migrate idempotency target: %w", ErrNotFound)
207 }
208 env, err := s.readBlobLocked(blobNameFor(meta), meta.Checksum)
209 if err != nil {
210 return fmt.Errorf("sessioninbox: migrate idempotency body: %w", err)
211 }
212 hash, err := idempotencyRequestHash(env)
213 if err != nil {
214 return fmt.Errorf("sessioninbox: migrate idempotency hash: %w", err)
215 }
216 man.IdempotencyHashes[key] = hash
217 }
218 man.SchemaVersion = SchemaVersion
219 }
220 // Cross-process recovery: another run left in-flight items.
221 recovered := 0
222 if man.RunID != "" && man.RunID != s.runID {
223 for i := range man.Items {
224 switch man.Items[i].State {
225 case StateRunning, StateSteerAccepted, StateSteerConsumed:
226 man.Items[i].State = StateUncertain
227 man.Items[i].UpdatedAt = time.Now().UTC()
228 recovered++
229 case StateQueued, StateBlocked, StateUncertain:
230 recovered++
231 }
232 }
233 if recovered > 0 || len(man.Items) > 0 {
234 man.Paused = true
235 man.Recovered = true
236 man.RecoveredN = recovered
237 }
238 }
239 man.RunID = s.runID
240 s.man = man
241 s.readonly = false
242 if recovered > 0 || migrated {
243 return s.commitManifestLocked(man)
244 }
245 // GC orphan blobs without holding callers longer than needed.
246 s.gcOrphansLocked()
247 return nil
248 }
249
250 // Snapshot returns a copy of current metadata.
251 func (s *Store) Snapshot() InboxSnapshot {
252 if s == nil {
253 return InboxSnapshot{}
254 }
255 s.mu.Lock()
256 defer s.mu.Unlock()
257 if release, err := s.beginDiskTransactionLocked(); err == nil {
258 release()
259 }
260 return s.snapshotLocked()
261 }
262
263 // CachedSnapshot returns the Store's current in-memory metadata without taking
264 // the cross-process disk lock. It is for owner-local admission decisions that
265 // must not add disk-lock latency; Snapshot remains the authoritative refresh.
266 func (s *Store) CachedSnapshot() InboxSnapshot {
267 if s == nil {
268 return InboxSnapshot{}
269 }
270 s.mu.Lock()
271 defer s.mu.Unlock()
272 return s.snapshotLocked()
273 }
274
275 // TryFreshSnapshot reads current metadata from disk without waiting for either
276 // the Store mutex or the cross-process transaction lock. It does not perform
277 // recovery, migration, cleanup, or any other durable mutation. Callers making
278 // latency-sensitive admission decisions should treat any error conservatively.
279 func (s *Store) TryFreshSnapshot() (InboxSnapshot, error) {
280 if s == nil {
281 return InboxSnapshot{}, ErrClosed
282 }
283 if !s.mu.TryLock() {
284 return InboxSnapshot{}, ErrSnapshotBusy
285 }
286 defer s.mu.Unlock()
287 if s.closed {
288 return InboxSnapshot{}, ErrClosed
289 }
290 if err := validatePrivateDir(s.dir); err != nil {
291 return InboxSnapshot{}, fmt.Errorf("sessioninbox: validate inbox directory: %w", err)
292 }
293 release, err := filelock.TryAcquire(filepath.Join(s.dir, diskLockName))
294 if err != nil {
295 if errors.Is(err, filelock.ErrHeld) {
296 return InboxSnapshot{}, ErrSnapshotBusy
297 }
298 return InboxSnapshot{}, fmt.Errorf("sessioninbox: acquire disk lock: %w", err)
299 }
300 defer release()
301 data, err := readRegularFile(filepath.Join(s.dir, manifestName), maxManifestBytes)
302 if errors.Is(err, os.ErrNotExist) {
303 return s.snapshotLocked(), nil
304 }
305 if err != nil {
306 return InboxSnapshot{}, fmt.Errorf("sessioninbox: read manifest: %w", err)
307 }
308 man, err := decodeManifest(data)
309 if err != nil {
310 return InboxSnapshot{}, fmt.Errorf("sessioninbox: decode manifest: %w", err)
311 }
312 return s.snapshotFromManifestLocked(man, man.SchemaVersion > SchemaVersion), nil
313 }
314
315 func (s *Store) snapshotLocked() InboxSnapshot {
316 m := s.man
317 if m == nil {
318 m = emptyManifest(s.runID)
319 }
320 return s.snapshotFromManifestLocked(m, s.readonly)
321 }
322
323 func (s *Store) snapshotFromManifestLocked(m *manifest, readonly bool) InboxSnapshot {
324 items := append([]InboxItemMeta(nil), m.Items...)
325 return InboxSnapshot{
326 SchemaVersion: m.SchemaVersion,
327 Revision: m.Revision,
328 Paused: m.Paused,
329 Recovered: m.Recovered,
330 RecoveredN: m.RecoveredN,
331 Readonly: readonly,
332 RunID: m.RunID,
333 SessionPath: s.session,
334 Items: items,
335 Capacity: Capacity{
336 Items: len(items),
337 MaxItems: s.limits.MaxItems,
338 Bytes: m.totalBytes(),
339 MaxBytes: s.limits.MaxTotalBytes,
340 MaxItemBytes: s.limits.MaxItemBytes,
341 },
342 }
343 }
344
345 // Enqueue durably appends an item. Only returns a receipt after blob+manifest
346 // commit succeed. Idempotent keys return the original item.
347 func (s *Store) Enqueue(req EnqueueRequest) (InboxReceipt, error) {
348 if s == nil {
349 return InboxReceipt{}, ErrClosed
350 }
351 env := completeEnqueueEnvelope(req.Envelope)
352 hasInvocation := env.Invocation != nil || len(env.Invocations) > 0
353 if strings.TrimSpace(env.SubmitText) == "" && strings.TrimSpace(env.DisplayText) == "" && strings.TrimSpace(env.RawText) == "" && !hasInvocation {
354 return InboxReceipt{}, ErrEmpty
355 }
356 intent := req.Intent
357 if intent != IntentSteer {
358 intent = IntentFollowup
359 }
360 idem := strings.TrimSpace(firstNonEmpty(req.Idempotency, env.Idempotency))
361 if idem != "" && !validIdempotencyKey(idem) {
362 return InboxReceipt{}, fmt.Errorf("sessioninbox: invalid idempotency key")
363 }
364 source := strings.TrimSpace(firstNonEmpty(req.Source, env.Source))
365
366 blobBytes, checksum, byteSize, err := encodeEnvelope(env)
367 if err != nil {
368 return InboxReceipt{}, err
369 }
370 requestHash, err := idempotencyRequestHash(env)
371 if err != nil {
372 return InboxReceipt{}, err
373 }
374
375 s.mu.Lock()
376 defer s.mu.Unlock()
377 release, err := s.beginDiskTransactionLocked()
378 if err != nil {
379 return InboxReceipt{}, err
380 }
381 defer release()
382 if s.closed {
383 return InboxReceipt{}, ErrClosed
384 }
385 if s.readonly {
386 return InboxReceipt{}, ErrSchemaReadonly
387 }
388 if receipt, found, err := s.idempotentReceiptLocked(idem, requestHash); err != nil || found {
389 return receipt, err
390 }
391 if byteSize > s.limits.MaxItemBytes {
392 return InboxReceipt{}, ErrItemTooLarge
393 }
394 if len(s.man.Items) >= s.limits.MaxItems {
395 return InboxReceipt{}, ErrCapacityItems
396 }
397 if s.man.totalBytes()+byteSize > s.limits.MaxTotalBytes {
398 return InboxReceipt{}, ErrCapacityBytes
399 }
400
401 id := newRandomID()
402 blobName := id
403 now := time.Now().UTC()
404 meta := InboxItemMeta{
405 ID: id,
406 SessionID: firstNonEmpty(req.SessionID, agentBranchID(s.session)),
407 Intent: intent,
408 State: StateQueued,
409 Revision: s.man.Revision + 1,
410 BlobName: blobName,
411 Source: source,
412 CreatedAt: now,
413 UpdatedAt: now,
414 Preview: PreviewText(env.DisplayText, DefaultPreviewRunes),
415 ByteSize: byteSize,
416 Checksum: checksum,
417 Idempotency: idem,
418 Refs: refSummaries(env.Refs),
419 RunID: s.runID,
420 }
421
422 // Transaction: write blob → commit manifest → receipt.
423 if err := s.writeBlobLocked(blobName, blobBytes); err != nil {
424 return InboxReceipt{}, err
425 }
426 next := s.man.clone()
427 next.Items = append(next.Items, meta)
428 bindIdempotency(next, idem, id, requestHash)
429 if err := s.commitManifestLocked(next); err != nil {
430 s.removeBlobLocked(blobName)
431 return InboxReceipt{}, err
432 }
433 snap := s.snapshotLocked()
434 s.notifyLocked(snap)
435 return InboxReceipt{
436 ItemID: id,
437 Disposition: DispositionQueuedFollowup,
438 Position: len(next.Items),
439 Paused: next.Paused,
440 Capacity: snap.Capacity,
441 }, nil
442 }
443
444 // ReadItem loads a full PromptEnvelope by ID.
445 func (s *Store) ReadItem(id string) (InboxItemMeta, PromptEnvelope, error) {
446 if s == nil {
447 return InboxItemMeta{}, PromptEnvelope{}, ErrClosed
448 }
449 id = strings.TrimSpace(id)
450 s.mu.Lock()
451 defer s.mu.Unlock()
452 release, err := s.beginDiskTransactionLocked()
453 if err != nil {
454 return InboxItemMeta{}, PromptEnvelope{}, err
455 }
456 defer release()
457 if s.closed {
458 return InboxItemMeta{}, PromptEnvelope{}, ErrClosed
459 }
460 meta, ok := s.man.item(id)
461 if !ok {
462 return InboxItemMeta{}, PromptEnvelope{}, ErrNotFound
463 }
464 env, err := s.readBlobLocked(blobNameFor(meta), meta.Checksum)
465 if err != nil {
466 return meta, PromptEnvelope{}, err
467 }
468 return meta, env, nil
469 }
470
471 // UpdateItem writes a new immutable blob, switches the manifest pointer, then
472 // deletes the old blob. Pre-commit crashes leave only a GC-able orphan; after
473 // commit the checksum always points at the new body.
474 func (s *Store) UpdateItem(id string, env PromptEnvelope) (InboxItemMeta, error) {
475 return s.UpdateItemWithIdempotency(id, env, "", PromptEnvelope{})
476 }
477
478 // UpdateItemWithIdempotency atomically updates an item and optionally binds an
479 // additional client idempotency key to it. aliasEnvelope is the original client
480 // request, not the merged body, so collect-mode redelivery remains deduplicated.
481 func (s *Store) UpdateItemWithIdempotency(id string, env PromptEnvelope, alias string, aliasEnvelope PromptEnvelope) (InboxItemMeta, error) {
482 return s.updateItem(id, env, alias, aliasEnvelope, "")
483 }
484
485 // UpdateItemWithIdempotencyIfVersion prevents collect-mode appends from
486 // overwriting an edit made while references were being prepared.
487 func (s *Store) UpdateItemWithIdempotencyIfVersion(id string, env PromptEnvelope, alias string, aliasEnvelope PromptEnvelope, version string) (InboxItemMeta, error) {
488 if version == "" {
489 return InboxItemMeta{}, ErrContentChanged
490 }
491 return s.updateItem(id, env, alias, aliasEnvelope, version)
492 }
493
494 func (s *Store) updateItem(id string, env PromptEnvelope, alias string, aliasEnvelope PromptEnvelope, version string) (InboxItemMeta, error) {
495 if s == nil {
496 return InboxItemMeta{}, ErrClosed
497 }
498 id = strings.TrimSpace(id)
499 alias = strings.TrimSpace(alias)
500 if alias != "" && !validIdempotencyKey(alias) {
501 return InboxItemMeta{}, fmt.Errorf("sessioninbox: invalid idempotency key")
502 }
503 if version == "" {
504 env = normalizeEnvelope(env)
505 }
506 if strings.TrimSpace(env.SubmitText) == "" && env.Invocation == nil && len(env.Invocations) == 0 {
507 return InboxItemMeta{}, ErrEmpty
508 }
509 blobBytes, checksum, byteSize, err := encodeEnvelope(env)
510 if err != nil {
511 return InboxItemMeta{}, err
512 }
513 aliasHash := ""
514 if alias != "" {
515 aliasEnvelope = completeEnqueueEnvelope(aliasEnvelope)
516 aliasHash, err = idempotencyRequestHash(aliasEnvelope)
517 if err != nil {
518 return InboxItemMeta{}, err
519 }
520 }
521 s.mu.Lock()
522 defer s.mu.Unlock()
523 release, err := s.beginDiskTransactionLocked()
524 if err != nil {
525 return InboxItemMeta{}, err
526 }
527 defer release()
528 if err := s.mutableLocked(); err != nil {
529 return InboxItemMeta{}, err
530 }
531 meta, ok := s.man.item(id)
532 if !ok {
533 return InboxItemMeta{}, ErrNotFound
534 }
535 if !isPendingState(meta.State) {
536 return InboxItemMeta{}, ErrInvalidState
537 }
538 replayed, err := s.idempotentAliasReplayLocked(alias, aliasHash, id)
539 if err != nil {
540 return InboxItemMeta{}, err
541 }
542 if replayed {
543 return meta, nil
544 }
545 if version != "" && ContentVersion(meta) != version {
546 return InboxItemMeta{}, ErrContentChanged
547 }
548 if byteSize > s.limits.MaxItemBytes {
549 return InboxItemMeta{}, ErrItemTooLarge
550 }
551 delta := byteSize - meta.ByteSize
552 if s.man.totalBytes()+delta > s.limits.MaxTotalBytes {
553 return InboxItemMeta{}, ErrCapacityBytes
554 }
555 oldBlob := blobNameFor(meta)
556 newBlob := id + "." + newRandomID()
557 if err := s.writeBlobLocked(newBlob, blobBytes); err != nil {
558 return InboxItemMeta{}, err
559 }
560 next := s.man.clone()
561 i := next.indexOf(id)
562 next.Items[i].BlobName = newBlob
563 next.Items[i].ByteSize = byteSize
564 next.Items[i].Checksum = checksum
565 next.Items[i].Preview = PreviewText(env.DisplayText, DefaultPreviewRunes)
566 next.Items[i].Refs = refSummaries(env.Refs)
567 next.Items[i].UpdatedAt = time.Now().UTC()
568 next.Items[i].Revision = next.Revision + 1
569 bindIdempotency(next, alias, id, aliasHash)
570 if next.Items[i].State == StateBlocked {
571 next.Items[i].State = StateQueued
572 next.Items[i].BlockReason = ""
573 }
574 if len(env.ReferenceErrors) > 0 {
575 // Preserve uncertain delivery until explicit retry, even after editing.
576 if next.Items[i].State != StateUncertain {
577 next.Items[i].State = StateBlocked
578 }
579 next.Items[i].BlockReason = strings.Join(env.ReferenceErrors, "; ")
580 next.Paused = true
581 }
582 if err := s.commitManifestLocked(next); err != nil {
583 s.removeBlobLocked(newBlob)
584 return InboxItemMeta{}, err
585 }
586 if oldBlob != newBlob {
587 s.removeBlobLocked(oldBlob)
588 }
589 updated := next.Items[i]
590 s.notifyLocked(s.snapshotLocked())
591 return updated, nil
592 }
593
594 func (s *Store) removeBlobLocked(blobName string) {
595 if err := validatePrivateDir(filepath.Join(s.dir, blobsDirName)); err != nil {
596 return
597 }
598 path, err := s.blobPath(blobName)
599 if err == nil {
600 _ = os.Remove(path)
601 }
602 }
603
603 lines GO