Compare commits
2
Commits
682f920114
...
9b7be2c76e
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9b7be2c76e | ||
|
|
56d4054118 |
@@ -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-hak` | `crucible hak` | ERF/HAK pack/unpack + hak manifests |
|
||||
| `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-wiki` | `crucible wiki` | render + deploy mechanical wiki pages |
|
||||
|
||||
|
||||
+36
-1
@@ -36,6 +36,7 @@ aliases.
|
||||
| `depot` | `pull` | Incremental verified pull of every referenced blob. |
|
||||
| `nwsync` | `emit` | Explode one artifact into NWSync blobs plus its own NSYM manifest. |
|
||||
| `nwsync` | `assemble` | Merge per-artifact NSYM manifests into one merged manifest. |
|
||||
| `nwsync` | `verify` | Decompress and hash a published manifest's blobs through the pull zone. |
|
||||
|
||||
`depot status` and `depot get` pick their backend either with `--out DIR`, a
|
||||
depot tree on disk, or with `--target bunny|cdn`, a remote backend. The two
|
||||
@@ -48,8 +49,9 @@ beside the artifact itself with the extension replaced, so `emit` and
|
||||
`assemble` agree on where it is without being told.
|
||||
|
||||
```
|
||||
nwsync emit [--as NAME] [--out DIR] [--jobs N] <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 verify [--sample N] [--base URL] [--jobs N] <manifest-sha1>
|
||||
```
|
||||
|
||||
`emit` is latency-bound, not CPU-bound: every blob costs an existence probe
|
||||
@@ -73,6 +75,28 @@ resolves it. `--tlk-key` has its own slot because the TLK shadows nothing.
|
||||
`--group-id` is per channel — 1 is current, 2 is testing, and 0 leaves the field
|
||||
out of the sidecar.
|
||||
|
||||
`nwsync verify` is the only check on a published blob upstream of a player's
|
||||
client. It reads the **pull zone**, not the storage API, and needs no
|
||||
credential: what matters is the bytes a client is served, edge behaviour
|
||||
included. Every blob is decompressed and hashed, and the zstd frame is asserted
|
||||
to declare its content size. Neither half is optional — a `Content-Length` check
|
||||
passes a byte-correct-looking object whose contents are short, and a round-trip
|
||||
check alone passes a frame the game client cannot decode but Go's decoder can.
|
||||
Failures are reported per blob as missing, malformed framing, size mismatch or
|
||||
hash mismatch, and the exit code is 1.
|
||||
|
||||
A full sweep of the live manifest is roughly 69,000 blobs and 15 GB, so
|
||||
`--sample N` exists to make verifying routine; the default is a full sweep.
|
||||
`--base URL` (or `NWSYNC_PULL_BASE`) overrides the public host.
|
||||
|
||||
`emit --verify` applies the same check where `emit` would otherwise skip. `emit`
|
||||
normally reads a blob's presence as proof of its contents, decided by a 1-byte
|
||||
range GET, so an object written truncated — or written by an emitter since found
|
||||
broken — is skipped by every later run forever and no backfill repairs it. With
|
||||
`--verify` the stored copy is read back, unwrapped, hashed against its own name,
|
||||
and replaced when it does not match. It costs a full GET per existing blob, so
|
||||
it is a repair pass, not the default.
|
||||
|
||||
`emit` uploads blobs first and the index last, so the presence of an index is
|
||||
the publication marker: an artifact whose emit died halfway leaves real blobs in
|
||||
the zone and no index. Blob names are content hashes, so re-running skips
|
||||
@@ -108,6 +132,17 @@ saving on hak content: a 250 MB hak emitted 2296 blobs totalling 59 MB on disk
|
||||
against 249 MB of resources, as recorded in the sidecar's `on_disk_bytes` and
|
||||
`total_bytes`.
|
||||
|
||||
The zstd frame always declares its `Frame_Content_Size`. The game client sizes
|
||||
its output buffer from that field and cannot decode a frame without one, but the
|
||||
Go encoder omits it below 256 bytes, so `emit` re-headers those frames into the
|
||||
shape reference libzstd emits: `Single_Segment_flag` set, `Window_Descriptor`
|
||||
dropped, and a one-byte content size in its place. `zstd -l <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
|
||||
|
||||
Existing scripts may continue using these names indefinitely. They are accepted
|
||||
|
||||
@@ -108,10 +108,11 @@ var Registry = []Builder{
|
||||
{
|
||||
Name: "nwsync",
|
||||
Bin: "crucible-nwsync",
|
||||
Summary: "publish NWSync blobs and manifests (emit/assemble)",
|
||||
Summary: "publish NWSync blobs and manifests (emit/assemble/verify)",
|
||||
Commands: []Command{
|
||||
{Name: "emit", Summary: "explode one artifact into blobs plus its own NSYM manifest", Usage: "crucible nwsync emit <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: "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,
|
||||
},
|
||||
|
||||
@@ -160,7 +160,7 @@ func TestCanonicalCommandSurface(t *testing.T) {
|
||||
"module": {"build", "extract", "validate", "compare", "manifest"},
|
||||
"topdata": {"validate", "build", "package", "compare", "convert"},
|
||||
"wiki": {"build", "deploy"},
|
||||
"nwsync": {"emit", "assemble"},
|
||||
"nwsync": {"emit", "assemble", "verify"},
|
||||
}
|
||||
for _, builder := range Registry {
|
||||
got := builder.subcommands()
|
||||
|
||||
@@ -3,6 +3,7 @@ package nwsync
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
|
||||
"github.com/klauspost/compress/zstd"
|
||||
@@ -31,6 +32,19 @@ var (
|
||||
blobDecoder, _ = zstd.NewReader(nil, zstd.WithDecoderConcurrency(1))
|
||||
)
|
||||
|
||||
// zstd frame header bits we care about. A frame starts with the magic, then a
|
||||
// one-byte Frame_Header_Descriptor: bits 7-6 size the Frame_Content_Size field,
|
||||
// bit 5 is Single_Segment_flag, bits 1-0 size the Dictionary_ID field.
|
||||
const (
|
||||
zstdFrameMagic = "\x28\xb5\x2f\xfd"
|
||||
frameSingleSegment = 1 << 5
|
||||
frameDictionaryMask = 0x03
|
||||
// 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.
|
||||
func compressBlob(data []byte) []byte {
|
||||
var out bytes.Buffer
|
||||
@@ -38,10 +52,95 @@ func compressBlob(data []byte) []byte {
|
||||
for _, field := range header {
|
||||
_ = 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()
|
||||
}
|
||||
|
||||
// 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
|
||||
// package wrote — or one upstream wrote — can be compared by its uncompressed
|
||||
// bytes, which is the only comparison that is meaningful across zstd
|
||||
|
||||
@@ -33,7 +33,9 @@ var skippedTypes = resTypes("nss", "ndb", "gic")
|
||||
// manifest quietly disagrees with. Bump it only when emitted bytes change — it
|
||||
// is deliberately not the build revision, which would invalidate every
|
||||
// published index on every unrelated commit.
|
||||
const emitterVersion = "1"
|
||||
// Version 2 declares Frame_Content_Size on every blob (#86); version 1 omitted
|
||||
// it below 256 bytes and no client could sync past such a blob.
|
||||
const emitterVersion = "2"
|
||||
|
||||
// serverTypes are loaded only server-side; a manifest holding nothing else
|
||||
// has no client contents. Mirrors upstream's GlobalResTypeServerList, whose
|
||||
@@ -78,6 +80,7 @@ type EmitOptions struct {
|
||||
As string // name override, for a TLK whose filename is not its published name
|
||||
OutDir string // write locally instead of uploading — the conformance path
|
||||
Jobs int // resources in flight at once; 0 means defaultEmitJobs
|
||||
Verify bool // hash what would be skipped instead of trusting presence
|
||||
Sink sink // test seam; nil means OutDir or the zone
|
||||
}
|
||||
|
||||
@@ -129,7 +132,7 @@ func Emit(options EmitOptions) (EmitResult, error) {
|
||||
if jobs < 1 {
|
||||
jobs = defaultEmitJobs
|
||||
}
|
||||
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs)
|
||||
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs, options.Verify)
|
||||
if err != nil {
|
||||
return EmitResult{}, err
|
||||
}
|
||||
@@ -199,7 +202,7 @@ func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name strin
|
||||
//
|
||||
// The returned entries are in artifact order whatever order the workers finish
|
||||
// in, because a manifest's bytes are promised deterministic by emitterVersion.
|
||||
func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jobs int) ([]Entry, int, int64, error) {
|
||||
func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jobs int, verify bool) ([]Entry, int, int64, error) {
|
||||
// A resref appearing twice inside one artifact resolves to the last one,
|
||||
// the way resman lets the last container added win.
|
||||
order := make([]Identity, 0, len(index))
|
||||
@@ -271,7 +274,7 @@ func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jo
|
||||
if duplicate {
|
||||
return
|
||||
}
|
||||
written, err := target.putBlob(fmt.Sprintf("%x", sum), func() []byte { return compressBlob(payload) })
|
||||
written, err := target.putBlob(fmt.Sprintf("%x", sum), verify, func() []byte { return compressBlob(payload) })
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if err != nil {
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
@@ -194,7 +195,14 @@ func marshalSidecar(sidecar Sidecar) ([]byte, error) {
|
||||
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 {
|
||||
return filepath.Join(root, "data", "sha1", sha1Hex[0:2], sha1Hex[2:4], sha1Hex)
|
||||
return filepath.Join(root, filepath.FromSlash(blobKey(sha1Hex)))
|
||||
}
|
||||
|
||||
@@ -501,3 +501,31 @@ func TestRunUsageErrors(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestEveryBlobDeclaresItsFrameContentSize guards the fault that stopped every
|
||||
// client sync (#86): klauspost/compress omits Frame_Content_Size for inputs
|
||||
// under 256 bytes, and the game client cannot decode a frame without it. This
|
||||
// asserts a frame property, not a round trip — the zstd CLI and Go's decoder
|
||||
// both stream such a frame happily, so round-tripping cannot see the defect.
|
||||
func TestEveryBlobDeclaresItsFrameContentSize(t *testing.T) {
|
||||
// 230 and 175 are real sizes from the manifest that failed to sync; 255/256
|
||||
// straddle the encoder's threshold.
|
||||
for _, size := range []int{1, 32, 175, 230, 255, 256, 257, 1024, 5000} {
|
||||
payload := make([]byte, size)
|
||||
for i := range payload {
|
||||
payload[i] = byte('a' + i%26)
|
||||
}
|
||||
blob := compressBlob(payload)
|
||||
if !frameDeclaresContentSize(blob[blobHeaderBytes:]) {
|
||||
t.Errorf("blob of %d bytes declares no frame content size (descriptor %#x)",
|
||||
size, blob[blobHeaderBytes+4])
|
||||
}
|
||||
got, err := decompressBlob(blob)
|
||||
if err != nil {
|
||||
t.Fatalf("decompress %d-byte blob: %v", size, err)
|
||||
}
|
||||
if !bytes.Equal(got, payload) {
|
||||
t.Errorf("%d-byte blob did not round trip", size)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+54
-1
@@ -12,12 +12,17 @@ import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
)
|
||||
|
||||
const (
|
||||
exitOK = 0
|
||||
exitUsage = 64
|
||||
exitInternal = 70
|
||||
// exitDrift says the command worked and the zone is wrong, which is a
|
||||
// different thing for CI to act on than the command failing. It matches
|
||||
// depot's code for the same meaning.
|
||||
exitDrift = 1
|
||||
)
|
||||
|
||||
// Run executes an nwsync subcommand. args[0] is the subcommand (emit|assemble);
|
||||
@@ -32,6 +37,8 @@ func Run(args []string, stdout, stderr io.Writer) int {
|
||||
return runEmit(args[1:], stderr)
|
||||
case "assemble":
|
||||
return runAssemble(args[1:], stderr)
|
||||
case "verify":
|
||||
return runVerify(args[1:], stdout, stderr, os.Getenv)
|
||||
case "-h", "--help", "help":
|
||||
printRunUsage(stdout)
|
||||
return exitOK
|
||||
@@ -44,17 +51,27 @@ func Run(args []string, stdout, stderr io.Writer) int {
|
||||
|
||||
func printRunUsage(w io.Writer) {
|
||||
fmt.Fprint(w, `usage:
|
||||
nwsync emit [--as NAME] [--out DIR] <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 verify [--sample N] [--base URL] <manifest-sha1>
|
||||
|
||||
emit explodes one .hak/.erf or one loose file (the TLK) into NWSync blobs plus
|
||||
a NSYM index covering only that artifact, and uploads both. assemble merges
|
||||
those indexes into one manifest, reading no bulk data. Artifact keys are depot
|
||||
keys; an index lives beside its artifact, with the extension replaced.
|
||||
|
||||
verify reads a published manifest and its blobs back through the public pull
|
||||
zone, with no credential, and decompresses and hashes every one. It is the only
|
||||
check on a published blob upstream of a player's client.
|
||||
|
||||
--verify makes emit hash what it would otherwise skip. emit normally treats a
|
||||
blob's presence as proof of its contents, so without this an object written
|
||||
truncated, or written by an emitter since found broken, is skipped forever.
|
||||
|
||||
--out DIR writes to a local repository tree instead of uploading, which is the
|
||||
conformance path against upstream nwn_nwsync_write. Without it, the zone comes
|
||||
from NWSYNC_STORAGE_ZONE, NWSYNC_STORAGE_PASSWORD and BUNNY_STORAGE_HOST.
|
||||
verify needs none of those; its base comes from --base or NWSYNC_PULL_BASE.
|
||||
`)
|
||||
}
|
||||
|
||||
@@ -82,6 +99,7 @@ func runEmit(args []string, stderr io.Writer) int {
|
||||
as := fs.String("as", "", "published name of the artifact, when it differs from the key")
|
||||
out := fs.String("out", "", "write to a local repository tree instead of uploading")
|
||||
jobs := fs.Int("jobs", defaultEmitJobs, "resources to hash, compress and store at once")
|
||||
verify := fs.Bool("verify", false, "read back and hash blobs that already exist instead of trusting their presence")
|
||||
positional, err := parseArgs(fs, args)
|
||||
if err != nil {
|
||||
return exitUsage
|
||||
@@ -101,6 +119,7 @@ func runEmit(args []string, stderr io.Writer) int {
|
||||
As: *as,
|
||||
OutDir: *out,
|
||||
Jobs: *jobs,
|
||||
Verify: *verify,
|
||||
})
|
||||
if err != nil {
|
||||
fmt.Fprintf(stderr, "nwsync emit: %v\n", err)
|
||||
@@ -111,6 +130,40 @@ func runEmit(args []string, stderr io.Writer) int {
|
||||
return exitOK
|
||||
}
|
||||
|
||||
func runVerify(args []string, stdout, stderr io.Writer, getenv func(string) string) int {
|
||||
fs := flag.NewFlagSet("verify", flag.ContinueOnError)
|
||||
fs.SetOutput(stderr)
|
||||
base := fs.String("base", getenv("NWSYNC_PULL_BASE"), "pull zone base URL to read through")
|
||||
sample := fs.Int("sample", 0, "check this many random blobs instead of all of them")
|
||||
jobs := fs.Int("jobs", defaultEmitJobs, "blobs to fetch and hash at once")
|
||||
positional, err := parseArgs(fs, args)
|
||||
if err != nil {
|
||||
return exitUsage
|
||||
}
|
||||
if len(positional) != 1 {
|
||||
fmt.Fprintf(stderr, "nwsync verify: exactly one <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.SetOutput(stderr)
|
||||
|
||||
+39
-10
@@ -20,11 +20,17 @@ import (
|
||||
// upstream's output and ours can be diffed on a developer machine.
|
||||
type sink interface {
|
||||
// putBlob stores one NWCompressedBuffer blob under its sha1 name and
|
||||
// returns the bytes stored, or 0 if the blob was already there. Blob names
|
||||
// are content hashes, so an existing name is existing content — which is
|
||||
// why body is a thunk: compression is the expensive part of emit and a
|
||||
// blob that is already stored must not pay for it.
|
||||
putBlob(sha1Hex string, body func() []byte) (int64, error)
|
||||
// returns the bytes stored, or 0 if a good copy was already there. Blob
|
||||
// names are content hashes, so an existing name is normally taken as
|
||||
// existing content — which is why body is a thunk: compression is the
|
||||
// expensive part of emit and a blob that is already stored must not pay
|
||||
// for it.
|
||||
//
|
||||
// verify stops trusting presence: the stored copy is read back, unwrapped
|
||||
// and hashed, and replaced when it is not what its name claims. Without it
|
||||
// an object written truncated, or written by an emitter since found broken,
|
||||
// is skipped by every later emit forever and no backfill can repair it.
|
||||
putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error)
|
||||
// putIndex stores a NSYM manifest and its sidecar under key, which is
|
||||
// either an artifact-derived object key or a local path.
|
||||
putIndex(key string, manifest, sidecar []byte) error
|
||||
@@ -37,9 +43,15 @@ type sink interface {
|
||||
// dirSink writes a local NWSync repository tree.
|
||||
type dirSink struct{ root string }
|
||||
|
||||
func (s dirSink) putBlob(sha1Hex string, body func() []byte) (int64, error) {
|
||||
func (s dirSink) putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error) {
|
||||
blob := blobPath(s.root, sha1Hex)
|
||||
if _, err := os.Stat(blob); err == nil {
|
||||
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
|
||||
}
|
||||
if err := os.MkdirAll(filepath.Dir(blob), 0o755); err != nil {
|
||||
@@ -91,8 +103,8 @@ type zoneSink struct {
|
||||
zone string
|
||||
}
|
||||
|
||||
func (s zoneSink) putBlob(sha1Hex string, body func() []byte) (int64, error) {
|
||||
key := path.Join("data", "sha1", sha1Hex[0:2], sha1Hex[2:4], sha1Hex)
|
||||
func (s zoneSink) putBlob(sha1Hex string, verify bool, body func() []byte) (int64, error) {
|
||||
key := blobKey(sha1Hex)
|
||||
// A throttled probe must never be read as "missing, re-upload" or as
|
||||
// "present, skip", so only a confirmed Present skips the upload.
|
||||
state, _, err := s.store.ProbeKey(s.ctx, key)
|
||||
@@ -100,7 +112,24 @@ func (s zoneSink) putBlob(sha1Hex string, body func() []byte) (int64, error) {
|
||||
return 0, fmt.Errorf("probe blob %s: %w", sha1Hex, err)
|
||||
}
|
||||
if state == depot.Present {
|
||||
return 0, nil
|
||||
if !verify {
|
||||
return 0, nil
|
||||
}
|
||||
// The probe only proved the object exists. Read it back and hold it to
|
||||
// its own name.
|
||||
//
|
||||
// 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()
|
||||
if err := s.put(key, data); err != nil {
|
||||
|
||||
@@ -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[:]
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -21,6 +21,13 @@ type fakeZone struct {
|
||||
objects map[string][]byte
|
||||
puts []string
|
||||
failOn func(key string) bool // when true, the PUT fails
|
||||
url string // base the same objects are readable at
|
||||
}
|
||||
|
||||
// pullZone reads the fake zone the way the public pull zone is read: plain
|
||||
// unauthenticated GETs, no storage API.
|
||||
func (z *fakeZone) pullZone() blobSource {
|
||||
return newPullZone(z.url)
|
||||
}
|
||||
|
||||
func newFakeZone(t *testing.T) (*fakeZone, func(string) string) {
|
||||
@@ -28,6 +35,7 @@ func newFakeZone(t *testing.T) (*fakeZone, func(string) string) {
|
||||
zone := &fakeZone{objects: map[string][]byte{}}
|
||||
server := httptest.NewServer(zone)
|
||||
t.Cleanup(server.Close)
|
||||
zone.url = server.URL + "/sow-nwsync"
|
||||
getenv := func(name string) string {
|
||||
switch name {
|
||||
case "NWSYNC_STORAGE_ZONE":
|
||||
|
||||
Reference in New Issue
Block a user