From de43a7192928df2b9fc3fc3bb43f2f1173eb435f Mon Sep 17 00:00:00 2001 From: ash Date: Mon, 7 Sep 2026 19:23:55 +0000 Subject: [PATCH 1/3] ops: Signal disk-full and eviction capacity pressure When the disk (or memory) tier is at cap or the volume returns ENOSPC, ops currently look like random misses with no clear "we are dropping data." Count those events as capacity_pressure_events on /metrics and log tier plus reason so operators can tell capacity pressure from a cold cache, without changing the existing evictions counter. Link: https://git.s1d3sw1ped.com/s1d3sw1ped/steamcache2/issues/36 --- README.md | 2 + steamcache/metrics/metrics.go | 132 +++++++++++++++++++---------- steamcache/metrics/metrics_test.go | 63 ++++++++++++++ steamcache/steamcache_test.go | 3 + vfs/disk/disk.go | 36 ++++---- vfs/disk/disk_test.go | 48 +++++++++++ vfs/disk/enospc_unix.go | 14 +++ vfs/disk/enospc_unix_test.go | 58 +++++++++++++ vfs/disk/enospc_windows.go | 17 ++++ vfs/disk/enospc_windows_test.go | 56 ++++++++++++ vfs/memory/memory.go | 20 ++--- vfs/memory/memory_test.go | 11 +++ 12 files changed, 382 insertions(+), 78 deletions(-) create mode 100644 steamcache/metrics/metrics_test.go create mode 100644 vfs/disk/enospc_unix.go create mode 100644 vfs/disk/enospc_unix_test.go create mode 100644 vfs/disk/enospc_windows.go create mode 100644 vfs/disk/enospc_windows_test.go 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..620c835 100644 --- a/steamcache/steamcache_test.go +++ b/steamcache/steamcache_test.go @@ -1114,6 +1114,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 diff --git a/vfs/disk/disk.go b/vfs/disk/disk.go index 3614142..6ed2550 100644 --- a/vfs/disk/disk.go +++ b/vfs/disk/disk.go @@ -403,11 +403,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 +440,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 +736,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 +787,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 +836,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 +890,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 +945,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..953f9d2 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). 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) { From fe04fb9fbbfb6b86beaf822ce6a7a705d72191f4 Mon Sep 17 00:00:00 2001 From: ash Date: Tue, 8 Sep 2026 15:16:59 +0000 Subject: [PATCH 2/3] vfs/disk: Harden pending-window test against empty-dir race Empty-dir disk init closes initDone almost immediately, so TestDiskTierSignalMixedPendingReady can observe DiskTierReady=0 then a heartbeat that already saw ready. Register a per-root init hold so the test can keep the pending window open without slowing production. Link: https://git.s1d3sw1ped.com/s1d3sw1ped/steamcache2/issues/36 --- steamcache/steamcache_test.go | 17 +++++++++-- vfs/disk/disk.go | 41 ++++++++++++++++++++++++-- vfs/disk/disk_test.go | 54 +++++++++++++++++++++++++++++++++++ 3 files changed, 108 insertions(+), 4 deletions(-) diff --git a/steamcache/steamcache_test.go b/steamcache/steamcache_test.go index 620c835..32c9428 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" @@ -1149,19 +1150,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) } @@ -1172,7 +1185,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 6ed2550..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). diff --git a/vfs/disk/disk_test.go b/vfs/disk/disk_test.go index 953f9d2..6aa0be4 100644 --- a/vfs/disk/disk_test.go +++ b/vfs/disk/disk_test.go @@ -662,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") + } +} From e3b2b8de1e9b337483fa23a1635117ccd0205b37 Mon Sep 17 00:00:00 2001 From: ash Date: Tue, 8 Sep 2026 15:30:02 +0000 Subject: [PATCH 3/3] test: Hold disk init in DiskOnlyDelayedAttach too Empty-dir attach can finish before the pending DiskTierReady assertion under CI load, same race MixedPendingReady already fixed with RegisterInitHold. Drop t.Parallel and hold the barrier through the pending checks so disk-only stays green. Link: https://git.s1d3sw1ped.com/s1d3sw1ped/steamcache2/issues/36 --- steamcache/steamcache_test.go | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/steamcache/steamcache_test.go b/steamcache/steamcache_test.go index 32c9428..b58e8ab 100644 --- a/steamcache/steamcache_test.go +++ b/steamcache/steamcache_test.go @@ -1050,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, "") @@ -1064,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) @@ -1077,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()