Compare commits

...

34 Commits

Author SHA1 Message Date
linus 8cebc1f96c ops: Signal disk-tier attach pending vs ready (#45)
CI / vulncheck (push) Successful in 15s
CI / check-and-test (push) Successful in 41s
Release Tag / release (push) Successful in 15s
disk_tier_ready + X-SteamCache-Disk-Tier. Closes #33.
2026-09-07 12:02:10 -05:00
pike 12ea3ee4f6 ops: Signal disk-tier attach pending vs ready
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 40s
Large disk caches can look memory-only/broken during async DiskFS attach.
Expose disk_tier_ready in /metrics, pending/ready logs for disk-only and
mixed modes, and X-SteamCache-Disk-Tier on /lancache-heartbeat. GetMetrics
skips blocking disk.Size() while attach is pending so /metrics stays usable.

Closes #33
2026-09-07 16:58:13 +00:00
linus a500f51f17 cache: Serve Range GET from disk when present (#44)
CI / vulncheck (push) Successful in 14s
CI / check-and-test (push) Successful in 41s
Range HIT/MISS as 206; range_cache/range_upstream metrics. Closes #38.
2026-09-07 11:37:06 -05:00
moss ac2d36f1ad cache: Serve Range GET from disk when present
CI / vulncheck (pull_request) Successful in 15s
CI / check-and-test (pull_request) Successful in 41s
Steam clients lean on Range requests. Hit/miss looked fine while downloads
still felt cold when a Range miss returned the full upstream body as 200.

Add range_cache / range_upstream metrics. HIT path already sliced 206 from
the cached full object; count range_cache there. On MISS, keep stripping
Range for the upstream fetch (full object still cached) but serve the
client's requested slice as 206 and count range_upstream.

Tests in steamcache/range_test.go cover Range HIT (local, no upstream),
Range MISS (206 + full object cached), and /metrics emission.
2026-09-07 16:32:36 +00:00
linus c43bfba568 Merge pull request 'docs: Real DNS + empty upstream Quick Start' (#43) from docs/quick-start-real-dns into main
docs: Real DNS + empty upstream Quick Start

Fixes #32
2026-09-07 10:00:58 -05:00
eva 3411defd51 docs: Real DNS + empty upstream Quick Start
CI / vulncheck (pull_request) Successful in 15s
CI / check-and-test (pull_request) Successful in 40s
2026-09-07 14:58:06 +00:00
linus 71d5106777 docs: Add README operator cache check (#30)
CI / vulncheck (push) Successful in 14s
CI / check-and-test (push) Successful in 41s
Release Tag / release (push) Successful in 16s
README Quick check + make validate-check using existing /metrics and heartbeat. Closes #29.
2026-09-03 00:19:53 -05:00
pike 8b1b229539 docs: Deduplicate validate-check helper
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 41s
Keep a single Makefile validate-check target that prints the full
/metrics dump, highlights hit/miss fields, and curls /lancache-heartbeat.
Drop the duplicate README troubleshooting item and align the validation
chapter plus validate-config.yaml comments with that behavior.

Link: #29
2026-09-03 05:18:19 +00:00
pike ff1ab31327 docs: Add README operator cache check
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 41s
The README buried hit/miss verification under the SteamPrefill validation
chapter, and advertised make validate-check without a Makefile target.
Operators need a short start → traffic → metrics path against the existing
/metrics and /lancache-heartbeat endpoints.

Link: #29
2026-09-03 05:16:14 +00:00
s1d3sw1ped_bot 85e14bc8af README: Replace --threads with real cobra concurrency flags
Release Tag / release (push) Successful in 17s
Docs sync: align README with current install/CLI behavior.

Co-authored-by: s1d3sw1ped_bot <s1d3sw1ped+giteabot@gmail.com>
Co-committed-by: s1d3sw1ped_bot <s1d3sw1ped+giteabot@gmail.com>
2026-09-02 10:43:52 -05:00
s1d3sw1ped_bot 3f5175b482 .gitea/workflows: Prefer job GITEA_TOKEN
CI / vulncheck (push) Successful in 15s
CI / check-and-test (push) Successful in 41s
Release Tag still used secrets.RELEASE_TOKEN. Prefer the job built-in token with contents: write, same path as vulture/teleport.

Link: #26
Co-authored-by: s1d3sw1ped_bot <s1d3sw1ped+giteabot@gmail.com>
Co-committed-by: s1d3sw1ped_bot <s1d3sw1ped+giteabot@gmail.com>
2026-09-02 10:08:00 -05:00
s1d3sw1ped_bot 04d1c6c368 Merge pull request 'Promote develop: vfs/cache promoteToFast race fix' (#24) from develop into main
CI / vulncheck (push) Successful in 15s
CI / check-and-test (push) Successful in 41s
Release Tag / release (push) Successful in 30s
Promote vfs/cache Stat snapshot + promoteToFast size fix.
2026-09-01 15:36:49 -05:00
s1d3sw1ped_bot 523a9a4782 Merge pull request 'vfs/cache: Fix promoteToFast race with WriteCloser Close' (#23) from vfs/cache-promote-race into develop
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 41s
Stat returns FileInfo snapshots; promoteToFast sizes from len(content).

#21
2026-09-01 15:35:20 -05:00
s1d3sw1ped_bot 8e09c89e24 vfs/cache: Fix promoteToFast race with WriteCloser Close
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 40s
promoteToFast read FileInfo.Size from Stat after the VFS lock was
released, while WriteCloser.Close updated Size on the same live
object. Concurrent Create/Open on overlapping keys tripped the race
detector.

Stat now returns a FileInfo snapshot, and promotion uses the
ReadAll length for the fast-tier Create size. Promotion stays
best-effort.

#21
2026-09-01 20:30:40 +00:00
s1d3sw1ped_bot a2ac13d317 vfs/disk: Fix EvictDiskVisibilityAndRecreateSafety flake
CI / vulncheck (push) Successful in 13s
CI / check-and-test (push) Failing after 40s
2026-09-01 15:23:10 -05:00
s1d3sw1ped_bot 5d006ac44f vfs/disk: Fix EvictDiskVisibilityAndRecreateSafety flake
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 41s
2026-09-01 15:21:05 -05:00
s1d3sw1ped_bot acd006d4a2 ci: Skip test workflow on markdown-only pushes
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 40s
Docs-only pushes to main (**.md / CONTRIBUTING.md) do not need the
full check-and-test job. Pull requests are unchanged.
2026-09-01 20:19:19 +00:00
s1d3sw1ped_bot 30a695458e vfs/disk: Fix EvictDiskVisibilityAndRecreateSafety flake
New() launches background calculateSizeAndPopulateIndex which scans
disk and calls insertBatch. Create does not wait on initDone, so a
file can be written and indexed, discovered by the scan, then removed
from d.info and disk by EvictLRU/EvictBySize. insertBatch then
re-inserted the stale discoveredFile without checking the path still
existed, so Stat succeeded from the index while os.Stat failed.

Re-stat under the lock and skip gone files so evicted keys are not
resurrected.

Fixes #18.
2026-09-01 20:19:14 +00:00
s1d3sw1ped_bot e8bcf0ddbd docs: Add CONTRIBUTING.md
CI / vulncheck (push) Successful in 14s
CI / check-and-test (push) Successful in 41s
Add short contribution guide on main from develop tip content.
2026-09-01 14:11:02 -05:00
s1d3sw1ped_bot f497e71ef0 docs: Add CONTRIBUTING.md
CI / vulncheck (pull_request) Successful in 16s
CI / check-and-test (pull_request) Successful in 43s
Copy the develop tip contribution guide onto main so the default
branch carries the same short CONTRIBUTING.md.
2026-09-01 19:10:15 +00:00
s1d3sw1ped_bot 0198e8990b docs: Add CONTRIBUTING.md
CI / vulncheck (pull_request) Successful in 31s
CI / check-and-test (pull_request) Failing after 57s
Contributors need a short guide for develop-targeted PRs, commit
subject form, and Gitea issue-closing rules.
2026-09-01 13:37:56 -05:00
s1d3sw1ped_bot 2a2cd8d393 Bump Go to 1.27.0 (#14)
CI / vulncheck (push) Successful in 8s
CI / check-and-test (push) Successful in 35s
2026-08-31 20:04:49 -05:00
s1d3sw1ped_bot 81b3a7df53 Restrict Host-based fetches to Steam CDN names (#13)
CI / vulncheck (push) Successful in 7s
CI / check-and-test (push) Successful in 27s
Allowlist steamcontent/steampowered/steamstatic Hosts; do not follow upstream redirects.
2026-08-31 18:58:53 -05:00
s1d3sw1ped_bot 8e8e877533 Restrict Host-based origin fetches to Steam CDN names.
CI / vulncheck (pull_request) Successful in 7s
CI / check-and-test (pull_request) Successful in 28s
When upstream is empty the cache used the client Host as the fetch URL, so any LAN client with a spoofed Steam User-Agent could proxy to literal IPs or arbitrary names. Reject those hosts, stop following upstream redirects, and keep path-only cache keys so real Steam CDNs still share entries.
2026-08-31 23:56:03 +00:00
s1d3sw1ped d63d7b4d3c Merge pull request 'Run CI on push to main' (#11) from ci/test-on-push into main
CI / vulncheck (push) Successful in 8s
CI / check-and-test (push) Successful in 28s
Release Tag / release (push) Successful in 14s
Reviewed-on: #11
2026-08-31 15:47:07 -05:00
s1d3sw1ped_bot 35d698a232 Keep v1 lint rules under golangci-lint v2 (no ST/QF, no G704/G705).
CI / vulncheck (pull_request) Successful in 7s
CI / check-and-test (pull_request) Successful in 38s
2026-08-31 20:30:36 +00:00
s1d3sw1ped_bot 36b613b6cd Drop empty gosec settings so golangci-lint v2 config verifies.
CI / vulncheck (pull_request) Successful in 7s
CI / check-and-test (pull_request) Failing after 13s
2026-08-31 20:28:54 +00:00
s1d3sw1ped_bot 79f02d9868 Fix CI on Go 1.26.7: golangci-lint-action v8 and vulncheck on 1.27.
CI / vulncheck (pull_request) Successful in 19s
CI / check-and-test (pull_request) Failing after 27s
v4 installs golangci-lint v1.64.8 (built with Go 1.24) which cannot lint a
1.26 module; v8 pulls v2.12. govulncheck@latest still fails to install on
1.26.7, so run that job on Go 1.27. Keep tests on push to main.
2026-08-31 20:27:35 +00:00
s1d3sw1ped_bot 19497eba0c Pin CI to Go 1.26.7 and stable action tags.
CI / vulncheck (pull_request) Failing after 13s
CI / check-and-test (pull_request) Failing after 23s
setup-go@main plus go-version-file 1.26.0 was flaky (version: not found) and still scanned an unpatched 1.26.0 stdlib.
2026-08-31 15:18:59 -05:00
s1d3sw1ped_bot c7a2312994 Raise module Go version from EOL 1.23.0 to 1.26.0.
CI / check-and-test (pull_request) Failing after 11s
CI / vulncheck (pull_request) Failing after 21s
govulncheck fails on stdlib crypto/tls and crypto/x509 findings (including GO-2025-4008 / CVE-2025-58189) that were never patched on 1.23. 1.26 is still a supported release.
2026-08-31 15:15:48 -05:00
s1d3sw1ped_bot fbb084d824 Use latest Go 1.23 patch in CI and run tests even if govulncheck fails.
CI / vulncheck (pull_request) Failing after 18s
CI / check-and-test (pull_request) Successful in 37s
setup-go was installing go.mod's 1.23.0 exactly, so govulncheck reported stdlib x509 findings and skipped tests. check-latest gets the patched 1.23, and vulncheck is its own job.
2026-08-31 15:11:43 -05:00
s1d3sw1ped_bot e7d4a19c3f Let govulncheck@latest fetch a newer Go toolchain.
CI / check-and-test (pull_request) Failing after 22s
CI was dying on: golang.org/x/vuln@v1.7.0 requires go >= 1.25.0 (running go 1.23.0; GOTOOLCHAIN=local).
2026-08-31 15:11:18 -05:00
s1d3sw1ped_bot 0c54ef3404 Let govulncheck@latest fetch a newer Go toolchain.
CI / check-and-test (pull_request) Failing after 24s
CI was dying on: golang.org/x/vuln@v1.7.0 requires go >= 1.25.0 (running go 1.23.0; GOTOOLCHAIN=local).
2026-08-31 15:08:11 -05:00
s1d3sw1ped_bot 3d3c74fdb2 Run CI on push to main as well as pull requests.
CI / check-and-test (pull_request) Failing after 1m1s
Direct pushes to main currently skip tests because the workflow only listened for pull_request. Scratchbox already tests every push; do the same here.
2026-08-31 15:05:28 -05:00
22 changed files with 1136 additions and 165 deletions
+5 -1
View File
@@ -7,6 +7,8 @@ on:
jobs:
release:
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@main
with:
@@ -21,4 +23,6 @@ jobs:
version: 'latest'
args: release
env:
GITEA_TOKEN: ${{secrets.RELEASE_TOKEN}}
GITEA_TOKEN: ${{ secrets.GITEA_TOKEN }}
GITHUB_TOKEN: ${{ secrets.GITEA_TOKEN }}
GORELEASER_FORCE_TOKEN: gitea
+22 -8
View File
@@ -1,24 +1,38 @@
name: PR Check
name: CI
on:
- pull_request
pull_request:
push:
branches:
- main
paths-ignore:
- '**.md'
- 'CONTRIBUTING.md'
jobs:
check-and-test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@main
- uses: actions/setup-go@main
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version-file: 'go.mod'
- run: go mod tidy
- run: go build ./...
- run: go vet ./...
- name: golangci-lint
uses: golangci/golangci-lint-action@v4
uses: golangci/golangci-lint-action@v8
with:
version: latest
version: v2.13.2
args: --timeout=5m
- run: go test -race -v -shuffle=on -coverprofile=coverage.out -timeout=5m ./...
- run: go tool cover -func=coverage.out | tail -10 # basic coverage report
vulncheck:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version-file: 'go.mod'
- run: go install golang.org/x/vuln/cmd/govulncheck@latest
- run: govulncheck ./...
- run: go test -race -v -shuffle=on -coverprofile=coverage.out -timeout=5m ./...
- run: go tool cover -func=coverage.out | tail -10 # basic coverage report
+74 -66
View File
@@ -1,86 +1,94 @@
# .golangci.yml - steamcache2 lint config
# .golangci.yml - steamcache2 lint config (golangci-lint v2)
# Philosophy: enable reasonable linters by default (golangci curated set + key additions)
# then use most specific suppressions possible (source //nosec with justification,
# _ = discard for errcheck on unavoidable client writes, narrow exclude-rules only for tests).
# This makes remaining accepted issues visible and actionable in the code.
# Run with: make lint (or golangci-lint run ./...)
# Install: go install github.com/golangci/golangci-lint/cmd/golangci-lint@latest
version: "2"
run:
timeout: 5m
modules-download-mode: readonly
linters:
# No disable-all: use golangci defaults (errcheck, govet, ineffassign, staticcheck, unused, gosimple, etc.)
# No default: none — use golangci defaults (errcheck, govet, ineffassign, staticcheck, unused, etc.)
# Explicitly enable the non-default linters we require for this LAN cache proxy.
enable:
- gosec # security checks (re-audited; see source //nosec for justified cases)
- misspell # documentation hygiene
- goimports # import formatting (enforced)
# gofmt covered via linter or goimports; errcheck/govet etc. from defaults
settings:
errcheck:
check-type-assertions: false
check-blank: false
# gosec: keep source-level //nosec for G104/G115/G301/G304/G306.
# G704/G705 are new taint-analysis rules (SSRF/XSS) not present in v1.64.8;
# a CDN cache proxy forwards upstream URLs and response bodies by design.
gosec:
excludes:
- G704
- G705
# v1 staticcheck checks: ["all"] meant SA* only. v2 merged stylecheck (ST*)
# and quickfix (QF*) into staticcheck; keep the previous SA*+gosimple set.
staticcheck:
checks:
- all
- "-ST*"
- "-QF*"
govet:
enable-all: true
disable:
- fieldalignment # performance tuning not a priority for this proxy appliance
- shadow # common idiomatic "err" redeclarations in error-handling chains (large ServeHTTP, root, parse funcs); enabling adds noise with no real bugs; would require scope refactor for little gain
exclusions:
generated: lax
paths:
- dist
- bin
rules:
- path: _test\.go
linters:
- errcheck
- gosec # tests often use weak patterns intentionally (e.g. error injection, temp files)
# NOTE: narrow SA9003 exclude retained only for the one remaining intentional empty branch in test (best-effort status check; main assert is metrics side-effect).
- path: steamcache/steamcache_test.go
linters:
- staticcheck
text: "SA9003: empty branch"
# Narrow gosec excludes for unavoidable classes after re-audit (LAN proxy threat model):
# - G115: int64<->uint casts in eviction/GC math (all sizes positive, guarded by capacity checks; API uses uint for bytesNeeded)
# - G304: path vars for Read/Open/Remove under trusted disk.root or user config file (sanitized keys, no traversal, no arbitrary inclusion from untrusted URLs)
# G306 for config WriteFile kept as source //nosec (one site).
# G301 fixed at source (0700 dirs). G104 addressed via errcheck fixes.
- path: vfs/memory/memory.go
linters:
- gosec
text: "G115"
- path: vfs/disk/disk.go
linters:
- gosec
text: "G115"
- path: vfs/gc/gc.go
linters:
- gosec
text: "G115"
- path: config/config.go
linters:
- gosec
text: "G304"
- path: vfs/disk/disk.go
linters:
- gosec
text: "G304"
linters-settings:
errcheck:
check-type-assertions: false
check-blank: false
gosec:
# Broad global excludes removed (G104/G115/G301/G304/G306).
# - G301 addressed by switching cache MkdirAll to 0700 (least privilege for CDN content).
# - Remaining justified cases documented with precise //nosec (or #nosec) + comments at the call sites.
# - G104 largely eliminated by errcheck + explicit _ = handling (or defer wrappers).
staticcheck:
checks: ["all"] # SA1019 exclusion removed (no deprecated API usages in tree)
govet:
enable-all: true
disable:
- fieldalignment # performance tuning not a priority for this proxy appliance
- shadow # common idiomatic "err" redeclarations in error-handling chains (large ServeHTTP, root, parse funcs); enabling adds noise with no real bugs; would require scope refactor for little gain
# Old global errcheck disable + aspirational "re-enable after refactors" comments deleted.
# errcheck is now on via defaults. Unavoidable cases handled at source with _ = or (rarely) narrow rules.
formatters:
enable:
- goimports
exclusions:
generated: lax
paths:
- dist
- bin
issues:
max-issues-per-linter: 0
max-same-issues: 0
exclude-use-default: false
exclude-dirs:
- dist
- bin
exclude-rules:
- path: _test\.go
linters:
- errcheck
- gosec # tests often use weak patterns intentionally (e.g. error injection, temp files)
# NOTE: narrow SA9003 exclude retained only for the one remaining intentional empty branch in test (best-effort status check; main assert is metrics side-effect).
# The config one was a truly redundant check (already errored above); deleted surgically in Fix Round 1 (Issue 1), eliminating its exclude-rule.
- path: steamcache/steamcache_test.go
linters:
- staticcheck
text: "SA9003: empty branch"
# Narrow gosec excludes for unavoidable classes after re-audit (LAN proxy threat model):
# - G115: int64<->uint casts in eviction/GC math (all sizes positive, guarded by capacity checks; API uses uint for bytesNeeded)
# - G304: path vars for Read/Open/Remove under trusted disk.root or user config file (sanitized keys, no traversal, no arbitrary inclusion from untrusted URLs)
# G306 for config WriteFile kept as source //nosec (one site).
# G301 fixed at source (0700 dirs). G104 addressed via errcheck fixes.
- path: vfs/memory/memory.go
linters:
- gosec
text: "G115"
- path: vfs/disk/disk.go
linters:
- gosec
text: "G115"
- path: vfs/gc/gc.go
linters:
- gosec
text: "G115"
- path: config/config.go
linters:
- gosec
text: "G304"
- path: vfs/disk/disk.go
linters:
- gosec
text: "G304"
# Predictive/* rules deleted: vfs/predictive/ removed in commit 0dbb2e0; rules were stale/dead.
# All other suppressions use source-level //nosec (gosec) or _= (errcheck) for precision and visibility.
+40
View File
@@ -0,0 +1,40 @@
# Contributing
## Propose changes
Open a pull request against `develop`. Keep the default branch for releases and
stable tips; land work on `develop` first.
Point at an existing issue when one fits. Prefer a short issue that states the
symptom or request before a large PR.
## Commits
Subject form:
```
area: Imperative summary
```
- **Area** is a real package, directory, or subsystem token (`ci:`, `docs:`,
Go package name). Not a lone filename.
- **Imperative** mood: Fix, Add, Remove — not "Fixed" or "This patch…".
- No trailing period. Aim ≤ ~7075 characters for the whole subject.
- Not conventional-commits (`feat:` / `fix:` / `chore:` as types).
Body explains **why**. Establish the problem, then say what you are doing.
One logical change per commit; split fix and cleanup.
## Pull requests
Title matches the primary commit subject.
- **What** changed
- **Why** (problem and impact)
- **Test** (concrete steps; "CI green" alone is weak)
## Issues and closing
Cite leftover issues by **full URL**. Gitea closes issues when `#N` appears in
merge text, so do not put `#N` in the merge message unless that issue is actually
done. Use `Fixes #N` / `Closes #N` only when the leftover work is finished.
+26 -1
View File
@@ -62,6 +62,30 @@ validate run-validation: build clean-disk ## Start steamcache2 on :80 with small
fi; \
exec "$$BINARY" --config docs/examples/validate-config.yaml --log-level info
validate-check: ## Curl local /metrics (full dump + hit/miss fields) and /lancache-heartbeat (default :80)
@echo "=== http://localhost/metrics ==="
@metrics=$$(curl -sf --max-time 5 http://localhost/metrics) || { \
echo "ERROR: could not fetch http://localhost/metrics"; \
echo "Is steamcache2 running on the default listen address :80?"; \
exit 1; \
}; \
printf '%s\n' "$$metrics"; \
echo ""; \
echo "=== hit/miss fields ==="; \
printf '%s\n' "$$metrics" | grep -E '^(total_requests|cache_hits|cache_misses|hit_rate|memory_cache_hits|disk_cache_hits|errors) ' || true; \
echo ""; \
echo "=== http://localhost/lancache-heartbeat (GET; expect 204 + X-LanCache-Processed-By: SteamCache2) ==="; \
hb=$$(curl -sD - -o /dev/null --max-time 5 http://localhost/lancache-heartbeat) || { \
echo "ERROR: could not fetch http://localhost/lancache-heartbeat"; \
echo "Is steamcache2 running on the default listen address :80?"; \
exit 1; \
}; \
printf '%s\n' "$$hb"; \
echo "$$hb" | grep -q '204' && echo "$$hb" | grep -qi 'X-LanCache-Processed-By' || { \
echo "ERROR: expected HTTP 204 and X-LanCache-Processed-By on /lancache-heartbeat"; \
exit 1; \
}
validate-kill: ## Kill leftover steamcache2 processes (safer, checks process name)
@echo "Looking for steamcache2 processes on common validation ports (80 is primary)..."
@for port in 80 8040 8080; do \
@@ -109,6 +133,7 @@ help: ## Show this help message
@echo " clean-disk Remove disk cache"
@echo " bench Run low-level VFS microbenchmarks"
@echo " validate / run-validation Start server on :80 (builds, auto-setcaps fresh binary, then runs as normal user, cleans disk cache first)"
@echo " validate-check Curl local /metrics (full dump + hit/miss fields) and /lancache-heartbeat (default :80)"
@echo " setcap Explicitly set cap on current build (for port 80 use outside validate)"
@echo " validate-kill Kill leftover steamcache2 processes (safer)"
@echo " prefill Download latest SteamPrefill into bin/steam-prefill/SteamPrefill (gitignored)"
@echo " prefill Download latest SteamPrefill into bin/steam-prefill/SteamPrefill (gitignored)"
+85 -18
View File
@@ -10,8 +10,8 @@ SteamCache2 is a blazing fast download cache for Steam, designed to reduce bandw
- Reduces bandwidth usage
- Easy to set up and configure aside from dns stuff to trick Steam into using it
- Supports multiple clients
- **NEW:** YAML configuration system with automatic config generation
- **NEW:** Simple Makefile for development workflow
- YAML configuration with automatic config generation on first run
- Makefile for development and validation workflows
- Cross-platform builds (Linux, macOS, Windows)
## Quick Start
@@ -34,7 +34,10 @@ SteamCache2 is a blazing fast download cache for Steam, designed to reduce bandw
The application will automatically create a `config.yaml` file with default settings and exit, allowing you to customize it.
3. **Edit the configuration** (`config.yaml`):
3. **Edit the configuration** (`config.yaml`) for a real Steam front-door:
Leave `upstream` empty. With an empty upstream, steamcache2 fetches from the request `Host` (the same pattern SteamPrefill / real Steam clients use when DNS points them at your cache). Do **not** paste a fake host like `https://steam.cdn.com` — that is not a Steam CDN and will not put you in front of Steam.
```yaml
listen_address: :80
cache:
@@ -45,14 +48,62 @@ SteamCache2 is a blazing fast download cache for Steam, designed to reduce bandw
size: 10GB
path: ./disk
gc_algorithm: hybrid
upstream: "https://steam.cdn.com" # Set your upstream server
# Empty upstream = use the client Host as the origin (Steam CDN names only).
upstream: ""
```
4. **Run the application again:**
4. **Point Steam (or SteamPrefill) at this cache** before you expect hits:
- **LAN DNS:** resolve `lancache.steamcontent.com` (and other Steam content names your clients use) to this server's LAN IP.
- **Single Windows PC:** add a hosts override — see [Windows Hosts File Override](#windows-hosts-file-override) (`<cache-ip> lancache.steamcontent.com`).
- Restart Steam (or your prefill tool) after DNS/hosts changes.
Empty-upstream direct fetch only allows Steam CDN host suffixes (`steamcontent.com`, `steampowered.com`, `steamstatic.com`). Literal IPs and unrelated hosts are rejected.
5. **Run the application again:**
```bash
make run # or ./steamcache2
```
### Quick check: is it caching?
After steamcache2 is running (default `listen_address: :80`) and has seen a little Steam traffic — a game download, a short SteamPrefill pass, or any cacheable request — confirm hits vs misses from the existing endpoints. You do not need a full benchmark or log diving.
```bash
make validate-check
# or, manually:
curl -s http://localhost/metrics
curl -s -i http://localhost/lancache-heartbeat
```
`make validate-check` prints the full `/metrics` dump, highlights hit/miss fields, and curls `/lancache-heartbeat`. Read these fields:
| Field | Meaning |
| --- | --- |
| `cache_hits` / `cache_misses` / `hit_rate` | Whether later requests were served from cache |
| `range_cache` / `range_upstream` | Range GETs served as 206 from a cached object vs after a full upstream fetch |
| `memory_cache_hits` / `disk_cache_hits` | Which tier served the hits |
| `total_requests` / `errors` | Volume and failures |
| `disk_tier_ready` | `0` while disk slow-tier attach pending; `1` when attached, or when no disk configured (N/A — not waiting) |
A first pass through new content is mostly misses (`hit_rate` near 0). Repeat the same content and `cache_hits` / `hit_rate` should rise.
Steam clients lean on Range requests. When an object is already cached, a Range GET is served locally as 206 from that full object (`range_cache`). On a Range miss the cache still fetches and stores the full upstream body, then returns the requested byte range as 206 (`range_upstream`).
To confirm the process is up (HTTP 204 and `X-LanCache-Processed-By: SteamCache2`):
```bash
curl -s -i http://localhost/lancache-heartbeat
```
Use GET (`curl -i`), not HEAD (`curl -I`): the server only accepts GET.
Heartbeat also returns `X-SteamCache-Disk-Tier: pending|ready|disabled` (`disabled` = memory-only / no disk; `pending`/`ready` = disk configured attach state).
These are the cache process's own `/metrics` and `/lancache-heartbeat` endpoints. There is no separate metrics daemon.
If you changed `listen_address`, point curl at that host:port instead. For a full SteamPrefill validation workflow (small caches, coalescing, GC), see [Validating Full Functionality](#validating-full-functionality-with-external-tools).
### Development Workflow
Use `make` for the majority of common development tasks. The Makefile handles running tests, linting, hygiene checks, building, running the application, and other routine boilerplate work.
@@ -96,13 +147,16 @@ When the server is running, point your external SteamPrefill (or other load gene
./SteamPrefill benchmark run ...
```
When finished, you can get a quick metrics summary with:
When finished, you can get a quick metrics + heartbeat report with:
```bash
make validate-check
# or, manually:
curl -s http://localhost/metrics
curl -s -i http://localhost/lancache-heartbeat
```
This is the recommended simple workflow. No automatic downloading or running of external tools.
See [Quick check: is it caching?](#quick-check-is-it-caching) for which fields to read. This is the recommended simple workflow. No automatic downloading or running of external tools.
#### Inspecting the Result
@@ -110,14 +164,17 @@ After a benchmark run you can ask for a quick report:
```bash
make validate-check
# or manually:
# or, manually:
curl -s http://localhost/metrics
curl -s -i http://localhost/lancache-heartbeat
```
Look for:
- High cache hit rate after the warmup pass
`make validate-check` prints the full `/metrics` dump, highlights hit/miss fields, and curls `/lancache-heartbeat`. Look for:
- High cache hit rate after the warmup pass (`cache_hits`, `hit_rate`, plus `memory_cache_hits` / `disk_cache_hits`)
- Non-zero `coalesced` and `disk` activity
- Zero unexpected errors
- Zero unexpected `errors`
Heartbeat should be HTTP 204 with `X-LanCache-Processed-By: SteamCache2`. Use GET (`curl -i`), not HEAD (`curl -I`).
#### The Validation Config
@@ -128,7 +185,7 @@ The recommended validation config is at [docs/examples/validate-config.yaml](doc
Running a realistic SteamPrefill benchmark workload through a built steamcache2 exercises the complete public surface that matters for production use:
- Steam User-Agent detection and depot/manifest/chunk URL patterns
- Full MISS → cache write → HIT (and HIT-COALESCED) paths
- Range request handling from cached full responses
- Range request handling from cached full responses (local 206 on HIT via `range_cache`; MISS fetches the full object then serves the requested slice as 206 via `range_upstream`)
- Request coalescing under concurrent load
- Memory tier + disk tier interaction (including async disk attach)
- Garbage collection and eviction under pressure
@@ -160,8 +217,9 @@ While most configuration is done via the YAML file, some runtime options are sti
# Set logging level
./steamcache2 --log-level debug --log-format json
# Set number of worker threads
./steamcache2 --threads 8
# Override concurrency from the CLI (0 = use config.yaml)
./steamcache2 --max-concurrent-requests 8
./steamcache2 --max-requests-per-client 4
# Show help
./steamcache2 --help
@@ -198,8 +256,9 @@ cache:
gc_algorithm: hybrid
# Upstream server configuration
# The upstream server to proxy requests to
upstream: "https://steam.cdn.com"
# Leave empty to fetch from the request Host (Steam CDN names only).
# Set only when chaining caches (table RAM cache -> room disk cache).
upstream: ""
```
#### Startup Validation
@@ -234,6 +293,9 @@ See `config.Validate()` and `steamcache.New` error paths. This ensures the LAN a
- The explicit startup guard (reduce size if pre-existing on-disk > cap) runs as the literal last step of bg init, before the barrier opens.
- Add a note for operators: very large disk caches (tens/hundreds GB with millions files) may show extended "memory-only or no-cache" behavior at startup (seconds to minutes depending on storage speed); this is by design for responsiveness.
- Godoc on `disk.New` and `DiskFS.Size` expanded with the barrier/attach behavior.
- 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.
#### Garbage Collection Algorithms
@@ -311,7 +373,7 @@ This will direct any requests to `lancache.steamcontent.com` to your SteamCache2
### Prerequisites
- Go 1.19 or later
- Go 1.27.0 or later
- Make (optional, but recommended)
### Build Commands
@@ -319,7 +381,7 @@ This will direct any requests to `lancache.steamcontent.com` to your SteamCache2
```bash
# Clone the repository
git clone <repository-url>
cd SteamCache2
cd steamcache2
# Download dependencies
make deps
@@ -371,6 +433,11 @@ make
- Consider using a different GC algorithm like `hybrid`
- Adjust the disk cache size to match available storage
6. **Not sure if it is caching**
- Do not start with the full SteamPrefill chapter. Use [Quick check: is it caching?](#quick-check-is-it-caching): `make validate-check` (full `/metrics`, hit/miss fields, and `/lancache-heartbeat`)
- A first pass is mostly `cache_misses`; repeating the same content should raise `cache_hits` / `hit_rate`
- Confirm the process is up with `curl -s -i http://localhost/lancache-heartbeat` (GET, not HEAD)
### Getting Help
- Check the logs for detailed error messages
+3
View File
@@ -27,7 +27,10 @@
# SteamPrefill benchmark run -c 20 ...
#
# After the benchmark run, inspect with:
# make validate-check # full /metrics + hit/miss fields + /lancache-heartbeat
# # or, manually:
# curl -s http://localhost/metrics
# curl -s -i http://localhost/lancache-heartbeat # GET, not HEAD
#
# Tweak sizes upward if you want to run very large workloads while still
# exercising the disk tier (workload >> RAM is ideal for real disk testing).
+1 -1
View File
@@ -1,6 +1,6 @@
module s1d3sw1ped/steamcache2
go 1.23.0
go 1.27.0
require (
github.com/docker/go-units v0.5.0
+7
View File
@@ -262,6 +262,13 @@ func (sc *SteamCache) streamCachedResponse(w http.ResponseWriter, r *http.Reques
// Send the range data
_, _ = w.Write(rangeData) // client write error ignored (disconnect during range body send is not actionable)
// Range served from cache: count the range-specific metric and the range
// bytes actually written (handleCacheHit skips full-blob byte counting for
// Range requests so these are not double-counted).
sc.metrics.IncrementRangeCache()
sc.metrics.AddBytesServed(int64(len(rangeData)))
sc.metrics.AddBytesSaved(int64(len(rangeData)))
logger.Logger.Info().
Str("cache_key", cacheKey).
Str("url", r.URL.String()).
+59 -6
View File
@@ -56,6 +56,15 @@ func (sc *SteamCache) handleSpecialEndpoints(w http.ResponseWriter, r *http.Requ
logger.Logger.Debug().
Str("client_ip", clientIP).
Msg("LanCache heartbeat request")
diskTier := "disabled"
if sc.disk != nil {
if sc.metrics.GetDiskTierReady() == 1 {
diskTier = "ready"
} else {
diskTier = "pending"
}
}
w.Header().Add("X-SteamCache-Disk-Tier", diskTier)
w.Header().Add("X-LanCache-Processed-By", "SteamCache2")
w.WriteHeader(http.StatusNoContent)
_, _ = w.Write(nil) // client write error ignored (heartbeat path; nil write is no-op)
@@ -109,8 +118,14 @@ func (sc *SteamCache) handleCacheHit(w http.ResponseWriter, r *http.Request, cac
// Track cache hit metrics
sc.metrics.IncrementCacheHits()
sc.metrics.AddResponseTime(time.Since(tstart))
sc.metrics.AddBytesServed(int64(len(cachedData)))
sc.metrics.AddBytesSaved(int64(len(cachedData)))
if r.Header.Get("Range") == "" {
// Full-object HIT: count the cached blob served. Range HITs skip
// this here — streamCachedResponse counts the range bytes actually
// written to the client instead, so BytesServed/Saved reflect the
// partial body and are not double-counted.
sc.metrics.AddBytesServed(int64(len(cachedData)))
sc.metrics.AddBytesSaved(int64(len(cachedData)))
}
sc.metrics.IncrementServiceRequests(service.Name)
logger.Logger.Debug().
@@ -345,6 +360,18 @@ func (sc *SteamCache) ServeHTTP(w http.ResponseWriter, r *http.Request) {
req.Host = r.Host
} else { // if no upstream server is configured, proxy the request to the host specified in the request
host := r.Host
if !hostAllowedForDirectFetch(host) {
logger.Logger.Warn().
Str("host", host).
Str("client_ip", clientIP).
Msg("Rejecting direct-fetch Host (not a Steam CDN name)")
sc.metrics.IncrementErrors()
if isNew {
coalescedReq.complete(nil, fmt.Errorf("host not allowed for direct fetch"))
}
http.Error(w, "Invalid URL", http.StatusBadRequest)
return
}
if r.Header.Get("X-Sls-Https") == "enable" {
host = "https://" + host
} else {
@@ -537,14 +564,40 @@ func (sc *SteamCache) ServeHTTP(w http.ResponseWriter, r *http.Request) {
w.Header().Set("X-LanCache-Status", "MISS")
w.Header().Set("X-LanCache-Processed-By", "SteamCache2")
// Stream the response body to client
w.WriteHeader(resp.StatusCode)
_, _ = w.Write(bodyData) // client write error ignored (disconnect during MISS body send is not actionable)
// Stream the response body to client.
// Range miss: the Range header was stripped for the upstream fetch (so the
// FULL object is cached below); serve the client's requested slice from the
// full body as 206, matching the HIT Range path.
if rangeHeader := r.Header.Get("Range"); rangeHeader != "" {
start, end, totalSize, rangeValid := parseRangeHeader(rangeHeader, int64(len(bodyData)))
if !rangeValid {
// Invalid range — 416 (consistent with the HIT Range path). Drop the
// upstream Content-Length: it describes the full body, which 416 does
// not send (a stale CL would hang clients waiting for a body).
w.Header().Del("Content-Length")
w.Header().Set("Content-Range", fmt.Sprintf("bytes */%d", len(bodyData)))
w.WriteHeader(http.StatusRequestedRangeNotSatisfiable)
} else {
rangeData := bodyData[start : end+1]
w.Header().Set("Content-Range", fmt.Sprintf("bytes %d-%d/%d", start, end, totalSize))
w.Header().Set("Content-Length", fmt.Sprintf("%d", len(rangeData)))
w.Header().Set("Accept-Ranges", "bytes")
w.WriteHeader(http.StatusPartialContent)
_, _ = w.Write(rangeData) // client write error ignored (disconnect during MISS range body send is not actionable)
// Range required an upstream fetch (full object) then served as 206
sc.metrics.IncrementRangeUpstream()
sc.metrics.AddBytesServed(int64(len(rangeData))) // range bytes only, not the full cached object
}
} else {
w.WriteHeader(resp.StatusCode)
_, _ = w.Write(bodyData) // client write error ignored (disconnect during MISS body send is not actionable)
sc.metrics.AddBytesServed(int64(len(bodyData)))
}
// Track cache miss metrics
sc.metrics.IncrementCacheMisses()
sc.metrics.AddResponseTime(time.Since(tstart))
sc.metrics.AddBytesServed(int64(len(bodyData)))
sc.metrics.IncrementServiceRequests(service.Name)
// Verify we received the complete file by checking Content-Length
+36
View File
@@ -16,6 +16,8 @@ type Metrics struct {
CacheHits int64
CacheMisses int64
CacheCoalesced int64
RangeCache int64 // Range requests served as 206 from an already-cached object (HIT)
RangeUpstream int64 // Range requests that required an upstream fetch (full object), served as 206
Errors int64
RateLimited int64
@@ -31,6 +33,7 @@ type Metrics struct {
DiskCacheHits int64
Promotions int64
Evictions int64
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
@@ -78,6 +81,17 @@ func (m *Metrics) IncrementCacheCoalesced() {
atomic.AddInt64(&m.CacheCoalesced, 1)
}
// IncrementRangeCache increments the Range-from-cache counter (HIT served as 206)
func (m *Metrics) IncrementRangeCache() {
atomic.AddInt64(&m.RangeCache, 1)
}
// IncrementRangeUpstream increments the Range-from-upstream counter (full object
// fetched upstream, requested slice served as 206)
func (m *Metrics) IncrementRangeUpstream() {
atomic.AddInt64(&m.RangeUpstream, 1)
}
// IncrementErrors increments the error counter
func (m *Metrics) IncrementErrors() {
atomic.AddInt64(&m.Errors, 1)
@@ -114,6 +128,17 @@ func (m *Metrics) SetDiskCacheSize(size int64) {
atomic.StoreInt64(&m.DiskCacheSize, size)
}
// SetDiskTierReady sets whether the disk slow tier is attached (1) or still pending (0).
// Memory-only (no disk) also uses 1 — meaning "not waiting on disk attach". Reset does not clear this.
func (m *Metrics) SetDiskTierReady(ready int64) {
atomic.StoreInt64(&m.DiskTierReady, ready)
}
// GetDiskTierReady returns 1 if disk tier is ready (or no disk configured), else 0 while attach pending.
func (m *Metrics) GetDiskTierReady() int64 {
return atomic.LoadInt64(&m.DiskTierReady)
}
// IncrementMemoryCacheHits increments memory cache hits
func (m *Metrics) IncrementMemoryCacheHits() {
atomic.AddInt64(&m.MemoryCacheHits, 1)
@@ -188,6 +213,8 @@ func (m *Metrics) GetStats() *Stats {
CacheHits: cacheHits,
CacheMisses: cacheMisses,
CacheCoalesced: atomic.LoadInt64(&m.CacheCoalesced),
RangeCache: atomic.LoadInt64(&m.RangeCache),
RangeUpstream: atomic.LoadInt64(&m.RangeUpstream),
Errors: atomic.LoadInt64(&m.Errors),
RateLimited: atomic.LoadInt64(&m.RateLimited),
HitRate: hitRate,
@@ -196,6 +223,7 @@ func (m *Metrics) GetStats() *Stats {
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),
@@ -215,6 +243,8 @@ func (m *Metrics) Reset() {
atomic.StoreInt64(&m.CacheHits, 0)
atomic.StoreInt64(&m.CacheMisses, 0)
atomic.StoreInt64(&m.CacheCoalesced, 0)
atomic.StoreInt64(&m.RangeCache, 0)
atomic.StoreInt64(&m.RangeUpstream, 0)
atomic.StoreInt64(&m.Errors, 0)
atomic.StoreInt64(&m.RateLimited, 0)
atomic.StoreInt64(&m.TotalResponseTime, 0)
@@ -244,6 +274,8 @@ type Stats struct {
CacheHits int64
CacheMisses int64
CacheCoalesced int64
RangeCache int64
RangeUpstream int64
Errors int64
RateLimited int64
HitRate float64
@@ -253,6 +285,7 @@ type Stats struct {
MemoryCacheSize int64
DiskCacheSize int64
DiskTierReady int64
MemoryCacheHits int64
DiskCacheHits int64
Promotions int64
@@ -276,6 +309,8 @@ func WriteText(w http.ResponseWriter, stats *Stats) {
_, _ = fmt.Fprintf(w, "cache_hits %d\n", stats.CacheHits)
_, _ = fmt.Fprintf(w, "cache_misses %d\n", stats.CacheMisses)
_, _ = fmt.Fprintf(w, "cache_coalesced %d\n", stats.CacheCoalesced)
_, _ = fmt.Fprintf(w, "range_cache %d\n", stats.RangeCache)
_, _ = fmt.Fprintf(w, "range_upstream %d\n", stats.RangeUpstream)
_, _ = fmt.Fprintf(w, "errors %d\n", stats.Errors)
_, _ = fmt.Fprintf(w, "rate_limited %d\n", stats.RateLimited)
_, _ = fmt.Fprintf(w, "upstream_errors %d\n", stats.UpstreamErrors)
@@ -297,5 +332,6 @@ func WriteText(w http.ResponseWriter, stats *Stats) {
_, _ = fmt.Fprintf(w, "memory_cache_size %d\n", stats.MemoryCacheSize)
_, _ = fmt.Fprintf(w, "disk_cache_size %d\n", stats.DiskCacheSize)
_, _ = fmt.Fprintf(w, "disk_tier_ready %d\n", stats.DiskTierReady)
_, _ = fmt.Fprintf(w, "uptime_seconds %.2f\n", stats.Uptime.Seconds())
}
+386
View File
@@ -0,0 +1,386 @@
// steamcache/range_test.go
package steamcache
import (
"io"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
)
// TestRangeHitServedLocally verifies that a Range GET against an already-cached
// object is served locally as 206 (no upstream re-fetch) and increments range_cache.
func TestRangeHitServedLocally(t *testing.T) {
body := []byte("0123456789abcdef") // 16 bytes
var upstreamCalls atomic.Int64
f := func(w http.ResponseWriter, r *http.Request) {
upstreamCalls.Add(1)
w.Header().Set("Content-Type", "application/octet-stream")
_, _ = w.Write(body)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
srv := newCacheServer(t, sc)
c := &http.Client{Timeout: 5 * time.Second}
// 1) Populate the cache with a full (non-Range) MISS.
req, err := http.NewRequest("GET", srv.URL+"/depot/rangetest/chunk", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
resp, err := c.Do(req)
if err != nil {
t.Fatalf("miss GET: %v", err)
}
if resp.StatusCode != http.StatusOK {
t.Fatalf("miss GET: expected 200, got %d", resp.StatusCode)
}
_, _ = io.Copy(io.Discard, resp.Body)
_ = resp.Body.Close()
if got := sc.GetMetrics().CacheMisses; got < 1 {
t.Fatalf("expected CacheMisses >= 1 after first GET, got %d", got)
}
// 2) Range GET against the same URL — must be a local 206 HIT.
req2, err := http.NewRequest("GET", srv.URL+"/depot/rangetest/chunk", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req2.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
req2.Header.Set("Range", "bytes=4-7")
resp2, err := c.Do(req2)
if err != nil {
t.Fatalf("range GET: %v", err)
}
data, err := io.ReadAll(resp2.Body)
_ = resp2.Body.Close()
if err != nil {
t.Fatalf("read range body: %v", err)
}
if resp2.StatusCode != http.StatusPartialContent {
t.Fatalf("range GET: expected 206, got %d", resp2.StatusCode)
}
if string(data) != "4567" {
t.Errorf("range GET: expected body %q, got %q", "4567", data)
}
if got := resp2.Header.Get("X-LanCache-Status"); got != "HIT" {
t.Errorf("range GET: expected X-LanCache-Status HIT, got %q", got)
}
if got := resp2.Header.Get("Content-Range"); got != "bytes 4-7/16" {
t.Errorf("range GET: expected Content-Range bytes 4-7/16, got %q", got)
}
if got := resp2.Header.Get("Accept-Ranges"); got != "bytes" {
t.Errorf("range GET: expected Accept-Ranges bytes, got %q", got)
}
// Upstream must NOT have been hit again.
if got := upstreamCalls.Load(); got != 1 {
t.Errorf("upstream hit %d times, want exactly 1 (range HIT must be local)", got)
}
// range_cache incremented, range_upstream untouched.
stats := sc.GetMetrics()
if stats.RangeCache != 1 {
t.Errorf("expected RangeCache == 1 after range HIT, got %d", stats.RangeCache)
}
if stats.RangeUpstream != 0 {
t.Errorf("expected RangeUpstream == 0 (no range miss yet), got %d", stats.RangeUpstream)
}
if stats.CacheHits < 1 {
t.Errorf("expected CacheHits >= 1 after range HIT, got %d", stats.CacheHits)
}
// BytesServed: 16 (full miss body) + 4 (range bytes) = 20.
if stats.TotalBytesServed != 20 {
t.Errorf("expected TotalBytesServed == 20 (16 + range 4), got %d", stats.TotalBytesServed)
}
}
// TestRangeMissServes206FromFullFetch verifies that a Range GET on a cold key fetches
// the FULL object from upstream (Range stripped), caches it, and serves the requested
// slice as 206 with range_upstream incremented.
func TestRangeMissServes206FromFullFetch(t *testing.T) {
body := []byte("0123456789abcdef") // 16 bytes
var upstreamCalls atomic.Int64
var upstreamSawRange atomic.Bool
f := func(w http.ResponseWriter, r *http.Request) {
upstreamCalls.Add(1)
if r.Header.Get("Range") != "" {
upstreamSawRange.Store(true)
}
w.Header().Set("Content-Type", "application/octet-stream")
_, _ = w.Write(body)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
srv := newCacheServer(t, sc)
c := &http.Client{Timeout: 5 * time.Second}
// Range GET on a cold key.
req, err := http.NewRequest("GET", srv.URL+"/depot/rangetest/chunk2", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
req.Header.Set("Range", "bytes=0-3")
resp, err := c.Do(req)
if err != nil {
t.Fatalf("range miss GET: %v", err)
}
data, err := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if err != nil {
t.Fatalf("read body: %v", err)
}
if resp.StatusCode != http.StatusPartialContent {
t.Fatalf("range miss GET: expected 206, got %d", resp.StatusCode)
}
if string(data) != "0123" {
t.Errorf("range miss GET: expected body %q, got %q", "0123", data)
}
if got := resp.Header.Get("X-LanCache-Status"); got != "MISS" {
t.Errorf("range miss GET: expected X-LanCache-Status MISS, got %q", got)
}
if got := resp.Header.Get("Content-Range"); got != "bytes 0-3/16" {
t.Errorf("range miss GET: expected Content-Range bytes 0-3/16, got %q", got)
}
if got := resp.Header.Get("Accept-Ranges"); got != "bytes" {
t.Errorf("range miss GET: expected Accept-Ranges bytes, got %q", got)
}
// Range must have been stripped for the upstream fetch (full file cached).
if upstreamSawRange.Load() {
t.Error("upstream received a Range header; Range must be stripped so the full object is cached")
}
if got := upstreamCalls.Load(); got != 1 {
t.Errorf("upstream hit %d times, want exactly 1", got)
}
// range_upstream incremented, range_cache untouched.
stats := sc.GetMetrics()
if stats.RangeUpstream != 1 {
t.Errorf("expected RangeUpstream == 1 after range MISS, got %d", stats.RangeUpstream)
}
if stats.RangeCache != 0 {
t.Errorf("expected RangeCache == 0 (no range hit yet), got %d", stats.RangeCache)
}
if stats.CacheMisses < 1 {
t.Errorf("expected CacheMisses >= 1, got %d", stats.CacheMisses)
}
// BytesServed for the range miss: only the 4 range bytes, not the 16-byte fetch.
if stats.TotalBytesServed != 4 {
t.Errorf("expected TotalBytesServed == 4 (range bytes only), got %d", stats.TotalBytesServed)
}
// The FULL object must have been cached: a subsequent full GET is a HIT with the
// complete 16-byte body.
req3, err := http.NewRequest("GET", srv.URL+"/depot/rangetest/chunk2", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req3.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
resp3, err := c.Do(req3)
if err != nil {
t.Fatalf("full GET after range miss: %v", err)
}
fullData, err := io.ReadAll(resp3.Body)
_ = resp3.Body.Close()
if err != nil {
t.Fatalf("read full body: %v", err)
}
if resp3.StatusCode != http.StatusOK {
t.Fatalf("full GET after range miss: expected 200, got %d", resp3.StatusCode)
}
if got := resp3.Header.Get("X-LanCache-Status"); got != "HIT" {
t.Errorf("full GET after range miss: expected X-LanCache-Status HIT, got %q", got)
}
if len(fullData) != len(body) || string(fullData) != string(body) {
t.Errorf("full GET after range miss: expected full %d-byte body, got %d bytes", len(body), len(fullData))
}
if got := upstreamCalls.Load(); got != 1 {
t.Errorf("upstream hit %d times after HIT, want still 1", got)
}
}
// TestRangeMissInvalidRange416 verifies that an unsatisfiable Range on a cold key
// fetches upstream, returns 416 (as on the HIT path), and does not count range_upstream.
func TestRangeMissInvalidRange416(t *testing.T) {
body := []byte("0123456789abcdef") // 16 bytes
var upstreamCalls atomic.Int64
f := func(w http.ResponseWriter, r *http.Request) {
upstreamCalls.Add(1)
_, _ = w.Write(body)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
srv := newCacheServer(t, sc)
c := &http.Client{Timeout: 5 * time.Second}
req, err := http.NewRequest("GET", srv.URL+"/depot/rangetest/chunk3", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
req.Header.Set("Range", "bytes=100-200")
resp, err := c.Do(req)
if err != nil {
t.Fatalf("invalid range GET: %v", err)
}
data, err := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if err != nil {
t.Fatalf("read body: %v", err)
}
if resp.StatusCode != http.StatusRequestedRangeNotSatisfiable {
t.Fatalf("invalid range GET: expected 416, got %d", resp.StatusCode)
}
if len(data) != 0 {
t.Errorf("invalid range GET: expected empty body, got %d bytes", len(data))
}
if got := resp.Header.Get("Content-Range"); got != "bytes */16" {
t.Errorf("invalid range GET: expected Content-Range bytes */16, got %q", got)
}
if got := upstreamCalls.Load(); got != 1 {
t.Errorf("upstream hit %d times, want exactly 1 (fetch happens, then 416 to client)", got)
}
stats := sc.GetMetrics()
if stats.RangeUpstream != 0 {
t.Errorf("expected RangeUpstream == 0 for unsatisfiable range, got %d", stats.RangeUpstream)
}
if stats.RangeCache != 0 {
t.Errorf("expected RangeCache == 0, got %d", stats.RangeCache)
}
}
// TestRangeMetricsWriteText verifies /metrics emits the range_cache and range_upstream
// lines with the expected values after a range HIT and a range MISS.
func TestRangeMetricsWriteText(t *testing.T) {
body := []byte("0123456789abcdef")
f := func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write(body)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
srv := newCacheServer(t, sc)
c := &http.Client{Timeout: 5 * time.Second}
get := func(path, rangeHeader string) int {
t.Helper()
req, err := http.NewRequest("GET", srv.URL+path, nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
if rangeHeader != "" {
req.Header.Set("Range", rangeHeader)
}
resp, err := c.Do(req)
if err != nil {
t.Fatalf("GET %s: %v", path, err)
}
_, _ = io.Copy(io.Discard, resp.Body)
_ = resp.Body.Close()
return resp.StatusCode
}
// Range MISS on cold key -> range_upstream; warm it; range HIT -> range_cache.
if code := get("/depot/rangetest/wt/1", "bytes=0-3"); code != http.StatusPartialContent {
t.Fatalf("range miss: expected 206, got %d", code)
}
if code := get("/depot/rangetest/wt/2", ""); code != http.StatusOK {
t.Fatalf("warm miss: expected 200, got %d", code)
}
if code := get("/depot/rangetest/wt/2", "bytes=8-11"); code != http.StatusPartialContent {
t.Fatalf("range hit: expected 206, got %d", code)
}
req, err := http.NewRequest("GET", srv.URL+"/metrics", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
resp, err := c.Do(req)
if err != nil {
t.Fatalf("GET /metrics: %v", err)
}
out, err := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if err != nil {
t.Fatalf("read /metrics: %v", err)
}
text := string(out)
if !strings.Contains(text, "range_cache 1\n") {
t.Errorf("/metrics missing 'range_cache 1' line:\n%s", text)
}
if !strings.Contains(text, "range_upstream 1\n") {
t.Errorf("/metrics missing 'range_upstream 1' line:\n%s", text)
}
}
// TestStreamCachedResponseRange206 is a focused unit test for streamCachedResponse:
// valid Range yields 206 with the exact slice + metrics; invalid Range yields 416
// with no range metrics.
func TestStreamCachedResponseRange206(t *testing.T) {
body := []byte("0123456789abcdef") // 16 bytes
raw := append([]byte("HTTP/1.1 200 OK\r\nContent-Type: application/octet-stream\r\n\r\n"), body...)
serialized, err := serializeRawResponse(raw)
if err != nil {
t.Fatalf("serialize cache file: %v", err)
}
cf, err := deserializeCacheFile(serialized)
if err != nil {
t.Fatalf("build cache file: %v", err)
}
sc, _ := newTestCacheWithFakeUpstream(t, func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte("x"))
}, "1MB", "0")
sc.ResetMetrics()
// Valid range -> 206 + slice + metrics.
rec := httptest.NewRecorder()
req := httptest.NewRequest("GET", "/depot/rangetest/chunk", nil)
req.Header.Set("Range", "bytes=4-7")
sc.streamCachedResponse(rec, req, cf, "steam/testkey", "127.0.0.1", time.Now())
if rec.Code != http.StatusPartialContent {
t.Fatalf("expected 206, got %d", rec.Code)
}
if rec.Body.String() != "4567" {
t.Errorf("expected body %q, got %q", "4567", rec.Body.String())
}
if got := rec.Header().Get("Content-Range"); got != "bytes 4-7/16" {
t.Errorf("expected Content-Range bytes 4-7/16, got %q", got)
}
if got := rec.Header().Get("X-LanCache-Status"); got != "HIT" {
t.Errorf("expected X-LanCache-Status HIT, got %q", got)
}
stats := sc.GetMetrics()
if stats.RangeCache != 1 {
t.Errorf("expected RangeCache == 1, got %d", stats.RangeCache)
}
if stats.TotalBytesServed != 4 {
t.Errorf("expected TotalBytesServed == 4 (range bytes), got %d", stats.TotalBytesServed)
}
if stats.TotalBytesSaved != 4 {
t.Errorf("expected TotalBytesSaved == 4 (range bytes), got %d", stats.TotalBytesSaved)
}
// Invalid range -> 416, no range metrics, no bytes served.
rec2 := httptest.NewRecorder()
req2 := httptest.NewRequest("GET", "/depot/rangetest/chunk", nil)
req2.Header.Set("Range", "bytes=100-200")
sc.streamCachedResponse(rec2, req2, cf, "steam/testkey", "127.0.0.1", time.Now())
if rec2.Code != http.StatusRequestedRangeNotSatisfiable {
t.Fatalf("expected 416, got %d", rec2.Code)
}
if got := rec2.Header().Get("Content-Range"); got != "bytes */16" {
t.Errorf("expected Content-Range bytes */16, got %q", got)
}
stats = sc.GetMetrics()
if stats.RangeCache != 1 {
t.Errorf("RangeCache must stay 1 after unsatisfiable range, got %d", stats.RangeCache)
}
if stats.TotalBytesServed != 4 {
t.Errorf("TotalBytesServed must stay 4 after 416, got %d", stats.TotalBytesServed)
}
}
+42
View File
@@ -5,6 +5,7 @@ import (
"crypto/sha256"
"encoding/hex"
"fmt"
"net"
"net/http"
"regexp"
"strings"
@@ -163,3 +164,44 @@ func generateServiceCacheKey(urlPath string, servicePrefix string) (string, erro
}
return servicePrefix + "/" + hash, nil
}
// requestHostName strips a port and brackets from an HTTP Host header.
func requestHostName(host string) string {
host = strings.TrimSpace(host)
if host == "" {
return ""
}
if h, _, err := net.SplitHostPort(host); err == nil {
host = h
}
return strings.Trim(host, "[]")
}
func hostIsLiteralIP(host string) bool {
return net.ParseIP(requestHostName(host)) != nil
}
// defaultDirectFetchSuffixes are CDN names Steam actually uses. Applied only when
// no configured upstream is set and the request Host is used as the fetch target.
var defaultDirectFetchSuffixes = []string{
"steamcontent.com",
"steampowered.com",
"steamstatic.com",
}
// hostAllowedForDirectFetch reports whether Host may be used as an origin when
// upstream is empty. Literal IPs are rejected (LAN/metadata SSRF). Names must
// be Steam CDN suffixes so a spoofed User-Agent cannot turn the cache into an
// open reverse proxy.
func hostAllowedForDirectFetch(host string) bool {
name := strings.ToLower(requestHostName(host))
if name == "" || hostIsLiteralIP(host) {
return false
}
for _, suf := range defaultDirectFetchSuffixes {
if name == suf || strings.HasSuffix(name, "."+suf) {
return true
}
}
return false
}
+18 -10
View File
@@ -200,23 +200,33 @@ func New(address string, memorySize string, diskSize string, diskPath, upstream,
if disksize == 0 && memorysize != 0 {
// memory only mode - no disk
sc.metrics.SetDiskTierReady(1) // no disk — N/A / not pending
c.SetSlow(mgc)
} else if disksize != 0 && memorysize == 0 {
// disk only mode: delay attach until disk ready (pure-proxy during scan; Create returns ErrNotFound until slow tier Set)
sc.metrics.SetDiskTierReady(0)
logger.Logger.Info().Msg("Disk slow tier attach pending; Size barrier in progress")
sc.wg.Add(1)
go func() {
defer sc.wg.Done()
t0 := time.Now()
_ = d.Size() // block on barrier per design (all Size callers during window do this; documented)
select {
case <-sc.shutdownCh:
return // Shutdown raced; do not attach or SetSlow after stop
default:
c.SetSlow(dgc)
sc.metrics.SetDiskTierReady(1)
logger.Logger.Info().
Dur("attach_delay", time.Since(t0)).
Msg("Disk slow tier attached (disk-only mode); prior traffic had no disk tier")
}
}()
} else if disksize != 0 && memorysize != 0 {
// memory and disk mode: fast mem immediate, disk delayed (mem-only during scan)
c.SetFast(mgc)
sc.metrics.SetDiskTierReady(0)
logger.Logger.Info().Msg("Disk slow tier attach pending; Size barrier in progress")
sc.wg.Add(1)
go func() {
defer sc.wg.Done()
@@ -227,6 +237,7 @@ func New(address string, memorySize string, diskSize string, diskPath, upstream,
return
default:
c.SetSlow(dgc)
sc.metrics.SetDiskTierReady(1)
logger.Logger.Info().
Dur("attach_delay", time.Since(t0)).
Msg("Disk slow tier attached (mixed mode); prior traffic was memory-only")
@@ -326,15 +337,13 @@ func (sc *SteamCache) Shutdown() {
// GetMetrics returns current metrics
func (sc *SteamCache) GetMetrics() *metrics.Stats {
// Update cache sizes
if sc.memory != nil {
sc.metrics.SetMemoryCacheSize(sc.memory.Size())
}
if sc.disk != nil {
// Note: blocks on initDone (post-eviction state) for accurate post-attach size during long disk init window.
// Skip disk.Size() while attach pending — Size() blocks on initDone and would hang /metrics.
if sc.disk != nil && sc.metrics.GetDiskTierReady() == 1 {
sc.metrics.SetDiskCacheSize(sc.disk.Size())
}
return sc.metrics.GetStats()
}
@@ -357,7 +366,7 @@ func newHTTPTransport() *http.Transport {
DialContext: (&net.Dialer{
Timeout: 10 * time.Second, // Faster connection timeout
KeepAlive: 60 * time.Second, // Longer keep-alive
DualStack: true, // Enable dual-stack (IPv4/IPv6)
// Dual-stack Happy Eyeballs is the default since Go 1.12 (DualStack is deprecated).
}).DialContext,
// Timeout optimizations
@@ -387,11 +396,10 @@ func newHTTPClient(transport *http.Transport) *http.Client {
Timeout: 60 * time.Second, // Optimized timeout for better responsiveness
// Add redirect policy for better performance
CheckRedirect: func(req *http.Request, via []*http.Request) error {
// Limit redirects to prevent infinite loops
if len(via) >= 10 {
return http.ErrUseLastResponse
}
return nil
// Do not follow redirects. Steam CDN chunk/manifest fetches are
// expected to be 200; following Location would let an origin send
// the cache at an arbitrary internal URL.
return http.ErrUseLastResponse
},
}
}
+149
View File
@@ -1057,6 +1057,12 @@ func TestDiskOnlyDelayedAttach(t *testing.T) {
t.Errorf("during init window, expected ErrNotFound from disk-only tiered Create (no slow), got %v", err)
}
// Disk tier is pending while the attach goroutine is in the Size barrier.
// GetMetrics must return quickly (it skips disk.Size() while pending) and report 0.
if got := sc.GetMetrics().DiskTierReady; got != 0 {
t.Errorf("during pending attach, DiskTierReady=%d, want 0", got)
}
// Wait the barrier (exercises the attach go's Size wait)
_ = sc.disk.Size()
@@ -1083,6 +1089,91 @@ func TestDiskOnlyDelayedAttach(t *testing.T) {
} else {
rc.Close()
}
// After attach, the disk tier must be marked ready (1)
if got := sc.GetMetrics().DiskTierReady; got != 1 {
t.Errorf("post-attach DiskTierReady=%d, want 1 (ready)", got)
}
// /metrics text output includes the disk_tier_ready line
rec := httptest.NewRecorder()
metrics.WriteText(rec, sc.GetMetrics())
if !bytes.Contains(rec.Body.Bytes(), []byte("disk_tier_ready 1")) {
t.Errorf("WriteText output missing \"disk_tier_ready 1\": %q", rec.Body.String())
}
}
// TestDiskTierSignalMemoryOnly covers memory-only mode: DiskTierReady=1 (N/A, not
// waiting on disk attach) and heartbeat header X-SteamCache-Disk-Tier: disabled.
func TestDiskTierSignalMemoryOnly(t *testing.T) {
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil)
if err != nil {
t.Fatalf("New memory-only: %v", err)
}
t.Cleanup(func() { sc.Shutdown() })
if got := sc.GetMetrics().DiskTierReady; got != 1 {
t.Errorf("DiskTierReady=%d, want 1 (memory-only = N/A/not pending)", got)
}
req := httptest.NewRequest("GET", "/lancache-heartbeat", nil)
rec := httptest.NewRecorder()
sc.ServeHTTP(rec, req)
if rec.Code != http.StatusNoContent {
t.Errorf("heartbeat status=%d, want 204", rec.Code)
}
if got := rec.Header().Get("X-SteamCache-Disk-Tier"); got != "disabled" {
t.Errorf("X-SteamCache-Disk-Tier=%q, want disabled", got)
}
if got := rec.Header().Get("X-LanCache-Processed-By"); got != "SteamCache2" {
t.Errorf("X-LanCache-Processed-By=%q, want SteamCache2", got)
}
}
// 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.
func TestDiskTierSignalMixedPendingReady(t *testing.T) {
td := t.TempDir()
diskPath := filepath.Join(td, "disk")
if err := os.MkdirAll(diskPath, 0755); err != nil {
t.Fatal(err)
}
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() })
// Immediately in the pending window
if got := sc.GetMetrics().DiskTierReady; got != 0 {
t.Errorf("immediate DiskTierReady=%d, want 0 (pending)", got)
}
req := httptest.NewRequest("GET", "/lancache-heartbeat", nil)
rec := httptest.NewRecorder()
sc.ServeHTTP(rec, req)
if got := rec.Header().Get("X-SteamCache-Disk-Tier"); got != "pending" {
t.Errorf("heartbeat header=%q, want pending", got)
}
// Wait the barrier, then retry until the attach goroutine flips the flag
_ = sc.disk.Size()
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if sc.GetMetrics().DiskTierReady == 1 {
break
}
time.Sleep(1 * time.Millisecond)
}
if got := sc.GetMetrics().DiskTierReady; got != 1 {
t.Fatalf("DiskTierReady=%d after barrier, want 1 (ready)", got)
}
rec2 := httptest.NewRecorder()
sc.ServeHTTP(rec2, httptest.NewRequest("GET", "/lancache-heartbeat", nil))
if got := rec2.Header().Get("X-SteamCache-Disk-Tier"); got != "ready" {
t.Errorf("heartbeat header=%q, want ready", got)
}
}
// --- Phase 2: narrow black-box tests for the new wrapper types ---
@@ -1165,3 +1256,61 @@ func TestClientRateLimiter_BlackBox(t *testing.T) {
t.Error("different clients must have distinct limiters")
}
}
func TestHostAllowedForDirectFetch(t *testing.T) {
allowed := []string{
"lancache.steamcontent.com",
"cache1-iad1.steamcontent.com:443",
"steamcontent.com",
"content.steampowered.com",
"cdn.steamstatic.com",
}
denied := []string{
"",
"127.0.0.1",
"127.0.0.1:80",
"[::1]:80",
"192.168.1.1",
"169.254.169.254",
"evil.example",
"example.com",
"notsteamcontent.com",
}
for _, h := range allowed {
if !hostAllowedForDirectFetch(h) {
t.Errorf("expected allowed: %q", h)
}
}
for _, h := range denied {
if hostAllowedForDirectFetch(h) {
t.Errorf("expected denied: %q", h)
}
}
}
func TestDirectFetchRejectsNonSteamHost(t *testing.T) {
td := t.TempDir()
sc, err := New("127.0.0.1:0", "1MB", "0", td, "", "lru", "lru", 200, 5, "0", nil)
if err != nil {
t.Fatalf("New: %v", err)
}
t.Cleanup(func() { sc.Shutdown() })
req := httptest.NewRequest("GET", "/depot/ssrf/chunk", nil)
req.Host = "127.0.0.1"
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
rec := httptest.NewRecorder()
sc.ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest {
t.Errorf("IP Host: expected 400, got %d", rec.Code)
}
req2 := httptest.NewRequest("GET", "/depot/ssrf/chunk2", nil)
req2.Host = "evil.example"
req2.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
rec2 := httptest.NewRecorder()
sc.ServeHTTP(rec2, req2)
if rec2.Code != http.StatusBadRequest {
t.Errorf("non-CDN Host: expected 400, got %d", rec2.Code)
}
}
+4 -2
View File
@@ -186,7 +186,7 @@ func (tc *TieredCache) Capacity() int64 {
func (tc *TieredCache) promoteToFast(key string, reader io.ReadCloser) {
defer func() { _ = reader.Close() }() // best-effort close; error secondary to promotion attempt (async best-effort path)
// Get file info from slow tier to determine size
// Size for the space/ReadAll guards comes from a Stat snapshot, not the live in-map FileInfo.
var size int64
if slow := tc.slow.Load(); slow != nil {
if vfs, ok := slow.(vfs.VFS); ok {
@@ -210,7 +210,7 @@ func (tc *TieredCache) promoteToFast(key string, reader io.ReadCloser) {
}
// Guard promotion ReadAll using already-fetched size (in addition to space check above)
if size > 0 && size > (1<<30) { // conservative 1GB hard limit on promotion reads (aligns with typical max_object_size)
if size > (1 << 30) { // conservative 1GB hard limit on promotion reads (aligns with typical max_object_size)
return
}
// Read the entire file content
@@ -218,6 +218,8 @@ func (tc *TieredCache) promoteToFast(key string, reader io.ReadCloser) {
if err != nil {
return // Skip promotion if read fails
}
// Create with the bytes we actually hold so we never reuse a live FileInfo.Size.
size = int64(len(content))
// Create the file in fast tier
if fast := tc.fast.Load(); fast != nil {
+24 -9
View File
@@ -248,15 +248,25 @@ func (d *DiskFS) calculateSizeAndPopulateIndex() {
// insertBatch populates info/LRU under lock for a bounded batch (follows maxEvictBatch pattern for short critical sections).
// Size is incremented here only for files actually added (prevents double-count vs. concurrent Create during window).
// Fail-closed: re-stat each path under d.mu and skip if the file is gone. Create does not wait on
// initDone, so a file the scanner observed can be Evict/Delete'd (info + os.Remove) before this
// insert runs. Inserting without a live-file check would resurrect the key in d.info and make
// Stat succeed while os.Stat fails.
func (d *DiskFS) insertBatch(batch []discoveredFile) {
d.mu.Lock()
for _, df := range batch {
if _, exists := d.info[df.key]; !exists {
fi := vfs.NewFileInfoFromOS(df.osInfo, df.key)
d.info[df.key] = fi
d.LRU.Add(df.key, fi)
d.size += df.size
if _, exists := d.info[df.key]; exists {
continue
}
path := d.pathForKey(df.key)
st, err := os.Stat(path)
if err != nil {
continue
}
fi := vfs.NewFileInfoFromOS(st, df.key)
d.info[df.key] = fi
d.LRU.Add(df.key, fi)
d.size += st.Size()
}
d.mu.Unlock()
}
@@ -602,7 +612,9 @@ func (d *DiskFS) Delete(key string) error {
return nil
}
// Stat returns file information with lazy discovery
// Stat returns a snapshot of file information with lazy discovery.
// The returned *FileInfo is not the live cache entry; Close may update Size
// on the in-map object under d.mu.
func (d *DiskFS) Stat(key string) (*vfs.FileInfo, error) {
if key == "" {
return nil, vfserror.ErrInvalidKey
@@ -617,9 +629,10 @@ func (d *DiskFS) Stat(key string) (*vfs.FileInfo, error) {
keyMu.RLock()
d.mu.RLock()
if fi, ok := d.info[key]; ok {
snap := fi.Clone()
d.mu.RUnlock()
keyMu.RUnlock()
return fi, nil
return snap, nil
}
d.mu.RUnlock()
keyMu.RUnlock()
@@ -639,8 +652,9 @@ func (d *DiskFS) Stat(key string) (*vfs.FileInfo, error) {
// Double-check after acquiring write lock
d.mu.Lock()
if fi, ok := d.info[key]; ok {
snap := fi.Clone()
d.mu.Unlock()
return fi, nil
return snap, nil
}
// Re-verify the file still exists on disk under the lock before inserting.
@@ -659,9 +673,10 @@ func (d *DiskFS) Stat(key string) (*vfs.FileInfo, error) {
fi.UpdateAccessBatched(d.timeUpdater)
// Note: size not updated on lazy discovery (preserves prior behavior; initial on-disk accounted via bg populate at New time,
// subsequent files come via Create which accounts size).
snap := fi.Clone()
d.mu.Unlock()
return fi, nil
return snap, nil
}
// EvictLRU evicts the least recently used files to free up space
+74 -41
View File
@@ -371,7 +371,8 @@ func testKey(i int) string {
// artifacts for victims are immediately gone (no resurrection via lazy discovery in Stat/Open),
// and that recreating the same key produces independent content that is not subject to any
// stale eviction unlinks. This exercises the coordinated WLock remove path for DiskFS.
// Uses tolerant checks suitable for raw DiskFS lazy discovery + bg size.
// Create does not wait on initDone, so this also covers insertBatch racing with eviction:
// gone files must not be re-indexed (Stat present / disk missing).
func TestDiskFS_EvictDiskVisibilityAndRecreateSafety(t *testing.T) {
t.Parallel()
td := t.TempDir()
@@ -400,47 +401,21 @@ func TestDiskFS_EvictDiskVisibilityAndRecreateSafety(t *testing.T) {
_ = d.EvictBySize(1024*1024, true)
}
// Consistency check: never have a key absent from Stat but with a file on disk (would indicate
// either resurrection risk or orphan). If Stat succeeds, file should exist.
// A few retries tolerate the documented lazy discovery + eviction coordination windows under
// artificial "force massive eviction then immediate audit" load (especially visible under -race).
for attempt := 0; attempt < 3; attempt++ {
bad := false
for _, k := range created {
p := d.pathForKey(k)
_, statErr := d.Stat(k)
_, diskErr := os.Stat(p)
if statErr != nil {
if !os.IsNotExist(diskErr) {
bad = true
}
} else {
if diskErr != nil {
bad = true
}
}
}
if !bad {
break
}
if attempt < 2 {
time.Sleep(10 * time.Millisecond)
} else {
// On final attempt, report the last observed state for the keys
for _, k := range created {
p := d.pathForKey(k)
_, statErr := d.Stat(k)
_, diskErr := os.Stat(p)
if statErr != nil {
if !os.IsNotExist(diskErr) {
t.Errorf("key %s absent via Stat but file lingers on disk at %s (resurrection risk)", k, p)
}
} else {
if diskErr != nil {
t.Errorf("key %s present via Stat but missing on disk: %v", k, diskErr)
}
}
// Drain bg population so insertBatch cannot still be in flight when we audit.
_ = d.Size()
// Consistency: Stat success iff the file exists on disk. insertBatch must not resurrect
// keys whose backing files were already evicted.
for _, k := range created {
p := d.pathForKey(k)
_, statErr := d.Stat(k)
_, diskErr := os.Stat(p)
if statErr != nil {
if !os.IsNotExist(diskErr) {
t.Errorf("key %s absent via Stat but file lingers on disk at %s (resurrection risk)", k, p)
}
} else if diskErr != nil {
t.Errorf("key %s present via Stat but missing on disk: %v", k, diskErr)
}
}
@@ -464,6 +439,64 @@ func TestDiskFS_EvictDiskVisibilityAndRecreateSafety(t *testing.T) {
}
}
// TestDiskFS_InsertBatchSkipsGoneFiles is the fail-closed contract for bg/lazy index
// insert: a discoveredFile whose path was removed (evicted) must not be re-inserted
// into d.info. That resurrection is what made Stat succeed while os.Stat failed.
func TestDiskFS_InsertBatchSkipsGoneFiles(t *testing.T) {
t.Parallel()
td := t.TempDir()
d, err := New(td, 10*1024*1024, nil)
if err != nil {
t.Fatal(err)
}
_ = d.Size() // finish constructor scan so it cannot also index these keys
liveKey := "live"
goneKey := "gone"
writeKey := func(key, body string) os.FileInfo {
t.Helper()
p := d.pathForKey(key)
if err := os.MkdirAll(filepath.Dir(p), 0700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(p, []byte(body), 0600); err != nil {
t.Fatal(err)
}
st, err := os.Stat(p)
if err != nil {
t.Fatal(err)
}
return st
}
liveInfo := writeKey(liveKey, "still-here")
goneInfo := writeKey(goneKey, "about-to-vanish")
if err := os.Remove(d.pathForKey(goneKey)); err != nil {
t.Fatal(err)
}
d.insertBatch([]discoveredFile{
{key: liveKey, size: liveInfo.Size(), osInfo: liveInfo},
{key: goneKey, size: goneInfo.Size(), osInfo: goneInfo},
})
d.mu.RLock()
_, liveExists := d.info[liveKey]
_, goneExists := d.info[goneKey]
d.mu.RUnlock()
if !liveExists {
t.Errorf("insertBatch skipped live key %s", liveKey)
}
if goneExists {
t.Errorf("insertBatch resurrected gone key %s", goneKey)
}
if _, err := d.Stat(goneKey); err == nil {
t.Errorf("Stat succeeded for gone key %s", goneKey)
}
if _, err := os.Stat(d.pathForKey(liveKey)); err != nil {
t.Errorf("live key %s missing on disk: %v", liveKey, err)
}
}
// TestDiskFS_EvictBoundedLargeN exercises the maxEvictBatch early-break logic (Idea #2)
// under a map size >> batch limit. Forces repeated eviction rounds via GC-style pressure
// and asserts progress + consistency (no resurrection/orphans). Covers bounded collection
+3 -2
View File
@@ -289,7 +289,8 @@ func (m *MemoryFS) Delete(key string) error {
return nil
}
// Stat returns file information
// Stat returns a snapshot of file information. The returned *FileInfo is not
// the live cache entry; Close may update Size on the in-map object under m.mu.
func (m *MemoryFS) Stat(key string) (*types.FileInfo, error) {
if key == "" {
return nil, vfserror.ErrInvalidKey
@@ -310,7 +311,7 @@ func (m *MemoryFS) Stat(key string) (*types.FileInfo, error) {
defer m.mu.RUnlock()
if fi, ok := m.info[key]; ok {
return fi, nil
return fi.Clone(), nil
}
return nil, vfserror.ErrNotFound
+39
View File
@@ -346,6 +346,45 @@ func TestMemoryFS_ConcurrentCloseAndEvict_RaceFree(t *testing.T) {
_ = m.LRU.Len()
}
func TestMemoryFS_StatReturnsSnapshot(t *testing.T) {
t.Parallel()
m, err := New(1024)
if err != nil {
t.Fatal(err)
}
w, err := m.Create("k", 10)
if err != nil {
t.Fatal(err)
}
if _, err := w.Write([]byte("hello")); err != nil {
t.Fatal(err)
}
if err := w.Close(); err != nil {
t.Fatal(err)
}
fi, err := m.Stat("k")
if err != nil {
t.Fatal(err)
}
if fi.Size != 5 {
t.Fatalf("size %d want 5", fi.Size)
}
fi.Size = 999
fi.AccessCount = 0
fi2, err := m.Stat("k")
if err != nil {
t.Fatal(err)
}
if fi2.Size != 5 {
t.Errorf("Stat returned live FileInfo; store size became %d", fi2.Size)
}
if fi2.AccessCount == 0 {
t.Error("Stat returned live FileInfo; AccessCount mutation leaked")
}
}
func TestMemoryFS_EvictVariantsAndErrors(t *testing.T) {
t.Parallel()
m, err := New(800)
+11
View File
@@ -27,6 +27,17 @@ func NewFileInfo(key string, size int64) *FileInfo {
}
}
// Clone returns a snapshot copy of fi. Stat returns Clone() so callers can
// read Size and other fields without racing Close/Open mutations of the
// in-map FileInfo.
func (fi *FileInfo) Clone() *FileInfo {
if fi == nil {
return nil
}
cp := *fi
return &cp
}
// NewFileInfoFromOS creates a FileInfo from os.FileInfo
func NewFileInfoFromOS(info os.FileInfo, key string) *FileInfo {
return &FileInfo{
+28
View File
@@ -16,6 +16,34 @@ func TestNewFileInfo(t *testing.T) {
}
}
func TestFileInfoClone(t *testing.T) {
t.Parallel()
fi := NewFileInfo("k", 42)
fi.AccessCount = 7
cp := fi.Clone()
if cp == fi {
t.Fatal("Clone returned the same pointer")
}
if cp.Key != fi.Key || cp.Size != fi.Size || cp.AccessCount != fi.AccessCount {
t.Errorf("Clone mismatch: %+v vs %+v", cp, fi)
}
if !cp.ATime.Equal(fi.ATime) || !cp.CTime.Equal(fi.CTime) {
t.Error("Clone timestamps mismatch")
}
cp.Size = 99
cp.AccessCount = 1
if fi.Size != 42 || fi.AccessCount != 7 {
t.Error("mutating Clone affected original")
}
if NewFileInfo("x", 1).Clone() == nil {
t.Error("Clone of non-nil was nil")
}
var none *FileInfo
if none.Clone() != nil {
t.Error("Clone of nil was non-nil")
}
}
func TestUpdateAccess(t *testing.T) {
t.Parallel()
fi := NewFileInfo("k", 1)