Compare commits

..
2 Commits
Author SHA1 Message Date
archvillainetteandClaude Opus 5 ba671a76ed fix(nwsync): hash through a section reader, pin the encoder claim (#76)
ci / ci (pull_request) Successful in 3m30s
Review follow-ups. checkArtifactKey drained the *os.File to EOF, so the
file offset was left at the end; correct only because everything after it
uses ReadAt. Hash a section reader instead, so no later sequential read
can silently see nothing.

The zstd concurrency comment claimed byte-identical output but nothing
asserted it, so add the test. erf.Read now says outright that Data must
not be mutated, since payloads alias one buffer.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-31 13:08:35 +02:00
archvillainetteandClaude Opus 5 59b233e3db fix(nwsync): stream emit so peak memory tracks the largest resource (#76)
Emit read the whole artifact into memory and erf.Read then allocated a
second full copy of every payload, so a hak cost about 5x its size in
RAM. The 7 GB runner was OOM-killed on any hak over ~1.4 GB, which
blocks the NWSync backfill and every release that rebuilds a large hak.

Emit now opens the artifact, hashes it by streaming for the key check,
parses only the header and resource table via erf.ReadIndex, and reads,
hashes, compresses and stores one payload at a time. erf.Read keeps its
old shape but returns payloads as subslices instead of fresh copies,
which removes the second copy for the other callers too.

The zstd encoder pool also held one window-sized history per CPU — about
200 MB of live heap on a 24-core runner. EncodeAll is single-threaded
per call, so concurrency 1 gives byte-identical blobs for far less
memory.

Peak heap is now flat at ~22 MB for both an 8 MB and a 64 MB hak, and a
regression test asserts it does not scale with artifact size.

Closes #76

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-31 13:06:29 +02:00
4 changed files with 13 additions and 246 deletions
+1 -10
View File
@@ -48,19 +48,10 @@ 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] <artifact-key> <file>
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...
```
`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
+12 -96
View File
@@ -12,7 +12,6 @@ import (
"sort"
"strconv"
"strings"
"sync"
"time"
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo"
@@ -63,21 +62,12 @@ 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
}
@@ -125,11 +115,7 @@ func Emit(options EmitOptions) (EmitResult, error) {
return EmitResult{}, err
}
jobs := options.Jobs
if jobs < 1 {
jobs = defaultEmitJobs
}
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target, jobs)
entries, blobs, onDiskBytes, err := emitResources(artifact, index, target)
if err != nil {
return EmitResult{}, err
}
@@ -190,16 +176,11 @@ 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 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) {
// 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) {
// 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))
@@ -227,95 +208,30 @@ func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink, jo
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))
entries := make([]Entry, 0, 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]
for _, identity := range order {
payload, err := erf.ReadPayload(artifact, latest[identity])
if err != nil {
mu.Lock()
if firstErr == nil {
firstErr = err
}
mu.Unlock()
return
return nil, 0, 0, err
}
sum := sha1.Sum(payload)
entries[i] = Entry{
entries = append(entries, 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 {
if firstErr == nil {
firstErr = err
}
return
return nil, 0, 0, err
}
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
}
-134
View File
@@ -1,134 +0,0 @@
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)
}
}
-6
View File
@@ -81,7 +81,6 @@ 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
@@ -90,17 +89,12 @@ func runEmit(args []string, stdout, stderr io.Writer) int {
fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> 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)