Compare commits

..
9 Commits
Author SHA1 Message Date
archvillainetteandClaude Opus 5 3bf2031c2e docs(nwsync): verify is what tells you which keys to purge
ci / ci (pull_request) Successful in 3m29s
Running the repair for sow-tools#88 disproved the advice #90 had just
landed. "Purge the zone, then believe verify" assumed the stale set was
unknowable. It is not: verify reads the edge, so a run straight after a
repair names every key the edge is still serving stale — a survey, not a
verdict. Purge those, re-run, and the second run is the verdict.

The measured numbers are the argument. The repair rewrote 2,603 blobs at
the origin; 8 were stale at the edge, all of them ones a failed player
sync had pulled ninety minutes earlier. The edge only caches what someone
fetched, so purging the whole zone would have cooled 69,169 objects to
fix 8.

Refs #88, #89, #75.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-01 09:00:20 +00:00
archvillainette 3cac6e9484 fix(project): a module resref names a file, not a resource (#93)
`module.resref` is the name of the built `.mod` on disk, so the 16-byte resref limit never applied to it — NWN:EE module file names are routinely longer. The blanket check rejected `ShadowsOverWestgate` (19 characters) and blocked sow-module#60:

```
crucible module build
  module.resref "ShadowsOverWestgate" exceeds 16 characters
```

## What changed

`internal/project/project.go` — `module.resref` is validated as a **file name** now, which still rejects a resref that is a path or empty.

The 16-character limit is kept where the value really does become a resref: a project with `paths.assets` and no `haks[]` names its single generated HAK after the module resref (`build.go:1144`), and a HAK name is a resref the engine loads. That case now says what to do about it instead of refusing every long module name.

## Verified

- 3 new tests in `internal/project`: a long module name validates and produces `ShadowsOverWestgate.mod`; a resref containing a path is rejected; a long resref that would name a generated HAK is still rejected.
- `make check` green.
- `crucible module build` in sow-module writes `module/ShadowsOverWestgate.mod` (70 resources).
- `crucible topdata validate` in sow-topdata still passes — its assets live under `topdata.assets`, not `paths.assets`, so the HAK guard does not bite.

Merge this **first**: sow-module's rename PR cannot go green in CI until this lands and its `flake.lock` is bumped.

Refs ShadowsOverWestgate/sow-module#60

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #93

Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-08-01 08:59:25 +00:00
archvillainette 3f78197f0a Purge the edge after a repair, or verify answers per PoP (#89) (#90)
build-binaries / build-binaries (push) Successful in 2m33s
#89 asked for a decision. This is it, and it is the laziest of the three options listed there: **purge the whole pull zone by hand after a repair, one call, documented in the repair procedure.**

Why not the other two:

- Purging from `emit --verify` needs a CDN credential `emit` deliberately does not hold, and `emit` reports how many blobs it wrote, never which ones — so it could not target the keys anyway.
- Waiting out the TTL means 30 days.

Whole-zone rather than per-key costs a cold cache on a zone whose objects are mostly cold, and a repair scatters thousands of keys across the tree regardless.

Also answers the question #89 left open: **the zone does not negative-cache.** A missing key answers 404 with `cache-control: no-cache` and `cdn-cache: MISS`, still MISS on an immediate retry (checked credential-free, 2026-08-01).

Docs only — `docs/command-surface.md` and `nwsync`'s usage text. The procedure itself lives in sow-platform's NWSync runbook, next to the zone it acts on.

Closes #89.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #90
Reviewed-by: xtul <mpiasecki720@protonmail.com>
Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 23:13:55 +00:00
archvillainette 7cc53aeb68 fix(nwsync): declare Frame_Content_Size on every blob, and verify what is published (#87)
build-binaries / build-binaries (push) Successful in 2m40s
Closes #86. Closes #85.

These land together on purpose. Fixing the encoder alone changes nothing for the blobs already in the zone, because `emit` skips whatever is already present.

## #86 — the framing fix

`klauspost/compress` omits the zstd `Frame_Content_Size` field for inputs under 256 bytes, which the format 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 — rejected roughly 6% of our blobs outright. Any single one stops a sync dead, so no client could complete a sync of the live manifest.

No encoder option changes this, so `compressBlob` re-headers the affected frames into the shape libzstd itself emits: `Single_Segment_flag` set, `Window_Descriptor` dropped, and the freed byte spent on a one-byte `Frame_Content_Size`. Same length in, same length out, and the same descriptor byte (`0x24`) the issue recorded from libzstd.

`compressBlob` then asserts its own output. An encoder upgrade that finds another way to omit the field would otherwise reproduce #86 in silence, and a blob is skipped by every later emit once written.

`emitter_version` goes to `2`, so `assemble` refuses to merge an index written by the encoder that omitted the field.

**Proved against the reference decoder, not just a round trip.** A real emitted 175-byte blob:

```
Frames  Skips  Compressed  Uncompressed  Ratio  Check  Filename
     1      0      48   B       175   B  3.646  XXH64  frame.zst
c59d6620d4ffd4bf3fe73df43b19b7afcfe8fea4  -            <- zstd -dc | sha1sum
c59d6620d4ffd4bf3fe73df43b19b7afcfe8fea4               <- the blob's own name
```

Before the fix that `Uncompressed` column was blank.

## #85 — `nwsync verify`

`crucible nwsync verify <manifest-sha1>` reads a manifest and its blobs back through the **public pull zone**, with no credential, because what matters is the bytes a client is served, edge behaviour included. Every distinct blob is decompressed and hashed; failures are reported per blob as missing / malformed framing / size mismatch / hash mismatch, and the exit code is 1.

- `--sample N` makes a routine check cheap against a manifest that is ~69,000 blobs and 15 GB; the default is a full sweep.
- `--base URL` / `NWSYNC_PULL_BASE` overrides the public host.
- The manifest is checked against its own sha1 before a single blob is fetched.
- `emit --verify` applies the same check where `emit` would otherwise trust presence, and replaces a stored blob that is not what its name claims. This is what makes the #86 blobs repairable.

## Why the existing checks missed this

Both new checks assert the **frame property**, not just a round trip. The conformance suite (#59) compares decompressed bytes, so a frame that decodes correctly passes regardless of its header; and the earlier zone audit decompressed 68 blobs with the `zstd` CLI, a *more* capable decoder than the client's, which certified exactly the blobs the client rejects.

## Checks

`make check` and `make smoke` green. Second commit is the fixes from a two-axis review of the first.

## Not in this PR

Three follow-ups, filed separately: the backfill has not been run, replacing a blob does not purge the pull-zone edge cache, and #85's runbook line belongs to `sow-platform`.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #87

Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 22:33:12 +00:00
archvillainette 682f920114 fix(module): run scripts/fetch-upstream-manifests (#84)
build-binaries / build-binaries (push) Successful in 2m30s
Closes #83. Part of #54; implements #63 Decision 5 ("Coordinated rename flip").

## What changed

One string literal in `internal/app/app.go:388`:

```go
[]string{"scripts", "fetch-hak-manifest"} -> []string{"scripts", "fetch-upstream-manifests"}
```

sow-module renamed its half in ShadowsOverWestgate/sow-module#57
(`scripts/fetch-hak-manifest.sh` -> `scripts/fetch-upstream-manifests.sh`,
extended to resolve the topdata channel as well). The script name is a
hardcoded path convention shared by the two repos, not config, so both sides
only work when they carry the same name.

No doc in this repo named the old script (`grep` over the tree found the one
call site only), so nothing else needed touching.

## No fallback, by decision

Per #63 Decision 5: no transitional symlink, no Go-side fallback. Both sides
flip and the short broken window is accepted, because the failure is loud and
unmistakable (`required project script is missing`). A fallback would keep
both names alive forever.

## Merge order matters

`sow-module/.gitea/workflows/release.yml` runs `nix flake update sow-tools`, so
it always builds against the latest crucible, unpinned. Every crucible module
build in sow-module fails between sow-module#57 merging and a crucible release
carrying this flip. **Merge sow-module#57 and cut a crucible release back to
back** to keep that gap short.

## Checks

- `go build ./...`, `go vet ./internal/app/`, and the full `go test ./...` suite pass.
- No test covers this call path, and none was added: it is a single hardcoded
  constant that has to match the other repo, so a test here could only assert
  the literal against itself. The real check is the sow-module build.

## Still open on the issue's "done when"

- [x] app.go runs `scripts/fetch-upstream-manifests`
- [x] no doc in this repo names `fetch-hak-manifest`
- [ ] a crucible release is cut, and the crucible module build succeeds in a
      sow-module checkout at #57's head — needs a release after this merges

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #84
Reviewed-by: xtul <mpiasecki720@protonmail.com>
Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 19:00:25 +00:00
archvillainette 2f860ca9e4 fix(nwsync): write emit and assemble summaries to stderr (#82)
build-binaries / build-binaries (push) Successful in 2m35s
Fixes #81.

`crucible nwsync emit` and `assemble` printed their summary line to **stdout**. Any caller that captures a script's stdout as a value gets the summary glued onto it — `pack-haks.sh` does `release_dir="$(...)"`, so all 11 emit summaries landed in `$release_dir` and `publish-release.sh` died with "release dir not found".

Both lines move to stderr, where `lib.sh`'s own `nwsync: emitted $key` log already goes. Neither line is a machine-readable contract: sow-topdata's contract tests grep an `EMIT_LOG` their own fake-crucible stub writes, not real stdout, so nothing parses these.

`runEmit`/`runAssemble` no longer take the stdout writer — a leak in those two functions is now impossible to write by accident. `Run` still passes stdout to `printRunUsage` for explicit `-h`/`help`, which is correct.

Adds `TestRunKeepsSummariesOffStdout`: runs both verbs end to end against a local tree and asserts stdout stays empty while the summary reaches stderr. It asserts emptiness, not wording, so the summary text stays free to change.

Full `go test ./...` green.

Follow-up, outside this repo: sow-assets-manifest needs a `flake.lock` bump, then a re-run of the v0.2.1-rc1 tag.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #82
Reviewed-by: xtul <mpiasecki720@protonmail.com>
Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 18:42:20 +00:00
archvillainette 1c2acc5530 feat(nwsync): emit blobs in parallel with a bounded worker pool (#80)
build-binaries / build-binaries (push) Successful in 2m21s
Closes #79.

Emit is latency-bound, not CPU-bound. Every blob costs two serial round-trips to the zone — a `ProbeKey` HEAD, then a `PutReader` PUT — so a hak with a few thousand resources pays a few thousand serialised latencies. Measured on the live sow-assets-manifest backfill: 26 s of CPU across 9.5 minutes of wall clock, on a 4-core host with 5 GB free and peak RSS of 51 MB.

`emit` now hashes, compresses and stores `--jobs N` resources at once, default 16 — matching `DEPOT_JOBS` and the transport's `MaxIdleConnsPerHost`, so a worker per connection needs no fresh TLS handshake. `--jobs 1` is exactly the old behaviour.

### Three properties had to survive

Each has a test in `internal/nwsync/jobs_test.go`:

- **Deterministic manifest bytes.** `emitterVersion` promises a manifest is a function of its artifact, so `entries` is index-addressed rather than appended to — a worker owns `entries[i]` alone and the slice comes back in artifact order whatever order uploads finish in. `TestEmitProducesTheSameIndexAtEveryJobCount` diffs the `.nsym` and its sidecar between `-jobs 1` and `-jobs 16`.
- **Index still lands last.** Any worker's failure aborts before a manifest is written. `TestEmitLeavesNoIndexWhenAParallelUploadFails` fails every blob PUT with 16 workers in flight and asserts no `.nsym` appears. Under `-race` it also covers the shared counters.
- **Identical content still shares one blob.** This one bit during development and is the reason to read the diff carefully: serially, the sink's existence check absorbed two resrefs with identical bytes. In parallel both workers probe, both miss, and both upload — `TestEmitWritesBlobsAndManifest` caught it as "wrote 2 blobs, want 1". Claiming the sha1 in-process restores the dedupe and skips a probe round-trip as well.

### Memory

Peak now tracks the resources in flight rather than one resource. The ceiling is `N` × the 15 MB `fileSizeLimit` plus its compressed copy — bounded by a constant this package enforces itself, and still not tracking the archive. `TestEmitPeakMemoryIsBoundedByJobCount` re-runs the #76 regression check at `-jobs 8`: a hak 8× bigger still costs the same.

This is only cheap because of #78. Before streaming emit, N workers would have meant N whole archives resident.

### Not done

Skipping the `ProbeKey` HEAD on a first-time emit would halve round-trips, but doubles uploaded bytes on a re-run — which is exactly what a backfill is. Noted in #79 so it is not rediscovered; parallelism is the better lever and this PR takes it.

Worker compression still serialises on `blobEncoder`, which is `WithEncoderConcurrency(1)` for the memory reason in `compressedbuf.go`. At 26 s of CPU per hak that is not worth trading memory for, but it is where to look if the numbers ever say otherwise.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #80

Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 15:25:33 +00:00
archvillainette fa32dd411f fix(nwsync): stream emit so peak memory tracks the largest resource (#76) (#78)
build-binaries / build-binaries (push) Successful in 2m13s
Closes #76. Part of #54.

`nwsync emit` held roughly 5× the artifact size in RAM, so ovh-main (7 GB, no swap) OOM-killed it on any hak over ~1.4 GB. That blocked the backfill in sow-assets-manifest and would have killed the next release rebuilding a large hak.

## What it does

- `erf.ReadIndex` / `erf.ReadPayload` — parse the header and resource table only, read one payload on demand. `erf.Read` keeps its shape but returns payloads as subslices of the buffer instead of fresh copies, which removes one full copy for the `pipeline` callers too. Callers must not mutate `Resource.Data`; the doc comment says so and no caller does.
- `nwsync.Emit` opens the artifact, hashes it by streaming for the key check (through a section reader, so the file offset stays put), then hashes, compresses and stores one resource at a time. The archive is never resident. Shadowed duplicate resrefs are now never read at all.
- Bounds checks in `ReadIndex` moved to `int64`, so key/resource-list offsets can no longer overflow.

## Not in the issue, but memory-motivated

The zstd blob encoder ran at the default concurrency, which is one encoder per CPU, each holding a window-sized history — about 200 MB of live heap doing nothing on a 24-core runner. `EncodeAll` is single-threaded per call, so concurrency 1 costs nothing. `TestSingleThreadedEncoderMatchesDefault` pins the claim that blobs come out byte-identical.

## Measured

Peak heap during emit, sampled 1 ms:

| hak | before | after |
|-----|--------|-------|
| 8 MB | 155 MB | 21 MB |
| 64 MB | 289 MB | 22 MB |

Flat, as the acceptance asks. `TestEmitPeakMemoryDoesNotScaleWithArtifactSize` fails if the 64 MB fixture costs more than the 8 MB one plus 24 MB of slack.

## Gaps

- The regression check measures Go heap, not RSS, and its largest fixture is 64 MB — a multi-GB run was not done here. A 2.15 GB hak now needs about the same ~22 MB the 64 MB one does, so the 5 GB budget is not close, but that is inference from the flat curve, not a measurement.
- mmap was suggested in the issue and skipped. Payload buffers are still anonymous, but they are one resource each (≤15 MiB), so making them file-backed buys nothing now.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #78
Reviewed-by: xtul <mpiasecki720@protonmail.com>
Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 11:16:23 +00:00
archvillainette 7437653f14 docs: describe the on-disk shape of an emitted NWSync tree (#77)
Documents what `nwsync emit --out DIR` actually produces, so a zone can be checked by hand.

The trap this removes: a blob's filename is the SHA-1 of the resource's *original* bytes, but the file on disk is NWCompressedBuffer framing — a 24-byte `NSYC` header then a zstd frame. `sha1sum <blob>` therefore never matches the name it is sitting under. The doc records the recipe that does match:

```
tail -c +25 <blob> | zstd -dc | sha1sum
```

Also notes the tree layout, that the uncompressed length is a little-endian `uint32` at offset 12, and the rough 4:1 compression ratio on hak content.

Verified two ways: against a real emit of a 250 MB hak (2296 blobs, 59 MB on disk against 249 MB of resources), and against `internal/nwsync/compressedbuf.go`, where the header is written as `[]uint32{blobMagic, blobVersion, algorithmZstd, uint32(len(data)), zstdHeaderVer, zstdDictionary}` — confirming field 3 at offset 12.

Docs only, no behaviour change. Falls out of the investigation in #76; that fix is not in this PR.

🤖 Generated with [Claude Code](https://claude.com/claude-code)Reviewed-on: #77
Reviewed-by: xtul <mpiasecki720@protonmail.com>
Co-authored-by: vickydotbat <vickydotbat@tutamail.com>
2026-07-31 11:00:49 +00:00
19 changed files with 1519 additions and 126 deletions
+1 -1
View File
@@ -19,7 +19,7 @@ Crucible is how the artifact repos turn source into artifacts.
| `crucible-depot` | `crucible depot` | content-addressed depot blob verify/move | | `crucible-depot` | `crucible depot` | content-addressed depot blob verify/move |
| `crucible-hak` | `crucible hak` | ERF/HAK pack/unpack + hak manifests | | `crucible-hak` | `crucible hak` | ERF/HAK pack/unpack + hak manifests |
| `crucible-module` | `crucible module` | build/extract/validate/compare the `.mod` | | `crucible-module` | `crucible module` | build/extract/validate/compare the `.mod` |
| `crucible-nwsync` | `crucible nwsync` | NWSync blob emit + manifest assemble | | `crucible-nwsync` | `crucible nwsync` | NWSync blob emit + manifest assemble + verify |
| `crucible-topdata` | `crucible topdata` | compile 2da/tlk topdata + packages | | `crucible-topdata` | `crucible topdata` | compile 2da/tlk topdata + packages |
| `crucible-wiki` | `crucible wiki` | render + deploy mechanical wiki pages | | `crucible-wiki` | `crucible wiki` | render + deploy mechanical wiki pages |
+91 -1
View File
@@ -36,6 +36,7 @@ aliases.
| `depot` | `pull` | Incremental verified pull of every referenced blob. | | `depot` | `pull` | Incremental verified pull of every referenced blob. |
| `nwsync` | `emit` | Explode one artifact into NWSync blobs plus its own NSYM manifest. | | `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` | `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 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 depot tree on disk, or with `--target bunny|cdn`, a remote backend. The two
@@ -48,10 +49,20 @@ beside the artifact itself with the extension replaced, so `emit` and
`assemble` agree on where it is without being told. `assemble` agree on where it is without being told.
``` ```
nwsync emit [--as NAME] [--out DIR] <artifact-key> <file> nwsync emit [--as NAME] [--out DIR] [--jobs N] [--verify] <artifact-key> <file>
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>... nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...
nwsync verify [--sample N] [--base URL] [--jobs N] <manifest-sha1>
``` ```
`emit` is latency-bound, not CPU-bound: every blob costs an existence probe
plus an upload, and a measured backfill spent 26 seconds of CPU across 9.5
minutes of wall clock. `--jobs N` (default 16) sets how many resources are in
flight at once. The manifest is byte-identical at any value — the number of
workers is never observable in the output. Peak memory is `N` times the
per-resource limit of 15 MB plus its compressed copy, so raising `N` far past
the default costs real memory for little gain: the transport keeps 16 idle
connections per host, and past that a worker pays a fresh TLS handshake.
Both verbs upload by default; nothing bulky is ever written to the runner's Both verbs upload by default; nothing bulky is ever written to the runner's
disk. `--out DIR` writes a local repository tree instead, which is the disk. `--out DIR` writes a local repository tree instead, which is the
conformance path against upstream `nwn_nwsync_write`. The zone comes from conformance path against upstream `nwn_nwsync_write`. The zone comes from
@@ -64,12 +75,91 @@ 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 `--group-id` is per channel — 1 is current, 2 is testing, and 0 leaves the field
out of the sidecar. 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.
**After a repair, `verify` is what tells you which keys to purge.** A repair is
the one thing that makes a key serve different bytes than it did before, and the
edge caches these objects for 30 days precisely because that normally cannot
happen. The two commands look at different copies on purpose: `emit --verify`
repairs the **origin**, `verify` reads the **edge**. So a `verify` run straight
after a repair is not a verdict — it is a survey, and every blob it still calls
bad is one the edge is serving stale. Purge exactly those, then re-run it; only
that second run is the verdict.
Purging the keys `verify` names beats purging the zone, because the edge only
ever cached what somebody actually fetched: the 2026-08-01 repair rewrote 2,603
blobs at the origin and left 8 stale at the edge. The purge belongs in the
repair procedure rather than in `emit`, which reports how many blobs it wrote
and never which ones — so it could not target one even with a CDN credential,
which it deliberately does not hold (#89; the procedure itself is in
sow-platform's NWSync runbook).
`emit` uploads blobs first and the index last, so the presence of an index is `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 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 the zone and no index. Blob names are content hashes, so re-running skips
whatever already landed, and `assemble` fails closed on an artifact with no whatever already landed, and `assemble` fails closed on an artifact with no
index rather than publishing a manifest that is missing a hak. index rather than publishing a manifest that is missing a hak.
### What an emitted tree looks like
`--out DIR` produces the same tree `emit` would upload, which makes it the way
to check a zone by hand without touching one:
```
<artifact-sha>.nsym binary index
<artifact-sha>.nsym.json the same index, readable
data/sha1/a7/4a/a74aa84a... one blob per resource, two-level fanout
```
A blob's name is the SHA-1 of the resource's **original** bytes, but the file on
disk is not those bytes: each blob is wrapped in NWCompressedBuffer framing, a
24-byte `NSYC` header followed by a zstd frame. Hashing the file directly will
not match its name, which is the obvious
first thing to try and the obvious first thing to be confused by. Strip the
header first:
```
tail -c +25 <blob> | zstd -dc | sha1sum # == the blob's filename
```
The header carries the uncompressed length as a little-endian `uint32` at offset
12, so the decompressed size is checkable without decompressing. Compression is
worth roughly a 4:1
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 <frame>` 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 ## Hidden compatibility aliases
Existing scripts may continue using these names indefinitely. They are accepted Existing scripts may continue using these names indefinitely. They are accepted
+1 -1
View File
@@ -385,7 +385,7 @@ func refreshBuildModuleManifest(ctx context, p *project.Project, progress func(s
} }
progress("Refreshing hak list from the latest published sow-assets manifest...") progress("Refreshing hak list from the latest published sow-assets manifest...")
if err := runProjectScript(ctx, p, []string{"scripts", "fetch-hak-manifest"}, manifestPath); err != nil { if err := runProjectScript(ctx, p, []string{"scripts", "fetch-upstream-manifests"}, manifestPath); err != nil {
return "", "", err return "", "", err
} }
if _, err := pipeline.ApplyHAKManifest(p, manifestPath); err != nil { if _, err := pipeline.ApplyHAKManifest(p, manifestPath); err != nil {
+2 -1
View File
@@ -108,10 +108,11 @@ var Registry = []Builder{
{ {
Name: "nwsync", Name: "nwsync",
Bin: "crucible-nwsync", Bin: "crucible-nwsync",
Summary: "publish NWSync blobs and manifests (emit/assemble)", Summary: "publish NWSync blobs and manifests (emit/assemble/verify)",
Commands: []Command{ Commands: []Command{
{Name: "emit", Summary: "explode one artifact into blobs plus its own NSYM manifest", Usage: "crucible nwsync emit <artifact> --out DIR"}, {Name: "emit", Summary: "explode one artifact into blobs plus its own NSYM manifest", Usage: "crucible nwsync emit <artifact> --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: "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 <manifest-sha1> [--sample N]"},
}, },
Wired: true, Wired: true,
}, },
+1 -1
View File
@@ -160,7 +160,7 @@ func TestCanonicalCommandSurface(t *testing.T) {
"module": {"build", "extract", "validate", "compare", "manifest"}, "module": {"build", "extract", "validate", "compare", "manifest"},
"topdata": {"validate", "build", "package", "compare", "convert"}, "topdata": {"validate", "build", "package", "compare", "convert"},
"wiki": {"build", "deploy"}, "wiki": {"build", "deploy"},
"nwsync": {"emit", "assemble"}, "nwsync": {"emit", "assemble", "verify"},
} }
for _, builder := range Registry { for _, builder := range Registry {
got := builder.subcommands() got := builder.subcommands()
+92 -45
View File
@@ -311,63 +311,110 @@ func Write(w io.Writer, archive Archive) error {
return nil return nil
} }
// IndexEntry locates one resource inside an archive without holding its
// payload. Streaming callers read one payload at a time from these, so peak
// memory tracks the largest resource instead of the whole archive.
type IndexEntry struct {
Name string
Type uint16
Offset int64
Size int64
}
// Index is the header plus the resource table of an ERF: everything except the
// payloads.
type Index struct {
FileType string
Version string
Entries []IndexEntry
}
// ReadIndex parses the tables of an ERF of the given size, reading only the
// header, the key list and the resource list.
func ReadIndex(r io.ReaderAt, size int64) (Index, error) {
if size < headerSize {
return Index{}, fmt.Errorf("erf file too small: %d bytes", size)
}
var hdr header
if err := binary.Read(io.NewSectionReader(r, 0, headerSize), binary.LittleEndian, &hdr); err != nil {
return Index{}, fmt.Errorf("decode erf header: %w", err)
}
if int64(hdr.KeyListOffset)+int64(hdr.EntryCount)*24 > size {
return Index{}, fmt.Errorf("erf key list exceeds file bounds")
}
keys := make([]keyEntry, hdr.EntryCount)
keyReader := io.NewSectionReader(r, int64(hdr.KeyListOffset), int64(hdr.EntryCount)*24)
if err := binary.Read(keyReader, binary.LittleEndian, &keys); err != nil {
return Index{}, fmt.Errorf("decode key list: %w", err)
}
if int64(hdr.ResourceListOffset)+int64(hdr.EntryCount)*8 > size {
return Index{}, fmt.Errorf("erf resource list exceeds file bounds")
}
entries := make([]resourceEntry, hdr.EntryCount)
entryReader := io.NewSectionReader(r, int64(hdr.ResourceListOffset), int64(hdr.EntryCount)*8)
if err := binary.Read(entryReader, binary.LittleEndian, &entries); err != nil {
return Index{}, fmt.Errorf("decode resource list: %w", err)
}
index := Index{
FileType: string(hdr.FileType[:]),
Version: string(hdr.Version[:]),
Entries: make([]IndexEntry, 0, hdr.EntryCount),
}
for position, key := range keys {
entry := entries[position]
if int64(entry.Offset)+int64(entry.Size) > size {
return Index{}, fmt.Errorf("resource %d exceeds file bounds", position)
}
index.Entries = append(index.Entries, IndexEntry{
Name: string(bytes.TrimRight(key.ResRef[:], "\x00")),
Type: key.ResourceType,
Offset: int64(entry.Offset),
Size: int64(entry.Size),
})
}
return index, nil
}
// ReadPayload returns one resource's bytes.
func ReadPayload(r io.ReaderAt, entry IndexEntry) ([]byte, error) {
payload := make([]byte, entry.Size)
if _, err := r.ReadAt(payload, entry.Offset); err != nil {
return nil, fmt.Errorf("read resource %q: %w", entry.Name, err)
}
return payload, nil
}
// Read materialises a whole archive. Payloads are subslices of the buffer the
// archive was read into, so nothing is copied twice: a caller must not mutate
// Data. Callers that only need one resource at a time should use ReadIndex
// instead, which never holds the archive at all.
func Read(r io.Reader) (Archive, error) { func Read(r io.Reader) (Archive, error) {
data, err := io.ReadAll(r) data, err := io.ReadAll(r)
if err != nil { if err != nil {
return Archive{}, fmt.Errorf("read erf: %w", err) return Archive{}, fmt.Errorf("read erf: %w", err)
} }
if len(data) < headerSize { index, err := ReadIndex(bytes.NewReader(data), int64(len(data)))
return Archive{}, fmt.Errorf("erf file too small: %d bytes", len(data)) if err != nil {
return Archive{}, err
} }
var hdr header resources := make([]Resource, 0, len(index.Entries))
if err := binary.Read(bytes.NewReader(data[:headerSize]), binary.LittleEndian, &hdr); err != nil { for _, entry := range index.Entries {
return Archive{}, fmt.Errorf("decode erf header: %w", err)
}
keyStart := int(hdr.KeyListOffset)
keyEnd := keyStart + int(hdr.EntryCount)*24
if keyEnd > len(data) {
return Archive{}, fmt.Errorf("erf key list exceeds file bounds")
}
keys := make([]keyEntry, hdr.EntryCount)
if err := binary.Read(bytes.NewReader(data[keyStart:keyEnd]), binary.LittleEndian, &keys); err != nil {
return Archive{}, fmt.Errorf("decode key list: %w", err)
}
resourceStart := int(hdr.ResourceListOffset)
resourceEnd := resourceStart + int(hdr.EntryCount)*8
if resourceEnd > len(data) {
return Archive{}, fmt.Errorf("erf resource list exceeds file bounds")
}
entries := make([]resourceEntry, hdr.EntryCount)
if err := binary.Read(bytes.NewReader(data[resourceStart:resourceEnd]), binary.LittleEndian, &entries); err != nil {
return Archive{}, fmt.Errorf("decode resource list: %w", err)
}
resources := make([]Resource, 0, hdr.EntryCount)
for index, key := range keys {
entry := entries[index]
start := int(entry.Offset)
end := start + int(entry.Size)
if end > len(data) {
return Archive{}, fmt.Errorf("resource %d exceeds file bounds", index)
}
resref := string(bytes.TrimRight(key.ResRef[:], "\x00"))
payload := make([]byte, entry.Size)
copy(payload, data[start:end])
resources = append(resources, Resource{ resources = append(resources, Resource{
Name: resref, Name: entry.Name,
Type: key.ResourceType, Type: entry.Type,
Data: payload, Data: data[entry.Offset : entry.Offset+entry.Size],
Size: int64(entry.Size), Size: entry.Size,
}) })
} }
return Archive{ return Archive{
FileType: string(hdr.FileType[:]), FileType: index.FileType,
Version: string(hdr.Version[:]), Version: index.Version,
Resources: resources, Resources: resources,
}, nil }, nil
} }
+106 -3
View File
@@ -3,6 +3,7 @@ package nwsync
import ( import (
"bytes" "bytes"
"encoding/binary" "encoding/binary"
"encoding/hex"
"fmt" "fmt"
"github.com/klauspost/compress/zstd" "github.com/klauspost/compress/zstd"
@@ -22,9 +23,26 @@ const (
blobHeaderBytes = 24 blobHeaderBytes = 24
) )
// EncodeAll/DecodeAll are single-threaded per call, so the default pool of one
// encoder per CPU only buys idle memory: each holds a window-sized history, so
// on a 24-core runner that is ~200 MB of live heap doing nothing. Concurrency 1
// produces byte-identical output.
var ( var (
blobEncoder, _ = zstd.NewWriter(nil) blobEncoder, _ = zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1))
blobDecoder, _ = zstd.NewReader(nil) 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
// oneByteContentSizeCeiling is the size above which a Frame_Content_Size no
// longer fits in one byte. Below it the field's size flag is 0, which is
// what lets klauspost/compress leave the field out entirely.
oneByteContentSizeCeiling = 256
) )
// compressBlob wraps data in NWCompressedBuffer framing. // compressBlob wraps data in NWCompressedBuffer framing.
@@ -34,10 +52,95 @@ func compressBlob(data []byte) []byte {
for _, field := range header { for _, field := range header {
_ = binary.Write(&out, binary.LittleEndian, field) _ = binary.Write(&out, binary.LittleEndian, field)
} }
out.Write(blobEncoder.EncodeAll(data, nil)) frame := declareFrameContentSize(blobEncoder.EncodeAll(data, nil), len(data))
// Fail closed rather than publish a blob no client can decode. An encoder
// upgrade that finds a new way to omit the field would otherwise reproduce
// #86 in silence, and a blob is skipped by every later emit once written.
if !frameDeclaresContentSize(frame) {
panic(fmt.Sprintf("nwsync: refusing to emit a %d-byte blob whose zstd frame declares no content size (descriptor %#x)",
len(data), frame[4]))
}
out.Write(frame)
return out.Bytes() 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(blob) > blobHeaderBytes && !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 >= oneByteContentSizeCeiling || 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 // decompressBlob unwraps NWCompressedBuffer framing. It exists so a blob this
// package wrote — or one upstream wrote — can be compared by its uncompressed // package wrote — or one upstream wrote — can be compared by its uncompressed
// bytes, which is the only comparison that is meaningful across zstd // bytes, which is the only comparison that is meaningful across zstd
+150 -49
View File
@@ -1,10 +1,10 @@
package nwsync package nwsync
import ( import (
"bytes"
"context" "context"
"crypto/sha1" "crypto/sha1"
"fmt" "fmt"
"io"
"os" "os"
"path" "path"
"path/filepath" "path/filepath"
@@ -12,6 +12,7 @@ import (
"sort" "sort"
"strconv" "strconv"
"strings" "strings"
"sync"
"time" "time"
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo" "git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo"
@@ -32,7 +33,9 @@ var skippedTypes = resTypes("nss", "ndb", "gic")
// manifest quietly disagrees with. Bump it only when emitted bytes change — it // manifest quietly disagrees with. Bump it only when emitted bytes change — it
// is deliberately not the build revision, which would invalidate every // is deliberately not the build revision, which would invalidate every
// published index on every unrelated commit. // 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 // serverTypes are loaded only server-side; a manifest holding nothing else
// has no client contents. Mirrors upstream's GlobalResTypeServerList, whose // has no client contents. Mirrors upstream's GlobalResTypeServerList, whose
@@ -62,12 +65,22 @@ type EmitResult struct {
BlobsWritten int BlobsWritten int
} }
// defaultEmitJobs is how many resources are hashed, compressed and stored at
// once. Emit is latency-bound, not CPU-bound: a blob costs a probe round-trip
// plus an upload round-trip, and a measured backfill spent 26 s of CPU across
// 9.5 minutes of wall clock. The figure matches depot's DEPOT_JOBS default and
// the transport's MaxIdleConnsPerHost, so a worker per connection needs no new
// TLS handshake.
const defaultEmitJobs = 16
// EmitOptions describes one emit run. // EmitOptions describes one emit run.
type EmitOptions struct { type EmitOptions struct {
ArtifactKey string // depot key of the artifact; the NSYM key is derived from it ArtifactKey string // depot key of the artifact; the NSYM key is derived from it
ArtifactPath string // the file on disk ArtifactPath string // the file on disk
As string // name override, for a TLK whose filename is not its published name As string // name override, for a TLK whose filename is not its published name
OutDir string // write locally instead of uploading — the conformance path 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 Sink sink // test seam; nil means OutDir or the zone
} }
@@ -79,11 +92,18 @@ type EmitOptions struct {
// real blobs in the zone and no index, which is unambiguous. Blob names are // real blobs in the zone and no index, which is unambiguous. Blob names are
// content hashes, so re-running skips whatever already landed. // content hashes, so re-running skips whatever already landed.
func Emit(options EmitOptions) (EmitResult, error) { func Emit(options EmitOptions) (EmitResult, error) {
artifact, err := os.ReadFile(options.ArtifactPath) artifact, err := os.Open(options.ArtifactPath)
if err != nil { if err != nil {
return EmitResult{}, fmt.Errorf("read artifact: %w", err) return EmitResult{}, fmt.Errorf("read artifact: %w", err)
} }
if err := checkArtifactKey(options.ArtifactKey, artifact); err != nil { defer artifact.Close()
info, err := artifact.Stat()
if err != nil {
return EmitResult{}, fmt.Errorf("read artifact: %w", err)
}
// A section reader, not the file itself: hashing must not move the file
// offset out from under everything that reads the artifact afterwards.
if err := checkArtifactKey(options.ArtifactKey, io.NewSectionReader(artifact, 0, info.Size())); err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
name := options.As name := options.As
@@ -98,7 +118,7 @@ func Emit(options EmitOptions) (EmitResult, error) {
return EmitResult{}, err return EmitResult{}, err
} }
resources, err := readArtifact(options.ArtifactPath, artifact, name) index, err := readArtifactIndex(options.ArtifactPath, artifact, info.Size(), name)
if err != nil { if err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
@@ -108,7 +128,11 @@ func Emit(options EmitOptions) (EmitResult, error) {
return EmitResult{}, err return EmitResult{}, err
} }
entries, blobs, onDiskBytes, err := emitResources(resources, target) jobs := options.Jobs
if jobs < 1 {
jobs = defaultEmitJobs
}
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs, options.Verify)
if err != nil { if err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
@@ -137,87 +161,164 @@ func openSink(outDir string, injected sink) (sink, error) {
return newZoneSink(context.Background(), os.Getenv) return newZoneSink(context.Background(), os.Getenv)
} }
// readArtifact returns the resources of an ERF/HAK/MOD, or the single resource // readArtifactIndex locates the resources of an ERF/HAK/MOD, or the single
// a loose file represents. Upstream's resman does the same dispatch on the // resource a loose file represents, without reading any payload. Upstream's
// file's first three bytes. name is the artifact's published name, which for a // resman does the same dispatch on the file's first three bytes. name is the
// loose file is also its resref. // artifact's published name, which for a loose file is also its resref.
func readArtifact(path string, data []byte, name string) ([]erf.Resource, error) { func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name string) ([]erf.IndexEntry, error) {
if len(data) >= 3 { magic := make([]byte, 3)
switch string(data[:3]) { if size >= 3 {
case "ERF", "HAK": if _, err := artifact.ReadAt(magic, 0); err != nil {
archive, err := erf.Read(bytes.NewReader(data)) return nil, fmt.Errorf("%s: %w", path, err)
if err != nil {
return nil, fmt.Errorf("%s: %w", path, err)
}
return archive.Resources, nil
case "MOD":
// A persistent world never publishes module contents, so the .mod
// contributes no bytes to a manifest — it only says which haks and
// which TLK the manifest covers.
return nil, fmt.Errorf("%s: a module is never emitted; a manifest is haks plus the TLK", path)
} }
} }
switch string(magic) {
case "ERF", "HAK":
index, err := erf.ReadIndex(artifact, size)
if err != nil {
return nil, fmt.Errorf("%s: %w", path, err)
}
return index.Entries, nil
case "MOD":
// A persistent world never publishes module contents, so the .mod
// contributes no bytes to a manifest — it only says which haks and
// which TLK the manifest covers.
return nil, fmt.Errorf("%s: a module is never emitted; a manifest is haks plus the TLK", path)
}
extension := filepath.Ext(filepath.Base(path)) extension := filepath.Ext(filepath.Base(path))
restype, ok := erf.ResourceTypeForExtension(extension) restype, ok := erf.ResourceTypeForExtension(extension)
if !ok { if !ok {
return nil, fmt.Errorf("%s: unknown resource type %q", path, extension) return nil, fmt.Errorf("%s: unknown resource type %q", path, extension)
} }
return []erf.Resource{{ return []erf.IndexEntry{{Name: name, Type: restype, Offset: 0, Size: size}}, nil
Name: name,
Type: restype,
Data: data,
}}, nil
} }
func emitResources(resources []erf.Resource, target sink) ([]Entry, int, int64, error) { // emitResources hashes, compresses and stores resources, reading each payload
// from the artifact only when its turn comes. Peak memory tracks the resources
// in flight, not the archive: a 2 GB hak must emit inside a runner's few spare
// GB. jobs of them are in flight at once, so the ceiling is jobs multiplied by
// fileSizeLimit and its compressed copy — bounded, and bounded by a constant
// this package enforces itself.
//
// 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, verify bool) ([]Entry, int, int64, error) {
// A resref appearing twice inside one artifact resolves to the last one, // A resref appearing twice inside one artifact resolves to the last one,
// the way resman lets the last container added win. // the way resman lets the last container added win.
order := make([]Identity, 0, len(resources)) order := make([]Identity, 0, len(index))
latest := make(map[Identity]erf.Resource, len(resources)) latest := make(map[Identity]erf.IndexEntry, len(index))
var tooBig []string var tooBig []string
for _, resource := range resources { for _, entry := range index {
if _, ok := erf.ExtensionForResourceType(resource.Type); !ok { if _, ok := erf.ExtensionForResourceType(entry.Type); !ok {
return nil, 0, 0, fmt.Errorf("resref %s is not resolvable (unknown restype %d)", resource.Name, resource.Type) return nil, 0, 0, fmt.Errorf("resref %s is not resolvable (unknown restype %d)", entry.Name, entry.Type)
} }
if slices.Contains(skippedTypes, resource.Type) { if slices.Contains(skippedTypes, entry.Type) {
continue continue
} }
if len(resource.Data) > fileSizeLimit { if entry.Size > fileSizeLimit {
tooBig = append(tooBig, fmt.Sprintf("%s: %d bytes > %d", resource.Name, len(resource.Data), fileSizeLimit)) tooBig = append(tooBig, fmt.Sprintf("%s: %d bytes > %d", entry.Name, entry.Size, fileSizeLimit))
continue continue
} }
identity := Identity{ResRef: strings.ToLower(resource.Name), ResType: resource.Type} identity := Identity{ResRef: strings.ToLower(entry.Name), ResType: entry.Type}
if _, seen := latest[identity]; !seen { if _, seen := latest[identity]; !seen {
order = append(order, identity) order = append(order, identity)
} }
latest[identity] = resource latest[identity] = entry
} }
if len(tooBig) > 0 { if len(tooBig) > 0 {
sort.Strings(tooBig) sort.Strings(tooBig)
return nil, 0, 0, fmt.Errorf("resources exceed the file size limit:\n %s", strings.Join(tooBig, "\n ")) return nil, 0, 0, fmt.Errorf("resources exceed the file size limit:\n %s", strings.Join(tooBig, "\n "))
} }
entries := make([]Entry, 0, len(order)) // Index-addressed, never appended to: a worker owns entries[i] alone, so
// the slice comes back in artifact order and needs no lock.
entries := make([]Entry, len(order))
var blobs int var blobs int
var onDiskBytes int64 var onDiskBytes int64
for _, identity := range order { var mu sync.Mutex
resource := latest[identity] var firstErr error
sum := sha1.Sum(resource.Data) // Two resrefs in one artifact can hold identical bytes, and therefore one
entries = append(entries, Entry{ // blob. Serially the sink's existence check absorbed that; in parallel both
// workers would probe, both miss, and both upload. Claiming the sha1 here
// restores the dedupe and skips the probe round-trip as well.
claimed := make(map[[20]byte]bool, len(order))
failed := func() bool {
mu.Lock()
defer mu.Unlock()
return firstErr != nil
}
store := func(i int) {
identity := order[i]
payload, err := erf.ReadPayload(artifact, latest[identity])
if err != nil {
mu.Lock()
if firstErr == nil {
firstErr = err
}
mu.Unlock()
return
}
sum := sha1.Sum(payload)
entries[i] = Entry{
SHA1: sum, SHA1: sum,
Size: uint32(len(resource.Data)), Size: uint32(len(payload)),
ResRef: identity.ResRef, ResRef: identity.ResRef,
ResType: identity.ResType, ResType: identity.ResType,
}) }
written, err := target.putBlob(fmt.Sprintf("%x", sum), func() []byte { return compressBlob(resource.Data) }) mu.Lock()
duplicate := claimed[sum]
claimed[sum] = true
mu.Unlock()
if duplicate {
return
}
written, err := target.putBlob(fmt.Sprintf("%x", sum), verify, func() []byte { return compressBlob(payload) })
mu.Lock()
defer mu.Unlock()
if err != nil { if err != nil {
return nil, 0, 0, err if firstErr == nil {
firstErr = err
}
return
} }
if written > 0 { if written > 0 {
blobs++ blobs++
onDiskBytes += written onDiskBytes += written
} }
} }
if jobs < 1 {
jobs = 1
}
work := make(chan int)
var wg sync.WaitGroup
for range jobs {
wg.Add(1)
go func() {
defer wg.Done()
for i := range work {
// After a failure the run is over — the caller discards
// everything and no index is written. Draining the rest of the
// channel cheaply, rather than returning, keeps the feeder from
// blocking on workers that have gone away.
if failed() {
continue
}
store(i)
}
}()
}
for i := range order {
work <- i
}
close(work)
wg.Wait()
if firstErr != nil {
return nil, 0, 0, firstErr
}
return entries, blobs, onDiskBytes, nil return entries, blobs, onDiskBytes, nil
} }
+134
View File
@@ -0,0 +1,134 @@
package nwsync
import (
"bytes"
"fmt"
"path/filepath"
"runtime/debug"
"strings"
"testing"
)
// manyResources is a hak body with enough distinct resources that a worker pool
// actually interleaves. Payloads differ so nothing is deduplicated away.
func manyResources(count int) map[string][]byte {
contents := make(map[string][]byte, count)
for i := range count {
contents[fmt.Sprintf("res%05d.tga", i)] = []byte(fmt.Sprintf("payload %d", i))
}
return contents
}
// TestEmitProducesTheSameIndexAtEveryJobCount is the contract that lets emit be
// parallel at all: emitterVersion promises a manifest's bytes are a function of
// its artifact, so the number of workers must not be observable in the output.
func TestEmitProducesTheSameIndexAtEveryJobCount(t *testing.T) {
// The sidecar stamps a wall-clock time unless this is set, which would make
// two runs differ for a reason that has nothing to do with job count.
t.Setenv("SOURCE_DATE_EPOCH", "1700000000")
dir := t.TempDir()
hak := filepath.Join(dir, "sow_test_01.hak")
writeHak(t, hak, manyResources(64))
key := artifactKey(t, hak)
emit := func(jobs int) (manifest, sidecar []byte, result EmitResult) {
out := filepath.Join(t.TempDir(), "out")
result, err := Emit(EmitOptions{
ArtifactKey: key,
ArtifactPath: hak,
OutDir: out,
Jobs: jobs,
})
if err != nil {
t.Fatalf("emit at -jobs %d: %v", jobs, err)
}
manifest, sidecar, err = dirSink{root: out}.getIndex(filepath.Base(result.ManifestPath))
if err != nil {
t.Fatalf("read index at -jobs %d: %v", jobs, err)
}
return manifest, sidecar, result
}
serialManifest, serialSidecar, serial := emit(1)
parallelManifest, parallelSidecar, parallel := emit(16)
if !bytes.Equal(serialManifest, parallelManifest) {
t.Errorf("manifest bytes differ between -jobs 1 and -jobs 16")
}
if !bytes.Equal(serialSidecar, parallelSidecar) {
t.Errorf("sidecar bytes differ between -jobs 1 and -jobs 16:\n %s\n %s", serialSidecar, parallelSidecar)
}
if serial.Entries != parallel.Entries || serial.BlobsWritten != parallel.BlobsWritten {
t.Errorf("-jobs 1 wrote %d entries/%d blobs, -jobs 16 wrote %d/%d",
serial.Entries, serial.BlobsWritten, parallel.Entries, parallel.BlobsWritten)
}
}
// TestEmitLeavesNoIndexWhenAParallelUploadFails is the fail-closed check with
// workers in flight: several uploads are in the air when the first one fails,
// and the index must still never appear. Run under -race this also covers the
// shared counters.
func TestEmitLeavesNoIndexWhenAParallelUploadFails(t *testing.T) {
fixture := newZoneFixture(t)
fixture.zone.failOn = func(key string) bool { return strings.HasPrefix(key, "data/sha1/") }
dir := t.TempDir()
hak := filepath.Join(dir, "sow_test_01.hak")
writeHak(t, hak, manyResources(64))
if _, err := Emit(EmitOptions{
ArtifactKey: artifactKey(t, hak),
ArtifactPath: hak,
Sink: fixture.sink,
Jobs: 16,
}); err == nil {
t.Fatal("emit reported success after an upload failed")
}
fixture.zone.mu.Lock()
defer fixture.zone.mu.Unlock()
for key := range fixture.zone.objects {
if strings.HasSuffix(key, ".nsym") {
t.Errorf("a half-emitted artifact published an index: %s", key)
}
}
}
// TestEmitPeakMemoryIsBoundedByJobCount pins the ceiling the parallel emit
// rests on. Peak still must not track the archive — it tracks the resources in
// flight, so a bigger hak at the same job count costs the same.
func TestEmitPeakMemoryIsBoundedByJobCount(t *testing.T) {
if testing.Short() {
t.Skip("writes a 64 MB fixture")
}
defer debug.SetGCPercent(debug.SetGCPercent(10))
measure := func(count, jobs int) uint64 {
dir := t.TempDir()
hak := filepath.Join(dir, "big.hak")
writeStreamedHak(t, hak, count)
options := EmitOptions{
ArtifactKey: artifactKey(t, hak),
ArtifactPath: hak,
As: filepath.Base(hak),
OutDir: filepath.Join(dir, "out"),
Jobs: jobs,
}
return peakHeapDuring(func() {
if _, err := Emit(options); err != nil {
t.Fatalf("emit %d resources at -jobs %d: %v", count, jobs, err)
}
})
}
const jobs = 8
small := measure(8, jobs) // 8 MB
large := measure(64, jobs) // 64 MB
// Each worker may hold one resourceSize payload plus its compressed copy,
// so the pool itself is the slack — not the archive.
const slack = 24 << 20
t.Logf("peak heap at -jobs %d: 8 MB hak %d bytes, 64 MB hak %d bytes", jobs, small, large)
if large > small+slack {
t.Fatalf("peak heap scaled with artifact size at -jobs %d: 8 MB hak peaked at %d bytes, 64 MB hak at %d", jobs, small, large)
}
}
+10 -2
View File
@@ -7,6 +7,7 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"io" "io"
"path"
"path/filepath" "path/filepath"
"sort" "sort"
"strings" "strings"
@@ -194,7 +195,14 @@ func marshalSidecar(sidecar Sidecar) ([]byte, error) {
return append(body, '\r', '\n'), nil return append(body, '\r', '\n'), nil
} }
// blobPath is the data store path for a blob, hash tree depth 2. // blobKey is where a blob lives in a zone, hash tree depth 2. emit writes it,
// verify reads it and the game client requests it, so the rule lives here and
// nowhere else.
func blobKey(sha1Hex string) string {
return path.Join("data", "sha1", sha1Hex[0:2], sha1Hex[2:4], sha1Hex)
}
// blobPath is the same location inside a local repository tree.
func blobPath(root, sha1Hex string) string { func blobPath(root, sha1Hex string) string {
return filepath.Join(root, "data", "sha1", sha1Hex[0:2], sha1Hex[2:4], sha1Hex) return filepath.Join(root, filepath.FromSlash(blobKey(sha1Hex)))
} }
+150
View File
@@ -0,0 +1,150 @@
package nwsync
import (
"bytes"
"fmt"
"math/rand"
"os"
"path/filepath"
"runtime"
"runtime/debug"
"testing"
"time"
"github.com/klauspost/compress/zstd"
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/erf"
)
// TestSingleThreadedEncoderMatchesDefault pins the claim the blob encoder's
// concurrency setting rests on: it saves memory only, and a published blob is
// the same bytes either way.
func TestSingleThreadedEncoderMatchesDefault(t *testing.T) {
standard, err := zstd.NewWriter(nil)
if err != nil {
t.Fatal(err)
}
defer standard.Close()
body := make([]byte, 4<<20)
random := rand.New(rand.NewSource(1))
random.Read(body[:len(body)/2])
for _, size := range []int{0, 1, 4 << 10, len(body)} {
if !bytes.Equal(blobEncoder.EncodeAll(body[:size], nil), standard.EncodeAll(body[:size], nil)) {
t.Fatalf("%d bytes compress differently at concurrency 1", size)
}
}
}
// resourceSize is one payload in the memory fixtures. Real haks hold a few MB
// per resource, and peak memory is meant to track that, not the archive.
const resourceSize = 1 << 20
// writeStreamedHak builds a hak of count resources without ever holding the
// archive in memory, so the fixture itself does not decide the measurement.
// Payloads are distinct, so no blob is deduplicated away.
func writeStreamedHak(t *testing.T, path string, count int) {
t.Helper()
payload := filepath.Join(t.TempDir(), "payload.bin")
body := make([]byte, resourceSize)
for index := range body {
body[index] = byte(index)
}
resources := make([]erf.Resource, 0, count)
for index := range count {
// A distinct first byte per resource is enough to give every payload
// its own sha1 while still streaming from one file per resource.
unique := filepath.Join(filepath.Dir(payload), fmt.Sprintf("p%d.bin", index))
body[0] = byte(index)
body[1] = byte(index >> 8)
if err := os.WriteFile(unique, body, 0o644); err != nil {
t.Fatal(err)
}
resources = append(resources, erf.Resource{
Name: fmt.Sprintf("res%05d", index),
Type: restype(t, "tga"),
SourcePath: unique,
Size: resourceSize,
})
}
file, err := os.Create(path)
if err != nil {
t.Fatal(err)
}
defer file.Close()
if err := erf.Write(file, erf.New("HAK", resources)); err != nil {
t.Fatalf("write hak: %v", err)
}
}
// peakHeapDuring runs work while sampling the heap, and returns the largest
// live heap it saw.
func peakHeapDuring(work func()) uint64 {
runtime.GC()
done := make(chan struct{})
peak := make(chan uint64, 1)
go func() {
var highest uint64
var stats runtime.MemStats
for {
select {
case <-done:
peak <- highest
return
default:
}
runtime.ReadMemStats(&stats)
if stats.HeapAlloc > highest {
highest = stats.HeapAlloc
}
time.Sleep(time.Millisecond)
}
}()
work()
close(done)
return <-peak
}
// TestEmitPeakMemoryDoesNotScaleWithArtifactSize is the regression check for
// the OOM kills on large haks: emit used to hold the whole archive (twice), so
// a 2 GB hak needed about 10 GB. Emitting an archive 8× bigger must not cost
// meaningfully more memory.
func TestEmitPeakMemoryDoesNotScaleWithArtifactSize(t *testing.T) {
if testing.Short() {
t.Skip("writes a 64 MB fixture")
}
// A lazy GC lets garbage pile up in proportion to the live heap, which
// hides the thing under test. Collecting eagerly makes the sampled heap
// track what emit actually holds.
defer debug.SetGCPercent(debug.SetGCPercent(10))
measure := func(count int) uint64 {
dir := t.TempDir()
hak := filepath.Join(dir, "big.hak")
writeStreamedHak(t, hak, count)
// The key is computed outside the measurement: the test helper reads
// the whole file to hash it, which emit itself no longer does.
options := EmitOptions{
ArtifactKey: artifactKey(t, hak),
ArtifactPath: hak,
As: filepath.Base(hak),
OutDir: filepath.Join(dir, "out"),
}
return peakHeapDuring(func() {
if _, err := Emit(options); err != nil {
t.Fatalf("emit %d resources: %v", count, err)
}
})
}
small := measure(8) // 8 MB
large := measure(64) // 64 MB
const slack = 24 << 20
t.Logf("peak heap: 8 MB hak %d bytes, 64 MB hak %d bytes", small, large)
if large > small+slack {
t.Fatalf("peak heap scaled with artifact size: 8 MB hak peaked at %d bytes, 64 MB hak at %d", small, large)
}
}
+55
View File
@@ -458,6 +458,33 @@ func TestAssembleFailsClosedOnMissingIndex(t *testing.T) {
} }
} }
// Callers capture a script's stdout as a value: `dir="$(pack-haks.sh)"`. A
// summary line on stdout gets glued onto that value, so both summaries belong
// on stderr.
func TestRunKeepsSummariesOffStdout(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "sow_top.hak")
writeHak(t, path, map[string][]byte{"appearance.2da": []byte("2da from sow_top")})
key := artifactKey(t, path)
out := filepath.Join(dir, "out")
for _, args := range [][]string{
{"emit", "--out", out, "--as", "sow_top.hak", key, path},
{"assemble", "--out", out, key},
} {
var stdout, stderr bytes.Buffer
if code := Run(args, &stdout, &stderr); code != exitOK {
t.Fatalf("Run(%v) exit=%d: %s", args, code, stderr.String())
}
if stdout.Len() != 0 {
t.Errorf("Run(%v) wrote to stdout: %q", args, stdout.String())
}
if stderr.Len() == 0 {
t.Errorf("Run(%v) reported no summary on stderr", args)
}
}
}
func TestRunUsageErrors(t *testing.T) { func TestRunUsageErrors(t *testing.T) {
cases := [][]string{ cases := [][]string{
nil, nil,
@@ -474,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)
}
}
}
+69 -7
View File
@@ -12,12 +12,17 @@ import (
"flag" "flag"
"fmt" "fmt"
"io" "io"
"os"
) )
const ( const (
exitOK = 0 exitOK = 0
exitUsage = 64 exitUsage = 64
exitInternal = 70 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); // Run executes an nwsync subcommand. args[0] is the subcommand (emit|assemble);
@@ -29,9 +34,11 @@ func Run(args []string, stdout, stderr io.Writer) int {
} }
switch args[0] { switch args[0] {
case "emit": case "emit":
return runEmit(args[1:], stdout, stderr) return runEmit(args[1:], stderr)
case "assemble": case "assemble":
return runAssemble(args[1:], stdout, stderr) return runAssemble(args[1:], stderr)
case "verify":
return runVerify(args[1:], stdout, stderr, os.Getenv)
case "-h", "--help", "help": case "-h", "--help", "help":
printRunUsage(stdout) printRunUsage(stdout)
return exitOK return exitOK
@@ -44,17 +51,30 @@ func Run(args []string, stdout, stderr io.Writer) int {
func printRunUsage(w io.Writer) { func printRunUsage(w io.Writer) {
fmt.Fprint(w, `usage: fmt.Fprint(w, `usage:
nwsync emit [--as NAME] [--out DIR] <artifact-key> <file> nwsync emit [--as NAME] [--out DIR] [--verify] <artifact-key> <file>
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>... nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...
nwsync verify [--sample N] [--base URL] <manifest-sha1>
emit explodes one .hak/.erf or one loose file (the TLK) into NWSync blobs plus 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 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 those indexes into one manifest, reading no bulk data. Artifact keys are depot
keys; an index lives beside its artifact, with the extension replaced. 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.
--verify repairs the storage zone, while verify reads the edge in front of it.
So a verify run right after a repair is a survey, not a verdict: it names the
keys the edge still serves stale. Purge those, then run it again.
--out DIR writes to a local repository tree instead of uploading, which is the --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 conformance path against upstream nwn_nwsync_write. Without it, the zone comes
from NWSYNC_STORAGE_ZONE, NWSYNC_STORAGE_PASSWORD and BUNNY_STORAGE_HOST. 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.
`) `)
} }
@@ -76,11 +96,13 @@ func parseArgs(fs *flag.FlagSet, args []string) ([]string, error) {
} }
} }
func runEmit(args []string, stdout, stderr io.Writer) int { func runEmit(args []string, stderr io.Writer) int {
fs := flag.NewFlagSet("emit", flag.ContinueOnError) fs := flag.NewFlagSet("emit", flag.ContinueOnError)
fs.SetOutput(stderr) fs.SetOutput(stderr)
as := fs.String("as", "", "published name of the artifact, when it differs from the key") 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") 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) positional, err := parseArgs(fs, args)
if err != nil { if err != nil {
return exitUsage return exitUsage
@@ -89,23 +111,63 @@ func runEmit(args []string, stdout, stderr io.Writer) int {
fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> are both required\n") fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> are both required\n")
return exitUsage return exitUsage
} }
if *jobs < 1 {
fmt.Fprintf(stderr, "nwsync emit: -jobs must be at least 1, got %d\n", *jobs)
return exitUsage
}
result, err := Emit(EmitOptions{ result, err := Emit(EmitOptions{
ArtifactKey: positional[0], ArtifactKey: positional[0],
ArtifactPath: positional[1], ArtifactPath: positional[1],
As: *as, As: *as,
OutDir: *out, OutDir: *out,
Jobs: *jobs,
Verify: *verify,
}) })
if err != nil { if err != nil {
fmt.Fprintf(stderr, "nwsync emit: %v\n", err) fmt.Fprintf(stderr, "nwsync emit: %v\n", err)
return exitInternal return exitInternal
} }
fmt.Fprintf(stdout, "emitted %s: %d resources, %d new blobs, index %s\n", fmt.Fprintf(stderr, "emitted %s: %d resources, %d new blobs, index %s\n",
result.Name, result.Entries, result.BlobsWritten, result.ManifestPath) result.Name, result.Entries, result.BlobsWritten, result.ManifestPath)
return exitOK return exitOK
} }
func runAssemble(args []string, stdout, stderr io.Writer) int { 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 <manifest-sha1> 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 := flag.NewFlagSet("assemble", flag.ContinueOnError)
fs.SetOutput(stderr) fs.SetOutput(stderr)
tlkKey := fs.String("tlk-key", "", "depot key of the TLK, which shadows nothing and merges last") tlkKey := fs.String("tlk-key", "", "depot key of the TLK, which shadows nothing and merges last")
@@ -134,7 +196,7 @@ func runAssemble(args []string, stdout, stderr io.Writer) int {
fmt.Fprintf(stderr, "nwsync assemble: %v\n", err) fmt.Fprintf(stderr, "nwsync assemble: %v\n", err)
return exitInternal return exitInternal
} }
fmt.Fprintf(stdout, "assembled manifest %s: %d resources, %s\n", fmt.Fprintf(stderr, "assembled manifest %s: %d resources, %s\n",
result.SHA1, result.Entries, result.ManifestPath) result.SHA1, result.Entries, result.ManifestPath)
return exitOK return exitOK
} }
+47 -13
View File
@@ -6,6 +6,7 @@ import (
"crypto/sha256" "crypto/sha256"
"encoding/hex" "encoding/hex"
"fmt" "fmt"
"io"
"os" "os"
"path" "path"
"path/filepath" "path/filepath"
@@ -19,11 +20,17 @@ import (
// upstream's output and ours can be diffed on a developer machine. // upstream's output and ours can be diffed on a developer machine.
type sink interface { type sink interface {
// putBlob stores one NWCompressedBuffer blob under its sha1 name and // putBlob stores one NWCompressedBuffer blob under its sha1 name and
// returns the bytes stored, or 0 if the blob was already there. Blob names // returns the bytes stored, or 0 if a good copy was already there. Blob
// are content hashes, so an existing name is existing content — which is // names are content hashes, so an existing name is normally taken as
// why body is a thunk: compression is the expensive part of emit and a // existing content — which is why body is a thunk: compression is the
// blob that is already stored must not pay for it. // expensive part of emit and a blob that is already stored must not pay
putBlob(sha1Hex string, body func() []byte) (int64, error) // 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 // putIndex stores a NSYM manifest and its sidecar under key, which is
// either an artifact-derived object key or a local path. // either an artifact-derived object key or a local path.
putIndex(key string, manifest, sidecar []byte) error putIndex(key string, manifest, sidecar []byte) error
@@ -36,9 +43,15 @@ type sink interface {
// dirSink writes a local NWSync repository tree. // dirSink writes a local NWSync repository tree.
type dirSink struct{ root string } 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) blob := blobPath(s.root, sha1Hex)
if _, err := os.Stat(blob); err == nil { if !verify {
// Stat, not read: the common path must not pay to open every blob that
// is already there.
if _, err := os.Stat(blob); err == nil {
return 0, nil
}
} else if stored, err := os.ReadFile(blob); err == nil && blobMatchesName(stored, sha1Hex) == nil {
return 0, nil return 0, nil
} }
if err := os.MkdirAll(filepath.Dir(blob), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(blob), 0o755); err != nil {
@@ -90,8 +103,8 @@ type zoneSink struct {
zone string 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) key := blobKey(sha1Hex)
// A throttled probe must never be read as "missing, re-upload" or as // A throttled probe must never be read as "missing, re-upload" or as
// "present, skip", so only a confirmed Present skips the upload. // "present, skip", so only a confirmed Present skips the upload.
state, _, err := s.store.ProbeKey(s.ctx, key) state, _, err := s.store.ProbeKey(s.ctx, key)
@@ -99,7 +112,24 @@ func (s zoneSink) putBlob(sha1Hex string, body func() []byte) (int64, error) {
return 0, fmt.Errorf("probe blob %s: %w", sha1Hex, err) return 0, fmt.Errorf("probe blob %s: %w", sha1Hex, err)
} }
if state == depot.Present { 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.
//
// This reads the storage API rather than the pull zone: emit holds the
// write credential, and a repair decision has to be made against the
// copy it is about to overwrite, not against an edge cache of it. A
// read that fails outright is a fault, not a verdict — treating it as
// "bad, re-upload" would turn a throttled zone into a full backfill.
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() data := body()
if err := s.put(key, data); err != nil { if err := s.put(key, data); err != nil {
@@ -193,14 +223,18 @@ func resolveIndexKey(artifactKey, outDir string) (string, error) {
// checkArtifactKey fails closed when the key's embedded digest is not the // checkArtifactKey fails closed when the key's embedded digest is not the
// digest of the bytes being emitted. Publishing an index under the wrong key // digest of the bytes being emitted. Publishing an index under the wrong key
// silently pairs a manifest with the wrong artifact. // silently pairs a manifest with the wrong artifact.
func checkArtifactKey(artifactKey string, artifact []byte) error { // artifact is hashed by streaming, so a multi-gigabyte hak is never resident.
func checkArtifactKey(artifactKey string, artifact io.Reader) error {
base := path.Base(artifactKey) base := path.Base(artifactKey)
digest := strings.TrimSuffix(base, path.Ext(base)) digest := strings.TrimSuffix(base, path.Ext(base))
if len(digest) != 64 { if len(digest) != 64 {
return fmt.Errorf("artifact key %q does not name a sha256", artifactKey) return fmt.Errorf("artifact key %q does not name a sha256", artifactKey)
} }
sum := sha256.Sum256(artifact) hash := sha256.New()
if got := hex.EncodeToString(sum[:]); got != digest { if _, err := io.Copy(hash, artifact); err != nil {
return fmt.Errorf("hash artifact: %w", err)
}
if got := hex.EncodeToString(hash.Sum(nil)); got != digest {
return fmt.Errorf("artifact key %q names digest %s but the file hashes to %s", artifactKey, digest, got) return fmt.Errorf("artifact key %q names digest %s but the file hashes to %s", artifactKey, digest, got)
} }
return nil return nil
+233
View File
@@ -0,0 +1,233 @@
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) {
// The argument is interpolated straight into a URL path, so it is checked
// rather than trusted: exactly 20 bytes of hex, nothing else.
if sum, err := hex.DecodeString(options.ManifestSHA1); err != nil || len(sum) != sha1.Size {
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 {
// 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(blobKey(entry.sha1Hex()))
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[:]
}
+285
View File
@@ -0,0 +1,285 @@
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
}
// keyOf is where a resource's blob lives, addressed by the sha1 of its
// uncompressed bytes — the same path the client requests.
func keyOf(body []byte) string {
sum := sha1.Sum(body)
return blobKey(hex.EncodeToString(sum[:]))
}
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)
}
// A default run is a full sweep, so it must reach every blob the manifest
// names — not some of them.
if result.Checked != result.Blobs || result.Blobs == 0 {
t.Errorf("checked %d of %d blobs; a full sweep must check all of them", result.Checked, result.Blobs)
}
}
func TestVerifyReportsAMissingBlob(t *testing.T) {
fixture := newVerifyFixture(t)
key := keyOf([]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 := keyOf([]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 := keyOf(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 := keyOf([]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 TestVerifyReportsAShortBlob(t *testing.T) {
fixture := newVerifyFixture(t)
key := keyOf([]byte("small"))
fixture.zone.mu.Lock()
// Well-formed all the way down and simply too short — the shape a killed
// upload leaves behind, and the one a Content-Length check would pass.
fixture.zone.objects[key] = compressBlob([]byte("sma"))
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, "size mismatch") {
t.Errorf("a short blob was not reported as a size 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 <= result.Checked {
t.Errorf("sampling %d of %d blobs is not a sample", result.Checked, 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 := keyOf(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)
}
}
+8
View File
@@ -21,6 +21,13 @@ type fakeZone struct {
objects map[string][]byte objects map[string][]byte
puts []string puts []string
failOn func(key string) bool // when true, the PUT fails 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) { 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{}} zone := &fakeZone{objects: map[string][]byte{}}
server := httptest.NewServer(zone) server := httptest.NewServer(zone)
t.Cleanup(server.Close) t.Cleanup(server.Close)
zone.url = server.URL + "/sow-nwsync"
getenv := func(name string) string { getenv := func(name string) string {
switch name { switch name {
case "NWSYNC_STORAGE_ZONE": case "NWSYNC_STORAGE_ZONE":
+12 -2
View File
@@ -560,8 +560,18 @@ func (p *Project) ValidateLayout() error {
if strings.TrimSpace(p.Config.Module.ResRef) == "" { if strings.TrimSpace(p.Config.Module.ResRef) == "" {
failures = append(failures, errors.New("module.resref is required")) failures = append(failures, errors.New("module.resref is required"))
} }
if len(p.Config.Module.ResRef) > 16 { // module.resref names the built .mod FILE, so the 16-byte resref limit does not
failures = append(failures, fmt.Errorf("module.resref %q exceeds 16 characters", p.Config.Module.ResRef)) // apply to it — NWN:EE module file names are routinely longer. It is validated as
// a file name instead. The limit still binds when the same value has to be a real
// resref: with no haks configured, an asset project names its single generated HAK
// after it, and a HAK name is a resref the engine loads.
if err := validateOutputFileName("module.resref", p.Config.Module.ResRef+".mod", ".mod"); err != nil {
failures = append(failures, err)
}
if len(p.Config.Module.ResRef) > 16 && strings.TrimSpace(p.Config.Paths.Assets) != "" && len(p.Config.HAKs) == 0 {
failures = append(failures, fmt.Errorf(
"module.resref %q exceeds 16 characters and would name this project's generated HAK; configure haks[] with a shorter name",
p.Config.Module.ResRef))
} }
if strings.TrimSpace(p.Config.Paths.Source) == "" && strings.TrimSpace(p.Config.Paths.Assets) == "" && !p.HasTopData() { if strings.TrimSpace(p.Config.Paths.Source) == "" && strings.TrimSpace(p.Config.Paths.Assets) == "" && !p.HasTopData() {
failures = append(failures, errors.New("at least one of paths.source, paths.assets, or topdata.source is required")) failures = append(failures, errors.New("at least one of paths.source, paths.assets, or topdata.source is required"))
+72
View File
@@ -1093,6 +1093,78 @@ func TestValidateLayoutAllowsMissingAssetsDir(t *testing.T) {
} }
} }
// module.resref names the built .mod FILE, not a resource inside an archive, so the
// 16-byte resref limit does not apply to it. NWN:EE module file names are commonly
// longer (ShadowsOverWestgate.mod is 19). The limit still binds everywhere a resref
// really is a resref — see TestValidateLayoutRejectsLongResRefWhenItNamesAHAK.
func TestValidateLayoutAllowsLongModuleResRef(t *testing.T) {
root := t.TempDir()
mkdirAll(t, filepath.Join(root, "src"))
mkdirAll(t, filepath.Join(root, "build"))
proj := &Project{
Root: root,
Config: Config{
Module: ModuleConfig{Name: "Shadows Over Westgate", ResRef: "ShadowsOverWestgate"},
Paths: PathConfig{Source: "src", Build: "build"},
},
}
if err := proj.ValidateLayout(); err != nil {
t.Fatalf("ValidateLayout rejected a 19-character module file name: %v", err)
}
if got, want := filepath.Base(proj.ModuleArchivePath()), "ShadowsOverWestgate.mod"; got != want {
t.Fatalf("ModuleArchivePath() = %q, want %q", got, want)
}
}
// A module.resref that is not a usable file name is still rejected.
func TestValidateLayoutRejectsModuleResRefThatIsAPath(t *testing.T) {
root := t.TempDir()
mkdirAll(t, filepath.Join(root, "src"))
proj := &Project{
Root: root,
Config: Config{
Module: ModuleConfig{Name: "Test", ResRef: "../escape/mod"},
Paths: PathConfig{Source: "src", Build: "build"},
},
}
err := proj.ValidateLayout()
if err == nil {
t.Fatal("ValidateLayout accepted a module.resref containing a path")
}
if !strings.Contains(err.Error(), "module.resref") {
t.Fatalf("error does not name the offending field: %v", err)
}
}
// When a project declares no haks, the module resref becomes the name of the single
// generated HAK — and a HAK name IS a resref the engine loads. The limit applies
// there, so a long name is only allowed for projects that build no HAKs.
func TestValidateLayoutRejectsLongResRefWhenItNamesAHAK(t *testing.T) {
root := t.TempDir()
mkdirAll(t, filepath.Join(root, "src"))
mkdirAll(t, filepath.Join(root, "assets"))
proj := &Project{
Root: root,
Config: Config{
Module: ModuleConfig{Name: "Shadows Over Westgate", ResRef: "ShadowsOverWestgate"},
Paths: PathConfig{Source: "src", Assets: "assets", Build: "build"},
},
}
err := proj.ValidateLayout()
if err == nil {
t.Fatal("ValidateLayout accepted a 19-character name for a generated HAK")
}
if !strings.Contains(err.Error(), "16") {
t.Fatalf("error does not explain the resref limit: %v", err)
}
}
// paths.build is an OUTPUT dir the builder creates (MkdirAll) before writing, so // paths.build is an OUTPUT dir the builder creates (MkdirAll) before writing, so
// a bare clone with no build dir yet must still validate/build with no pre-step // a bare clone with no build dir yet must still validate/build with no pre-step
// (R2/parity). Only a build path that exists but is not a directory is an error. // (R2/parity). Only a build path that exists but is not a directory is an error.