Compare commits

..
2 Commits
Author SHA1 Message Date
archvillainette 978e2913f6 fix(nwsync): pin emitter version, ship the shim, harden manifest reads
ci / ci (pull_request) Successful in 3m28s
Review follow-ups on the same work:

- emitter_version is its own sidecar field, not the build revision. Keying
  the merge check on created_with invalidated every published index on every
  unrelated crucible commit, forcing a re-emit of the whole corpus.
- flake.nix builds cmd/crucible-nwsync, so Nix consumers get the binary the
  docs promise.
- readManifest uses io.ReadFull: bytes.Reader.Read can return a short count
  with no error, so a truncated manifest parsed into zero-padded garbage.
- emit refuses a .mod outright and names why, rather than failing on an
  unknown extension.
- SOURCE_DATE_EPOCH pins the sidecar timestamp; the manifest was already
  deterministic.
- Entry.identity replaces the resref/restype key struct duplicated in emit
  and assemble.

Refs #53.
2026-07-29 12:17:08 +02:00
archvillainette 0de1c1e280 feat(nwsync): emit blobs and per-artifact NSYM, assemble merged manifests
Adds `crucible nwsync emit` and `crucible nwsync assemble`, the split
replacement for upstream nwn_nwsync_write, which cannot be used because it
wants every hak, the TLK and the module present in one run on one disk.

emit explodes one artifact — a .hak/.erf or a loose file such as the TLK —
into NWSync blobs (NWCompressedBuffer framing, magic NSYC, zstd) plus a NSYM
v3 manifest describing only that artifact, with the same .json sidecar
upstream writes. Blob names are the sha1 of the uncompressed bytes, restypes
nss/ndb/gic are always skipped, an unresolvable restype is a hard error, a
resource over 15 MB fails closed, and no latest or .origin file is ever
written.

assemble merges the per-artifact manifests into one, reading no bulk data.
The merge rule is resref shadowing, not concatenation: a resref present in
more than one artifact resolves to the earliest artifact in --order, the way
the game resolves it. It refuses to merge across mismatched emitter versions,
and takes --group-id from the caller (1 current, 2 testing; 0 is absent).

