diff --git a/README.md b/README.md index 18c46f6..deae975 100644 --- a/README.md +++ b/README.md @@ -87,6 +87,7 @@ curl -s -i http://localhost/lancache-heartbeat | `total_requests` / `errors` | Volume and failures | | `upstream_errors` / `cache_write_failures` / `rate_limited` | Upstream pipe, cache write, and rate-limit pressure (Quick check highlights these next to hit/miss) | | `disk_tier_ready` | `0` while disk slow-tier attach pending; `1` when attached, or when no disk configured (N/A — not waiting) | +| `capacity_pressure_events` | Soft eviction under the memory or disk cap, and/or disk Create/Write/Mkdir hitting ENOSPC (volume full). Distinct from cold-cache misses and from the existing `evictions` counter. Logs `tier` (memory or disk) and `reason` (eviction or enospc). | A first pass through new content is mostly misses (`hit_rate` near 0). Repeat the same content and `cache_hits` / `hit_rate` should rise. @@ -306,6 +307,7 @@ See `config.Validate()` and `steamcache.New` error paths. This ensures the LAN a - Startup logs: Info "Disk slow tier attach pending..." then later "Disk slow tier attached (...)" for disk-only and mixed modes. - `/metrics` exposes `disk_tier_ready` 0/1 and stays responsive during attach (GetMetrics does not block on Size while pending). - `/lancache-heartbeat` header `X-SteamCache-Disk-Tier` mirrors that state. +- `/metrics` `capacity_pressure_events` counts times the cache dropped data under capacity pressure (soft eviction at the memory or disk cap, or disk Create/Write/Mkdir returning ENOSPC). Logs include `tier=memory|disk` and `reason=eviction|enospc` so operators can grep and tell this apart from a cold cache. The existing `evictions` counter is unchanged. #### Garbage Collection Algorithms diff --git a/steamcache/metrics/metrics.go b/steamcache/metrics/metrics.go index 7936e74..c224205 100644 --- a/steamcache/metrics/metrics.go +++ b/steamcache/metrics/metrics.go @@ -7,6 +7,8 @@ import ( "sync" "sync/atomic" "time" + + "s1d3sw1ped/steamcache2/steamcache/logger" ) // Metrics tracks various performance and operational metrics @@ -28,13 +30,14 @@ type Metrics struct { TotalBytesSaved int64 // bytes served from cache instead of being re-downloaded from upstream // Cache metrics - MemoryCacheSize int64 - DiskCacheSize int64 - MemoryCacheHits int64 - DiskCacheHits int64 - Promotions int64 - Evictions int64 - DiskTierReady int64 // 0=pending (or unset), 1=ready or no-disk (N/A) + MemoryCacheSize int64 + DiskCacheSize int64 + MemoryCacheHits int64 + DiskCacheHits int64 + Promotions int64 + Evictions int64 + CapacityPressureEvents int64 // soft eviction under cap and/or disk ENOSPC + DiskTierReady int64 // 0=pending (or unset), 1=ready or no-disk (N/A) // Expanded observability (upstream breakdowns, cache write failures, per-service errors) UpstreamErrors int64 @@ -169,8 +172,39 @@ func (m *Metrics) GetServiceRequests(service string) int64 { return m.ServiceRequests[service] } -func (m *Metrics) IncrementPromotions() { atomic.AddInt64(&m.Promotions, 1) } -func (m *Metrics) IncrementEvictions() { atomic.AddInt64(&m.Evictions, 1) } +func (m *Metrics) IncrementPromotions() { atomic.AddInt64(&m.Promotions, 1) } +func (m *Metrics) IncrementEvictions() { atomic.AddInt64(&m.Evictions, 1) } +func (m *Metrics) IncrementCapacityPressureEvents() { atomic.AddInt64(&m.CapacityPressureEvents, 1) } + +// NoteSoftEviction records one cap-pressure eviction batch that freed bytes. +// Keeps the existing evictions counter and also increments capacity_pressure_events. +// Nil m is safe: the log still fires so ops can grep without metrics wired. +func NoteSoftEviction(m *Metrics, tier string, evicted uint) { + if evicted == 0 { + return + } + if m != nil { + m.IncrementEvictions() + m.IncrementCapacityPressureEvents() + } + logger.Logger.Info(). + Str("tier", tier). + Str("reason", "eviction"). + Uint("bytes_evicted", evicted). + Msg("cache capacity pressure") +} + +// NoteNoSpace records a disk Create/Write/Mkdir ENOSPC (or equivalent) event. +func NoteNoSpace(m *Metrics, err error) { + if m != nil { + m.IncrementCapacityPressureEvents() + } + logger.Logger.Warn(). + Str("tier", "disk"). + Str("reason", "enospc"). + Err(err). + Msg("cache capacity pressure") +} // Additional observability counters func (m *Metrics) IncrementUpstreamErrors() { atomic.AddInt64(&m.UpstreamErrors, 1) } @@ -215,32 +249,33 @@ func (m *Metrics) GetStats() *Stats { } return &Stats{ - TotalRequests: totalRequests, - CacheHits: cacheHits, - CacheMisses: cacheMisses, - CacheCoalesced: atomic.LoadInt64(&m.CacheCoalesced), - NegativeCacheHits: atomic.LoadInt64(&m.NegativeCacheHits), - RangeCache: atomic.LoadInt64(&m.RangeCache), - RangeUpstream: atomic.LoadInt64(&m.RangeUpstream), - Errors: atomic.LoadInt64(&m.Errors), - RateLimited: atomic.LoadInt64(&m.RateLimited), - HitRate: hitRate, - AvgResponseTime: avgResponseTime, - TotalBytesServed: atomic.LoadInt64(&m.TotalBytesServed), - TotalBytesSaved: atomic.LoadInt64(&m.TotalBytesSaved), - MemoryCacheSize: atomic.LoadInt64(&m.MemoryCacheSize), - DiskCacheSize: atomic.LoadInt64(&m.DiskCacheSize), - DiskTierReady: atomic.LoadInt64(&m.DiskTierReady), - MemoryCacheHits: atomic.LoadInt64(&m.MemoryCacheHits), - DiskCacheHits: atomic.LoadInt64(&m.DiskCacheHits), - Promotions: atomic.LoadInt64(&m.Promotions), - Evictions: atomic.LoadInt64(&m.Evictions), - ServiceRequests: serviceRequests, - UpstreamErrors: atomic.LoadInt64(&m.UpstreamErrors), - CacheWriteFailures: atomic.LoadInt64(&m.CacheWriteFailures), - ServiceErrors: serviceErrors, - Uptime: time.Since(m.StartTime), - LastResetTime: m.LastResetTime, + TotalRequests: totalRequests, + CacheHits: cacheHits, + CacheMisses: cacheMisses, + CacheCoalesced: atomic.LoadInt64(&m.CacheCoalesced), + NegativeCacheHits: atomic.LoadInt64(&m.NegativeCacheHits), + RangeCache: atomic.LoadInt64(&m.RangeCache), + RangeUpstream: atomic.LoadInt64(&m.RangeUpstream), + Errors: atomic.LoadInt64(&m.Errors), + RateLimited: atomic.LoadInt64(&m.RateLimited), + HitRate: hitRate, + AvgResponseTime: avgResponseTime, + TotalBytesServed: atomic.LoadInt64(&m.TotalBytesServed), + TotalBytesSaved: atomic.LoadInt64(&m.TotalBytesSaved), + MemoryCacheSize: atomic.LoadInt64(&m.MemoryCacheSize), + DiskCacheSize: atomic.LoadInt64(&m.DiskCacheSize), + DiskTierReady: atomic.LoadInt64(&m.DiskTierReady), + MemoryCacheHits: atomic.LoadInt64(&m.MemoryCacheHits), + DiskCacheHits: atomic.LoadInt64(&m.DiskCacheHits), + Promotions: atomic.LoadInt64(&m.Promotions), + Evictions: atomic.LoadInt64(&m.Evictions), + CapacityPressureEvents: atomic.LoadInt64(&m.CapacityPressureEvents), + ServiceRequests: serviceRequests, + UpstreamErrors: atomic.LoadInt64(&m.UpstreamErrors), + CacheWriteFailures: atomic.LoadInt64(&m.CacheWriteFailures), + ServiceErrors: serviceErrors, + Uptime: time.Since(m.StartTime), + LastResetTime: m.LastResetTime, } } @@ -262,6 +297,7 @@ func (m *Metrics) Reset() { atomic.StoreInt64(&m.DiskCacheHits, 0) atomic.StoreInt64(&m.Promotions, 0) atomic.StoreInt64(&m.Evictions, 0) + atomic.StoreInt64(&m.CapacityPressureEvents, 0) atomic.StoreInt64(&m.UpstreamErrors, 0) atomic.StoreInt64(&m.CacheWriteFailures, 0) @@ -293,18 +329,19 @@ type Stats struct { TotalBytesSaved int64 MemoryCacheSize int64 - DiskCacheSize int64 - DiskTierReady int64 - MemoryCacheHits int64 - DiskCacheHits int64 - Promotions int64 - Evictions int64 - UpstreamErrors int64 - CacheWriteFailures int64 - ServiceErrors map[string]int64 - ServiceRequests map[string]int64 - Uptime time.Duration - LastResetTime time.Time + DiskCacheSize int64 + DiskTierReady int64 + MemoryCacheHits int64 + DiskCacheHits int64 + Promotions int64 + Evictions int64 + CapacityPressureEvents int64 + UpstreamErrors int64 + CacheWriteFailures int64 + ServiceErrors map[string]int64 + ServiceRequests map[string]int64 + Uptime time.Duration + LastResetTime time.Time } // WriteText emits the Prometheus-style text metrics to the ResponseWriter. @@ -329,6 +366,7 @@ func WriteText(w http.ResponseWriter, stats *Stats) { _, _ = fmt.Fprintf(w, "disk_cache_hits %d\n", stats.DiskCacheHits) _, _ = fmt.Fprintf(w, "promotions %d\n", stats.Promotions) _, _ = fmt.Fprintf(w, "evictions %d\n", stats.Evictions) + _, _ = fmt.Fprintf(w, "capacity_pressure_events %d\n", stats.CapacityPressureEvents) for svc, cnt := range stats.ServiceErrors { _, _ = fmt.Fprintf(w, "service_errors{service=%q} %d\n", svc, cnt) } diff --git a/steamcache/metrics/metrics_test.go b/steamcache/metrics/metrics_test.go new file mode 100644 index 0000000..89ee3ab --- /dev/null +++ b/steamcache/metrics/metrics_test.go @@ -0,0 +1,63 @@ +package metrics + +import ( + "bytes" + "errors" + "net/http/httptest" + "testing" +) + +func TestCapacityPressureEventsWriteTextAndReset(t *testing.T) { + t.Parallel() + m := NewMetrics() + if got := m.GetStats().CapacityPressureEvents; got != 0 { + t.Fatalf("initial CapacityPressureEvents=%d, want 0", got) + } + + NoteSoftEviction(m, "memory", 0) + if got := m.GetStats().CapacityPressureEvents; got != 0 { + t.Fatalf("zero-byte eviction counted: %d", got) + } + + NoteSoftEviction(m, "memory", 128) + st := m.GetStats() + if st.CapacityPressureEvents != 1 { + t.Fatalf("after memory eviction, CapacityPressureEvents=%d, want 1", st.CapacityPressureEvents) + } + if st.Evictions != 1 { + t.Fatalf("after memory eviction, Evictions=%d, want 1 (existing counter kept)", st.Evictions) + } + + NoteSoftEviction(m, "disk", 64) + NoteNoSpace(m, errors.New("no space left on device")) + st = m.GetStats() + if st.CapacityPressureEvents != 3 { + t.Fatalf("after disk eviction + ENOSPC, CapacityPressureEvents=%d, want 3", st.CapacityPressureEvents) + } + if st.Evictions != 2 { + t.Fatalf("ENOSPC must not increment evictions; Evictions=%d, want 2", st.Evictions) + } + + rec := httptest.NewRecorder() + WriteText(rec, st) + body := rec.Body.Bytes() + if !bytes.Contains(body, []byte("capacity_pressure_events 3")) { + t.Errorf("WriteText missing capacity_pressure_events 3: %q", rec.Body.String()) + } + if !bytes.Contains(body, []byte("evictions 2")) { + t.Errorf("WriteText missing evictions 2: %q", rec.Body.String()) + } + + m.Reset() + st = m.GetStats() + if st.CapacityPressureEvents != 0 || st.Evictions != 0 { + t.Errorf("after Reset, CapacityPressureEvents=%d Evictions=%d, want 0", st.CapacityPressureEvents, st.Evictions) + } +} + +func TestNoteSoftEvictionNilMetrics(t *testing.T) { + t.Parallel() + // Must not panic when metrics are not wired (unit tests / early init). + NoteSoftEviction(nil, "memory", 10) + NoteNoSpace(nil, errors.New("ENOSPC")) +} diff --git a/steamcache/steamcache_test.go b/steamcache/steamcache_test.go index 8f29c93..b58e8ab 100644 --- a/steamcache/steamcache_test.go +++ b/steamcache/steamcache_test.go @@ -14,6 +14,7 @@ import ( "path/filepath" "runtime" "s1d3sw1ped/steamcache2/steamcache/metrics" + "s1d3sw1ped/steamcache2/vfs/disk" "s1d3sw1ped/steamcache2/vfs/eviction" "s1d3sw1ped/steamcache2/vfs/memory" "s1d3sw1ped/steamcache2/vfs/vfserror" @@ -1049,13 +1050,21 @@ func TestP1_03_EvictionAlgorithmsDistinct(t *testing.T) { // TestDiskOnlyDelayedAttach covers pure disk-only mode (mem=0 + disk>0) hitting the exact delayed attach path. // During init window (pre Size barrier), TieredCache has no slow tier so Create returns ErrNotFound (proxy semantics, no disk caching). // Post-barrier + attach, Create succeeds. Uses real temp dir. +// An init hold keeps the empty-dir attach from finishing before the pending assertions (CI race). func TestDiskOnlyDelayedAttach(t *testing.T) { - t.Parallel() td := t.TempDir() diskPath := filepath.Join(td, "disk") if err := os.MkdirAll(diskPath, 0755); err != nil { t.Fatal(err) } + hold := make(chan struct{}) + var holdOnce sync.Once + closeHold := func() { holdOnce.Do(func() { close(hold) }) } + disk.RegisterInitHold(diskPath, hold) + t.Cleanup(func() { + closeHold() + disk.ClearInitHold(diskPath) + }) // mem=0, disk>0 -> pure disk delayed path (go func) sc, err := New("localhost:0", "0", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "") @@ -1063,6 +1072,7 @@ func TestDiskOnlyDelayedAttach(t *testing.T) { t.Fatalf("New disk-only: %v", err) } t.Cleanup(func() { sc.Shutdown() }) + t.Cleanup(closeHold) // before Shutdown: attach is blocked in Size() until the hold closes // Immediately in window: no slow tier attached yet -> Create must ErrNotFound (proxy, no disk write) _, err = sc.vfs.Create("during-init-key", 100) @@ -1076,6 +1086,7 @@ func TestDiskOnlyDelayedAttach(t *testing.T) { t.Errorf("during pending attach, DiskTierReady=%d, want 0", got) } + closeHold() // Wait the barrier (exercises the attach go's Size wait) _ = sc.disk.Size() @@ -1114,6 +1125,9 @@ func TestDiskOnlyDelayedAttach(t *testing.T) { if !bytes.Contains(rec.Body.Bytes(), []byte("disk_tier_ready 1")) { t.Errorf("WriteText output missing \"disk_tier_ready 1\": %q", rec.Body.String()) } + if !bytes.Contains(rec.Body.Bytes(), []byte("capacity_pressure_events")) { + t.Errorf("WriteText output missing capacity_pressure_events: %q", rec.Body.String()) + } } // TestDiskTierSignalMemoryOnly covers memory-only mode: DiskTierReady=1 (N/A, not @@ -1146,19 +1160,31 @@ func TestDiskTierSignalMemoryOnly(t *testing.T) { // TestDiskTierSignalMixedPendingReady covers mixed mode: DiskTierReady=0 (header // pending) while the disk attach is in the Size barrier, then DiskTierReady=1 // (header ready) after the barrier opens and the attach goroutine sets SetSlow. +// An init hold keeps the empty-dir attach from finishing between the pending +// metric check and the heartbeat, which otherwise races under CI load. func TestDiskTierSignalMixedPendingReady(t *testing.T) { td := t.TempDir() diskPath := filepath.Join(td, "disk") if err := os.MkdirAll(diskPath, 0755); err != nil { t.Fatal(err) } + hold := make(chan struct{}) + var holdOnce sync.Once + closeHold := func() { holdOnce.Do(func() { close(hold) }) } + disk.RegisterInitHold(diskPath, hold) + t.Cleanup(func() { + closeHold() + disk.ClearInitHold(diskPath) + }) + sc, err := New("127.0.0.1:0", "1MB", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "") if err != nil { t.Fatalf("New mixed: %v", err) } t.Cleanup(func() { sc.Shutdown() }) + t.Cleanup(closeHold) // before Shutdown: attach is blocked in Size() until the hold closes - // Immediately in the pending window + // Pending window is held open until closeHold; attach cannot finish. if got := sc.GetMetrics().DiskTierReady; got != 0 { t.Errorf("immediate DiskTierReady=%d, want 0 (pending)", got) } @@ -1169,7 +1195,7 @@ func TestDiskTierSignalMixedPendingReady(t *testing.T) { t.Errorf("heartbeat header=%q, want pending", got) } - // Wait the barrier, then retry until the attach goroutine flips the flag + closeHold() _ = sc.disk.Size() deadline := time.Now().Add(2 * time.Second) for time.Now().Before(deadline) { diff --git a/vfs/disk/disk.go b/vfs/disk/disk.go index 3614142..f6c9858 100644 --- a/vfs/disk/disk.go +++ b/vfs/disk/disk.go @@ -45,6 +45,25 @@ type DiskFS struct { initCloseOnce sync.Once startupEvict func(vfs.VFS, uint) uint // passed to New (via gc.GetGCAlgorithm); invoked as last step of bg init if over cap (no post-ctor race) metrics *metrics.Metrics + // initHold, if non-nil, is received on before closing initDone (test pending-window hold). + initHold <-chan struct{} +} + +// initHolds is a per-root registry of optional init holds (root path -> <-chan struct{}). +// Tests call RegisterInitHold before New so that instance copies the channel and waits +// before closing initDone; production never registers, so init is unchanged. +var initHolds sync.Map + +// RegisterInitHold registers a channel that DiskFS.New for this root copies onto that +// instance. calculateSizeAndPopulateIndex receives on it before closing initDone, so +// tests can observe the pending-attach window. Other DiskFS roots are unaffected. +func RegisterInitHold(root string, ch <-chan struct{}) { + initHolds.Store(root, ch) +} + +// ClearInitHold removes a previously registered hold for root. +func ClearInitHold(root string) { + initHolds.Delete(root) } // shardPath converts a Steam cache key to a sharded directory path to reduce inode pressure @@ -129,6 +148,12 @@ func New(root string, capacity int64, evict func(vfs.VFS, uint) uint) (*DiskFS, startupEvict: evict, } + if v, ok := initHolds.Load(root); ok { + if ch, ok := v.(<-chan struct{}); ok { + d.initHold = ch + } + } + d.initDone = make(chan struct{}) // Launch heavy population asynchronously so New returns fast (scans millions of files without blocking ctor or using O(N) temp RAM). // The initDone barrier ensures first Size() and subsequent ops (including late tier attach) see fully populated + post-eviction state. @@ -152,7 +177,7 @@ func (d *DiskFS) calculateSizeAndPopulateIndex() { if r := recover(); r != nil { logger.Logger.Error().Interface("recovered_panic", r).Msg("calculateSizeAndPopulateIndex panicked; ensuring initDone closed to unblock Size waiters and prevent hang") } - d.initCloseOnce.Do(func() { close(d.initDone) }) + d.closeInitDone() }() tstart := time.Now() @@ -243,7 +268,19 @@ func (d *DiskFS) calculateSizeAndPopulateIndex() { // Signal readiness: Size() and callers (late tier attach + Evict*) now see correct populated + post-eviction state. // Use Once (recover path also uses it) to guarantee exactly one close even under panic. - d.initCloseOnce.Do(func() { close(d.initDone) }) + d.closeInitDone() +} + +// closeInitDone receives on a copied test hold (if any) then closes initDone once. +// Both the normal end of calculateSizeAndPopulateIndex and the panic-recovery defer +// call this so Size() waiters unblock in either path. +func (d *DiskFS) closeInitDone() { + d.initCloseOnce.Do(func() { + if d.initHold != nil { + <-d.initHold + } + close(d.initDone) + }) } // insertBatch populates info/LRU under lock for a bounded batch (follows maxEvictBatch pattern for short critical sections). @@ -403,11 +440,13 @@ func (d *DiskFS) Create(key string, size int64) (io.WriteCloser, error) { dir := filepath.Dir(path) // 0700 (not 0755): per-shard cache dirs hold untrusted CDN content; restrict to owner only (G301 addressed). if err := os.MkdirAll(dir, 0700); err != nil { + d.recordIfNoSpace(err) return nil, err } file, err := os.Create(path) // #nosec G304 -- path built by pathForKey from sanitized (Clean, no ..) hash-derived key under trusted disk.root; no untrusted file inclusion if err != nil { + d.recordIfNoSpace(err) return nil, err } @@ -438,7 +477,19 @@ type diskWriteCloser struct { } func (dwc *diskWriteCloser) Write(p []byte) (n int, err error) { - return dwc.file.Write(p) + n, err = dwc.file.Write(p) + if err != nil { + dwc.disk.recordIfNoSpace(err) + } + return n, err +} + +// recordIfNoSpace increments capacity_pressure_events and logs when err is ENOSPC (or Windows disk-full). +func (d *DiskFS) recordIfNoSpace(err error) { + if !isNoSpaceError(err) { + return + } + metrics.NoteNoSpace(d.metrics, err) } func (dwc *diskWriteCloser) Close() error { @@ -722,9 +773,7 @@ func (d *DiskFS) EvictLRU(bytesNeeded uint) uint { } d.mu.Unlock() - if d.metrics != nil && evicted > 0 { - d.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(d.metrics, "disk", evicted) return evicted } @@ -775,9 +824,7 @@ func (d *DiskFS) EvictBySize(bytesNeeded uint, ascending bool) uint { } d.mu.Unlock() - if d.metrics != nil && evicted > 0 { - d.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(d.metrics, "disk", evicted) return evicted } @@ -826,9 +873,7 @@ func (d *DiskFS) EvictFIFO(bytesNeeded uint) uint { } d.mu.Unlock() - if d.metrics != nil && evicted > 0 { - d.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(d.metrics, "disk", evicted) return evicted } @@ -882,9 +927,7 @@ func (d *DiskFS) EvictLFU(bytesNeeded uint) uint { } d.mu.Unlock() - if d.metrics != nil && evicted > 0 { - d.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(d.metrics, "disk", evicted) return evicted } @@ -939,8 +982,6 @@ func (d *DiskFS) EvictHybrid(bytesNeeded uint) uint { } d.mu.Unlock() - if d.metrics != nil && evicted > 0 { - d.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(d.metrics, "disk", evicted) return evicted } diff --git a/vfs/disk/disk_test.go b/vfs/disk/disk_test.go index f77fc7a..6aa0be4 100644 --- a/vfs/disk/disk_test.go +++ b/vfs/disk/disk_test.go @@ -11,6 +11,7 @@ import ( "testing" "time" + "s1d3sw1ped/steamcache2/steamcache/metrics" "s1d3sw1ped/steamcache2/vfs" ) @@ -122,6 +123,42 @@ func TestDiskFS_InitPopulatesIndexOnRestart(t *testing.T) { } } +func TestDiskFS_CapacityPressureOnEvict(t *testing.T) { + t.Parallel() + td := t.TempDir() + d, err := New(td, 500, nil) + if err != nil { + t.Fatal(err) + } + _ = d.Size() + met := metrics.NewMetrics() + d.SetMetrics(met) + for i := 0; i < 3; i++ { + k := "f" + string(rune('0'+i)) + w, cerr := d.Create(k, 200) + if cerr != nil { + t.Fatal(cerr) + } + if _, werr := w.Write(make([]byte, 200)); werr != nil { + t.Fatal(werr) + } + if cerr := w.Close(); cerr != nil { + t.Fatal(cerr) + } + } + evicted := d.EvictLRU(100) + if evicted == 0 { + t.Fatalf("expected eviction under cap, size=%d cap=%d", d.Size(), d.Capacity()) + } + st := met.GetStats() + if st.Evictions == 0 { + t.Error("evictions counter not incremented under disk cap pressure") + } + if st.CapacityPressureEvents == 0 { + t.Error("capacity_pressure_events not incremented under disk cap pressure") + } +} + func TestDiskFS_EvictAndLazyStat(t *testing.T) { t.Parallel() td := t.TempDir() @@ -138,10 +175,21 @@ func TestDiskFS_EvictAndLazyStat(t *testing.T) { w.Write(make([]byte, 120)) w.Close() } + met := metrics.NewMetrics() + d.SetMetrics(met) ev := d.EvictLRU(200) if ev == 0 { t.Log("no evict (size calc async or snapshot tolerance?)") } + if ev > 0 { + st := met.GetStats() + if st.Evictions == 0 { + t.Error("evictions counter not incremented after disk EvictLRU freed bytes") + } + if st.CapacityPressureEvents == 0 { + t.Error("capacity_pressure_events not incremented after disk EvictLRU freed bytes") + } + } // Explicit post-evict consistency checks: for any key no longer visible via Stat, its on-disk // file must be absent (verifies coordinated unlink + no resurrection via lazy discovery). // Keys still present after this small evict are allowed (accounting tolerance in raw DiskFS). @@ -614,3 +662,57 @@ func TestDiskFS_NewMkdirError(t *testing.T) { t.Errorf("expected mkdir failure error for file-as-dir, got: %v", err) } } + +// TestDiskFS_InitHoldBlocksOnlyRegisteredRoot covers the per-root init hold: +// Size() stays blocked while the hold is open, and a DiskFS on a different root +// does not wait on that hold. +func TestDiskFS_InitHoldBlocksOnlyRegisteredRoot(t *testing.T) { + td := t.TempDir() + hold := make(chan struct{}) + var holdOnce sync.Once + closeHold := func() { holdOnce.Do(func() { close(hold) }) } + RegisterInitHold(td, hold) + t.Cleanup(func() { + closeHold() + ClearInitHold(td) + }) + + d, err := New(td, 10*1024*1024, nil) + if err != nil { + t.Fatal(err) + } + + blocked := make(chan struct{}) + go func() { + _ = d.Size() + close(blocked) + }() + select { + case <-blocked: + t.Fatal("Size returned while init hold still open") + case <-time.After(50 * time.Millisecond): + } + + td2 := t.TempDir() + d2, err := New(td2, 10*1024*1024, nil) + if err != nil { + t.Fatal(err) + } + other := make(chan struct{}) + go func() { + _ = d2.Size() + close(other) + }() + select { + case <-other: + case <-time.After(2 * time.Second): + t.Fatal("unrelated DiskFS Size hung; init hold leaked across roots") + } + + closeHold() + select { + case <-blocked: + case <-time.After(2 * time.Second): + t.Fatal("Size did not return after init hold released") + } +} diff --git a/vfs/disk/enospc_unix.go b/vfs/disk/enospc_unix.go new file mode 100644 index 0000000..cda5560 --- /dev/null +++ b/vfs/disk/enospc_unix.go @@ -0,0 +1,14 @@ +//go:build !windows + +package disk + +import ( + "errors" + + "golang.org/x/sys/unix" +) + +// isNoSpaceError reports whether err is ENOSPC (or wraps it). +func isNoSpaceError(err error) bool { + return err != nil && errors.Is(err, unix.ENOSPC) +} diff --git a/vfs/disk/enospc_unix_test.go b/vfs/disk/enospc_unix_test.go new file mode 100644 index 0000000..e383cfc --- /dev/null +++ b/vfs/disk/enospc_unix_test.go @@ -0,0 +1,58 @@ +//go:build !windows + +package disk + +import ( + "io" + "os" + "testing" + + "golang.org/x/sys/unix" + + "s1d3sw1ped/steamcache2/steamcache/metrics" +) + +func TestIsNoSpaceError(t *testing.T) { + t.Parallel() + if isNoSpaceError(nil) { + t.Error("nil must not be ENOSPC") + } + if isNoSpaceError(io.EOF) { + t.Error("EOF must not be ENOSPC") + } + if !isNoSpaceError(unix.ENOSPC) { + t.Error("unix.ENOSPC should match") + } + wrapped := &os.PathError{Op: "write", Path: "x", Err: unix.ENOSPC} + if !isNoSpaceError(wrapped) { + t.Error("PathError wrapping ENOSPC should match") + } +} + +func TestDiskFS_ENOSPCCapacityPressure(t *testing.T) { + t.Parallel() + d, err := New(t.TempDir(), 1024, nil) + if err != nil { + t.Fatal(err) + } + met := metrics.NewMetrics() + d.SetMetrics(met) + + d.recordIfNoSpace(io.EOF) + if got := met.GetStats().CapacityPressureEvents; got != 0 { + t.Fatalf("non-ENOSPC counted: %d", got) + } + + d.recordIfNoSpace(unix.ENOSPC) + if got := met.GetStats().CapacityPressureEvents; got != 1 { + t.Fatalf("unix.ENOSPC: CapacityPressureEvents=%d, want 1", got) + } + if got := met.GetStats().Evictions; got != 0 { + t.Fatalf("ENOSPC must not increment evictions, got %d", got) + } + + d.recordIfNoSpace(&os.PathError{Op: "write", Path: "p", Err: unix.ENOSPC}) + if got := met.GetStats().CapacityPressureEvents; got != 2 { + t.Fatalf("wrapped ENOSPC: CapacityPressureEvents=%d, want 2", got) + } +} diff --git a/vfs/disk/enospc_windows.go b/vfs/disk/enospc_windows.go new file mode 100644 index 0000000..2607eaa --- /dev/null +++ b/vfs/disk/enospc_windows.go @@ -0,0 +1,17 @@ +//go:build windows + +package disk + +import ( + "errors" + + "golang.org/x/sys/windows" +) + +// isNoSpaceError reports whether err is a Windows disk-full equivalent of ENOSPC. +func isNoSpaceError(err error) bool { + if err == nil { + return false + } + return errors.Is(err, windows.ERROR_DISK_FULL) || errors.Is(err, windows.ERROR_HANDLE_DISK_FULL) +} diff --git a/vfs/disk/enospc_windows_test.go b/vfs/disk/enospc_windows_test.go new file mode 100644 index 0000000..763512a --- /dev/null +++ b/vfs/disk/enospc_windows_test.go @@ -0,0 +1,56 @@ +//go:build windows + +package disk + +import ( + "io" + "os" + "testing" + + "golang.org/x/sys/windows" + + "s1d3sw1ped/steamcache2/steamcache/metrics" +) + +func TestIsNoSpaceError(t *testing.T) { + t.Parallel() + if isNoSpaceError(nil) { + t.Error("nil must not be disk-full") + } + if isNoSpaceError(io.EOF) { + t.Error("EOF must not be disk-full") + } + if !isNoSpaceError(windows.ERROR_DISK_FULL) { + t.Error("ERROR_DISK_FULL should match") + } + if !isNoSpaceError(windows.ERROR_HANDLE_DISK_FULL) { + t.Error("ERROR_HANDLE_DISK_FULL should match") + } + wrapped := &os.PathError{Op: "write", Path: "x", Err: windows.ERROR_DISK_FULL} + if !isNoSpaceError(wrapped) { + t.Error("PathError wrapping ERROR_DISK_FULL should match") + } +} + +func TestDiskFS_ENOSPCCapacityPressure(t *testing.T) { + t.Parallel() + d, err := New(t.TempDir(), 1024, nil) + if err != nil { + t.Fatal(err) + } + met := metrics.NewMetrics() + d.SetMetrics(met) + + d.recordIfNoSpace(io.EOF) + if got := met.GetStats().CapacityPressureEvents; got != 0 { + t.Fatalf("non-ENOSPC counted: %d", got) + } + + d.recordIfNoSpace(windows.ERROR_DISK_FULL) + if got := met.GetStats().CapacityPressureEvents; got != 1 { + t.Fatalf("ERROR_DISK_FULL: CapacityPressureEvents=%d, want 1", got) + } + if got := met.GetStats().Evictions; got != 0 { + t.Fatalf("disk-full must not increment evictions, got %d", got) + } +} diff --git a/vfs/memory/memory.go b/vfs/memory/memory.go index 0ce6ac7..4c2f61d 100644 --- a/vfs/memory/memory.go +++ b/vfs/memory/memory.go @@ -358,9 +358,7 @@ func (m *MemoryFS) EvictLRU(bytesNeeded uint) uint { } m.mu.Unlock() - if m.metrics != nil && evicted > 0 { - m.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(m.metrics, "memory", evicted) return evicted } @@ -413,9 +411,7 @@ func (m *MemoryFS) EvictBySize(bytesNeeded uint, ascending bool) uint { } m.mu.Unlock() - if m.metrics != nil && evicted > 0 { - m.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(m.metrics, "memory", evicted) return evicted } @@ -464,9 +460,7 @@ func (m *MemoryFS) EvictFIFO(bytesNeeded uint) uint { } m.mu.Unlock() - if m.metrics != nil && evicted > 0 { - m.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(m.metrics, "memory", evicted) return evicted } @@ -520,9 +514,7 @@ func (m *MemoryFS) EvictLFU(bytesNeeded uint) uint { } m.mu.Unlock() - if m.metrics != nil && evicted > 0 { - m.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(m.metrics, "memory", evicted) return evicted } @@ -578,8 +570,6 @@ func (m *MemoryFS) EvictHybrid(bytesNeeded uint) uint { } m.mu.Unlock() - if m.metrics != nil && evicted > 0 { - m.metrics.IncrementEvictions() - } + metrics.NoteSoftEviction(m.metrics, "memory", evicted) return evicted } diff --git a/vfs/memory/memory_test.go b/vfs/memory/memory_test.go index 0c3d730..4c3cd87 100644 --- a/vfs/memory/memory_test.go +++ b/vfs/memory/memory_test.go @@ -8,6 +8,8 @@ import ( "sync/atomic" "testing" "time" + + "s1d3sw1ped/steamcache2/steamcache/metrics" ) func TestMemoryFS_Basic(t *testing.T) { @@ -60,6 +62,8 @@ func TestMemoryFS_EvictUnderPressure(t *testing.T) { if err != nil { t.Fatal(err) } + met := metrics.NewMetrics() + m.SetMetrics(met) // create 3x200 = 600 >500, should trigger internal? but direct evict call for i := 0; i < 3; i++ { w, _ := m.Create("f"+string(rune('0'+i)), 200) @@ -71,6 +75,13 @@ func TestMemoryFS_EvictUnderPressure(t *testing.T) { if evicted == 0 || m.Size() > 500 { t.Errorf("evict failed: evicted=%d size=%d", evicted, m.Size()) } + st := met.GetStats() + if st.Evictions == 0 { + t.Error("evictions counter not incremented under memory cap pressure") + } + if st.CapacityPressureEvents == 0 { + t.Error("capacity_pressure_events not incremented under memory cap pressure") + } } func TestMemoryFS_SizeNeverExceedsAfterEvict(t *testing.T) {