| 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 |