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>
366 lines
12 KiB
Go
366 lines
12 KiB
Go
package nwsync
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha1"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path"
|
|
"path/filepath"
|
|
"slices"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo"
|
|
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/erf"
|
|
)
|
|
|
|
// fileSizeLimit matches upstream's --limit-file-size default of 15 MB. A
|
|
// resource over it is a hard failure, not a skip: upstream quit(1)s and so do
|
|
// we. Our largest resource today is 13.66 MiB, so the headroom is thin.
|
|
const fileSizeLimit = 15 * 1024 * 1024
|
|
|
|
// skippedTypes are never published, matching upstream's GobalResTypeSkipList.
|
|
var skippedTypes = resTypes("nss", "ndb", "gic")
|
|
|
|
// emitterVersion identifies the blob/manifest byte format this package
|
|
// produces. assemble refuses to merge indexes that disagree on it, because two
|
|
// producers of blobs mean a skewed emitter can otherwise write blobs the merged
|
|
// 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.
|
|
// 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
|
|
// trailing 0 is RESTYPE_INVALID.
|
|
var serverTypes = append(resTypes(
|
|
"are", "dlg", "fac", "gic", "git", "ifo", "itp", "jrl", "ncs", "ndb",
|
|
"nss", "ptm", "utc", "utd", "ute", "uti", "utm", "utp", "uts", "utt", "utw",
|
|
), 0)
|
|
|
|
func resTypes(extensions ...string) []uint16 {
|
|
types := make([]uint16, 0, len(extensions))
|
|
for _, extension := range extensions {
|
|
restype, ok := erf.ResourceTypeForExtension(extension)
|
|
if !ok {
|
|
panic("nwsync: unknown restype " + extension)
|
|
}
|
|
types = append(types, restype)
|
|
}
|
|
return types
|
|
}
|
|
|
|
// EmitResult reports what one emit run produced.
|
|
type EmitResult struct {
|
|
Name string // artifact name, without extension
|
|
ManifestPath string
|
|
Entries 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.
|
|
type EmitOptions struct {
|
|
ArtifactKey string // depot key of the artifact; the NSYM key is derived from it
|
|
ArtifactPath string // the file on disk
|
|
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
|
|
}
|
|
|
|
// Emit explodes one artifact — a .hak/.erf or a loose file such as the TLK —
|
|
// into NWSync blobs plus a NSYM manifest describing only that artifact.
|
|
//
|
|
// Blobs go up as they are produced and the index lands last, so the presence of
|
|
// an index is the publication marker: an artifact whose emit died halfway has
|
|
// real blobs in the zone and no index, which is unambiguous. Blob names are
|
|
// content hashes, so re-running skips whatever already landed.
|
|
func Emit(options EmitOptions) (EmitResult, error) {
|
|
artifact, err := os.Open(options.ArtifactPath)
|
|
if err != nil {
|
|
return EmitResult{}, fmt.Errorf("read artifact: %w", err)
|
|
}
|
|
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
|
|
}
|
|
name := options.As
|
|
if name == "" {
|
|
name = path.Base(options.ArtifactKey)
|
|
}
|
|
extension := path.Ext(name)
|
|
name = strings.TrimSuffix(name, extension)
|
|
|
|
key, err := resolveIndexKey(options.ArtifactKey, options.OutDir)
|
|
if err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
|
|
index, err := readArtifactIndex(options.ArtifactPath, artifact, info.Size(), name)
|
|
if err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
|
|
target, err := openSink(options.OutDir, options.Sink)
|
|
if err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
|
|
jobs := options.Jobs
|
|
if jobs < 1 {
|
|
jobs = defaultEmitJobs
|
|
}
|
|
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs, options.Verify)
|
|
if err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
if len(entries) == 0 {
|
|
return EmitResult{}, fmt.Errorf("%s: nothing to index (no publishable resources)", options.ArtifactPath)
|
|
}
|
|
|
|
data, err := writeManifest(entries)
|
|
if err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
if err := putManifestPair(target, key, data, entries, onDiskBytes, Sidecar{ModuleName: name}); err != nil {
|
|
return EmitResult{}, err
|
|
}
|
|
return EmitResult{Name: name, ManifestPath: target.describe(key), Entries: len(entries), BlobsWritten: blobs}, nil
|
|
}
|
|
|
|
// openSink returns the zone sink, or a local directory when outDir is set.
|
|
func openSink(outDir string, injected sink) (sink, error) {
|
|
if injected != nil {
|
|
return injected, nil
|
|
}
|
|
if outDir != "" {
|
|
return dirSink{root: outDir}, nil
|
|
}
|
|
return newZoneSink(context.Background(), os.Getenv)
|
|
}
|
|
|
|
// readArtifactIndex locates the resources of an ERF/HAK/MOD, or the single
|
|
// resource a loose file represents, without reading any payload. Upstream's
|
|
// resman does the same dispatch on the file's first three bytes. name is the
|
|
// artifact's published name, which for a loose file is also its resref.
|
|
func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name string) ([]erf.IndexEntry, error) {
|
|
magic := make([]byte, 3)
|
|
if size >= 3 {
|
|
if _, err := artifact.ReadAt(magic, 0); err != nil {
|
|
return nil, fmt.Errorf("%s: %w", path, err)
|
|
}
|
|
}
|
|
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))
|
|
restype, ok := erf.ResourceTypeForExtension(extension)
|
|
if !ok {
|
|
return nil, fmt.Errorf("%s: unknown resource type %q", path, extension)
|
|
}
|
|
return []erf.IndexEntry{{Name: name, Type: restype, Offset: 0, Size: size}}, nil
|
|
}
|
|
|
|
// 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,
|
|
// the way resman lets the last container added win.
|
|
order := make([]Identity, 0, len(index))
|
|
latest := make(map[Identity]erf.IndexEntry, len(index))
|
|
var tooBig []string
|
|
for _, entry := range index {
|
|
if _, ok := erf.ExtensionForResourceType(entry.Type); !ok {
|
|
return nil, 0, 0, fmt.Errorf("resref %s is not resolvable (unknown restype %d)", entry.Name, entry.Type)
|
|
}
|
|
if slices.Contains(skippedTypes, entry.Type) {
|
|
continue
|
|
}
|
|
if entry.Size > fileSizeLimit {
|
|
tooBig = append(tooBig, fmt.Sprintf("%s: %d bytes > %d", entry.Name, entry.Size, fileSizeLimit))
|
|
continue
|
|
}
|
|
identity := Identity{ResRef: strings.ToLower(entry.Name), ResType: entry.Type}
|
|
if _, seen := latest[identity]; !seen {
|
|
order = append(order, identity)
|
|
}
|
|
latest[identity] = entry
|
|
}
|
|
if len(tooBig) > 0 {
|
|
sort.Strings(tooBig)
|
|
return nil, 0, 0, fmt.Errorf("resources exceed the file size limit:\n %s", strings.Join(tooBig, "\n "))
|
|
}
|
|
|
|
// 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 onDiskBytes int64
|
|
var mu sync.Mutex
|
|
var firstErr error
|
|
// Two resrefs in one artifact can hold identical bytes, and therefore one
|
|
// 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,
|
|
Size: uint32(len(payload)),
|
|
ResRef: identity.ResRef,
|
|
ResType: identity.ResType,
|
|
}
|
|
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 firstErr == nil {
|
|
firstErr = err
|
|
}
|
|
return
|
|
}
|
|
if written > 0 {
|
|
blobs++
|
|
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
|
|
}
|
|
|
|
// created is the sidecar timestamp. SOURCE_DATE_EPOCH pins it so a build can
|
|
// be reproduced byte for byte; the manifest itself is deterministic already.
|
|
func created() int64 {
|
|
if raw := os.Getenv("SOURCE_DATE_EPOCH"); raw != "" {
|
|
if seconds, err := strconv.ParseInt(raw, 10, 64); err == nil {
|
|
return seconds
|
|
}
|
|
}
|
|
return time.Now().Unix()
|
|
}
|
|
|
|
// putManifestPair stores a NSYM manifest and its .json sidecar at key. The
|
|
// caller supplies the sidecar fields it knows; the rest are derived from the
|
|
// entries. data must be the serialised form of entries.
|
|
func putManifestPair(target sink, key string, data []byte, entries []Entry, onDiskBytes int64, sidecar Sidecar) error {
|
|
var totalBytes int64
|
|
clientContents := false
|
|
for _, entry := range entries {
|
|
totalBytes += int64(entry.Size)
|
|
if !slices.Contains(serverTypes, entry.ResType) {
|
|
clientContents = true
|
|
}
|
|
}
|
|
sidecar.Version = manifestVersion
|
|
sidecar.SHA1 = fmt.Sprintf("%x", sha1.Sum(data))
|
|
sidecar.HashTreeDepth = hashTreeDepth
|
|
sidecar.IncludesModuleContents = false
|
|
sidecar.IncludesClientContents = clientContents
|
|
sidecar.TotalFiles = len(entries)
|
|
sidecar.TotalBytes = totalBytes
|
|
sidecar.OnDiskBytes = onDiskBytes
|
|
sidecar.Created = created()
|
|
sidecar.CreatedWith = buildinfo.String()
|
|
sidecar.EmitterVersion = emitterVersion
|
|
|
|
body, err := marshalSidecar(sidecar)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return target.putIndex(key, data, body)
|
|
}
|