From f69121a9e81894290ac689e4600105f88cddb444 Mon Sep 17 00:00:00 2001 From: vickydotbat Date: Fri, 31 Jul 2026 16:28:45 +0200 Subject: [PATCH] feat(nwsync): emit blobs in parallel with a bounded worker pool (#79) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Emit is latency-bound, not CPU-bound. Every blob costs two serial HTTP round-trips to the zone — an existence probe, then an upload — so a hak with a few thousand resources pays a few thousand serialised latencies. A measured backfill spent 26 seconds of CPU across 9.5 minutes of wall clock, on a host with three of four cores idle and 5 GB free. Emit now hashes, compresses and stores `--jobs N` resources at once (default 16, matching DEPOT_JOBS and the transport's idle connections per host). Three properties had to survive, and each has a test: - The manifest's bytes are promised deterministic by emitterVersion, 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 the uploads finish in. - The index is still the publication marker, so any worker's failure aborts the run before a manifest is written. Workers drain the rest of the channel instead of returning, which keeps the feeder from blocking on workers that have gone away. - Two resrefs holding identical bytes still share one blob. Serially the sink's existence check absorbed that; in parallel both workers would probe, both miss, and both upload. Claiming the sha1 in-process restores the dedupe and skips a probe round-trip as well. Peak memory is now the resources in flight rather than one resource, so the ceiling is N times the 15 MB per-resource limit and its compressed copy — bounded by a constant this package enforces itself, and still nowhere near tracking the archive. Co-Authored-By: Claude Opus 5 (1M context) --- docs/command-surface.md | 11 ++- internal/nwsync/emit.go | 108 ++++++++++++++++++++++++---- internal/nwsync/jobs_test.go | 134 +++++++++++++++++++++++++++++++++++ internal/nwsync/run.go | 6 ++ 4 files changed, 246 insertions(+), 13 deletions(-) create mode 100644 internal/nwsync/jobs_test.go diff --git a/docs/command-surface.md b/docs/command-surface.md index 33162a5..ae19ce2 100644 --- a/docs/command-surface.md +++ b/docs/command-surface.md @@ -48,10 +48,19 @@ 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] +nwsync emit [--as NAME] [--out DIR] [--jobs N] nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] ... ``` +`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 disk. `--out DIR` writes a local repository tree instead, which is the conformance path against upstream `nwn_nwsync_write`. The zone comes from diff --git a/internal/nwsync/emit.go b/internal/nwsync/emit.go index 7574e70..b1c6e10 100644 --- a/internal/nwsync/emit.go +++ b/internal/nwsync/emit.go @@ -12,6 +12,7 @@ import ( "sort" "strconv" "strings" + "sync" "time" "git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo" @@ -62,12 +63,21 @@ type EmitResult struct { 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 Sink sink // test seam; nil means OutDir or the zone } @@ -115,7 +125,11 @@ func Emit(options EmitOptions) (EmitResult, error) { return EmitResult{}, err } - entries, blobs, onDiskBytes, err := emitResources(artifact, index, target) + jobs := options.Jobs + if jobs < 1 { + jobs = defaultEmitJobs + } + entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs) if err != nil { return EmitResult{}, err } @@ -176,11 +190,16 @@ func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name strin return []erf.IndexEntry{{Name: name, Type: restype, Offset: 0, Size: size}}, nil } -// emitResources hashes, compresses and stores one resource at a time, reading -// each payload from the artifact only when its turn comes. Peak memory -// therefore tracks the largest single resource, not the archive: a 2 GB hak -// must emit inside a runner's few spare GB. -func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, 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) ([]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)) @@ -208,30 +227,95 @@ func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink) ([ 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 onDiskBytes int64 - for _, identity := range order { + 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 { - return nil, 0, 0, err + mu.Lock() + if firstErr == nil { + firstErr = err + } + mu.Unlock() + return } sum := sha1.Sum(payload) - entries = append(entries, Entry{ + 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), func() []byte { return compressBlob(payload) }) + mu.Lock() + defer mu.Unlock() if err != nil { - return nil, 0, 0, err + 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 } diff --git a/internal/nwsync/jobs_test.go b/internal/nwsync/jobs_test.go new file mode 100644 index 0000000..19e6b1e --- /dev/null +++ b/internal/nwsync/jobs_test.go @@ -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) + } +} diff --git a/internal/nwsync/run.go b/internal/nwsync/run.go index 40dffd2..de2a88c 100644 --- a/internal/nwsync/run.go +++ b/internal/nwsync/run.go @@ -81,6 +81,7 @@ func runEmit(args []string, stdout, stderr io.Writer) int { fs.SetOutput(stderr) 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") positional, err := parseArgs(fs, args) if err != nil { return exitUsage @@ -89,12 +90,17 @@ func runEmit(args []string, stdout, stderr io.Writer) int { fmt.Fprintf(stderr, "nwsync emit: and are both required\n") 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{ ArtifactKey: positional[0], ArtifactPath: positional[1], As: *as, OutDir: *out, + Jobs: *jobs, }) if err != nil { fmt.Fprintf(stderr, "nwsync emit: %v\n", err) -- 2.54.0