diff --git a/docs/command-surface.md b/docs/command-surface.md index ae19ce2..9e17d20 100644 --- a/docs/command-surface.md +++ b/docs/command-surface.md @@ -36,6 +36,7 @@ aliases. | `depot` | `pull` | Incremental verified pull of every referenced blob. | | `nwsync` | `emit` | Explode one artifact into NWSync blobs plus its own NSYM manifest. | | `nwsync` | `assemble` | Merge per-artifact NSYM manifests into one merged manifest. | +| `nwsync` | `verify` | Decompress and hash a published manifest's blobs through the pull zone. | `depot status` and `depot get` pick their backend either with `--out DIR`, a depot tree on disk, or with `--target bunny|cdn`, a remote backend. The two @@ -48,8 +49,9 @@ beside the artifact itself with the extension replaced, so `emit` and `assemble` agree on where it is without being told. ``` -nwsync emit [--as NAME] [--out DIR] [--jobs N] +nwsync emit [--as NAME] [--out DIR] [--jobs N] [--verify] nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] ... +nwsync verify [--sample N] [--base URL] [--jobs N] ``` `emit` is latency-bound, not CPU-bound: every blob costs an existence probe @@ -73,6 +75,28 @@ resolves it. `--tlk-key` has its own slot because the TLK shadows nothing. `--group-id` is per channel — 1 is current, 2 is testing, and 0 leaves the field out of the sidecar. +`nwsync verify` is the only check on a published blob upstream of a player's +client. It reads the **pull zone**, not the storage API, and needs no +credential: what matters is the bytes a client is served, edge behaviour +included. Every blob is decompressed and hashed, and the zstd frame is asserted +to declare its content size. Neither half is optional — a `Content-Length` check +passes a byte-correct-looking object whose contents are short, and a round-trip +check alone passes a frame the game client cannot decode but Go's decoder can. +Failures are reported per blob as missing, malformed framing, size mismatch or +hash mismatch, and the exit code is 1. + +A full sweep of the live manifest is roughly 69,000 blobs and 15 GB, so +`--sample N` exists to make verifying routine; the default is a full sweep. +`--base URL` (or `NWSYNC_PULL_BASE`) overrides the public host. + +`emit --verify` applies the same check where `emit` would otherwise skip. `emit` +normally reads a blob's presence as proof of its contents, decided by a 1-byte +range GET, so an object written truncated — or written by an emitter since found +broken — is skipped by every later run forever and no backfill repairs it. With +`--verify` the stored copy is read back, unwrapped, hashed against its own name, +and replaced when it does not match. It costs a full GET per existing blob, so +it is a repair pass, not the default. + `emit` uploads blobs first and the index last, so the presence of an index is the publication marker: an artifact whose emit died halfway leaves real blobs in the zone and no index. Blob names are content hashes, so re-running skips @@ -108,6 +132,17 @@ saving on hak content: a 250 MB hak emitted 2296 blobs totalling 59 MB on disk against 249 MB of resources, as recorded in the sidecar's `on_disk_bytes` and `total_bytes`. +The zstd frame always declares its `Frame_Content_Size`. The game client sizes +its output buffer from that field and cannot decode a frame without one, but the +Go encoder omits it below 256 bytes, so `emit` re-headers those frames into the +shape reference libzstd emits: `Single_Segment_flag` set, `Window_Descriptor` +dropped, and a one-byte content size in its place. `zstd -l ` must print a +decompressed size; a blank column there is the fault, and it is invisible to any +check that only decompresses, because both `zstd -dc` and Go's decoder stream +such a frame happily. This is what the sidecar's `emitter_version` counts: +version 1 omitted the field and no client could sync past such a blob, version 2 +declares it. `assemble` refuses to merge indexes that disagree. + ## Hidden compatibility aliases Existing scripts may continue using these names indefinitely. They are accepted diff --git a/internal/dispatch/dispatch.go b/internal/dispatch/dispatch.go index e370caf..cf7b7e5 100644 --- a/internal/dispatch/dispatch.go +++ b/internal/dispatch/dispatch.go @@ -108,10 +108,11 @@ var Registry = []Builder{ { Name: "nwsync", Bin: "crucible-nwsync", - Summary: "publish NWSync blobs and manifests (emit/assemble)", + Summary: "publish NWSync blobs and manifests (emit/assemble/verify)", Commands: []Command{ {Name: "emit", Summary: "explode one artifact into blobs plus its own NSYM manifest", Usage: "crucible nwsync emit --out DIR"}, {Name: "assemble", Summary: "merge per-artifact NSYM manifests into one", Usage: "crucible nwsync assemble --order NAMES --entries DIR --out DIR [--group-id N]"}, + {Name: "verify", Summary: "read a published manifest's blobs back through the pull zone and hash them", Usage: "crucible nwsync verify [--sample N]"}, }, Wired: true, }, diff --git a/internal/dispatch/dispatch_test.go b/internal/dispatch/dispatch_test.go index 1a459eb..3182fcc 100644 --- a/internal/dispatch/dispatch_test.go +++ b/internal/dispatch/dispatch_test.go @@ -160,7 +160,7 @@ func TestCanonicalCommandSurface(t *testing.T) { "module": {"build", "extract", "validate", "compare", "manifest"}, "topdata": {"validate", "build", "package", "compare", "convert"}, "wiki": {"build", "deploy"}, - "nwsync": {"emit", "assemble"}, + "nwsync": {"emit", "assemble", "verify"}, } for _, builder := range Registry { got := builder.subcommands() diff --git a/internal/nwsync/compressedbuf.go b/internal/nwsync/compressedbuf.go index af708d3..e66267c 100644 --- a/internal/nwsync/compressedbuf.go +++ b/internal/nwsync/compressedbuf.go @@ -3,6 +3,7 @@ package nwsync import ( "bytes" "encoding/binary" + "encoding/hex" "fmt" "github.com/klauspost/compress/zstd" @@ -31,6 +32,18 @@ var ( blobDecoder, _ = zstd.NewReader(nil, zstd.WithDecoderConcurrency(1)) ) +// zstd frame header bits we care about. A frame starts with the magic, then a +// one-byte Frame_Header_Descriptor: bits 7-6 size the Frame_Content_Size field, +// bit 5 is Single_Segment_flag, bits 1-0 size the Dictionary_ID field. +const ( + zstdFrameMagic = "\x28\xb5\x2f\xfd" + frameSingleSegment = 1 << 5 + frameDictionaryMask = 0x03 + // The size below which klauspost/compress writes no Frame_Content_Size, + // matching the field's own encoding threshold. + frameSizeThreshold = 256 +) + // compressBlob wraps data in NWCompressedBuffer framing. func compressBlob(data []byte) []byte { var out bytes.Buffer @@ -38,10 +51,87 @@ func compressBlob(data []byte) []byte { for _, field := range header { _ = binary.Write(&out, binary.LittleEndian, field) } - out.Write(blobEncoder.EncodeAll(data, nil)) + out.Write(declareFrameContentSize(blobEncoder.EncodeAll(data, nil), len(data))) return out.Bytes() } +// inspectBlob unwraps a stored blob the way the game client reads it, and is the +// only reader that should be trusted to judge a published blob. +// +// It asserts the frame property on top of the round trip. Go's decoder — like +// the zstd CLI — streams a frame that declares no content size, so a check that +// only decompresses and hashes is a *more* capable decoder than the client's: it +// certifies exactly the blobs the client rejects, which is how #86 reached +// production and survived an audit. +func inspectBlob(blob []byte) ([]byte, error) { + data, err := decompressBlob(blob) + if err != nil { + return nil, fmt.Errorf("malformed framing: %w", err) + } + if len(data) > 0 && !frameDeclaresContentSize(blob[blobHeaderBytes:]) { + return nil, fmt.Errorf("malformed framing: the zstd frame declares no content size, which the game client cannot decode") + } + return data, nil +} + +// blobMatchesName holds a stored blob to its own file name: a blob is named +// after the sha1 of its uncompressed bytes, so the name is a complete statement +// about the contents and nothing else is needed to check it. +func blobMatchesName(blob []byte, sha1Hex string) error { + data, err := inspectBlob(blob) + if err != nil { + return err + } + if got := hex.EncodeToString(sha1Sum(data)); got != sha1Hex { + return fmt.Errorf("blob %s holds the contents of %s", sha1Hex, got) + } + return nil +} + +// frameDeclaresContentSize reports whether a zstd frame states how many bytes it +// decompresses to. A frame with a zero-sized Frame_Content_Size field declares +// one only when Single_Segment_flag is set; otherwise the size is unknown. +func frameDeclaresContentSize(frame []byte) bool { + if len(frame) < 5 || string(frame[:4]) != zstdFrameMagic { + return false + } + descriptor := frame[4] + return descriptor>>6 != 0 || descriptor&frameSingleSegment != 0 +} + +// declareFrameContentSize rewrites a frame that does not declare its +// Frame_Content_Size so that it does, and returns any other frame unchanged. +// +// klauspost/compress omits the field for inputs under 256 bytes, which the spec +// permits. Reference libzstd never does, so the NWN client — which sizes its +// output buffer from ZSTD_getFrameContentSize and has therefore never met a +// frame without one — rejects the blob outright with an empty "potential +// compression error" (#86). No encoder option changes this, so the frame is +// re-headered here. +// +// The result is the shape libzstd itself emits for the same input: setting +// Single_Segment_flag drops the Window_Descriptor byte, and the freed byte pays +// for a one-byte Frame_Content_Size. Window_Size then equals the content size, +// which is sound because the content is under 256 bytes and every match in it +// therefore falls inside that window. Same length in, same length out. +func declareFrameContentSize(frame []byte, size int) []byte { + if size <= 0 || size >= frameSizeThreshold || len(frame) < 6 || string(frame[:4]) != zstdFrameMagic { + return frame + } + descriptor := frame[4] + // Rewrite only the exact shape a small input produces: no declared size, no + // single segment, no dictionary. Anything else either declares a size + // already or is not a frame this reinterpretation is safe on. + if descriptor>>6 != 0 || descriptor&frameSingleSegment != 0 || descriptor&frameDictionaryMask != 0 { + return frame + } + reframed := make([]byte, len(frame)) + copy(reframed, frame) + reframed[4] = descriptor | frameSingleSegment + reframed[5] = byte(size) // replaces Window_Descriptor + return reframed +} + // decompressBlob unwraps NWCompressedBuffer framing. It exists so a blob this // package wrote — or one upstream wrote — can be compared by its uncompressed // bytes, which is the only comparison that is meaningful across zstd diff --git a/internal/nwsync/emit.go b/internal/nwsync/emit.go index b1c6e10..826e262 100644 --- a/internal/nwsync/emit.go +++ b/internal/nwsync/emit.go @@ -33,7 +33,9 @@ var skippedTypes = resTypes("nss", "ndb", "gic") // manifest quietly disagrees with. Bump it only when emitted bytes change — it // is deliberately not the build revision, which would invalidate every // published index on every unrelated commit. -const emitterVersion = "1" +// Version 2 declares Frame_Content_Size on every blob (#86); version 1 omitted +// it below 256 bytes and no client could sync past such a blob. +const emitterVersion = "2" // serverTypes are loaded only server-side; a manifest holding nothing else // has no client contents. Mirrors upstream's GlobalResTypeServerList, whose @@ -78,6 +80,7 @@ type EmitOptions struct { As string // name override, for a TLK whose filename is not its published name OutDir string // write locally instead of uploading — the conformance path Jobs int // resources in flight at once; 0 means defaultEmitJobs + Verify bool // hash what would be skipped instead of trusting presence Sink sink // test seam; nil means OutDir or the zone } @@ -129,7 +132,7 @@ func Emit(options EmitOptions) (EmitResult, error) { if jobs < 1 { jobs = defaultEmitJobs } - entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs) + entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs, options.Verify) if err != nil { return EmitResult{}, err } @@ -199,7 +202,7 @@ func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name strin // // The returned entries are in artifact order whatever order the workers finish // in, because a manifest's bytes are promised deterministic by emitterVersion. -func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jobs int) ([]Entry, int, int64, error) { +func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jobs int, verify bool) ([]Entry, int, int64, error) { // A resref appearing twice inside one artifact resolves to the last one, // the way resman lets the last container added win. order := make([]Identity, 0, len(index)) @@ -271,7 +274,7 @@ func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jo if duplicate { return } - written, err := target.putBlob(fmt.Sprintf("%x", sum), func() []byte { return compressBlob(payload) }) + written, err := target.putBlob(fmt.Sprintf("%x", sum), verify, func() []byte { return compressBlob(payload) }) mu.Lock() defer mu.Unlock() if err != nil { diff --git a/internal/nwsync/nwsync_test.go b/internal/nwsync/nwsync_test.go index 0673f25..0934b86 100644 --- a/internal/nwsync/nwsync_test.go +++ b/internal/nwsync/nwsync_test.go @@ -501,3 +501,31 @@ func TestRunUsageErrors(t *testing.T) { } } } + +// TestEveryBlobDeclaresItsFrameContentSize guards the fault that stopped every +// client sync (#86): klauspost/compress omits Frame_Content_Size for inputs +// under 256 bytes, and the game client cannot decode a frame without it. This +// asserts a frame property, not a round trip — the zstd CLI and Go's decoder +// both stream such a frame happily, so round-tripping cannot see the defect. +func TestEveryBlobDeclaresItsFrameContentSize(t *testing.T) { + // 230 and 175 are real sizes from the manifest that failed to sync; 255/256 + // straddle the encoder's threshold. + for _, size := range []int{1, 32, 175, 230, 255, 256, 257, 1024, 5000} { + payload := make([]byte, size) + for i := range payload { + payload[i] = byte('a' + i%26) + } + blob := compressBlob(payload) + if !frameDeclaresContentSize(blob[blobHeaderBytes:]) { + t.Errorf("blob of %d bytes declares no frame content size (descriptor %#x)", + size, blob[blobHeaderBytes+4]) + } + got, err := decompressBlob(blob) + if err != nil { + t.Fatalf("decompress %d-byte blob: %v", size, err) + } + if !bytes.Equal(got, payload) { + t.Errorf("%d-byte blob did not round trip", size) + } + } +} diff --git a/internal/nwsync/run.go b/internal/nwsync/run.go index a25b519..81d324c 100644 --- a/internal/nwsync/run.go +++ b/internal/nwsync/run.go @@ -12,12 +12,17 @@ import ( "flag" "fmt" "io" + "os" ) const ( exitOK = 0 exitUsage = 64 exitInternal = 70 + // exitDrift says the command worked and the zone is wrong, which is a + // different thing for CI to act on than the command failing. It matches + // depot's code for the same meaning. + exitDrift = 1 ) // Run executes an nwsync subcommand. args[0] is the subcommand (emit|assemble); @@ -32,6 +37,8 @@ func Run(args []string, stdout, stderr io.Writer) int { return runEmit(args[1:], stderr) case "assemble": return runAssemble(args[1:], stderr) + case "verify": + return runVerify(args[1:], stdout, stderr, os.Getenv) case "-h", "--help", "help": printRunUsage(stdout) return exitOK @@ -44,17 +51,27 @@ func Run(args []string, stdout, stderr io.Writer) int { func printRunUsage(w io.Writer) { fmt.Fprint(w, `usage: - nwsync emit [--as NAME] [--out DIR] + nwsync emit [--as NAME] [--out DIR] [--verify] nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] ... + nwsync verify [--sample N] [--base URL] emit explodes one .hak/.erf or one loose file (the TLK) into NWSync blobs plus a NSYM index covering only that artifact, and uploads both. assemble merges those indexes into one manifest, reading no bulk data. Artifact keys are depot keys; an index lives beside its artifact, with the extension replaced. +verify reads a published manifest and its blobs back through the public pull +zone, with no credential, and decompresses and hashes every one. It is the only +check on a published blob upstream of a player's client. + +--verify makes emit hash what it would otherwise skip. emit normally treats a +blob's presence as proof of its contents, so without this an object written +truncated, or written by an emitter since found broken, is skipped forever. + --out DIR writes to a local repository tree instead of uploading, which is the conformance path against upstream nwn_nwsync_write. Without it, the zone comes from NWSYNC_STORAGE_ZONE, NWSYNC_STORAGE_PASSWORD and BUNNY_STORAGE_HOST. +verify needs none of those; its base comes from --base or NWSYNC_PULL_BASE. `) } @@ -82,6 +99,7 @@ func runEmit(args []string, stderr io.Writer) int { as := fs.String("as", "", "published name of the artifact, when it differs from the key") out := fs.String("out", "", "write to a local repository tree instead of uploading") jobs := fs.Int("jobs", defaultEmitJobs, "resources to hash, compress and store at once") + verify := fs.Bool("verify", false, "read back and hash blobs that already exist instead of trusting their presence") positional, err := parseArgs(fs, args) if err != nil { return exitUsage @@ -101,6 +119,7 @@ func runEmit(args []string, stderr io.Writer) int { As: *as, OutDir: *out, Jobs: *jobs, + Verify: *verify, }) if err != nil { fmt.Fprintf(stderr, "nwsync emit: %v\n", err) @@ -111,6 +130,40 @@ func runEmit(args []string, stderr io.Writer) int { return exitOK } +func runVerify(args []string, stdout, stderr io.Writer, getenv func(string) string) int { + fs := flag.NewFlagSet("verify", flag.ContinueOnError) + fs.SetOutput(stderr) + base := fs.String("base", getenv("NWSYNC_PULL_BASE"), "pull zone base URL to read through") + sample := fs.Int("sample", 0, "check this many random blobs instead of all of them") + jobs := fs.Int("jobs", defaultEmitJobs, "blobs to fetch and hash at once") + positional, err := parseArgs(fs, args) + if err != nil { + return exitUsage + } + if len(positional) != 1 { + fmt.Fprintf(stderr, "nwsync verify: exactly one is required\n") + return exitUsage + } + + result, err := Verify(VerifyOptions{ + ManifestSHA1: positional[0], + Base: *base, + Sample: *sample, + Jobs: *jobs, + Log: stderr, + }) + if err != nil { + fmt.Fprintf(stderr, "nwsync verify: %v\n", err) + return exitInternal + } + fmt.Fprintf(stdout, "verified %d of %d blobs behind %d resources: %d failures, %d bytes checked\n", + result.Checked, result.Blobs, result.Entries, result.Failures, result.Bytes) + if result.Failures > 0 { + return exitDrift + } + return exitOK +} + func runAssemble(args []string, stderr io.Writer) int { fs := flag.NewFlagSet("assemble", flag.ContinueOnError) fs.SetOutput(stderr) diff --git a/internal/nwsync/sink.go b/internal/nwsync/sink.go index f552737..d8f5930 100644 --- a/internal/nwsync/sink.go +++ b/internal/nwsync/sink.go @@ -20,11 +20,17 @@ import ( // upstream's output and ours can be diffed on a developer machine. type sink interface { // putBlob stores one NWCompressedBuffer blob under its sha1 name and - // returns the bytes stored, or 0 if the blob was already there. Blob names - // are content hashes, so an existing name is existing content — which is - // why body is a thunk: compression is the expensive part of emit and a - // blob that is already stored must not pay for it. - putBlob(sha1Hex string, body func() []byte) (int64, error) + // returns the bytes stored, or 0 if a good copy was already there. Blob + // names are content hashes, so an existing name is normally taken as + // existing content — which is why body is a thunk: compression is the + // expensive part of emit and a blob that is already stored must not pay + // for it. + // + // verify stops trusting presence: the stored copy is read back, unwrapped + // and hashed, and replaced when it is not what its name claims. Without it + // an object written truncated, or written by an emitter since found broken, + // is skipped by every later emit forever and no backfill can repair it. + putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error) // putIndex stores a NSYM manifest and its sidecar under key, which is // either an artifact-derived object key or a local path. putIndex(key string, manifest, sidecar []byte) error @@ -37,10 +43,12 @@ type sink interface { // dirSink writes a local NWSync repository tree. type dirSink struct{ root string } -func (s dirSink) putBlob(sha1Hex string, body func() []byte) (int64, error) { +func (s dirSink) putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error) { blob := blobPath(s.root, sha1Hex) - if _, err := os.Stat(blob); err == nil { - return 0, nil + if stored, err := os.ReadFile(blob); err == nil { + if !verify || blobMatchesName(stored, sha1Hex) == nil { + return 0, nil + } } if err := os.MkdirAll(filepath.Dir(blob), 0o755); err != nil { return 0, fmt.Errorf("create blob directory: %w", err) @@ -91,7 +99,7 @@ type zoneSink struct { zone string } -func (s zoneSink) putBlob(sha1Hex string, body func() []byte) (int64, error) { +func (s zoneSink) putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error) { key := path.Join("data", "sha1", sha1Hex[0:2], sha1Hex[2:4], sha1Hex) // A throttled probe must never be read as "missing, re-upload" or as // "present, skip", so only a confirmed Present skips the upload. @@ -100,7 +108,19 @@ func (s zoneSink) putBlob(sha1Hex string, body func() []byte) (int64, error) { return 0, fmt.Errorf("probe blob %s: %w", sha1Hex, err) } if state == depot.Present { - return 0, nil + if !verify { + return 0, nil + } + // The probe only proved the object exists. Read it back and hold it to + // its own name; a failure here means re-upload, not abort, because + // repairing what is there is the whole point of verifying. + stored, err := s.store.GetKey(s.ctx, key) + if err != nil { + return 0, fmt.Errorf("read back blob %s: %w", sha1Hex, err) + } + if blobMatchesName(stored, sha1Hex) == nil { + return 0, nil + } } data := body() if err := s.put(key, data); err != nil { diff --git a/internal/nwsync/verify.go b/internal/nwsync/verify.go new file mode 100644 index 0000000..8c8fdb9 --- /dev/null +++ b/internal/nwsync/verify.go @@ -0,0 +1,232 @@ +package nwsync + +import ( + "crypto/sha1" + "encoding/hex" + "errors" + "fmt" + "io" + "math/rand/v2" + "net/http" + "path" + "sort" + "strconv" + "sync" + "time" + + "git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/erf" +) + +// defaultPullBase is the public NWSync host, which is a Bunny pull zone fronting +// the storage zone. Verify reads through it rather than through the storage API +// on purpose: what matters is the bytes a client is served, edge behaviour +// included, not what the origin believes it holds. +const defaultPullBase = "https://nwsync.westgate.pw" + +// errBlobMissing marks an object the zone does not serve at all, as distinct +// from one it serves badly. +var errBlobMissing = errors.New("missing") + +// blobSource reads one object out of the zone by key. Verify never writes and +// never authenticates, so this is deliberately narrower than sink. +type blobSource interface { + get(key string) ([]byte, error) + describe(key string) string +} + +// pullZone reads the zone over plain HTTP, with no credential. +type pullZone struct { + base string + client *http.Client +} + +func newPullZone(base string) blobSource { + if base == "" { + base = defaultPullBase + } + return pullZone{ + base: base, + // A full sweep is tens of thousands of small requests, so connections + // have to be reused; the default transport does that already. + client: &http.Client{Timeout: 60 * time.Second}, + } +} + +func (z pullZone) describe(key string) string { return z.base + "/" + key } + +func (z pullZone) get(key string) ([]byte, error) { + resp, err := z.client.Get(z.describe(key)) + if err != nil { + return nil, err + } + defer resp.Body.Close() + if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusGone { + _, _ = io.Copy(io.Discard, resp.Body) + return nil, errBlobMissing + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + _, _ = io.Copy(io.Discard, resp.Body) + return nil, fmt.Errorf("unexpected status %d", resp.StatusCode) + } + return io.ReadAll(resp.Body) +} + +// VerifyOptions describes one verify run. +type VerifyOptions struct { + ManifestSHA1 string // the merged manifest to verify + Base string // pull zone base URL; empty means defaultPullBase + Sample int // check this many random blobs; 0 means all of them + Jobs int // blobs in flight at once; 0 means defaultEmitJobs + Source blobSource // test seam; nil means the pull zone at Base + Log io.Writer // per-blob failures land here; nil discards them +} + +// VerifyResult reports what one verify run found. +type VerifyResult struct { + Entries int // resources the manifest names + Blobs int // distinct blobs behind those resources + Checked int // blobs actually fetched + Failures int // blobs that failed a check + Bytes int64 // uncompressed bytes verified +} + +// Verify reads a published manifest and its blobs the way a client reads them, +// and reports every blob that is not what the manifest says it is. +// +// Presence is not correctness. emit skips an object that already exists on the +// strength of a 1-byte range GET, so a truncated or wrongly framed object is +// skipped by every later emit forever and the backfill cannot repair it. This is +// the only thing upstream of a player's client that can tell that has happened. +// +// Every blob is decompressed and hashed. A Content-Length check would pass the +// exact failure mode being hunted — a byte-correct-looking object whose contents +// are wrong — and a round-trip check alone would pass a frame that omits its +// content size, because Go's decoder is more capable than the client's (#86). +func Verify(options VerifyOptions) (VerifyResult, error) { + if len(options.ManifestSHA1) != 40 { + return VerifyResult{}, fmt.Errorf("%q is not a manifest sha1", options.ManifestSHA1) + } + source := options.Source + if source == nil { + source = newPullZone(options.Base) + } + log := options.Log + if log == nil { + log = io.Discard + } + + manifestKey := path.Join("manifests", options.ManifestSHA1) + data, err := source.get(manifestKey) + if err != nil { + return VerifyResult{}, fmt.Errorf("%s: %w", source.describe(manifestKey), err) + } + // A manifest is named after its own sha1, so this catches the zone serving + // a different manifest — or a truncated one — before any blob is fetched. + if got := hex.EncodeToString(sha1Sum(data)); got != options.ManifestSHA1 { + return VerifyResult{}, fmt.Errorf("%s hashes to %s, not the manifest asked for", + source.describe(manifestKey), got) + } + entries, err := readManifest(data) + if err != nil { + return VerifyResult{}, fmt.Errorf("%s: %w", source.describe(manifestKey), err) + } + + // A manifest names one blob many times over: mappings share a sha1, and so + // do resrefs with identical contents. Fetch each blob once. + blobs := make([]Entry, 0, len(entries)) + seen := make(map[[20]byte]bool, len(entries)) + for _, entry := range entries { + if seen[entry.SHA1] { + continue + } + seen[entry.SHA1] = true + blobs = append(blobs, entry) + } + + result := VerifyResult{Entries: len(entries), Blobs: len(blobs)} + checking := blobs + if options.Sample > 0 && options.Sample < len(blobs) { + // A full sweep of the live manifest is ~69,000 objects and ~15 GB, so + // sampling is what makes verifying a routine act rather than an event. + picks := rand.Perm(len(blobs))[:options.Sample] + checking = make([]Entry, 0, options.Sample) + for _, i := range picks { + checking = append(checking, blobs[i]) + } + } + result.Checked = len(checking) + + jobs := options.Jobs + if jobs < 1 { + jobs = defaultEmitJobs + } + var ( + mu sync.Mutex + failures []string + ) + work := make(chan Entry) + var wg sync.WaitGroup + for range jobs { + wg.Go(func() { + for entry := range work { + fault := checkEntry(source, entry) + mu.Lock() + if fault != "" { + failures = append(failures, fault) + } else { + result.Bytes += int64(entry.Size) + } + mu.Unlock() + } + }) + } + for _, entry := range checking { + work <- entry + } + close(work) + wg.Wait() + + // Workers finish in any order; a report an operator can diff must not. + sort.Strings(failures) + for _, fault := range failures { + fmt.Fprintln(log, fault) + } + result.Failures = len(failures) + return result, nil +} + +// checkEntry fetches one blob and returns a one-line fault, or "" if it is +// exactly what the manifest entry says it is. +func checkEntry(source blobSource, entry Entry) string { + key := path.Join("data", "sha1", entry.sha1Hex()[0:2], entry.sha1Hex()[2:4], entry.sha1Hex()) + // Name the resource, not just the hash: an operator has to find the thing + // in a hak, and a bare sha1 says nothing about where to look. + extension, ok := erf.ExtensionForResourceType(entry.ResType) + if !ok { + extension = strconv.Itoa(int(entry.ResType)) + } + where := fmt.Sprintf("%s (%s.%s)", entry.sha1Hex(), entry.ResRef, extension) + blob, err := source.get(key) + if err != nil { + if errors.Is(err, errBlobMissing) { + return where + ": missing" + } + return where + ": unreadable: " + err.Error() + } + data, err := inspectBlob(blob) + if err != nil { + return where + ": " + err.Error() + } + if uint32(len(data)) != entry.Size { + return fmt.Sprintf("%s: size mismatch: %d bytes, manifest says %d", where, len(data), entry.Size) + } + if sha1.Sum(data) != entry.SHA1 { + return fmt.Sprintf("%s: hash mismatch: contents hash to %x", where, sha1.Sum(data)) + } + return "" +} + +func sha1Sum(data []byte) []byte { + sum := sha1.Sum(data) + return sum[:] +} diff --git a/internal/nwsync/verify_test.go b/internal/nwsync/verify_test.go new file mode 100644 index 0000000..b44c142 --- /dev/null +++ b/internal/nwsync/verify_test.go @@ -0,0 +1,263 @@ +package nwsync + +import ( + "bytes" + "crypto/sha1" + "encoding/hex" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/klauspost/compress/zstd" +) + +// verifyFixture emits two haks and a TLK into a fake zone, assembles them, and +// hands back a verifier reading that zone the way a client would. +type verifyFixture struct { + *zoneSinkFixture + manifestSHA1 string +} + +func newVerifyFixture(t *testing.T) *verifyFixture { + t.Helper() + zone := newZoneFixture(t) + dir := t.TempDir() + hak := filepath.Join(dir, "sow_test_01.hak") + // A payload under 256 bytes is the one the frame-header check exists for. + writeHak(t, hak, map[string][]byte{ + "bloodstain1.tga": []byte("small"), + "appearance.2da": bytes.Repeat([]byte("2DA V2.0\n"), 200), + }) + tlk := filepath.Join(dir, "sow_tlk.tlk") + if err := os.WriteFile(tlk, []byte("TLK V3.0 payload"), 0o644); err != nil { + t.Fatal(err) + } + zone.emit(t, hak) + if _, err := Emit(EmitOptions{ + ArtifactKey: artifactKey(t, tlk), ArtifactPath: tlk, As: "sow_tlk.tlk", Sink: zone.sink, + }); err != nil { + t.Fatalf("emit tlk: %v", err) + } + assembled, err := Assemble(AssembleOptions{ + ArtifactKeys: []string{artifactKey(t, hak)}, TLKKey: artifactKey(t, tlk), Sink: zone.sink, + }) + if err != nil { + t.Fatalf("assemble: %v", err) + } + return &verifyFixture{zoneSinkFixture: zone, manifestSHA1: assembled.SHA1} +} + +func (f *verifyFixture) verify(t *testing.T, sample int) (VerifyResult, string, error) { + t.Helper() + var log bytes.Buffer + result, err := Verify(VerifyOptions{ + ManifestSHA1: f.manifestSHA1, + Sample: sample, + Source: f.zone.pullZone(), + Log: &log, + }) + return result, log.String(), err +} + +// blobKey is where a resource's blob lives, addressed by the sha1 of its +// uncompressed bytes — the same path the client requests. +func blobKey(body []byte) string { + sum := sha1.Sum(body) + hexed := hex.EncodeToString(sum[:]) + return "data/sha1/" + hexed[0:2] + "/" + hexed[2:4] + "/" + hexed +} + +func TestVerifyPassesACleanZone(t *testing.T) { + fixture := newVerifyFixture(t) + result, log, err := fixture.verify(t, 0) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Failures != 0 { + t.Errorf("verify reported %d failures on a clean zone: %s", result.Failures, log) + } + if result.Checked != 3 { + t.Errorf("checked %d blobs, want 3", result.Checked) + } +} + +func TestVerifyReportsAMissingBlob(t *testing.T) { + fixture := newVerifyFixture(t) + key := blobKey([]byte("small")) + fixture.zone.mu.Lock() + delete(fixture.zone.objects, key) + fixture.zone.mu.Unlock() + + result, log, err := fixture.verify(t, 0) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Failures != 1 { + t.Fatalf("reported %d failures, want 1: %s", result.Failures, log) + } + if !strings.Contains(log, "missing") { + t.Errorf("a deleted blob was not reported as missing: %s", log) + } +} + +func TestVerifyReportsATruncatedBlob(t *testing.T) { + fixture := newVerifyFixture(t) + key := blobKey([]byte("small")) + fixture.zone.mu.Lock() + fixture.zone.objects[key] = fixture.zone.objects[key][:blobHeaderBytes+4] + fixture.zone.mu.Unlock() + + result, log, err := fixture.verify(t, 0) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Failures != 1 { + t.Fatalf("reported %d failures, want 1: %s", result.Failures, log) + } + if !strings.Contains(log, "framing") { + t.Errorf("a truncated blob was not reported as malformed framing: %s", log) + } +} + +// TestVerifyRejectsABlobWithNoDeclaredFrameContentSize is the check that #86 +// slipped past: the blob decompresses to exactly the right bytes, so a verifier +// that only round-trips certifies it, yet the client cannot decode it. +func TestVerifyRejectsABlobWithNoDeclaredFrameContentSize(t *testing.T) { + fixture := newVerifyFixture(t) + body := []byte("small") + key := blobKey(body) + + fixture.zone.mu.Lock() + good := fixture.zone.objects[key] + encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) + if err != nil { + t.Fatal(err) + } + bad := append(append([]byte{}, good[:blobHeaderBytes]...), encoder.EncodeAll(body, nil)...) + fixture.zone.objects[key] = bad + fixture.zone.mu.Unlock() + + if frameDeclaresContentSize(bad[blobHeaderBytes:]) { + t.Fatal("the fixture blob declares a content size; it cannot exercise the check") + } + if got, err := decompressBlob(bad); err != nil || !bytes.Equal(got, body) { + t.Fatalf("the fixture blob must round trip, or it proves nothing: %v", err) + } + + result, log, err := fixture.verify(t, 0) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Failures != 1 { + t.Fatalf("reported %d failures, want 1: %s", result.Failures, log) + } + if !strings.Contains(log, "content size") { + t.Errorf("undeclared frame content size was not the reported reason: %s", log) + } +} + +func TestVerifyReportsWrongContents(t *testing.T) { + fixture := newVerifyFixture(t) + key := blobKey([]byte("small")) + fixture.zone.mu.Lock() + // Valid framing, valid zstd, wrong bytes: only decompressing and hashing + // can see this, which is why Content-Length is not enough. + fixture.zone.objects[key] = compressBlob([]byte("wrong")) + fixture.zone.mu.Unlock() + + result, log, err := fixture.verify(t, 0) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Failures != 1 { + t.Fatalf("reported %d failures, want 1: %s", result.Failures, log) + } + if !strings.Contains(log, "hash mismatch") { + t.Errorf("wrong contents were not reported as a hash mismatch: %s", log) + } +} + +func TestVerifyFailsWhenTheManifestIsNotTheOneAsked(t *testing.T) { + fixture := newVerifyFixture(t) + fixture.zone.mu.Lock() + fixture.zone.objects["manifests/"+fixture.manifestSHA1] = []byte("NSYM garbage") + fixture.zone.mu.Unlock() + + if _, _, err := fixture.verify(t, 0); err == nil { + t.Fatal("verify accepted a manifest that is not the one requested") + } +} + +func TestVerifySampleChecksFewerBlobs(t *testing.T) { + fixture := newVerifyFixture(t) + result, log, err := fixture.verify(t, 1) + if err != nil { + t.Fatalf("verify: %v", err) + } + if result.Checked != 1 { + t.Errorf("--sample 1 checked %d blobs, want 1: %s", result.Checked, log) + } + if result.Blobs != 3 { + t.Errorf("reported %d distinct blobs, want 3", result.Blobs) + } +} + +// TestEmitVerifyReplacesABlobThatIsNotItsName covers the reason #86 could not be +// fixed by the encoder alone: emit skips whatever is already present, so every +// blob published by the broken encoder stays broken until emit stops trusting +// presence. +func TestEmitVerifyReplacesABlobThatIsNotItsName(t *testing.T) { + fixture := newZoneFixture(t) + dir := t.TempDir() + hak := filepath.Join(dir, "sow_test_01.hak") + body := []byte("blood") + writeHak(t, hak, map[string][]byte{"bloodstain1.tga": body}) + + emit := func(verify bool) EmitResult { + t.Helper() + result, err := Emit(EmitOptions{ + ArtifactKey: artifactKey(t, hak), + ArtifactPath: hak, + Sink: fixture.sink, + Verify: verify, + }) + if err != nil { + t.Fatalf("emit (verify=%v): %v", verify, err) + } + return result + } + + emit(false) + key := blobKey(body) + encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) + if err != nil { + t.Fatal(err) + } + fixture.zone.mu.Lock() + good := fixture.zone.objects[key] + fixture.zone.objects[key] = append(append([]byte{}, good[:blobHeaderBytes]...), encoder.EncodeAll(body, nil)...) + fixture.zone.mu.Unlock() + + if plain := emit(false); plain.BlobsWritten != 0 { + t.Fatalf("a plain re-emit wrote %d blobs; it is supposed to trust presence", plain.BlobsWritten) + } + if verified := emit(true); verified.BlobsWritten != 1 { + t.Fatalf("--verify wrote %d blobs, want 1 (the bad copy must be replaced)", verified.BlobsWritten) + } + + fixture.zone.mu.Lock() + repaired := fixture.zone.objects[key] + fixture.zone.mu.Unlock() + if !bytes.Equal(repaired, good) { + t.Error("the replaced blob is not what the current encoder produces") + } + if _, err := inspectBlob(repaired); err != nil { + t.Errorf("the replaced blob still fails inspection: %v", err) + } + + // A second verifying run has nothing left to repair. + if again := emit(true); again.BlobsWritten != 0 { + t.Errorf("--verify rewrote %d good blobs", again.BlobsWritten) + } +} diff --git a/internal/nwsync/zone_test.go b/internal/nwsync/zone_test.go index eeb4df3..e17f9ae 100644 --- a/internal/nwsync/zone_test.go +++ b/internal/nwsync/zone_test.go @@ -21,6 +21,13 @@ type fakeZone struct { objects map[string][]byte puts []string failOn func(key string) bool // when true, the PUT fails + url string // base the same objects are readable at +} + +// pullZone reads the fake zone the way the public pull zone is read: plain +// unauthenticated GETs, no storage API. +func (z *fakeZone) pullZone() blobSource { + return newPullZone(z.url) } func newFakeZone(t *testing.T) (*fakeZone, func(string) string) { @@ -28,6 +35,7 @@ func newFakeZone(t *testing.T) (*fakeZone, func(string) string) { zone := &fakeZone{objects: map[string][]byte{}} server := httptest.NewServer(zone) t.Cleanup(server.Close) + zone.url = server.URL + "/sow-nwsync" getenv := func(name string) string { switch name { case "NWSYNC_STORAGE_ZONE":