返回 DeepSeek-Reasonix
ops.go
1 package sessioninbox
2
3 import (
4 "strings"
5 "time"
6 )
7
8 // DeleteItem removes metadata first, then the blob (crash may leave orphan).
9 func (s *Store) DeleteItem(id string) error {
10 return s.deleteItem(id, false)
11 }
12
13 // DeletePendingOrAcceptedItem atomically withdraws a queued item or an
14 // accepted-but-unconsumed steer. A concurrent consumed transition wins by
15 // making the delete fail with ErrInvalidState.
16 func (s *Store) DeletePendingOrAcceptedItem(id string) error {
17 return s.deleteItem(id, true)
18 }
19
20 func (s *Store) deleteItem(id string, allowAcceptedSteer bool) error {
21 if s == nil {
22 return ErrClosed
23 }
24 id = strings.TrimSpace(id)
25 s.mu.Lock()
26 defer s.mu.Unlock()
27 release, err := s.beginDiskTransactionLocked()
28 if err != nil {
29 return err
30 }
31 defer release()
32 if err := s.mutableLocked(); err != nil {
33 return err
34 }
35 meta, ok := s.man.item(id)
36 if !ok {
37 return ErrNotFound
38 }
39 if !isPendingState(meta.State) && !(allowAcceptedSteer && meta.State == StateSteerAccepted) {
40 return ErrInvalidState
41 }
42 next := s.man.clone()
43 keys := next.idempotencyKeysFor(id)
44 next.rememberReceipt(keys, id, Disposition("deleted"), time.Now().UTC())
45 removed, _ := next.removeItem(id)
46 clearPauseIfEmpty(next)
47 if err := s.commitManifestLocked(next); err != nil {
48 return err
49 }
50 s.removeBlobLocked(blobNameFor(removed))
51 s.notifyLocked(s.snapshotLocked())
52 return nil
53 }
54
55 // DiscardPendingItems removes the named, not-yet-admitted items in one
56 // manifest transaction. Missing IDs are treated as already consumed so a
57 // frontend may safely cancel from a slightly stale metadata snapshot. Items
58 // that have crossed the admission boundary are never deleted here.
59 func (s *Store) DiscardPendingItems(ids []string) error {
60 return s.DiscardPendingItemsOwned(ids, "")
61 }
62
63 // DiscardPendingItemsOwned atomically removes pending IDs belonging to source.
64 // Foreign-source IDs are ignored so one frontend cannot cancel another.
65 func (s *Store) DiscardPendingItemsOwned(ids []string, source string) error {
66 _, err := s.discardPendingItemsOwnedResult(ids, source, true)
67 return err
68 }
69
70 // DiscardPendingItemsOwnedResult atomically removes cancellable IDs belonging
71 // to source and returns exactly the IDs committed as discarded. Items that
72 // already crossed the durable delivery boundary are ignored instead of making
73 // a mixed batch fail as a whole.
74 func (s *Store) DiscardPendingItemsOwnedResult(ids []string, source string) ([]string, error) {
75 return s.discardPendingItemsOwnedResult(ids, source, false)
76 }
77
78 func (s *Store) discardPendingItemsOwnedResult(ids []string, source string, strict bool) ([]string, error) {
79 if s == nil {
80 return nil, ErrClosed
81 }
82 wanted := make(map[string]struct{}, len(ids))
83 for _, id := range ids {
84 if id = strings.TrimSpace(id); id != "" {
85 wanted[id] = struct{}{}
86 }
87 }
88 if len(wanted) == 0 {
89 return []string{}, nil
90 }
91
92 s.mu.Lock()
93 defer s.mu.Unlock()
94 release, err := s.beginDiskTransactionLocked()
95 if err != nil {
96 return nil, err
97 }
98 defer release()
99 if err := s.mutableLocked(); err != nil {
100 return nil, err
101 }
102 for _, item := range s.man.Items {
103 if _, ok := wanted[item.ID]; !ok {
104 continue
105 }
106 if source != "" && item.Source != source {
107 continue
108 }
109 switch item.State {
110 case StateQueued, StateBlocked, StateUncertain:
111 case StateSteerAccepted:
112 if strict {
113 return nil, ErrInvalidState
114 }
115 case StateRunning, StateSteerConsumed:
116 if strict {
117 return nil, ErrInvalidState
118 }
119 default:
120 return nil, ErrInvalidState
121 }
122 }
123
124 next := s.man.clone()
125 removed := make([]InboxItemMeta, 0, len(wanted))
126 kept := next.Items[:0]
127 for _, item := range next.Items {
128 _, selected := wanted[item.ID]
129 owned := source == "" || item.Source == source
130 cancellable := item.State == StateQueued || item.State == StateBlocked || item.State == StateUncertain || item.State == StateSteerAccepted
131 if selected && owned && cancellable {
132 removed = append(removed, item)
133 continue
134 }
135 kept = append(kept, item)
136 }
137 if len(removed) == 0 {
138 return []string{}, nil
139 }
140 next.Items = kept
141 now := time.Now().UTC()
142 for _, item := range removed {
143 keys := next.idempotencyKeysFor(item.ID)
144 next.rememberReceipt(keys, item.ID, Disposition("discarded"), now)
145 for _, key := range keys {
146 delete(next.Idempotency, key)
147 delete(next.IdempotencyHashes, key)
148 }
149 }
150 clearPauseIfEmpty(next)
151 if err := s.commitManifestLocked(next); err != nil {
152 return nil, err
153 }
154 for _, item := range removed {
155 s.removeBlobLocked(blobNameFor(item))
156 }
157 s.notifyLocked(s.snapshotLocked())
158 discarded := make([]string, 0, len(removed))
159 for _, item := range removed {
160 discarded = append(discarded, item.ID)
161 }
162 return discarded, nil
163 }
164
165 // MoveItem reorders the queue. toIndex is 0-based; values past the end append.
166 func (s *Store) MoveItem(id string, toIndex int) error {
167 if s == nil {
168 return ErrClosed
169 }
170 id = strings.TrimSpace(id)
171 s.mu.Lock()
172 defer s.mu.Unlock()
173 release, err := s.beginDiskTransactionLocked()
174 if err != nil {
175 return err
176 }
177 defer release()
178 if err := s.mutableLocked(); err != nil {
179 return err
180 }
181 next := s.man.clone()
182 from := next.indexOf(id)
183 if from < 0 {
184 return ErrNotFound
185 }
186 if !isPendingState(next.Items[from].State) {
187 return ErrInvalidState
188 }
189 if toIndex < 0 {
190 toIndex = 0
191 }
192 if toIndex >= len(next.Items) {
193 toIndex = len(next.Items) - 1
194 }
195 if from == toIndex {
196 return nil
197 }
198 it := next.Items[from]
199 next.Items = append(next.Items[:from], next.Items[from+1:]...)
200 if toIndex > len(next.Items) {
201 toIndex = len(next.Items)
202 }
203 next.Items = append(next.Items[:toIndex], append([]InboxItemMeta{it}, next.Items[toIndex:]...)...)
204 if err := s.commitManifestLocked(next); err != nil {
205 return err
206 }
207 s.notifyLocked(s.snapshotLocked())
208 return nil
209 }
210
211 // SetPaused toggles the recovery/inspection pause flag.
212 func (s *Store) SetPaused(paused bool) error {
213 if s == nil {
214 return ErrClosed
215 }
216 s.mu.Lock()
217 defer s.mu.Unlock()
218 release, err := s.beginDiskTransactionLocked()
219 if err != nil {
220 return err
221 }
222 defer release()
223 if err := s.mutableLocked(); err != nil {
224 return err
225 }
226 if s.man.Paused == paused {
227 return nil
228 }
229 next := s.man.clone()
230 next.Paused = paused
231 if !paused {
232 next.Recovered = false
233 next.RecoveredN = 0
234 }
235 if err := s.commitManifestLocked(next); err != nil {
236 return err
237 }
238 s.notifyLocked(s.snapshotLocked())
239 return nil
240 }
241
242 // PauseIfPending pauses dispatch only when the inbox still contains work.
243 func (s *Store) PauseIfPending() error {
244 if s == nil {
245 return ErrClosed
246 }
247 s.mu.Lock()
248 defer s.mu.Unlock()
249 release, err := s.beginDiskTransactionLocked()
250 if err != nil {
251 return err
252 }
253 defer release()
254 if err := s.mutableLocked(); err != nil {
255 return err
256 }
257 if len(s.man.Items) == 0 || s.man.Paused {
258 return nil
259 }
260 next := s.man.clone()
261 next.Paused = true
262 if err := s.commitManifestLocked(next); err != nil {
263 return err
264 }
265 s.notifyLocked(s.snapshotLocked())
266 return nil
267 }
268
269 // SetState transitions one item's durable state.
270 func (s *Store) SetState(id string, state InboxState, blockReason string) error {
271 if s == nil {
272 return ErrClosed
273 }
274 id = strings.TrimSpace(id)
275 s.mu.Lock()
276 defer s.mu.Unlock()
277 release, err := s.beginDiskTransactionLocked()
278 if err != nil {
279 return err
280 }
281 defer release()
282 if err := s.mutableLocked(); err != nil {
283 return err
284 }
285 next := s.man.clone()
286 i := next.indexOf(id)
287 if i < 0 {
288 return ErrNotFound
289 }
290 next.Items[i].State = state
291 next.Items[i].BlockReason = blockReason
292 next.Items[i].UpdatedAt = time.Now().UTC()
293 if err := s.commitManifestLocked(next); err != nil {
294 return err
295 }
296 s.notifyLocked(s.snapshotLocked())
297 return nil
298 }
299
300 // MarkSteerConsumed is the durable steer delivery boundary. A loader must
301 // commit this transition before returning the instruction to the agent. If a
302 // concurrent cancellation removed the accepted item first, the loader fails
303 // closed and the instruction is not applied.
304 func (s *Store) MarkSteerConsumed(id string) error {
305 return s.transitionAcceptedSteer(id, StateSteerConsumed, "", true)
306 }
307
308 // MarkAcceptedSteerUncertain preserves an accepted steer that left the agent
309 // queue without being applied. It refuses to overwrite a consumed item.
310 func (s *Store) MarkAcceptedSteerUncertain(id, reason string) error {
311 return s.transitionAcceptedSteer(id, StateUncertain, reason, false)
312 }
313
314 func (s *Store) transitionAcceptedSteer(id string, target InboxState, blockReason string, consumedIdempotent bool) error {
315 if s == nil {
316 return ErrClosed
317 }
318 id = strings.TrimSpace(id)
319 s.mu.Lock()
320 defer s.mu.Unlock()
321 release, err := s.beginDiskTransactionLocked()
322 if err != nil {
323 return err
324 }
325 defer release()
326 if err := s.mutableLocked(); err != nil {
327 return err
328 }
329 next := s.man.clone()
330 i := next.indexOf(id)
331 if i < 0 {
332 return ErrNotFound
333 }
334 if consumedIdempotent && next.Items[i].State == StateSteerConsumed {
335 return nil
336 }
337 if next.Items[i].State != StateSteerAccepted {
338 return ErrInvalidState
339 }
340 next.Items[i].State = target
341 next.Items[i].BlockReason = blockReason
342 next.Items[i].UpdatedAt = time.Now().UTC()
343 if err := s.commitManifestLocked(next); err != nil {
344 return err
345 }
346 s.notifyLocked(s.snapshotLocked())
347 return nil
348 }
349
350 // ClaimItem atomically transitions one queued item to running. It is the
351 // durable admission boundary for both asynchronous and synchronous frontends.
352 func (s *Store) ClaimItem(id string) error {
353 if s == nil {
354 return ErrClosed
355 }
356 id = strings.TrimSpace(id)
357 s.mu.Lock()
358 defer s.mu.Unlock()
359 release, err := s.beginDiskTransactionLocked()
360 if err != nil {
361 return err
362 }
363 defer release()
364 if err := s.mutableLocked(); err != nil {
365 return err
366 }
367 if s.man.Paused {
368 return ErrPaused
369 }
370 next := s.man.clone()
371 i := next.indexOf(id)
372 if i < 0 {
373 return ErrNotFound
374 }
375 if next.Items[i].State != StateQueued {
376 return ErrInvalidState
377 }
378 next.Items[i].State = StateRunning
379 next.Items[i].BlockReason = ""
380 next.Items[i].UpdatedAt = time.Now().UTC()
381 if err := s.commitManifestLocked(next); err != nil {
382 return err
383 }
384 s.notifyLocked(s.snapshotLocked())
385 return nil
386 }
387
388 // ConvertIntent changes followup ↔ steer while keeping the item queued.
389 func (s *Store) ConvertIntent(id string, intent InboxIntent) error {
390 if s == nil {
391 return ErrClosed
392 }
393 if intent != IntentSteer {
394 intent = IntentFollowup
395 }
396 s.mu.Lock()
397 defer s.mu.Unlock()
398 release, err := s.beginDiskTransactionLocked()
399 if err != nil {
400 return err
401 }
402 defer release()
403 if err := s.mutableLocked(); err != nil {
404 return err
405 }
406 next := s.man.clone()
407 i := next.indexOf(id)
408 if i < 0 {
409 return ErrNotFound
410 }
411 if !isPendingState(next.Items[i].State) {
412 return ErrInvalidState
413 }
414 next.Items[i].Intent = intent
415 next.Items[i].UpdatedAt = time.Now().UTC()
416 if err := s.commitManifestLocked(next); err != nil {
417 return err
418 }
419 s.notifyLocked(s.snapshotLocked())
420 return nil
421 }
422
423 // AckDequeue removes a running/consumed item after durable transcript commit.
424 func (s *Store) AckDequeue(id string) error {
425 if s == nil {
426 return ErrClosed
427 }
428 id = strings.TrimSpace(id)
429 s.mu.Lock()
430 defer s.mu.Unlock()
431 release, err := s.beginDiskTransactionLocked()
432 if err != nil {
433 return err
434 }
435 defer release()
436 if err := s.mutableLocked(); err != nil {
437 return err
438 }
439 next := s.man.clone()
440 keys := next.idempotencyKeysFor(id)
441 next.rememberReceipt(keys, id, Disposition("acknowledged"), time.Now().UTC())
442 removed, ok := next.removeItem(id)
443 if !ok {
444 return ErrNotFound
445 }
446 switch removed.State {
447 case StateRunning, StateSteerAccepted, StateSteerConsumed:
448 default:
449 return ErrInvalidState
450 }
451 clearPauseIfEmpty(next)
452 if err := s.commitManifestLocked(next); err != nil {
453 return err
454 }
455 s.removeBlobLocked(blobNameFor(removed))
456 s.notifyLocked(s.snapshotLocked())
457 return nil
458 }
459
460 // RetryItem resets uncertain/blocked items to queued.
461 func (s *Store) RetryItem(id string) error {
462 if s == nil {
463 return ErrClosed
464 }
465 id = strings.TrimSpace(id)
466 s.mu.Lock()
467 defer s.mu.Unlock()
468 release, err := s.beginDiskTransactionLocked()
469 if err != nil {
470 return err
471 }
472 defer release()
473 if err := s.mutableLocked(); err != nil {
474 return err
475 }
476 next := s.man.clone()
477 i := next.indexOf(id)
478 if i < 0 {
479 return ErrNotFound
480 }
481 switch next.Items[i].State {
482 case StateUncertain, StateBlocked:
483 next.Items[i].State = StateQueued
484 next.Items[i].BlockReason = ""
485 next.Items[i].UpdatedAt = time.Now().UTC()
486 default:
487 return ErrInvalidState
488 }
489 if err := s.commitManifestLocked(next); err != nil {
490 return err
491 }
492 s.notifyLocked(s.snapshotLocked())
493 return nil
494 }
495
496 // NextQueued returns the first FIFO queued follow-up (or rejected steer kept as
497 // follow-up) when the inbox is not paused.
498 func (s *Store) NextQueued() (InboxItemMeta, bool) {
499 if s == nil {
500 return InboxItemMeta{}, false
501 }
502 s.mu.Lock()
503 defer s.mu.Unlock()
504 if release, err := s.beginDiskTransactionLocked(); err == nil {
505 release()
506 }
507 if s.man == nil || s.man.Paused || s.readonly {
508 return InboxItemMeta{}, false
509 }
510 for _, it := range s.man.Items {
511 if it.State == StateQueued && it.Intent == IntentFollowup {
512 return it, true
513 }
514 // Rejected steers that remain intent=steer but queued are still follow-ups
515 // for the dispatcher after ConvertIntent; only followup intent is admitted.
516 }
517 // Also admit steer-intent items that are still queued (user wants them as turns).
518 for _, it := range s.man.Items {
519 if it.State == StateQueued {
520 return it, true
521 }
522 }
523 return InboxItemMeta{}, false
524 }
525
526 // Pause marks paused=true without requiring a mutation check beyond schema.
527 func (s *Store) ForcePause(reasonRecovered bool, n int) error {
528 if s == nil {
529 return ErrClosed
530 }
531 s.mu.Lock()
532 defer s.mu.Unlock()
533 release, err := s.beginDiskTransactionLocked()
534 if err != nil {
535 return err
536 }
537 defer release()
538 if s.closed || s.readonly {
539 if s.readonly {
540 return nil
541 }
542 return ErrClosed
543 }
544 next := s.man.clone()
545 next.Paused = true
546 if reasonRecovered {
547 next.Recovered = true
548 if n > 0 {
549 next.RecoveredN = n
550 }
551 }
552 if err := s.commitManifestLocked(next); err != nil {
553 return err
554 }
555 s.notifyLocked(s.snapshotLocked())
556 return nil
557 }
558
559 func clearPauseIfEmpty(m *manifest) {
560 if m == nil || len(m.Items) > 0 {
561 return
562 }
563 m.Paused = false
564 m.Recovered = false
565 m.RecoveredN = 0
566 }
567
568 func isPendingState(state InboxState) bool {
569 switch state {
570 case StateQueued, StateBlocked, StateUncertain:
571 return true
572 default:
573 return false
574 }
575 }
576
576 lines GO