返回 DeepSeek-Reasonix
queue_mutations.go
根目录 / internal / sessioninbox / queue_mutations.go
1 package sessioninbox
2
3 import (
4 "crypto/sha256"
5 "errors"
6 "fmt"
7 "slices"
8 "time"
9 )
10
11 var (
12 ErrContentChanged = errors.New("inbox content changed")
13 ErrOrderChanged = errors.New("inbox order changed")
14 ErrAnchorMissing = errors.New("inbox reorder anchor missing")
15 )
16
17 // ContentVersion identifies an immutable body, including an ABA replacement.
18 // It exposes neither the private blob path nor a new persisted schema field.
19 func ContentVersion(meta InboxItemMeta) string {
20 return fmt.Sprintf("%x", sha256.Sum256([]byte(meta.ID+"\x00"+blobNameFor(meta)+"\x00"+meta.Checksum)))
21 }
22
23 // UpdateItemIfVersion commits body, reference failure and lifecycle together.
24 func (s *Store) UpdateItemIfVersion(id string, env PromptEnvelope, version string) (InboxItemMeta, error) {
25 if version == "" {
26 return InboxItemMeta{}, ErrContentChanged
27 }
28 return s.updateItem(id, env, "", PromptEnvelope{}, version)
29 }
30
31 // MoveItemBefore changes only pending slots; accepted/running items stay put.
32 func (s *Store) MoveItemBefore(id string, before *string, revision int64) error {
33 if s == nil {
34 return ErrClosed
35 }
36 s.mu.Lock()
37 defer s.mu.Unlock()
38 release, err := s.beginDiskTransactionLocked()
39 if err != nil {
40 return err
41 }
42 defer release()
43 if err := s.mutableLocked(); err != nil {
44 return err
45 }
46 if s.man.Revision != revision {
47 return ErrOrderChanged
48 }
49 meta, ok := s.man.item(id)
50 if !ok {
51 return ErrNotFound
52 }
53 if !isPendingState(meta.State) {
54 return ErrInvalidState
55 }
56 if before != nil {
57 anchor, found := s.man.item(*before)
58 if !found || !isPendingState(anchor.State) {
59 return ErrAnchorMissing
60 }
61 if *before == id {
62 return nil
63 }
64 }
65 var pending []InboxItemMeta
66 var original []string
67 for _, item := range s.man.Items {
68 if isPendingState(item.State) {
69 original = append(original, item.ID)
70 if item.ID != id {
71 pending = append(pending, item)
72 }
73 }
74 }
75 index := len(pending)
76 if before != nil {
77 index = slices.IndexFunc(pending, func(item InboxItemMeta) bool { return item.ID == *before })
78 }
79 pending = slices.Insert(pending, index, meta)
80 changed := false
81 for i := range pending {
82 changed = changed || pending[i].ID != original[i]
83 }
84 if !changed {
85 return nil
86 }
87 next := s.man.clone()
88 j := 0
89 for i := range next.Items {
90 if isPendingState(next.Items[i].State) {
91 next.Items[i] = pending[j]
92 j++
93 }
94 }
95 if err := s.commitManifestLocked(next); err != nil {
96 return err
97 }
98 s.notifyLocked(s.snapshotLocked())
99 return nil
100 }
101
102 // TransitionPrepared validates what was prepared outside the disk transaction.
103 // A successful edit or move before this boundary must affect actual execution.
104 func (s *Store) TransitionPrepared(id, version string, target InboxState, reason string, requireHead bool) error {
105 if s == nil {
106 return ErrClosed
107 }
108 s.mu.Lock()
109 defer s.mu.Unlock()
110 release, err := s.beginDiskTransactionLocked()
111 if err != nil {
112 return err
113 }
114 defer release()
115 if err := s.mutableLocked(); err != nil {
116 return err
117 }
118 meta, ok := s.man.item(id)
119 if !ok {
120 return ErrNotFound
121 }
122 if meta.State != StateQueued {
123 return ErrInvalidState
124 }
125 if ContentVersion(meta) != version {
126 return ErrContentChanged
127 }
128 if target != StateBlocked && s.man.Paused {
129 return ErrPaused
130 }
131 if requireHead {
132 for _, item := range s.man.Items {
133 if isPendingState(item.State) {
134 if item.ID != id {
135 return ErrOrderChanged
136 }
137 break
138 }
139 }
140 }
141 next := s.man.clone()
142 i := next.indexOf(id)
143 next.Items[i].State = target
144 next.Items[i].BlockReason = reason
145 next.Items[i].UpdatedAt = time.Now().UTC()
146 if target == StateBlocked {
147 next.Paused = true
148 }
149 if err := s.commitManifestLocked(next); err != nil {
150 return err
151 }
152 s.notifyLocked(s.snapshotLocked())
153 return nil
154 }
155
155 lines GO