Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1c2acc5530 |
+10
-1
@@ -48,10 +48,19 @@ beside the artifact itself with the extension replaced, so `emit` and
|
|||||||
`assemble` agree on where it is without being told.
|
`assemble` agree on where it is without being told.
|
||||||
|
|
||||||
```
|
```
|
||||||
nwsync emit [--as NAME] [--out DIR] <artifact-key> <file>
|
nwsync emit [--as NAME] [--out DIR] [--jobs N] <artifact-key> <file>
|
||||||
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...
|
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...
|
||||||
```
|
```
|
||||||
|
|
||||||
|
`emit` is latency-bound, not CPU-bound: every blob costs an existence probe
|
||||||
|
plus an upload, and a measured backfill spent 26 seconds of CPU across 9.5
|
||||||
|
minutes of wall clock. `--jobs N` (default 16) sets how many resources are in
|
||||||
|
flight at once. The manifest is byte-identical at any value — the number of
|
||||||
|
workers is never observable in the output. Peak memory is `N` times the
|
||||||
|
per-resource limit of 15 MB plus its compressed copy, so raising `N` far past
|
||||||
|
the default costs real memory for little gain: the transport keeps 16 idle
|
||||||
|
connections per host, and past that a worker pays a fresh TLS handshake.
|
||||||
|
|
||||||
Both verbs upload by default; nothing bulky is ever written to the runner's
|
Both verbs upload by default; nothing bulky is ever written to the runner's
|
||||||
disk. `--out DIR` writes a local repository tree instead, which is the
|
disk. `--out DIR` writes a local repository tree instead, which is the
|
||||||
conformance path against upstream `nwn_nwsync_write`. The zone comes from
|
conformance path against upstream `nwn_nwsync_write`. The zone comes from
|
||||||
|
|||||||
+96
-12
@@ -12,6 +12,7 @@ import (
|
|||||||
"sort"
|
"sort"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo"
|
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/buildinfo"
|
||||||
@@ -62,12 +63,21 @@ type EmitResult struct {
|
|||||||
BlobsWritten int
|
BlobsWritten int
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// defaultEmitJobs is how many resources are hashed, compressed and stored at
|
||||||
|
// once. Emit is latency-bound, not CPU-bound: a blob costs a probe round-trip
|
||||||
|
// plus an upload round-trip, and a measured backfill spent 26 s of CPU across
|
||||||
|
// 9.5 minutes of wall clock. The figure matches depot's DEPOT_JOBS default and
|
||||||
|
// the transport's MaxIdleConnsPerHost, so a worker per connection needs no new
|
||||||
|
// TLS handshake.
|
||||||
|
const defaultEmitJobs = 16
|
||||||
|
|
||||||
// EmitOptions describes one emit run.
|
// EmitOptions describes one emit run.
|
||||||
type EmitOptions struct {
|
type EmitOptions struct {
|
||||||
ArtifactKey string // depot key of the artifact; the NSYM key is derived from it
|
ArtifactKey string // depot key of the artifact; the NSYM key is derived from it
|
||||||
ArtifactPath string // the file on disk
|
ArtifactPath string // the file on disk
|
||||||
As string // name override, for a TLK whose filename is not its published name
|
As string // name override, for a TLK whose filename is not its published name
|
||||||
OutDir string // write locally instead of uploading — the conformance path
|
OutDir string // write locally instead of uploading — the conformance path
|
||||||
|
Jobs int // resources in flight at once; 0 means defaultEmitJobs
|
||||||
Sink sink // test seam; nil means OutDir or the zone
|
Sink sink // test seam; nil means OutDir or the zone
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -115,7 +125,11 @@ func Emit(options EmitOptions) (EmitResult, error) {
|
|||||||
return EmitResult{}, err
|
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 {
|
if err != nil {
|
||||||
return EmitResult{}, err
|
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
|
return []erf.IndexEntry{{Name: name, Type: restype, Offset: 0, Size: size}}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// emitResources hashes, compresses and stores one resource at a time, reading
|
// emitResources hashes, compresses and stores resources, reading each payload
|
||||||
// each payload from the artifact only when its turn comes. Peak memory
|
// from the artifact only when its turn comes. Peak memory tracks the resources
|
||||||
// therefore tracks the largest single resource, not the archive: a 2 GB hak
|
// in flight, not the archive: a 2 GB hak must emit inside a runner's few spare
|
||||||
// must emit inside a runner's few spare GB.
|
// GB. jobs of them are in flight at once, so the ceiling is jobs multiplied by
|
||||||
func emitResources(artifact io.ReaderAt, index []erf.IndexEntry, target sink) ([]Entry, int, int64, error) {
|
// 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,
|
// A resref appearing twice inside one artifact resolves to the last one,
|
||||||
// the way resman lets the last container added win.
|
// the way resman lets the last container added win.
|
||||||
order := make([]Identity, 0, len(index))
|
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 "))
|
return nil, 0, 0, fmt.Errorf("resources exceed the file size limit:\n %s", strings.Join(tooBig, "\n "))
|
||||||
}
|
}
|
||||||
|
|
||||||
entries := make([]Entry, 0, len(order))
|
// Index-addressed, never appended to: a worker owns entries[i] alone, so
|
||||||
|
// the slice comes back in artifact order and needs no lock.
|
||||||
|
entries := make([]Entry, len(order))
|
||||||
var blobs int
|
var blobs int
|
||||||
var onDiskBytes int64
|
var onDiskBytes int64
|
||||||
for _, identity := range order {
|
var mu sync.Mutex
|
||||||
|
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])
|
payload, err := erf.ReadPayload(artifact, latest[identity])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, 0, 0, err
|
mu.Lock()
|
||||||
|
if firstErr == nil {
|
||||||
|
firstErr = err
|
||||||
|
}
|
||||||
|
mu.Unlock()
|
||||||
|
return
|
||||||
}
|
}
|
||||||
sum := sha1.Sum(payload)
|
sum := sha1.Sum(payload)
|
||||||
entries = append(entries, Entry{
|
entries[i] = Entry{
|
||||||
SHA1: sum,
|
SHA1: sum,
|
||||||
Size: uint32(len(payload)),
|
Size: uint32(len(payload)),
|
||||||
ResRef: identity.ResRef,
|
ResRef: identity.ResRef,
|
||||||
ResType: identity.ResType,
|
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) })
|
written, err := target.putBlob(fmt.Sprintf("%x", sum), func() []byte { return compressBlob(payload) })
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, 0, 0, err
|
if firstErr == nil {
|
||||||
|
firstErr = err
|
||||||
|
}
|
||||||
|
return
|
||||||
}
|
}
|
||||||
if written > 0 {
|
if written > 0 {
|
||||||
blobs++
|
blobs++
|
||||||
onDiskBytes += written
|
onDiskBytes += written
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if jobs < 1 {
|
||||||
|
jobs = 1
|
||||||
|
}
|
||||||
|
work := make(chan int)
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for range jobs {
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for i := range work {
|
||||||
|
// After a failure the run is over — the caller discards
|
||||||
|
// everything and no index is written. Draining the rest of the
|
||||||
|
// channel cheaply, rather than returning, keeps the feeder from
|
||||||
|
// blocking on workers that have gone away.
|
||||||
|
if failed() {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
store(i)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
for i := range order {
|
||||||
|
work <- i
|
||||||
|
}
|
||||||
|
close(work)
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
if firstErr != nil {
|
||||||
|
return nil, 0, 0, firstErr
|
||||||
|
}
|
||||||
return entries, blobs, onDiskBytes, nil
|
return entries, blobs, onDiskBytes, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -81,6 +81,7 @@ func runEmit(args []string, stdout, stderr io.Writer) int {
|
|||||||
fs.SetOutput(stderr)
|
fs.SetOutput(stderr)
|
||||||
as := fs.String("as", "", "published name of the artifact, when it differs from the key")
|
as := fs.String("as", "", "published name of the artifact, when it differs from the key")
|
||||||
out := fs.String("out", "", "write to a local repository tree instead of uploading")
|
out := fs.String("out", "", "write to a local repository tree instead of uploading")
|
||||||
|
jobs := fs.Int("jobs", defaultEmitJobs, "resources to hash, compress and store at once")
|
||||||
positional, err := parseArgs(fs, args)
|
positional, err := parseArgs(fs, args)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return exitUsage
|
return exitUsage
|
||||||
@@ -89,12 +90,17 @@ func runEmit(args []string, stdout, stderr io.Writer) int {
|
|||||||
fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> are both required\n")
|
fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> are both required\n")
|
||||||
return exitUsage
|
return exitUsage
|
||||||
}
|
}
|
||||||
|
if *jobs < 1 {
|
||||||
|
fmt.Fprintf(stderr, "nwsync emit: -jobs must be at least 1, got %d\n", *jobs)
|
||||||
|
return exitUsage
|
||||||
|
}
|
||||||
|
|
||||||
result, err := Emit(EmitOptions{
|
result, err := Emit(EmitOptions{
|
||||||
ArtifactKey: positional[0],
|
ArtifactKey: positional[0],
|
||||||
ArtifactPath: positional[1],
|
ArtifactPath: positional[1],
|
||||||
As: *as,
|
As: *as,
|
||||||
OutDir: *out,
|
OutDir: *out,
|
||||||
|
Jobs: *jobs,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Fprintf(stderr, "nwsync emit: %v\n", err)
|
fmt.Fprintf(stderr, "nwsync emit: %v\n", err)
|
||||||
|
|||||||
Reference in New Issue
Block a user