返回 DeepSeek-Reasonix
service.go
根目录 / internal / skill / skillwatch / service.go
1 // Package skillwatch shares physical skill-directory watches across stores.
2 // Healthy roots use coalesced native events. Failed roots use bounded backoff
3 // scans until registration recovers. Windows isolates blocking filesystem APIs
4 // in a helper process. Subscribe completes registration before the caller's
5 // first catalog scan, subject to the helper timeout.
6 package skillwatch
7
8 import (
9 "context"
10 "crypto/sha256"
11 "fmt"
12 "io"
13 "path/filepath"
14 "sync"
15 "time"
16 )
17
18 // scanBackoffs is the degraded-mode cadence replacing the retired 250ms
19 // polling generation.
20 var scanBackoffs = []time.Duration{2 * time.Second, 5 * time.Second, 15 * time.Second, 30 * time.Second}
21
22 const (
23 coalesceWindow = 100 * time.Millisecond
24 coalesceMaxDelay = 500 * time.Millisecond
25 // helperControlTimeout bounds one helper control round trip; a slower
26 // helper is torn down and rebuilt. helperRestartWindow / helperRestartLimit
27 // cap automatic rebuilds before the service degrades to scanning.
28 helperControlTimeout = 5 * time.Second
29 helperRestartWindow = 30 * time.Second
30 helperRestartLimit = 2
31 // helperSubscribeWait bounds how long Subscribe waits for the helper to
32 // confirm a registration before letting the caller's first scan proceed.
33 // Expiry is non-fatal: the registration stays in flight under the control
34 // timeout and late confirmation still arms the watches.
35 helperSubscribeWait = 2 * time.Second
36 )
37
38 // ScopeFunc returns the directories discovery can visit under root for the
39 // given depth. It includes the nearest existing ancestor of a missing root so
40 // later creation invalidates the snapshot, and excludes subtrees discovery
41 // skips (scripts/assets/references bodies) so their content churn cannot
42 // trigger catalog rebuilds.
43 type ScopeFunc func(ctx context.Context, root string, maxDepth int) (dirs []string, complete bool)
44
45 // HashFunc summarizes one root's watched tree for degraded-mode scanning and
46 // reports how many entries the scan visited. ok is false when the scan could
47 // not complete (cancelled, IO error).
48 type HashFunc func(ctx context.Context, root string, maxDepth int) (sum [sha256.Size]byte, entries int, ok bool)
49
50 // Options configure the service.
51 type Options struct {
52 // Stderr receives diagnostic warnings; nil defaults to os.Stderr.
53 Stderr io.Writer
54 // HelperCommand overrides the helper process used by Windows and tests.
55 HelperCommand func(ctx context.Context) (helperProcess, error)
56 // ForceHelper routes every platform through the helper backend. Test-only.
57 ForceHelper bool
58 // ScanOnly serves subscriptions from backoff scanning and never starts a
59 // helper. Test-only: it lets a caller hold a real, closable service without
60 // spawning a child (which in a test binary is the test binary itself).
61 ScanOnly bool
62 }
63
64 // Diagnostics snapshots the resource counters. Healthy idle state keeps scans
65 // at zero after initialization and never grows physical watches while the
66 // subscription set is constant.
67 type Diagnostics struct {
68 PhysicalWatches uint64 `json:"physicalWatches"`
69 LogicalSubscriptions uint64 `json:"logicalSubscriptions"`
70 Scans uint64 `json:"scans"`
71 ScannedEntries uint64 `json:"scannedEntries"`
72 EventsReceived uint64 `json:"eventsReceived"`
73 Notifications uint64 `json:"notifications"`
74 DegradedRoots uint64 `json:"degradedRoots"`
75 HelperRestarts uint64 `json:"helperRestarts"`
76 }
77
78 // Subscription is one store's handle on a physical root. Release is safe to
79 // call twice and never blocks on backend IO: the logical subscription dies
80 // immediately and late events for it are dropped.
81 type Subscription struct {
82 svc *Service
83 root string // canonical root path (resolved)
84 maxDepth int
85 onChange func(reason string)
86
87 mu sync.Mutex
88 released bool
89 }
90
91 // Release detaches the subscription. The last subscription on a root tears the
92 // physical watches down; no historical project watches are retained.
93 func (sub *Subscription) Release() {
94 if sub == nil {
95 return
96 }
97 sub.mu.Lock()
98 if sub.released {
99 sub.mu.Unlock()
100 return
101 }
102 sub.released = true
103 sub.mu.Unlock()
104 sub.svc.dropSubscription(sub)
105 }
106
107 type rootState struct {
108 canonical string
109 gen uint64 // bumped on every (re)registration; late events dropped
110 dirs []string
111 watchID uint64 // backend registration carrying the current gen
112 maxDepth int
113 subs map[*Subscription]struct{}
114 scope ScopeFunc
115 hash HashFunc
116
117 scopeDirty bool
118
119 // regDone closes when the current registration attempt settles, bounding
120 // the Subscribe-side wait for the register-before-scan ordering.
121 regDone chan struct{}
122
123 // coalescing
124 pending bool
125 firstSeen time.Time
126 timer *time.Timer
127
128 // degraded scanning
129 degraded bool
130 scanStop chan struct{}
131 scanDone chan struct{}
132 lastHash [sha256.Size]byte
133 hasHash bool
134 backoffIx int
135 }
136
137 // Service owns physical watches for any number of subscription roots.
138 type Service struct {
139 stderr io.Writer
140
141 mu sync.Mutex
142 roots map[string]*rootState
143 nextID uint64
144 closed bool
145 stats Diagnostics
146 backend backend
147 backendKind string
148 helper *helperClient
149 }
150
151 // NewService builds the host watch service. The backend is chosen once: the
152 // helper process where available (Windows by default), otherwise in-process
153 // native watches.
154 func NewService(opts Options) *Service {
155 stderr := opts.Stderr
156 if stderr == nil {
157 stderr = io.Discard
158 }
159 svc := &Service{
160 stderr: stderr,
161 roots: map[string]*rootState{},
162 // Registration ids start at 1: zero means "no previous registration"
163 // in the cancel-on-reregister bookkeeping.
164 nextID: 1,
165 }
166 if opts.ScanOnly {
167 // No platform backend, so no helper process on any platform.
168 svc.backend, svc.backendKind = scanOnlyBackend{}, "scan-only"
169 return svc
170 }
171 svc.backend, svc.backendKind, svc.helper = newPlatformBackend(svc, opts)
172 return svc
173 }
174
175 // Subscribe adds one logical subscription for root. It registers the physical
176 // watches before returning (bounded on the helper path) so the caller's first
177 // scan cannot race the watch into missing changes. A registration failure
178 // degrades the root to backoff scanning instead of failing the store.
179 func (s *Service) Subscribe(root string, maxDepth int, scope ScopeFunc, hash HashFunc, onChange func(reason string)) *Subscription {
180 if abs, err := filepath.Abs(root); err == nil {
181 root = abs
182 }
183 if resolved, err := filepath.EvalSymlinks(root); err == nil {
184 root = resolved
185 }
186 sub := &Subscription{
187 svc: s,
188 root: root,
189 maxDepth: maxDepth,
190 onChange: onChange,
191 }
192 s.mu.Lock()
193 if s.closed {
194 s.mu.Unlock()
195 // A closed service hands out dead subscriptions; stores built during
196 // teardown keep working from their initial scan.
197 return sub
198 }
199 state := s.roots[root]
200 if state == nil {
201 state = &rootState{
202 canonical: root, subs: map[*Subscription]struct{}{},
203 maxDepth: maxDepth, scope: scope, hash: hash,
204 }
205 s.roots[root] = state
206 }
207 state.subs[sub] = struct{}{}
208 raised := maxDepth > state.maxDepth
209 if raised {
210 state.maxDepth = maxDepth
211 }
212 s.stats.LogicalSubscriptions++
213 if len(state.dirs) == 0 || raised || state.scopeDirty {
214 s.registerRootLocked(state)
215 }
216 regDone := state.regDone
217 kind := s.backendKind
218 s.mu.Unlock()
219 // Preserve register-before-scan ordering. Native and scan-only registration
220 // are local and settle before Subscribe returns; helper confirmation travels
221 // over a pipe, so only that path needs a timeout for a wedged child.
222 if regDone != nil && (kind == "native" || kind == "scan-only") {
223 <-regDone
224 }
225 if regDone != nil && kind == "helper" {
226 timer := time.NewTimer(helperSubscribeWait)
227 defer timer.Stop()
228 select {
229 case <-regDone:
230 case <-timer.C:
231 }
232 }
233 return sub
234 }
235
236 // registerRootLocked rebuilds the watch scope for one root and hands it to the
237 // backend. Caller holds s.mu. Registration is generation fenced: any result
238 // for an older generation is discarded. The superseded registration is
239 // cancelled once the new one settles, so a root keeps at most one live
240 // backend registration and physical watches never accumulate.
241 func (s *Service) registerRootLocked(state *rootState) {
242 state.gen++
243 gen := state.gen
244 dirs, _ := state.scope(context.Background(), state.canonical, state.maxDepth)
245 state.dirs = dirs
246 state.scopeDirty = false
247 id := s.nextID
248 s.nextID++
249 oldID := state.watchID
250 state.watchID = id
251 done := make(chan struct{})
252 state.regDone = done
253 go func() {
254 defer close(done)
255 err := s.backend.register(id, gen, state.canonical, dirs)
256 if oldID != 0 {
257 s.backend.cancel(oldID)
258 }
259 s.mu.Lock()
260 if state.gen != gen {
261 // Superseded by a newer registration, which owns cancelling id.
262 s.mu.Unlock()
263 return
264 }
265 if len(state.subs) == 0 {
266 // Fully released while registering: reclaim our own registration.
267 s.mu.Unlock()
268 s.backend.cancel(id)
269 return
270 }
271 if err != nil {
272 s.degradeRootLocked(state, err)
273 s.mu.Unlock()
274 return
275 }
276 if state.degraded {
277 s.stopScanLocked(state)
278 }
279 s.mu.Unlock()
280 }()
281 }
282
283 func (s *Service) dropSubscription(sub *Subscription) {
284 s.mu.Lock()
285 state := s.roots[sub.root]
286 if state == nil {
287 s.mu.Unlock()
288 return
289 }
290 if _, ok := state.subs[sub]; !ok {
291 s.mu.Unlock()
292 return
293 }
294 delete(state.subs, sub)
295 if s.stats.LogicalSubscriptions > 0 {
296 s.stats.LogicalSubscriptions--
297 }
298 if len(state.subs) > 0 {
299 s.mu.Unlock()
300 return
301 }
302 // Last subscription gone: release the physical watches. Reference counting
303 // to zero is what keeps historical project roots from accumulating.
304 id := state.watchID
305 var scanDone chan struct{}
306 if state.degraded {
307 scanDone = state.scanDone
308 s.stopScanLocked(state) // closes scanStop
309 }
310 s.cancelTimerLocked(state)
311 delete(s.roots, sub.root)
312 s.mu.Unlock()
313
314 // Backend IO happens outside the lock: cancel never waits on registration.
315 if id != 0 {
316 s.backend.cancel(id)
317 }
318 if scanDone != nil {
319 <-scanDone
320 }
321 }
322
323 // eventArrived receives one raw backend event. Events for unknown, released or
324 // superseded registrations are dropped without producing a notification.
325 func (s *Service) eventArrived(id, rootGen uint64, op Op) {
326 s.mu.Lock()
327 defer s.mu.Unlock()
328 if s.closed {
329 return
330 }
331 var state *rootState
332 for _, r := range s.roots {
333 if r.watchID == id {
334 state = r
335 break
336 }
337 }
338 if state == nil || state.gen != rootGen || len(state.subs) == 0 {
339 return
340 }
341 s.stats.EventsReceived++
342 if state.pending {
343 // Extend the window unless that would exceed the maximum delay.
344 if time.Since(state.firstSeen)+coalesceWindow <= coalesceMaxDelay {
345 state.timer.Reset(coalesceWindow)
346 if op&(OpCreate|OpRename) != 0 {
347 state.scopeDirty = true
348 }
349 return
350 }
351 s.fireRootLocked(state)
352 }
353 state.pending = true
354 state.firstSeen = time.Now()
355 state.timer = time.AfterFunc(coalesceWindow, func() {
356 s.mu.Lock()
357 defer s.mu.Unlock()
358 if state.pending {
359 s.fireRootLocked(state)
360 }
361 })
362 if op&(OpCreate|OpRename) != 0 {
363 // A create/rename can introduce a nested directory, a symlink target or
364 // a previously missing root; rebuild subscriptions after notification.
365 state.scopeDirty = true
366 }
367 }
368
369 // fireRootLocked delivers the coalesced notification and applies a deferred
370 // scope rebuild. Caller holds s.mu.
371 func (s *Service) fireRootLocked(state *rootState) {
372 s.cancelTimerLocked(state)
373 state.pending = false
374 s.stats.Notifications++
375 if state.scopeDirty {
376 s.registerRootLocked(state)
377 }
378 for sub := range state.subs {
379 sub.onChange("filesystem changed")
380 }
381 }
382
383 func (s *Service) cancelTimerLocked(state *rootState) {
384 if state.timer != nil {
385 state.timer.Stop()
386 state.timer = nil
387 }
388 }
389
390 // degradeRootLocked switches a root to bounded signature scanning. The helper
391 // restart budget is spent by the helper client before this point; degraded
392 // roots recover automatically when a later registration succeeds.
393 // Caller holds s.mu.
394 func (s *Service) degradeRootLocked(state *rootState, cause error) {
395 if state.degraded {
396 return
397 }
398 state.degraded = true
399 s.stats.DegradedRoots++
400 s.warnf("skillwatch: root %s degraded to scan fallback: %v", state.canonical, cause)
401 state.scanStop = make(chan struct{})
402 state.scanDone = make(chan struct{})
403 stop, done := state.scanStop, state.scanDone
404 go func() {
405 s.scanLoop(state, stop, done)
406 }()
407 }
408
409 func (s *Service) stopScanLocked(state *rootState) {
410 if !state.degraded {
411 return
412 }
413 state.degraded = false
414 if s.stats.DegradedRoots > 0 {
415 s.stats.DegradedRoots--
416 }
417 if state.scanStop != nil {
418 close(state.scanStop)
419 state.scanStop = nil
420 }
421 state.backoffIx = 0
422 state.hasHash = false
423 }
424
425 // scanLoop is the degraded-mode replacement for the retired 250ms polling
426 // generation: one scan at a time per root, growing 2/5/15/30s backoff while
427 // the watch backend stays unavailable. A signature diff notifies subscribers
428 // exactly like a native event would, and every pass retries the backend so a
429 // recovered helper or filesystem returns the root to native watching.
430 func (s *Service) scanLoop(state *rootState, stop, done chan struct{}) {
431 defer close(done)
432 for {
433 interval := s.nextScanInterval(state)
434 select {
435 case <-stop:
436 return
437 case <-time.After(interval):
438 }
439 s.mu.Lock()
440 if s.closed || !state.degraded || len(state.subs) == 0 {
441 s.mu.Unlock()
442 return
443 }
444 gen, hash, depth := state.gen, state.hash, state.maxDepth
445 s.mu.Unlock()
446
447 sum, entries, ok := hash(context.Background(), state.canonical, depth)
448
449 s.mu.Lock()
450 s.stats.Scans++
451 s.stats.ScannedEntries += uint64(entries)
452 if state.gen != gen || !state.degraded {
453 s.mu.Unlock()
454 return
455 }
456 if ok && (!state.hasHash || sum != state.lastHash) {
457 state.lastHash, state.hasHash = sum, true
458 for sub := range state.subs {
459 sub.onChange("filesystem changed (scan fallback)")
460 }
461 }
462 // Every pass retries the watch backend: a recovered helper or
463 // filesystem flips the root back to native watching, and
464 // registerRootLocked itself stops the scan on success.
465 s.registerRootLocked(state)
466 s.mu.Unlock()
467 }
468 }
469
470 func (s *Service) nextScanInterval(state *rootState) time.Duration {
471 s.mu.Lock()
472 defer s.mu.Unlock()
473 ix := state.backoffIx
474 if ix >= len(scanBackoffs) {
475 ix = len(scanBackoffs) - 1
476 }
477 state.backoffIx++
478 return scanBackoffs[ix]
479 }
480
481 // helperDied tells the service the helper exhausted its restart budget (or
482 // never started): every live root re-registers, fails, and degrades to its
483 // scan fallback. Recovery from degraded state needs a working backend, which
484 // a degraded host regains only through the per-pass re-registration in
485 // scanLoop once the helper serves again.
486 func (s *Service) helperDied() {
487 s.mu.Lock()
488 defer s.mu.Unlock()
489 if s.closed {
490 return
491 }
492 states := make([]*rootState, 0, len(s.roots))
493 for _, state := range s.roots {
494 states = append(states, state)
495 }
496 for _, state := range states {
497 if len(state.subs) > 0 {
498 s.registerRootLocked(state)
499 }
500 }
501 }
502
503 func (s *Service) helperRestarted() {
504 s.mu.Lock()
505 s.stats.HelperRestarts++
506 s.mu.Unlock()
507 }
508
509 func (s *Service) warnf(format string, args ...any) {
510 // Diagnostics only: counts, root paths and failure causes, never contents.
511 _, _ = fmt.Fprintf(s.stderr, format+"\n", args...)
512 }
513
514 // Diagnostics returns a snapshot of the resource counters.
515 func (s *Service) Diagnostics() Diagnostics {
516 s.mu.Lock()
517 defer s.mu.Unlock()
518 out := s.stats
519 out.PhysicalWatches = s.backend.physicalWatches()
520 out.DegradedRoots = 0
521 for _, state := range s.roots {
522 if state.degraded {
523 out.DegradedRoots++
524 }
525 }
526 return out
527 }
528
529 // Close tears the service down: every subscription dies, physical watches and
530 // the helper process are reclaimed. It is idempotent and safe to call during
531 // application exit; it never waits on uninterruptible backend registration.
532 func (s *Service) Close() error {
533 s.mu.Lock()
534 if s.closed {
535 s.mu.Unlock()
536 return nil
537 }
538 s.closed = true
539 type teardown struct {
540 id uint64
541 done chan struct{}
542 }
543 steps := make([]teardown, 0, len(s.roots)*2)
544 for _, state := range s.roots {
545 s.cancelTimerLocked(state)
546 if state.degraded && state.scanStop != nil {
547 close(state.scanStop)
548 state.scanStop = nil
549 steps = append(steps, teardown{done: state.scanDone})
550 }
551 if state.watchID != 0 {
552 steps = append(steps, teardown{id: state.watchID})
553 }
554 }
555 s.roots = map[string]*rootState{}
556 s.stats.LogicalSubscriptions = 0
557 s.mu.Unlock()
558 for _, st := range steps {
559 if st.id != 0 {
560 s.backend.cancel(st.id)
561 }
562 if st.done != nil {
563 <-st.done
564 }
565 }
566 return s.backend.close()
567 }
568
568 lines GO