Refs #53.
2026-07-29 12:12:22 +02:00
14 changed files with 290 additions and 1437 deletions
+4 -64
View File
@@ -43,70 +43,10 @@ flags are mutually exclusive.
`nwsync emit` runs where an artifact is born (a `.hak`/`.erf`, or a loose file `nwsync emit` runs where an artifact is born (a `.hak`/`.erf`, or a loose file
such as the TLK); `nwsync assemble` runs at module release and reads only the such as the TLK); `nwsync assemble` runs at module release and reads only the
small per-artifact indexes. Both take **depot keys**: an artifact's index lives small per-artifact manifests. `--order` lists artifact names highest priority
beside the artifact itself with the extension replaced, so `emit` and first: a resref in more than one artifact resolves to the earliest one, the way
`assemble` agree on where it is without being told. the game resolves it. `--group-id` is per channel — 1 is current, 2 is testing,
and 0 leaves the field out of the sidecar.
```
nwsync emit [--as NAME] [--out DIR] [--jobs N] <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
`NWSYNC_STORAGE_ZONE` and `NWSYNC_STORAGE_PASSWORD`, with the host from
`BUNNY_STORAGE_HOST` — NWSync data is a separate zone from the asset depot.
`assemble`'s artifact keys are in `Mod_HakList` order, highest priority first: a
resref in more than one artifact resolves to the earliest one, the way the game
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.
`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
whatever already landed, and `assemble` fails closed on an artifact with no
index rather than publishing a manifest that is missing a hak.
### What an emitted tree looks like
`--out DIR` produces the same tree `emit` would upload, which makes it the way
to check a zone by hand without touching one:
```
<artifact-sha>.nsym binary index
<artifact-sha>.nsym.json the same index, readable
data/sha1/a7/4a/a74aa84a... one blob per resource, two-level fanout
```
A blob's name is the SHA-1 of the resource's **original** bytes, but the file on
disk is not those bytes: each blob is wrapped in NWCompressedBuffer framing, a
24-byte `NSYC` header followed by a zstd frame. Hashing the file directly will
not match its name, which is the obvious
first thing to try and the obvious first thing to be confused by. Strip the
header first:
```
tail -c +25 <blob> | zstd -dc | sha1sum # == the blob's filename
```
The header carries the uncompressed length as a little-endian `uint32` at offset
12, so the decompressed size is checkable without decompressing. Compression is
worth roughly a 4:1
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`.
## Hidden compatibility aliases ## Hidden compatibility aliases
+1 -1
View File
@@ -385,7 +385,7 @@ func refreshBuildModuleManifest(ctx context, p *project.Project, progress func(s
} }
progress("Refreshing hak list from the latest published sow-assets manifest...") progress("Refreshing hak list from the latest published sow-assets manifest...")
if err := runProjectScript(ctx, p, []string{"scripts", "fetch-upstream-manifests"}, manifestPath); err != nil { if err := runProjectScript(ctx, p, []string{"scripts", "fetch-hak-manifest"}, manifestPath); err != nil {
return "", "", err return "", "", err
} }
if _, err := pipeline.ApplyHAKManifest(p, manifestPath); err != nil { if _, err := pipeline.ApplyHAKManifest(p, manifestPath); err != nil {
-121
View File
@@ -1,121 +0,0 @@
package depot
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"os"
"strings"
)
// KeyStore is the zone addressed by object key rather than by depot sha. The
// depot names every object after the sha256 of its contents; NWSync does not —
// a blob is named after the sha1 of its *uncompressed* bytes while the body
// uploaded is the compressed form, and a per-artifact index is named after its
// artifact. Both addressing modes want the same transport, retry and probe
// discipline, so the sha-addressed Backend rides on this rather than the other
// way round.
type KeyStore interface {
// ProbeKey returns the existence state of one key. transient=true means a
// retry might change the answer — never read it as "missing, re-upload".
ProbeKey(ctx context.Context, key string) (state ProbeState, transient bool, err error)
// PutReader uploads size bytes read from r to key. checksum is the
// uppercase hex sha256 of those bytes, which Bunny verifies server-side.
PutReader(ctx context.Context, key string, r io.Reader, size int64, checksum string) error
// GetKey fetches the whole object at key. Small objects only — it holds
// the body in memory and does no hash check, because a key is not always
// a content hash.
GetKey(ctx context.Context, key string) ([]byte, error)
}
// NewKeyStore returns a KeyStore for cfg's storage zone. Fails closed on a
// missing host or read key, matching NewBackend.
func NewKeyStore(cfg Config) (KeyStore, error) {
if cfg.StorageHost == "" {
return nil, errors.New("storage backend requires a storage host")
}
if cfg.StorageZone == "" {
return nil, errors.New("storage backend requires a storage zone")
}
if cfg.ReadKey == "" {
return nil, errors.New("storage backend requires a read key")
}
return &httpBackend{name: "bunny", client: newHTTPClient(cfg), cfg: cfg}, nil
}
// keyURL is the storage URL of one object key.
func (b *httpBackend) keyURL(key string) string {
host := b.cfg.StorageHost
// StorageHost is normally a bare host ("storage.bunnycdn.com"); allow a
// full scheme (used by tests against httptest.NewServer) to pass through
// unchanged.
if !strings.Contains(host, "://") {
host = "https://" + host
}
return fmt.Sprintf("%s/%s/%s", strings.TrimSuffix(host, "/"), b.cfg.StorageZone, key)
}
func (b *httpBackend) ProbeKey(ctx context.Context, key string) (ProbeState, bool, error) {
return b.rangeProbe(ctx, b.keyURL(key), map[string]string{"AccessKey": b.cfg.ReadKey})
}
func (b *httpBackend) PutReader(ctx context.Context, key string, r io.Reader, size int64, checksum string) error {
if b.name == "cdn" {
return errors.New("cdn backend is read-only")
}
if b.cfg.WriteKey == "" {
return errors.New("storage backend requires a write key to write")
}
req, err := http.NewRequestWithContext(ctx, http.MethodPut, b.keyURL(key), r)
if err != nil {
return err
}
req.ContentLength = size
req.Header.Set("AccessKey", b.cfg.WriteKey)
// Bunny defines Checksum as sha256 of the body and rejects a mismatch, so
// this is server-side integrity checking, not decoration.
req.Header.Set("Checksum", strings.ToUpper(checksum))
resp, err := b.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
_, _ = io.Copy(io.Discard, resp.Body)
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("put %s: unexpected status %d", key, resp.StatusCode)
}
return nil
}
func (b *httpBackend) GetKey(ctx context.Context, key string) ([]byte, error) {
resp, err := b.get(ctx, b.keyURL(key), map[string]string{"AccessKey": b.cfg.ReadKey})
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
_, _ = io.Copy(io.Discard, resp.Body)
return nil, fmt.Errorf("get %s: unexpected status %d", key, resp.StatusCode)
}
return io.ReadAll(resp.Body)
}
// putFile uploads the file at src to key, streaming it. checksum is the
// uppercase hex sha256 of the file's bytes.
func (b *httpBackend) putFile(ctx context.Context, key, src, checksum string) error {
f, err := os.Open(src)
if err != nil {
return err
}
defer f.Close()
info, err := f.Stat()
if err != nil {
return err
}
return b.PutReader(ctx, key, f, info.Size(), checksum)
}
+36 -4
View File
@@ -10,6 +10,7 @@ import (
"net/http" "net/http"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"time" "time"
) )
@@ -70,7 +71,16 @@ type httpBackend struct {
func (b *httpBackend) Name() string { return b.name } func (b *httpBackend) Name() string { return b.name }
func (b *httpBackend) storageURL(sha string) string { return b.keyURL(BlobKey(sha)) } func (b *httpBackend) storageURL(sha string) string {
host := b.cfg.StorageHost
// StorageHost is normally a bare host ("storage.bunnycdn.com"); allow a
// full scheme (used by tests against httptest.NewServer) to pass through
// unchanged.
if strings.Contains(host, "://") {
return fmt.Sprintf("%s/%s/%s", strings.TrimSuffix(host, "/"), b.cfg.StorageZone, BlobKey(sha))
}
return fmt.Sprintf("https://%s/%s/%s", host, b.cfg.StorageZone, BlobKey(sha))
}
func (b *httpBackend) cdnURL(sha string) string { func (b *httpBackend) cdnURL(sha string) string {
return fmt.Sprintf("%s/%s", b.cfg.CDNBase, BlobKey(sha)) return fmt.Sprintf("%s/%s", b.cfg.CDNBase, BlobKey(sha))
@@ -130,8 +140,7 @@ func (b *httpBackend) Probe(ctx context.Context, sha string) (ProbeState, bool,
return b.rangeProbe(ctx, b.storageURL(sha), map[string]string{"AccessKey": b.cfg.ReadKey}) return b.rangeProbe(ctx, b.storageURL(sha), map[string]string{"AccessKey": b.cfg.ReadKey})
} }
// Put uploads src for sha. cdn is read-only. A depot object is named after the // Put uploads src for sha. cdn is read-only.
// sha256 of its own bytes, so the key's sha doubles as the Checksum header.
func (b *httpBackend) Put(ctx context.Context, sha, src string) error { func (b *httpBackend) Put(ctx context.Context, sha, src string) error {
if b.name == "cdn" { if b.name == "cdn" {
return errors.New("cdn backend is read-only") return errors.New("cdn backend is read-only")
@@ -148,7 +157,30 @@ func (b *httpBackend) Put(ctx context.Context, sha, src string) error {
return nil return nil
} }
return b.putFile(ctx, BlobKey(sha), src, sha) f, err := os.Open(src)
if err != nil {
return err
}
defer f.Close()
req, err := http.NewRequestWithContext(ctx, http.MethodPut, b.storageURL(sha), f)
if err != nil {
return err
}
req.Header.Set("AccessKey", b.cfg.WriteKey)
req.Header.Set("Checksum", strings.ToUpper(sha))
resp, err := b.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
_, _ = io.Copy(io.Discard, resp.Body)
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("bunny put %s: unexpected status %d", sha, resp.StatusCode)
}
return nil
} }
// Get fetches sha into dest via temp file + rename, re-hashing and deleting // Get fetches sha into dest via temp file + rename, re-hashing and deleting
+45 -92
View File
@@ -311,110 +311,63 @@ func Write(w io.Writer, archive Archive) error {
return nil return nil
} }
// IndexEntry locates one resource inside an archive without holding its
// payload. Streaming callers read one payload at a time from these, so peak
// memory tracks the largest resource instead of the whole archive.
type IndexEntry struct {
Name string
Type uint16
Offset int64
Size int64
}
// Index is the header plus the resource table of an ERF: everything except the
// payloads.
type Index struct {
FileType string
Version string
Entries []IndexEntry
}
// ReadIndex parses the tables of an ERF of the given size, reading only the
// header, the key list and the resource list.
func ReadIndex(r io.ReaderAt, size int64) (Index, error) {
if size < headerSize {
return Index{}, fmt.Errorf("erf file too small: %d bytes", size)
}
var hdr header
if err := binary.Read(io.NewSectionReader(r, 0, headerSize), binary.LittleEndian, &hdr); err != nil {
return Index{}, fmt.Errorf("decode erf header: %w", err)
}
if int64(hdr.KeyListOffset)+int64(hdr.EntryCount)*24 > size {
return Index{}, fmt.Errorf("erf key list exceeds file bounds")
}
keys := make([]keyEntry, hdr.EntryCount)
keyReader := io.NewSectionReader(r, int64(hdr.KeyListOffset), int64(hdr.EntryCount)*24)
if err := binary.Read(keyReader, binary.LittleEndian, &keys); err != nil {
return Index{}, fmt.Errorf("decode key list: %w", err)
}
if int64(hdr.ResourceListOffset)+int64(hdr.EntryCount)*8 > size {
return Index{}, fmt.Errorf("erf resource list exceeds file bounds")
}
entries := make([]resourceEntry, hdr.EntryCount)
entryReader := io.NewSectionReader(r, int64(hdr.ResourceListOffset), int64(hdr.EntryCount)*8)
if err := binary.Read(entryReader, binary.LittleEndian, &entries); err != nil {
return Index{}, fmt.Errorf("decode resource list: %w", err)
}
index := Index{
FileType: string(hdr.FileType[:]),
Version: string(hdr.Version[:]),
Entries: make([]IndexEntry, 0, hdr.EntryCount),
}
for position, key := range keys {
entry := entries[position]
if int64(entry.Offset)+int64(entry.Size) > size {
return Index{}, fmt.Errorf("resource %d exceeds file bounds", position)
}
index.Entries = append(index.Entries, IndexEntry{
Name: string(bytes.TrimRight(key.ResRef[:], "\x00")),
Type: key.ResourceType,
Offset: int64(entry.Offset),
Size: int64(entry.Size),
})
}
return index, nil
}
// ReadPayload returns one resource's bytes.
func ReadPayload(r io.ReaderAt, entry IndexEntry) ([]byte, error) {
payload := make([]byte, entry.Size)
if _, err := r.ReadAt(payload, entry.Offset); err != nil {
return nil, fmt.Errorf("read resource %q: %w", entry.Name, err)
}
return payload, nil
}
// Read materialises a whole archive. Payloads are subslices of the buffer the
// archive was read into, so nothing is copied twice: a caller must not mutate
// Data. Callers that only need one resource at a time should use ReadIndex
// instead, which never holds the archive at all.
func Read(r io.Reader) (Archive, error) { func Read(r io.Reader) (Archive, error) {
data, err := io.ReadAll(r) data, err := io.ReadAll(r)
if err != nil { if err != nil {
return Archive{}, fmt.Errorf("read erf: %w", err) return Archive{}, fmt.Errorf("read erf: %w", err)
} }
index, err := ReadIndex(bytes.NewReader(data), int64(len(data))) if len(data) < headerSize {
if err != nil { return Archive{}, fmt.Errorf("erf file too small: %d bytes", len(data))
return Archive{}, err
} }
resources := make([]Resource, 0, len(index.Entries)) var hdr header
for _, entry := range index.Entries { if err := binary.Read(bytes.NewReader(data[:headerSize]), binary.LittleEndian, &hdr); err != nil {
return Archive{}, fmt.Errorf("decode erf header: %w", err)
}
keyStart := int(hdr.KeyListOffset)
keyEnd := keyStart + int(hdr.EntryCount)*24
if keyEnd > len(data) {
return Archive{}, fmt.Errorf("erf key list exceeds file bounds")
}
keys := make([]keyEntry, hdr.EntryCount)
if err := binary.Read(bytes.NewReader(data[keyStart:keyEnd]), binary.LittleEndian, &keys); err != nil {
return Archive{}, fmt.Errorf("decode key list: %w", err)
}
resourceStart := int(hdr.ResourceListOffset)
resourceEnd := resourceStart + int(hdr.EntryCount)*8
if resourceEnd > len(data) {
return Archive{}, fmt.Errorf("erf resource list exceeds file bounds")
}
entries := make([]resourceEntry, hdr.EntryCount)
if err := binary.Read(bytes.NewReader(data[resourceStart:resourceEnd]), binary.LittleEndian, &entries); err != nil {
return Archive{}, fmt.Errorf("decode resource list: %w", err)
}
resources := make([]Resource, 0, hdr.EntryCount)
for index, key := range keys {
entry := entries[index]
start := int(entry.Offset)
end := start + int(entry.Size)
if end > len(data) {
return Archive{}, fmt.Errorf("resource %d exceeds file bounds", index)
}
resref := string(bytes.TrimRight(key.ResRef[:], "\x00"))
payload := make([]byte, entry.Size)
copy(payload, data[start:end])
resources = append(resources, Resource{ resources = append(resources, Resource{
Name: entry.Name, Name: resref,
Type: entry.Type, Type: key.ResourceType,
Data: data[entry.Offset : entry.Offset+entry.Size], Data: payload,
Size: entry.Size, Size: int64(entry.Size),
}) })
} }
return Archive{ return Archive{
FileType: index.FileType, FileType: string(hdr.FileType[:]),
Version: index.Version, Version: string(hdr.Version[:]),
Resources: resources, Resources: resources,
}, nil }, nil
} }
+26 -43
View File
@@ -4,21 +4,18 @@ import (
"crypto/sha1" "crypto/sha1"
"encoding/json" "encoding/json"
"fmt" "fmt"
"path" "os"
"path/filepath"
) )
// AssembleOptions describes one merged manifest. // AssembleOptions describes one merged manifest.
type AssembleOptions struct { type AssembleOptions struct {
// ArtifactKeys are the depot keys of the artifacts to merge, in Order []string // artifact names, highest priority first
// Mod_HakList order — highest priority first. Each one's index is read EntriesDir string // directory holding <name>.nsym and <name>.nsym.json
// from the key beside it. OutDir string // repository root; the manifest lands in <out>/manifests
ArtifactKeys []string GroupID int // 1 = current, 2 = testing; 0 means absent
TLKKey string // the TLK's key, if the manifest carries one ModuleName string
OutDir string // write locally instead of uploading — the conformance path Description string
GroupID int // 1 = current, 2 = testing; 0 means absent
ModuleName string
Description string
Sink sink // test seam; nil means OutDir or the zone
} }
// AssembleResult reports what one assemble run produced. // AssembleResult reports what one assemble run produced.
@@ -36,43 +33,25 @@ type AssembleResult struct {
// the game resolves it (upstream's resman adds haks in reverse and lets the // the game resolves it (upstream's resman adds haks in reverse and lets the
// last one win). Get this backwards and the wrong texture ships silently. // last one win). Get this backwards and the wrong texture ships silently.
func Assemble(options AssembleOptions) (AssembleResult, error) { func Assemble(options AssembleOptions) (AssembleResult, error) {
if len(options.ArtifactKeys) == 0 { if len(options.Order) == 0 {
return AssembleResult{}, fmt.Errorf("assemble: no artifact keys given") return AssembleResult{}, fmt.Errorf("assemble: --order names no artifacts")
}
target, err := openSink(options.OutDir, options.Sink)
if err != nil {
return AssembleResult{}, err
}
// The TLK carries no precedence — it is not a hak and shadows nothing —
// so it merges last, after every hak has had its say.
keys := append([]string{}, options.ArtifactKeys...)
if options.TLKKey != "" {
keys = append(keys, options.TLKKey)
} }
merged := make([]Entry, 0, 1024) merged := make([]Entry, 0, 1024)
winner := make(map[Identity]bool, 1024) winner := make(map[Identity]bool, 1024)
var onDiskBytes int64 var onDiskBytes int64
for _, artifactKey := range keys { for _, name := range options.Order {
key, err := resolveIndexKey(artifactKey, options.OutDir) manifestPath := filepath.Join(options.EntriesDir, name+".nsym")
data, err := os.ReadFile(manifestPath)
if err != nil { if err != nil {
return AssembleResult{}, err return AssembleResult{}, fmt.Errorf("assemble: no index for %q: %w", name, err)
}
data, sidecarBody, err := target.getIndex(key)
if err != nil {
// An artifact with no index is an artifact whose emit never
// finished. Publishing a manifest without it would ship a release
// missing a hak, so this fails closed.
return AssembleResult{}, fmt.Errorf("assemble: no index for %s: %w", artifactKey, err)
} }
entries, err := readManifest(data) entries, err := readManifest(data)
if err != nil { if err != nil {
return AssembleResult{}, fmt.Errorf("%s: %w", target.describe(key), err) return AssembleResult{}, fmt.Errorf("%s: %w", manifestPath, err)
} }
sidecar, err := parseSidecar(target.describe(key), sidecarBody) sidecar, err := readSidecar(manifestPath + ".json")
if err != nil { if err != nil {
return AssembleResult{}, err return AssembleResult{}, err
} }
@@ -82,7 +61,7 @@ func Assemble(options AssembleOptions) (AssembleResult, error) {
if sidecar.EmitterVersion != emitterVersion { if sidecar.EmitterVersion != emitterVersion {
return AssembleResult{}, fmt.Errorf( return AssembleResult{}, fmt.Errorf(
"assemble: emitter version mismatch: %s was emitted by emitter %q, this is emitter %q", "assemble: emitter version mismatch: %s was emitted by emitter %q, this is emitter %q",
artifactKey, sidecar.EmitterVersion, emitterVersion) name, sidecar.EmitterVersion, emitterVersion)
} }
// on_disk_bytes overcounts by the handful of cross-artifact // on_disk_bytes overcounts by the handful of cross-artifact
// duplicates. It is a display statistic; no dedupe pass for it. // duplicates. It is a display statistic; no dedupe pass for it.
@@ -107,22 +86,26 @@ func Assemble(options AssembleOptions) (AssembleResult, error) {
return AssembleResult{}, err return AssembleResult{}, err
} }
sha1Hex := fmt.Sprintf("%x", sha1.Sum(data)) sha1Hex := fmt.Sprintf("%x", sha1.Sum(data))
manifestKey := path.Join("manifests", sha1Hex) manifestPath := filepath.Join(options.OutDir, "manifests", sha1Hex)
sidecar := Sidecar{ sidecar := Sidecar{
ModuleName: options.ModuleName, ModuleName: options.ModuleName,
Description: options.Description, Description: options.Description,
GroupID: options.GroupID, GroupID: options.GroupID,
} }
if err := putManifestPair(target, manifestKey, data, merged, onDiskBytes, sidecar); err != nil { if err := writeManifestPair(manifestPath, data, merged, onDiskBytes, sidecar); err != nil {
return AssembleResult{}, err return AssembleResult{}, err
} }
return AssembleResult{SHA1: sha1Hex, ManifestPath: target.describe(manifestKey), Entries: len(merged)}, nil return AssembleResult{SHA1: sha1Hex, ManifestPath: manifestPath, Entries: len(merged)}, nil
} }
func parseSidecar(where string, data []byte) (Sidecar, error) { func readSidecar(path string) (Sidecar, error) {
data, err := os.ReadFile(path)
if err != nil {
return Sidecar{}, fmt.Errorf("assemble: missing sidecar: %w", err)
}
var sidecar Sidecar var sidecar Sidecar
if err := json.Unmarshal(data, &sidecar); err != nil { if err := json.Unmarshal(data, &sidecar); err != nil {
return Sidecar{}, fmt.Errorf("%s: %w", where, err) return Sidecar{}, fmt.Errorf("%s: %w", path, err)
} }
return sidecar, nil return sidecar, nil
} }
+2 -6
View File
@@ -22,13 +22,9 @@ const (
blobHeaderBytes = 24 blobHeaderBytes = 24
) )
// EncodeAll/DecodeAll are single-threaded per call, so the default pool of one
// encoder per CPU only buys idle memory: each holds a window-sized history, so
// on a 24-core runner that is ~200 MB of live heap doing nothing. Concurrency 1
// produces byte-identical output.
var ( var (
blobEncoder, _ = zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) blobEncoder, _ = zstd.NewWriter(nil)
blobDecoder, _ = zstd.NewReader(nil, zstd.WithDecoderConcurrency(1)) blobDecoder, _ = zstd.NewReader(nil)
) )
// compressBlob wraps data in NWCompressedBuffer framing. // compressBlob wraps data in NWCompressedBuffer framing.
+91 -206
View File
@@ -1,18 +1,15 @@
package nwsync package nwsync
import ( import (
"context" "bytes"
"crypto/sha1" "crypto/sha1"
"fmt" "fmt"
"io"
"os" "os"
"path"
"path/filepath" "path/filepath"
"slices" "slices"
"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"
@@ -63,259 +60,119 @@ type EmitResult struct {
BlobsWritten int BlobsWritten int
} }
// defaultEmitJobs is how many resources are hashed, compressed and stored at // Emit explodes one artifact — a .hak/.erf/.mod or a loose file such as the
// once. Emit is latency-bound, not CPU-bound: a blob costs a probe round-trip // TLK — into NWSync blobs plus a NSYM manifest describing only that artifact.
// plus an upload round-trip, and a measured backfill spent 26 s of CPU across func Emit(artifactPath, outDir string) (EmitResult, error) {
// 9.5 minutes of wall clock. The figure matches depot's DEPOT_JOBS default and resources, err := readArtifact(artifactPath)
// 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
}
// 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 { if err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
name := strings.TrimSuffix(filepath.Base(artifactPath), filepath.Ext(artifactPath))
index, err := readArtifactIndex(options.ArtifactPath, artifact, info.Size(), name) entries, blobs, onDiskBytes, err := emitResources(resources, outDir)
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)
if err != nil { if err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
if len(entries) == 0 { if len(entries) == 0 {
return EmitResult{}, fmt.Errorf("%s: nothing to index (no publishable resources)", options.ArtifactPath) return EmitResult{}, fmt.Errorf("%s: nothing to index (no publishable resources)", artifactPath)
} }
data, err := writeManifest(entries) data, err := writeManifest(entries)
if err != nil { if err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
if err := putManifestPair(target, key, data, entries, onDiskBytes, Sidecar{ModuleName: name}); err != nil { manifestPath := filepath.Join(outDir, name+".nsym")
if err := writeManifestPair(manifestPath, data, entries, onDiskBytes, Sidecar{ModuleName: name}); err != nil {
return EmitResult{}, err return EmitResult{}, err
} }
return EmitResult{Name: name, ManifestPath: target.describe(key), Entries: len(entries), BlobsWritten: blobs}, nil return EmitResult{Name: name, ManifestPath: manifestPath, Entries: len(entries), BlobsWritten: blobs}, nil
} }
// openSink returns the zone sink, or a local directory when outDir is set. // readArtifact returns the resources of an ERF/HAK/MOD, or the single resource
func openSink(outDir string, injected sink) (sink, error) { // a loose file represents. Upstream's resman does the same dispatch on the
if injected != nil { // file's first three bytes.
return injected, nil func readArtifact(path string) ([]erf.Resource, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("read artifact: %w", err)
} }
if outDir != "" { if len(data) >= 3 {
return dirSink{root: outDir}, nil switch string(data[:3]) {
} case "ERF", "HAK":
return newZoneSink(context.Background(), os.Getenv) archive, err := erf.Read(bytes.NewReader(data))
} if err != nil {
return nil, fmt.Errorf("%s: %w", path, err)
// readArtifactIndex locates the resources of an ERF/HAK/MOD, or the single }
// resource a loose file represents, without reading any payload. Upstream's return archive.Resources, nil
// resman does the same dispatch on the file's first three bytes. name is the case "MOD":
// artifact's published name, which for a loose file is also its resref. // A persistent world never publishes module contents, so the .mod
func readArtifactIndex(path string, artifact io.ReaderAt, size int64, name string) ([]erf.IndexEntry, error) { // contributes no bytes to a manifest — it only says which haks and
magic := make([]byte, 3) // which TLK the manifest covers.
if size >= 3 { return nil, fmt.Errorf("%s: a module is never emitted; a manifest is haks plus the TLK", path)
if _, err := artifact.ReadAt(magic, 0); err != nil {
return nil, fmt.Errorf("%s: %w", path, err)
} }
} }
switch string(magic) { base := filepath.Base(path)
case "ERF", "HAK": extension := filepath.Ext(base)
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) restype, ok := erf.ResourceTypeForExtension(extension)
if !ok { if !ok {
return nil, fmt.Errorf("%s: unknown resource type %q", path, extension) return nil, fmt.Errorf("%s: unknown resource type %q", path, extension)
} }
return []erf.IndexEntry{{Name: name, Type: restype, Offset: 0, Size: size}}, nil return []erf.Resource{{
Name: strings.TrimSuffix(base, extension),
Type: restype,
Data: data,
}}, nil
} }
// emitResources hashes, compresses and stores resources, reading each payload func emitResources(resources []erf.Resource, outDir string) ([]Entry, int, int64, error) {
// 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, // 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(resources))
latest := make(map[Identity]erf.IndexEntry, len(index)) latest := make(map[Identity]erf.Resource, len(resources))
var tooBig []string var tooBig []string
for _, entry := range index { for _, resource := range resources {
if _, ok := erf.ExtensionForResourceType(entry.Type); !ok { if _, ok := erf.ExtensionForResourceType(resource.Type); !ok {
return nil, 0, 0, fmt.Errorf("resref %s is not resolvable (unknown restype %d)", entry.Name, entry.Type) return nil, 0, 0, fmt.Errorf("resref %s is not resolvable (unknown restype %d)", resource.Name, resource.Type)
} }
if slices.Contains(skippedTypes, entry.Type) { if slices.Contains(skippedTypes, resource.Type) {
continue continue
} }
if entry.Size > fileSizeLimit { if len(resource.Data) > fileSizeLimit {
tooBig = append(tooBig, fmt.Sprintf("%s: %d bytes > %d", entry.Name, entry.Size, fileSizeLimit)) tooBig = append(tooBig, fmt.Sprintf("%s: %d bytes > %d", resource.Name, len(resource.Data), fileSizeLimit))
continue continue
} }
identity := Identity{ResRef: strings.ToLower(entry.Name), ResType: entry.Type} identity := Identity{ResRef: strings.ToLower(resource.Name), ResType: resource.Type}
if _, seen := latest[identity]; !seen { if _, seen := latest[identity]; !seen {
order = append(order, identity) order = append(order, identity)
} }
latest[identity] = entry latest[identity] = resource
} }
if len(tooBig) > 0 { if len(tooBig) > 0 {
sort.Strings(tooBig) sort.Strings(tooBig)
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 "))
} }
// Index-addressed, never appended to: a worker owns entries[i] alone, so entries := make([]Entry, 0, len(order))
// 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
var mu sync.Mutex for _, identity := range order {
var firstErr error resource := latest[identity]
// Two resrefs in one artifact can hold identical bytes, and therefore one sum := sha1.Sum(resource.Data)
// blob. Serially the sink's existence check absorbed that; in parallel both entries = append(entries, Entry{
// 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, SHA1: sum,
Size: uint32(len(payload)), Size: uint32(len(resource.Data)),
ResRef: identity.ResRef, ResRef: identity.ResRef,
ResType: identity.ResType, ResType: identity.ResType,
} })
mu.Lock() written, err := writeBlob(outDir, sum, resource.Data)
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 err != nil {
if firstErr == nil { return nil, 0, 0, err
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
} }
@@ -330,10 +187,35 @@ func created() int64 {
return time.Now().Unix() return time.Now().Unix()
} }
// putManifestPair stores a NSYM manifest and its .json sidecar at key. The // writeBlob writes one NWCompressedBuffer blob and returns its size on disk,
// caller supplies the sidecar fields it knows; the rest are derived from the // or 0 if the blob already existed. Blob names are content hashes, so an
// entries. data must be the serialised form of entries. // existing name is existing content.
func putManifestPair(target sink, key string, data []byte, entries []Entry, onDiskBytes int64, sidecar Sidecar) error { func writeBlob(outDir string, sum [20]byte, data []byte) (int64, error) {
path := blobPath(outDir, fmt.Sprintf("%x", sum))
if _, err := os.Stat(path); err == nil {
return 0, nil
}
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return 0, fmt.Errorf("create blob directory: %w", err)
}
blob := compressBlob(data)
if err := os.WriteFile(path, blob, 0o644); err != nil {
return 0, fmt.Errorf("write blob: %w", err)
}
return int64(len(blob)), nil
}
// writeManifestPair writes a NSYM manifest and its .json sidecar. The caller
// supplies the sidecar fields it knows; the rest are derived from the entries.
// data must be the serialised form of entries.
func writeManifestPair(manifestPath string, data []byte, entries []Entry, onDiskBytes int64, sidecar Sidecar) error {
if err := os.MkdirAll(filepath.Dir(manifestPath), 0o755); err != nil {
return fmt.Errorf("create manifest directory: %w", err)
}
if err := os.WriteFile(manifestPath, data, 0o644); err != nil {
return fmt.Errorf("write manifest: %w", err)
}
var totalBytes int64 var totalBytes int64
clientContents := false clientContents := false
for _, entry := range entries { for _, entry := range entries {
@@ -358,5 +240,8 @@ func putManifestPair(target sink, key string, data []byte, entries []Entry, onDi
if err != nil { if err != nil {
return err return err
} }
return target.putIndex(key, data, body) if err := os.WriteFile(manifestPath+".json", body, 0o644); err != nil {
return fmt.Errorf("write sidecar: %w", err)
}
return 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)
}
}
-150
View File
@@ -1,150 +0,0 @@
package nwsync
import (
"bytes"
"fmt"
"math/rand"
"os"
"path/filepath"
"runtime"
"runtime/debug"
"testing"
"time"
"github.com/klauspost/compress/zstd"
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/erf"
)
// TestSingleThreadedEncoderMatchesDefault pins the claim the blob encoder's
// concurrency setting rests on: it saves memory only, and a published blob is
// the same bytes either way.
func TestSingleThreadedEncoderMatchesDefault(t *testing.T) {
standard, err := zstd.NewWriter(nil)
if err != nil {
t.Fatal(err)
}
defer standard.Close()
body := make([]byte, 4<<20)
random := rand.New(rand.NewSource(1))
random.Read(body[:len(body)/2])
for _, size := range []int{0, 1, 4 << 10, len(body)} {
if !bytes.Equal(blobEncoder.EncodeAll(body[:size], nil), standard.EncodeAll(body[:size], nil)) {
t.Fatalf("%d bytes compress differently at concurrency 1", size)
}
}
}
// resourceSize is one payload in the memory fixtures. Real haks hold a few MB
// per resource, and peak memory is meant to track that, not the archive.
const resourceSize = 1 << 20
// writeStreamedHak builds a hak of count resources without ever holding the
// archive in memory, so the fixture itself does not decide the measurement.
// Payloads are distinct, so no blob is deduplicated away.
func writeStreamedHak(t *testing.T, path string, count int) {
t.Helper()
payload := filepath.Join(t.TempDir(), "payload.bin")
body := make([]byte, resourceSize)
for index := range body {
body[index] = byte(index)
}
resources := make([]erf.Resource, 0, count)
for index := range count {
// A distinct first byte per resource is enough to give every payload
// its own sha1 while still streaming from one file per resource.
unique := filepath.Join(filepath.Dir(payload), fmt.Sprintf("p%d.bin", index))
body[0] = byte(index)
body[1] = byte(index >> 8)
if err := os.WriteFile(unique, body, 0o644); err != nil {
t.Fatal(err)
}
resources = append(resources, erf.Resource{
Name: fmt.Sprintf("res%05d", index),
Type: restype(t, "tga"),
SourcePath: unique,
Size: resourceSize,
})
}
file, err := os.Create(path)
if err != nil {
t.Fatal(err)
}
defer file.Close()
if err := erf.Write(file, erf.New("HAK", resources)); err != nil {
t.Fatalf("write hak: %v", err)
}
}
// peakHeapDuring runs work while sampling the heap, and returns the largest
// live heap it saw.
func peakHeapDuring(work func()) uint64 {
runtime.GC()
done := make(chan struct{})
peak := make(chan uint64, 1)
go func() {
var highest uint64
var stats runtime.MemStats
for {
select {
case <-done:
peak <- highest
return
default:
}
runtime.ReadMemStats(&stats)
if stats.HeapAlloc > highest {
highest = stats.HeapAlloc
}
time.Sleep(time.Millisecond)
}
}()
work()
close(done)
return <-peak
}
// TestEmitPeakMemoryDoesNotScaleWithArtifactSize is the regression check for
// the OOM kills on large haks: emit used to hold the whole archive (twice), so
// a 2 GB hak needed about 10 GB. Emitting an archive 8× bigger must not cost
// meaningfully more memory.
func TestEmitPeakMemoryDoesNotScaleWithArtifactSize(t *testing.T) {
if testing.Short() {
t.Skip("writes a 64 MB fixture")
}
// A lazy GC lets garbage pile up in proportion to the live heap, which
// hides the thing under test. Collecting eagerly makes the sampled heap
// track what emit actually holds.
defer debug.SetGCPercent(debug.SetGCPercent(10))
measure := func(count int) uint64 {
dir := t.TempDir()
hak := filepath.Join(dir, "big.hak")
writeStreamedHak(t, hak, count)
// The key is computed outside the measurement: the test helper reads
// the whole file to hash it, which emit itself no longer does.
options := EmitOptions{
ArtifactKey: artifactKey(t, hak),
ArtifactPath: hak,
As: filepath.Base(hak),
OutDir: filepath.Join(dir, "out"),
}
return peakHeapDuring(func() {
if _, err := Emit(options); err != nil {
t.Fatalf("emit %d resources: %v", count, err)
}
})
}
small := measure(8) // 8 MB
large := measure(64) // 64 MB
const slack = 24 << 20
t.Logf("peak heap: 8 MB hak %d bytes, 64 MB hak %d bytes", small, large)
if large > small+slack {
t.Fatalf("peak heap scaled with artifact size: 8 MB hak peaked at %d bytes, 64 MB hak at %d", small, large)
}
}
+40 -93
View File
@@ -3,7 +3,6 @@ package nwsync
import ( import (
"bytes" "bytes"
"crypto/sha1" "crypto/sha1"
"crypto/sha256"
"encoding/binary" "encoding/binary"
"encoding/hex" "encoding/hex"
"encoding/json" "encoding/json"
@@ -46,30 +45,6 @@ func writeHak(t *testing.T, path string, contents map[string][]byte) {
} }
} }
// artifactKey is the depot key a file would be published under: the sha256 of
// its bytes, hash-tree depth 2, keeping the extension.
func artifactKey(t *testing.T, path string) string {
t.Helper()
body, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
sum := sha256.Sum256(body)
digest := hex.EncodeToString(sum[:])
return "artifacts/haks/sha256/" + digest[0:2] + "/" + digest[2:4] + "/" + digest + filepath.Ext(path)
}
// emitLocal emits one artifact into a local tree, the conformance path.
func emitLocal(t *testing.T, path, out string) (EmitResult, error) {
t.Helper()
return Emit(EmitOptions{
ArtifactKey: artifactKey(t, path),
ArtifactPath: path,
As: filepath.Base(path),
OutDir: out,
})
}
func TestBlobFramingRoundTrips(t *testing.T) { func TestBlobFramingRoundTrips(t *testing.T) {
data := []byte("the quick brown fox jumps over the lazy dog, repeatedly and at length") data := []byte("the quick brown fox jumps over the lazy dog, repeatedly and at length")
blob := compressBlob(data) blob := compressBlob(data)
@@ -168,7 +143,7 @@ func TestEmitWritesBlobsAndManifest(t *testing.T) {
}) })
out := filepath.Join(dir, "out") out := filepath.Join(dir, "out")
result, err := emitLocal(t, hak, out) result, err := Emit(hak, out)
if err != nil { if err != nil {
t.Fatalf("emit: %v", err) t.Fatalf("emit: %v", err)
} }
@@ -196,7 +171,7 @@ func TestEmitWritesBlobsAndManifest(t *testing.T) {
t.Errorf("blob decompressed to %q, want %q", got, body) t.Errorf("blob decompressed to %q, want %q", got, body)
} }
entries := readEmitted(t, result.ManifestPath) entries := readEmitted(t, out, "sow_test_01")
for _, entry := range entries { for _, entry := range entries {
if entry.ResType == restype(t, "nss") || entry.ResType == restype(t, "ndb") || entry.ResType == restype(t, "gic") { if entry.ResType == restype(t, "nss") || entry.ResType == restype(t, "ndb") || entry.ResType == restype(t, "gic") {
t.Errorf("skipped restype leaked into the manifest: %+v", entry) t.Errorf("skipped restype leaked into the manifest: %+v", entry)
@@ -225,9 +200,9 @@ func TestEmitWritesBlobsAndManifest(t *testing.T) {
} }
} }
func readEmitted(t *testing.T, indexPath string) []Entry { func readEmitted(t *testing.T, dir, name string) []Entry {
t.Helper() t.Helper()
data, err := os.ReadFile(indexPath) data, err := os.ReadFile(filepath.Join(dir, name+".nsym"))
if err != nil { if err != nil {
t.Fatalf("read emitted manifest: %v", err) t.Fatalf("read emitted manifest: %v", err)
} }
@@ -243,7 +218,7 @@ func TestEmitFailsClosedOnOversizeResource(t *testing.T) {
hak := filepath.Join(dir, "big.hak") hak := filepath.Join(dir, "big.hak")
writeHak(t, hak, map[string][]byte{"huge1.tga": make([]byte, fileSizeLimit+1)}) writeHak(t, hak, map[string][]byte{"huge1.tga": make([]byte, fileSizeLimit+1)})
if _, err := emitLocal(t, hak, filepath.Join(dir, "out")); err == nil { if _, err := Emit(hak, filepath.Join(dir, "out")); err == nil {
t.Fatal("emit accepted a resource over the 15 MB limit") t.Fatal("emit accepted a resource over the 15 MB limit")
} }
} }
@@ -255,11 +230,11 @@ func TestEmitLooseFile(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
out := filepath.Join(dir, "out") out := filepath.Join(dir, "out")
result, err := emitLocal(t, tlk, out) result, err := Emit(tlk, out)
if err != nil { if err != nil {
t.Fatalf("emit tlk: %v", err) t.Fatalf("emit tlk: %v", err)
} }
entries := readEmitted(t, result.ManifestPath) entries := readEmitted(t, out, "sow_tlk")
if len(entries) != 1 || entries[0].ResRef != "sow_tlk" || entries[0].ResType != restype(t, "tlk") { if len(entries) != 1 || entries[0].ResRef != "sow_tlk" || entries[0].ResType != restype(t, "tlk") {
t.Fatalf("tlk emitted as %+v", entries) t.Fatalf("tlk emitted as %+v", entries)
} }
@@ -270,35 +245,34 @@ func TestEmitLooseFile(t *testing.T) {
// emitFixture emits two haks that share a resref, so the merge rule is // emitFixture emits two haks that share a resref, so the merge rule is
// observable: "top" holds the winning body, "assets" the shadowed one. // observable: "top" holds the winning body, "assets" the shadowed one.
func emitFixture(t *testing.T) (out string, keys map[string]string, topBody, assetBody []byte) { func emitFixture(t *testing.T) (string, []byte, []byte) {
t.Helper() t.Helper()
dir := t.TempDir() dir := t.TempDir()
topBody = []byte("2da from sow_top") topBody := []byte("2da from sow_top")
assetBody = []byte("2da from the asset hak") assetBody := []byte("2da from the asset hak")
writeHak(t, filepath.Join(dir, "sow_top.hak"), map[string][]byte{"appearance.2da": topBody}) writeHak(t, filepath.Join(dir, "sow_top.hak"), map[string][]byte{"appearance.2da": topBody})
writeHak(t, filepath.Join(dir, "sow_core_01.hak"), map[string][]byte{ writeHak(t, filepath.Join(dir, "sow_core_01.hak"), map[string][]byte{
"appearance.2da": assetBody, "appearance.2da": assetBody,
"bloodstain1.tga": []byte("blood"), "bloodstain1.tga": []byte("blood"),
}) })
out = filepath.Join(dir, "out") out := filepath.Join(dir, "out")
keys = map[string]string{}
for _, name := range []string{"sow_top", "sow_core_01"} { for _, name := range []string{"sow_top", "sow_core_01"} {
path := filepath.Join(dir, name+".hak") if _, err := Emit(filepath.Join(dir, name+".hak"), out); err != nil {
keys[name] = artifactKey(t, path)
if _, err := emitLocal(t, path, out); err != nil {
t.Fatalf("emit %s: %v", name, err) t.Fatalf("emit %s: %v", name, err)
} }
} }
return out, keys, topBody, assetBody return out, topBody, assetBody
} }
func TestAssembleShadowsByOrder(t *testing.T) { func TestAssembleShadowsByOrder(t *testing.T) {
entriesDir, keys, topBody, assetBody := emitFixture(t) entriesDir, topBody, assetBody := emitFixture(t)
out := t.TempDir()
result, err := Assemble(AssembleOptions{ result, err := Assemble(AssembleOptions{
ArtifactKeys: []string{keys["sow_top"], keys["sow_core_01"]}, Order: []string{"sow_top", "sow_core_01"},
OutDir: entriesDir, EntriesDir: entriesDir,
GroupID: 2, OutDir: out,
GroupID: 2,
}) })
if err != nil { if err != nil {
t.Fatalf("assemble: %v", err) t.Fatalf("assemble: %v", err)
@@ -351,11 +325,13 @@ func TestAssembleShadowsByOrder(t *testing.T) {
} }
func TestAssembleReversedOrderPicksTheOtherHak(t *testing.T) { func TestAssembleReversedOrderPicksTheOtherHak(t *testing.T) {
entriesDir, keys, topBody, assetBody := emitFixture(t) entriesDir, topBody, assetBody := emitFixture(t)
out := t.TempDir()
result, err := Assemble(AssembleOptions{ result, err := Assemble(AssembleOptions{
ArtifactKeys: []string{keys["sow_core_01"], keys["sow_top"]}, Order: []string{"sow_core_01", "sow_top"},
OutDir: entriesDir, EntriesDir: entriesDir,
OutDir: out,
}) })
if err != nil { if err != nil {
t.Fatalf("assemble: %v", err) t.Fatalf("assemble: %v", err)
@@ -374,13 +350,9 @@ func TestAssembleReversedOrderPicksTheOtherHak(t *testing.T) {
} }
func TestAssembleRefusesMismatchedEmitterVersions(t *testing.T) { func TestAssembleRefusesMismatchedEmitterVersions(t *testing.T) {
entriesDir, keys, _, _ := emitFixture(t) entriesDir, _, _ := emitFixture(t)
index, err := resolveIndexKey(keys["sow_core_01"], entriesDir) path := filepath.Join(entriesDir, "sow_core_01.nsym.json")
if err != nil {
t.Fatal(err)
}
path := filepath.Join(entriesDir, index+".json")
var sidecar Sidecar var sidecar Sidecar
body, err := os.ReadFile(path) body, err := os.ReadFile(path)
if err != nil { if err != nil {
@@ -399,8 +371,9 @@ func TestAssembleRefusesMismatchedEmitterVersions(t *testing.T) {
} }
_, err = Assemble(AssembleOptions{ _, err = Assemble(AssembleOptions{
ArtifactKeys: []string{keys["sow_top"], keys["sow_core_01"]}, Order: []string{"sow_top", "sow_core_01"},
OutDir: entriesDir, EntriesDir: entriesDir,
OutDir: t.TempDir(),
}) })
if err == nil || !strings.Contains(err.Error(), "emitter version mismatch") { if err == nil || !strings.Contains(err.Error(), "emitter version mismatch") {
t.Fatalf("assemble merged across emitter versions: %v", err) t.Fatalf("assemble merged across emitter versions: %v", err)
@@ -412,7 +385,7 @@ func TestEmitHonoursSourceDateEpoch(t *testing.T) {
dir := t.TempDir() dir := t.TempDir()
hak := filepath.Join(dir, "pinned.hak") hak := filepath.Join(dir, "pinned.hak")
writeHak(t, hak, map[string][]byte{"one1.tga": []byte("body")}) writeHak(t, hak, map[string][]byte{"one1.tga": []byte("body")})
result, err := emitLocal(t, hak, filepath.Join(dir, "out")) result, err := Emit(hak, filepath.Join(dir, "out"))
if err != nil { if err != nil {
t.Fatalf("emit: %v", err) t.Fatalf("emit: %v", err)
} }
@@ -441,58 +414,32 @@ func TestEmitRejectsAModule(t *testing.T) {
if err := os.WriteFile(path, out.Bytes(), 0o644); err != nil { if err := os.WriteFile(path, out.Bytes(), 0o644); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if _, err := emitLocal(t, path, filepath.Join(dir, "out")); err == nil { if _, err := Emit(path, filepath.Join(dir, "out")); err == nil {
t.Fatal("emit accepted a .mod; a manifest never carries module contents") t.Fatal("emit accepted a .mod; a manifest never carries module contents")
} }
} }
func TestAssembleFailsClosedOnMissingIndex(t *testing.T) { func TestAssembleFailsClosedOnMissingIndex(t *testing.T) {
entriesDir, keys, _, _ := emitFixture(t) entriesDir, _, _ := emitFixture(t)
missing := "artifacts/haks/sha256/00/11/" + strings.Repeat("0", 64) + ".hak"
_, err := Assemble(AssembleOptions{ _, err := Assemble(AssembleOptions{
ArtifactKeys: []string{keys["sow_top"], missing}, Order: []string{"sow_top", "sow_never_published"},
OutDir: entriesDir, EntriesDir: entriesDir,
OutDir: t.TempDir(),
}) })
if err == nil || !strings.Contains(err.Error(), missing) { if err == nil || !strings.Contains(err.Error(), "sow_never_published") {
t.Fatalf("assemble did not fail closed and name the missing artifact: %v", err) t.Fatalf("assemble did not fail closed and name the missing artifact: %v", err)
} }
} }
// Callers capture a script's stdout as a value: `dir="$(pack-haks.sh)"`. A
// summary line on stdout gets glued onto that value, so both summaries belong
// on stderr.
func TestRunKeepsSummariesOffStdout(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "sow_top.hak")
writeHak(t, path, map[string][]byte{"appearance.2da": []byte("2da from sow_top")})
key := artifactKey(t, path)
out := filepath.Join(dir, "out")
for _, args := range [][]string{
{"emit", "--out", out, "--as", "sow_top.hak", key, path},
{"assemble", "--out", out, key},
} {
var stdout, stderr bytes.Buffer
if code := Run(args, &stdout, &stderr); code != exitOK {
t.Fatalf("Run(%v) exit=%d: %s", args, code, stderr.String())
}
if stdout.Len() != 0 {
t.Errorf("Run(%v) wrote to stdout: %q", args, stdout.String())
}
if stderr.Len() == 0 {
t.Errorf("Run(%v) reported no summary on stderr", args)
}
}
}
func TestRunUsageErrors(t *testing.T) { func TestRunUsageErrors(t *testing.T) {
cases := [][]string{ cases := [][]string{
nil, nil,
{"nope"}, {"nope"},
{"emit"}, {"emit"},
{"emit", "artifact-key.hak"}, {"emit", "artifact.hak"},
{"emit", "a", "b", "c"}, {"assemble", "--entries", "x", "--out", "y"},
{"assemble", "--out", "y"}, {"assemble", "--order", "a", "--out", "y"},
{"assemble", "--order", "a", "--entries", "x"},
} }
for _, args := range cases { for _, args := range cases {
var out, errw bytes.Buffer var out, errw bytes.Buffer
+45 -62
View File
@@ -12,6 +12,7 @@ import (
"flag" "flag"
"fmt" "fmt"
"io" "io"
"strings"
) )
const ( const (
@@ -29,9 +30,9 @@ func Run(args []string, stdout, stderr io.Writer) int {
} }
switch args[0] { switch args[0] {
case "emit": case "emit":
return runEmit(args[1:], stderr) return runEmit(args[1:], stdout, stderr)
case "assemble": case "assemble":
return runAssemble(args[1:], stderr) return runAssemble(args[1:], stdout, stderr)
case "-h", "--help", "help": case "-h", "--help", "help":
printRunUsage(stdout) printRunUsage(stdout)
return exitOK return exitOK
@@ -44,103 +45,85 @@ func Run(args []string, stdout, stderr io.Writer) int {
func printRunUsage(w io.Writer) { func printRunUsage(w io.Writer) {
fmt.Fprint(w, `usage: fmt.Fprint(w, `usage:
nwsync emit [--as NAME] [--out DIR] <artifact-key> <file> nwsync emit <artifact> --out DIR
nwsync assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>... nwsync assemble --order NAMES --entries DIR --out DIR [--group-id N]
emit explodes one .hak/.erf or one loose file (the TLK) into NWSync blobs plus 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 a NSYM manifest covering only that artifact. assemble merges those per-artifact
those indexes into one manifest, reading no bulk data. Artifact keys are depot manifests into one, reading no bulk data.
keys; an index lives beside its artifact, with the extension replaced.
--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.
`) `)
} }
// parseArgs parses flags that may appear before, after or between positionals. func runEmit(args []string, stdout, stderr io.Writer) int {
// Go's flag package stops at the first non-flag argument, which turns
// `emit <key> <file> --out DIR` into a confusing arity error.
func parseArgs(fs *flag.FlagSet, args []string) ([]string, error) {
var positional []string
for {
if err := fs.Parse(args); err != nil {
return nil, err
}
rest := fs.Args()
if len(rest) == 0 {
return positional, nil
}
positional = append(positional, rest[0])
args = rest[1:]
}
}
func runEmit(args []string, stderr io.Writer) int {
fs := flag.NewFlagSet("emit", flag.ContinueOnError) fs := flag.NewFlagSet("emit", flag.ContinueOnError)
fs.SetOutput(stderr) fs.SetOutput(stderr)
as := fs.String("as", "", "published name of the artifact, when it differs from the key") out := fs.String("out", "", "output directory (blobs plus the per-artifact manifest)")
out := fs.String("out", "", "write to a local repository tree instead of uploading") if err := fs.Parse(args); err != nil {
jobs := fs.Int("jobs", defaultEmitJobs, "resources to hash, compress and store at once")
positional, err := parseArgs(fs, args)
if err != nil {
return exitUsage return exitUsage
} }
if len(positional) != 2 { if fs.NArg() != 1 {
fmt.Fprintf(stderr, "nwsync emit: <artifact-key> and <file> are both required\n") fmt.Fprintf(stderr, "nwsync emit: exactly one artifact is required\n")
return exitUsage return exitUsage
} }
if *jobs < 1 { if *out == "" {
fmt.Fprintf(stderr, "nwsync emit: -jobs must be at least 1, got %d\n", *jobs) fmt.Fprintf(stderr, "nwsync emit: --out is required\n")
return exitUsage return exitUsage
} }
result, err := Emit(EmitOptions{ result, err := Emit(fs.Arg(0), *out)
ArtifactKey: positional[0],
ArtifactPath: positional[1],
As: *as,
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)
return exitInternal return exitInternal
} }
fmt.Fprintf(stderr, "emitted %s: %d resources, %d new blobs, index %s\n", fmt.Fprintf(stdout, "emitted %s: %d resources, %d new blobs, manifest %s\n",
result.Name, result.Entries, result.BlobsWritten, result.ManifestPath) result.Name, result.Entries, result.BlobsWritten, result.ManifestPath)
return exitOK return exitOK
} }
func runAssemble(args []string, stderr io.Writer) int { func runAssemble(args []string, stdout, stderr io.Writer) int {
fs := flag.NewFlagSet("assemble", flag.ContinueOnError) fs := flag.NewFlagSet("assemble", flag.ContinueOnError)
fs.SetOutput(stderr) fs.SetOutput(stderr)
tlkKey := fs.String("tlk-key", "", "depot key of the TLK, which shadows nothing and merges last") order := fs.String("order", "", "comma-separated artifact names, highest priority first")
out := fs.String("out", "", "write to a local repository tree instead of uploading") entries := fs.String("entries", "", "directory holding the per-artifact .nsym files")
out := fs.String("out", "", "repository root; the manifest lands in <out>/manifests")
groupID := fs.Int("group-id", 0, "NWSync group id (1 = current, 2 = testing; 0 omits it)") groupID := fs.Int("group-id", 0, "NWSync group id (1 = current, 2 = testing; 0 omits it)")
moduleName := fs.String("module-name", "", "module name recorded in the sidecar") moduleName := fs.String("module-name", "", "module name recorded in the sidecar")
description := fs.String("description", "", "description recorded in the sidecar") description := fs.String("description", "", "description recorded in the sidecar")
positional, err := parseArgs(fs, args) if err := fs.Parse(args); err != nil {
if err != nil {
return exitUsage return exitUsage
} }
if len(positional) == 0 { switch {
fmt.Fprintf(stderr, "nwsync assemble: at least one artifact key is required\n") case *order == "":
fmt.Fprintf(stderr, "nwsync assemble: --order is required\n")
return exitUsage
case *entries == "":
fmt.Fprintf(stderr, "nwsync assemble: --entries is required\n")
return exitUsage
case *out == "":
fmt.Fprintf(stderr, "nwsync assemble: --out is required\n")
return exitUsage return exitUsage
} }
names := make([]string, 0, 8)
for _, name := range strings.Split(*order, ",") {
if name = strings.TrimSpace(name); name != "" {
names = append(names, name)
}
}
result, err := Assemble(AssembleOptions{ result, err := Assemble(AssembleOptions{
ArtifactKeys: positional, Order: names,
TLKKey: *tlkKey, EntriesDir: *entries,
OutDir: *out, OutDir: *out,
GroupID: *groupID, GroupID: *groupID,
ModuleName: *moduleName, ModuleName: *moduleName,
Description: *description, Description: *description,
}) })
if err != nil { if err != nil {
fmt.Fprintf(stderr, "nwsync assemble: %v\n", err) fmt.Fprintf(stderr, "nwsync assemble: %v\n", err)
return exitInternal return exitInternal
} }
fmt.Fprintf(stderr, "assembled manifest %s: %d resources, %s\n", fmt.Fprintf(stdout, "assembled %s: %d resources from %d artifacts\n",
result.SHA1, result.Entries, result.ManifestPath) result.SHA1, result.Entries, len(names))
return exitOK return exitOK
} }
-212
View File
@@ -1,212 +0,0 @@
package nwsync
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"io"
"os"
"path"
"path/filepath"
"strings"
"git.westgate.pw/ShadowsOverWestgate/sow-tools/internal/depot"
)
// sink is where an emit or assemble run puts what it produces. The zone is the
// production sink; a local directory exists only as the conformance path, so
// 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)
// 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
// getIndex reads back a NSYM manifest and its sidecar.
getIndex(key string) (manifest, sidecar []byte, err error)
// describe names the sink for messages.
describe(key string) string
}
// dirSink writes a local NWSync repository tree.
type dirSink struct{ root string }
func (s dirSink) putBlob(sha1Hex string, body func() []byte) (int64, error) {
blob := blobPath(s.root, sha1Hex)
if _, err := os.Stat(blob); err == nil {
return 0, nil
}
if err := os.MkdirAll(filepath.Dir(blob), 0o755); err != nil {
return 0, fmt.Errorf("create blob directory: %w", err)
}
data := body()
if err := os.WriteFile(blob, data, 0o644); err != nil {
return 0, fmt.Errorf("write blob: %w", err)
}
return int64(len(data)), nil
}
func (s dirSink) putIndex(key string, manifest, sidecar []byte) error {
target := filepath.Join(s.root, filepath.FromSlash(key))
if err := os.MkdirAll(filepath.Dir(target), 0o755); err != nil {
return fmt.Errorf("create manifest directory: %w", err)
}
if err := os.WriteFile(target, manifest, 0o644); err != nil {
return fmt.Errorf("write manifest: %w", err)
}
if err := os.WriteFile(target+".json", sidecar, 0o644); err != nil {
return fmt.Errorf("write sidecar: %w", err)
}
return nil
}
func (s dirSink) getIndex(key string) ([]byte, []byte, error) {
target := filepath.Join(s.root, filepath.FromSlash(key))
manifest, err := os.ReadFile(target)
if err != nil {
return nil, nil, fmt.Errorf("read index: %w", err)
}
sidecar, err := os.ReadFile(target + ".json")
if err != nil {
return nil, nil, fmt.Errorf("read sidecar: %w", err)
}
return manifest, sidecar, nil
}
func (s dirSink) describe(key string) string {
return filepath.Join(s.root, filepath.FromSlash(key))
}
// zoneSink uploads straight to the NWSync storage zone. Nothing bulky is ever
// written to the runner's disk: the working set is one resource at a time.
type zoneSink struct {
store depot.KeyStore
ctx context.Context
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)
// 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)
if err != nil {
return 0, fmt.Errorf("probe blob %s: %w", sha1Hex, err)
}
if state == depot.Present {
return 0, nil
}
data := body()
if err := s.put(key, data); err != nil {
return 0, fmt.Errorf("upload blob %s: %w", sha1Hex, err)
}
return int64(len(data)), nil
}
func (s zoneSink) putIndex(key string, manifest, sidecar []byte) error {
// The manifest lands last: its presence is the publication marker, so it
// must never appear before the blobs it names.
if err := s.put(key+".json", sidecar); err != nil {
return fmt.Errorf("upload sidecar %s: %w", key, err)
}
if err := s.put(key, manifest); err != nil {
return fmt.Errorf("upload index %s: %w", key, err)
}
return nil
}
func (s zoneSink) getIndex(key string) ([]byte, []byte, error) {
manifest, err := s.store.GetKey(s.ctx, key)
if err != nil {
return nil, nil, fmt.Errorf("read index: %w", err)
}
sidecar, err := s.store.GetKey(s.ctx, key+".json")
if err != nil {
return nil, nil, fmt.Errorf("read sidecar: %w", err)
}
return manifest, sidecar, nil
}
func (s zoneSink) describe(key string) string { return s.zone + "/" + key }
func (s zoneSink) put(key string, body []byte) error {
sum := sha256.Sum256(body)
return s.store.PutReader(s.ctx, key, bytes.NewReader(body), int64(len(body)), hex.EncodeToString(sum[:]))
}
// newZoneSink builds the upload sink from the environment. NWSync data lives
// in its own zone, separate from the asset depot, so it has its own zone and
// credential; only the host is shared, and Crucible has no default host.
func newZoneSink(ctx context.Context, getenv func(string) string) (sink, error) {
cfg := depot.LoadConfig(getenv)
cfg.StorageZone = getenv("NWSYNC_STORAGE_ZONE")
cfg.WriteKey = getenv("NWSYNC_STORAGE_PASSWORD")
cfg.ReadKey = cfg.WriteKey
if cfg.StorageZone == "" {
return nil, fmt.Errorf("NWSYNC_STORAGE_ZONE is unset (or pass --out DIR to write locally)")
}
if cfg.WriteKey == "" {
return nil, fmt.Errorf("NWSYNC_STORAGE_PASSWORD is unset (or pass --out DIR to write locally)")
}
if cfg.StorageHost == "" {
return nil, fmt.Errorf("BUNNY_STORAGE_HOST is unset")
}
store, err := depot.NewKeyStore(cfg)
if err != nil {
return nil, err
}
return zoneSink{store: store, ctx: ctx, zone: cfg.StorageZone}, nil
}
// indexKey is where an artifact's NSYM lives: beside the artifact itself, with
// the final extension replaced. emit and assemble must agree on this one rule,
// so it lives here and nowhere else.
//
// artifacts/haks/sha256/30/46/3046….hak -> artifacts/haks/sha256/30/46/3046….nsym
func indexKey(artifactKey string) (string, error) {
extension := path.Ext(artifactKey)
if extension == "" {
return "", fmt.Errorf("artifact key %q has no extension", artifactKey)
}
return strings.TrimSuffix(artifactKey, extension) + ".nsym", nil
}
// resolveIndexKey is where emit writes an artifact's index and where assemble
// reads it from. On the zone that is beside the artifact; locally the indexes
// sit flat beside the data tree, so upstream's output and ours diff directly.
func resolveIndexKey(artifactKey, outDir string) (string, error) {
key, err := indexKey(artifactKey)
if err != nil {
return "", err
}
if outDir != "" {
return path.Base(key), nil
}
return key, nil
}
// checkArtifactKey fails closed when the key's embedded digest is not the
// digest of the bytes being emitted. Publishing an index under the wrong key
// silently pairs a manifest with the wrong artifact.
// artifact is hashed by streaming, so a multi-gigabyte hak is never resident.
func checkArtifactKey(artifactKey string, artifact io.Reader) error {
base := path.Base(artifactKey)
digest := strings.TrimSuffix(base, path.Ext(base))
if len(digest) != 64 {
return fmt.Errorf("artifact key %q does not name a sha256", artifactKey)
}
hash := sha256.New()
if _, err := io.Copy(hash, artifact); err != nil {
return fmt.Errorf("hash artifact: %w", err)
}
if got := hex.EncodeToString(hash.Sum(nil)); got != digest {
return fmt.Errorf("artifact key %q names digest %s but the file hashes to %s", artifactKey, digest, got)
}
return nil
}
-249
View File
@@ -1,249 +0,0 @@
package nwsync
import (
"crypto/sha1"
"crypto/sha256"
"encoding/hex"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"sync"
"testing"
)
// fakeZone is a Bunny-shaped object store: PUT stores, GET reads, and the
// Checksum header is verified the way Bunny verifies it.
type fakeZone struct {
mu sync.Mutex
objects map[string][]byte
puts []string
failOn func(key string) bool // when true, the PUT fails
}
func newFakeZone(t *testing.T) (*fakeZone, func(string) string) {
t.Helper()
zone := &fakeZone{objects: map[string][]byte{}}
server := httptest.NewServer(zone)
t.Cleanup(server.Close)
getenv := func(name string) string {
switch name {
case "NWSYNC_STORAGE_ZONE":
return "sow-nwsync"
case "NWSYNC_STORAGE_PASSWORD":
return "write-key"
case "BUNNY_STORAGE_HOST":
return server.URL
}
return ""
}
return zone, getenv
}
func (z *fakeZone) ServeHTTP(w http.ResponseWriter, r *http.Request) {
key := strings.TrimPrefix(r.URL.Path, "/sow-nwsync/")
switch r.Method {
case http.MethodPut:
if z.failOn != nil && z.failOn(key) {
http.Error(w, "boom", http.StatusInternalServerError)
return
}
body, err := io.ReadAll(r.Body)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
sum := sha256.Sum256(body)
if want := strings.ToUpper(hex.EncodeToString(sum[:])); r.Header.Get("Checksum") != want {
http.Error(w, "checksum mismatch", http.StatusBadRequest)
return
}
z.mu.Lock()
z.objects[key] = body
z.puts = append(z.puts, key)
z.mu.Unlock()
w.WriteHeader(http.StatusCreated)
case http.MethodGet:
z.mu.Lock()
body, ok := z.objects[key]
z.mu.Unlock()
if !ok {
http.Error(w, "not found", http.StatusNotFound)
return
}
_, _ = w.Write(body)
default:
http.Error(w, "unsupported", http.StatusMethodNotAllowed)
}
}
func (z *zoneSinkFixture) emit(t *testing.T, path string) EmitResult {
t.Helper()
result, err := Emit(EmitOptions{
ArtifactKey: artifactKey(t, path),
ArtifactPath: path,
Sink: z.sink,
})
if err != nil {
t.Fatalf("emit %s: %v", path, err)
}
return result
}
type zoneSinkFixture struct {
zone *fakeZone
sink sink
}
func newZoneFixture(t *testing.T) *zoneSinkFixture {
t.Helper()
zone, getenv := newFakeZone(t)
target, err := newZoneSink(t.Context(), getenv)
if err != nil {
t.Fatalf("zone sink: %v", err)
}
return &zoneSinkFixture{zone: zone, sink: target}
}
func sha1Of(body []byte) [20]byte { return sha1.Sum(body) }
func TestEmitUploadsBlobsThenIndex(t *testing.T) {
fixture := newZoneFixture(t)
dir := t.TempDir()
hak := filepath.Join(dir, "sow_test_01.hak")
body := []byte("texture bytes")
writeHak(t, hak, map[string][]byte{"bloodstain1.tga": body, "copy1.txi": body})
key := artifactKey(t, hak)
result := fixture.emit(t, hak)
index, err := indexKey(key)
if err != nil {
t.Fatal(err)
}
if _, ok := fixture.zone.objects[index]; !ok {
t.Fatalf("no index at %s; zone holds %v", index, fixture.zone.puts)
}
if result.BlobsWritten != 1 {
t.Errorf("uploaded %d blobs, want 1 (identical content shares a blob)", result.BlobsWritten)
}
// The index is the publication marker, so it must land after every blob it
// names — including its own sidecar.
last := fixture.zone.puts[len(fixture.zone.puts)-1]
if last != index {
t.Errorf("index landed at position %d of %d; it must be last", len(fixture.zone.puts), len(fixture.zone.puts))
}
for _, key := range fixture.zone.puts[:len(fixture.zone.puts)-1] {
if strings.HasPrefix(key, "data/sha1/") || key == index+".json" {
continue
}
t.Errorf("unexpected object uploaded before the index: %s", key)
}
}
func TestEmitSkipsBlobsAlreadyInTheZone(t *testing.T) {
fixture := newZoneFixture(t)
dir := t.TempDir()
hak := filepath.Join(dir, "sow_test_01.hak")
writeHak(t, hak, map[string][]byte{"bloodstain1.tga": []byte("blood")})
first := fixture.emit(t, hak)
if first.BlobsWritten != 1 {
t.Fatalf("first emit uploaded %d blobs, want 1", first.BlobsWritten)
}
second := fixture.emit(t, hak)
if second.BlobsWritten != 0 {
t.Errorf("re-emit uploaded %d blobs, want 0 (a blob name is its content)", second.BlobsWritten)
}
}
func TestEmitLeavesNoIndexWhenAnUploadFails(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, map[string][]byte{"bloodstain1.tga": []byte("blood")})
_, err := Emit(EmitOptions{
ArtifactKey: artifactKey(t, hak),
ArtifactPath: hak,
Sink: fixture.sink,
})
if err == nil {
t.Fatal("emit reported success after an upload failed")
}
for key := range fixture.zone.objects {
if strings.HasSuffix(key, ".nsym") {
t.Errorf("a half-emitted artifact published an index: %s", key)
}
}
}
func TestEmitRejectsAKeyThatDoesNotMatchTheFile(t *testing.T) {
fixture := newZoneFixture(t)
dir := t.TempDir()
hak := filepath.Join(dir, "sow_test_01.hak")
writeHak(t, hak, map[string][]byte{"bloodstain1.tga": []byte("blood")})
wrong := "artifacts/haks/sha256/00/11/" + strings.Repeat("0", 64) + ".hak"
_, err := Emit(EmitOptions{ArtifactKey: wrong, ArtifactPath: hak, Sink: fixture.sink})
if err == nil || !strings.Contains(err.Error(), "hashes to") {
t.Fatalf("emit published under a key that names another artifact: %v", err)
}
}
func TestAssembleReadsIndexesFromTheZone(t *testing.T) {
fixture := newZoneFixture(t)
dir := t.TempDir()
topBody := []byte("2da from sow_top")
assetBody := []byte("2da from the asset hak")
top := filepath.Join(dir, "sow_top.hak")
core := filepath.Join(dir, "sow_core_01.hak")
writeHak(t, top, map[string][]byte{"appearance.2da": topBody})
writeHak(t, core, map[string][]byte{"appearance.2da": assetBody, "bloodstain1.tga": []byte("blood")})
tlkPath := filepath.Join(dir, "sow_tlk.tlk")
if err := os.WriteFile(tlkPath, []byte("TLK V3.0 payload"), 0o644); err != nil {
t.Fatal(err)
}
fixture.emit(t, top)
fixture.emit(t, core)
if _, err := Emit(EmitOptions{
ArtifactKey: artifactKey(t, tlkPath),
ArtifactPath: tlkPath,
As: "sow_tlk.tlk",
Sink: fixture.sink,
}); err != nil {
t.Fatalf("emit tlk: %v", err)
}
result, err := Assemble(AssembleOptions{
ArtifactKeys: []string{artifactKey(t, top), artifactKey(t, core)},
TLKKey: artifactKey(t, tlkPath),
GroupID: 2,
Sink: fixture.sink,
})
if err != nil {
t.Fatalf("assemble: %v", err)
}
if result.Entries != 3 {
t.Fatalf("merged %d entries, want 3 (appearance.2da is shadowed, the TLK adds one)", result.Entries)
}
manifest, ok := fixture.zone.objects["manifests/"+result.SHA1]
if !ok {
t.Fatalf("no merged manifest in the zone; it holds %v", fixture.zone.puts)
}
entries, err := readManifest(manifest)
if err != nil {
t.Fatalf("parse merged manifest: %v", err)
}
for _, entry := range entries {
if entry.ResRef == "appearance" && entry.SHA1 != sha1Of(topBody) {
t.Errorf("appearance.2da resolved to the shadowed hak, not the first one given")
}
}
}