返回 DeepSeek-Reasonix
recovery.go
根目录 / internal / sessioninbox / recovery.go
1 package sessioninbox
2
3 import "time"
4
5 // RecoverOrphanedInFlight converts admitted items that no live Controller owns
6 // into reviewable pending work. The transition is atomic so a crash cannot
7 // leave only part of a multi-item active set recoverable.
8 func (s *Store) RecoverOrphanedInFlight(ownedIDs []string) (int, error) {
9 owned := make(map[string]struct{}, len(ownedIDs))
10 for _, id := range ownedIDs {
11 if id != "" {
12 owned[id] = struct{}{}
13 }
14 }
15 return s.RecoverOrphanedInFlightOwnedBy(func(id string) bool {
16 _, ok := owned[id]
17 return ok
18 })
19 }
20
21 // RecoverOrphanedInFlightOwnedBy resolves live ownership only after the Store
22 // transaction is current. The callback must be lock-free and must not call
23 // Store methods; Controller uses sync.Map-backed ownership so a newly admitted
24 // item cannot be recovered from a stale pre-transaction snapshot.
25 func (s *Store) RecoverOrphanedInFlightOwnedBy(ownedBy func(string) bool) (int, error) {
26 if s == nil {
27 return 0, ErrClosed
28 }
29 isOwned := func(id string) bool {
30 return ownedBy != nil && ownedBy(id)
31 }
32
33 s.mu.Lock()
34 defer s.mu.Unlock()
35 needsRecovery := false
36 for i := range s.man.Items {
37 if isOwned(s.man.Items[i].ID) {
38 continue
39 }
40 switch s.man.Items[i].State {
41 case StateRunning, StateSteerAccepted, StateSteerConsumed:
42 needsRecovery = true
43 }
44 }
45 if !needsRecovery {
46 return 0, nil
47 }
48 release, err := s.beginDiskTransactionLocked()
49 if err != nil {
50 return 0, err
51 }
52 defer release()
53 if err := s.mutableLocked(); err != nil {
54 return 0, err
55 }
56
57 next := s.man.clone()
58 now := time.Now().UTC()
59 recovered := 0
60 for i := range next.Items {
61 if isOwned(next.Items[i].ID) {
62 continue
63 }
64 switch next.Items[i].State {
65 case StateRunning, StateSteerAccepted, StateSteerConsumed:
66 next.Items[i].State = StateUncertain
67 next.Items[i].BlockReason = "in-flight owner is no longer active"
68 next.Items[i].UpdatedAt = now
69 recovered++
70 }
71 }
72 if recovered == 0 {
73 return 0, nil
74 }
75 next.Paused = true
76 next.Recovered = true
77 next.RecoveredN = min(len(next.Items), next.RecoveredN+recovered)
78 if err := s.commitManifestLocked(next); err != nil {
79 return 0, err
80 }
81 s.notifyLocked(s.snapshotLocked())
82 return recovered, nil
83 }
84
84 lines GO