feat(nwsync): upload sink, key-addressed CLI, fail-closed publication marker
ci / ci (pull_request) Successful in 3m21s
ci / ci (pull_request) Successful in 3m21s
Resolves the build half of sow-tools#60 and closes the CLI gap sow-tools#53's sweep found: PR #71 shipped #53's original surface, not the one #56, #62 and #65 settled. - `depot.KeyStore` (`ProbeKey`/`PutReader`/`GetKey`) addresses the zone by object key instead of by depot sha, reusing the existing IPv4-pinned transport, retry and tri-state probe. The sha-addressed `Backend` now rides on it; no new HTTP client. - `nwsync` gains a sink: the zone by default, a local tree with `--out DIR` as the conformance path. Blobs upload as they are produced, the index lands last, and a blob already in the zone is skipped without paying for compression. - CLI is now `emit [--as NAME] [--out DIR] <artifact-key> <file>` and `assemble --group-id N [--tlk-key KEY] [--out DIR] <artifact-key>...`. Indexes live beside their artifact with the extension replaced, derived in one place. Flags may follow positionals, which cost a run during sow-tools#59. - Fail-closed: an artifact key whose digest does not match the file is refused, and `assemble` refuses an artifact with no index rather than publishing a manifest missing a hak. Conformance re-checked through the new CLI against upstream nwn_nwsync_write 2.1.2 on sow_vfxs_01.hak: the manifest is still byte-identical. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,121 @@
|
||||
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)
|
||||
}
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -71,16 +70,7 @@ type httpBackend struct {
|
||||
|
||||
func (b *httpBackend) Name() string { return b.name }
|
||||
|
||||
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) storageURL(sha string) string { return b.keyURL(BlobKey(sha)) }
|
||||
|
||||
func (b *httpBackend) cdnURL(sha string) string {
|
||||
return fmt.Sprintf("%s/%s", b.cfg.CDNBase, BlobKey(sha))
|
||||
@@ -140,7 +130,8 @@ 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})
|
||||
}
|
||||
|
||||
// Put uploads src for sha. cdn is read-only.
|
||||
// Put uploads src for sha. cdn is read-only. A depot object is named after the
|
||||
// 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 {
|
||||
if b.name == "cdn" {
|
||||
return errors.New("cdn backend is read-only")
|
||||
@@ -157,30 +148,7 @@ func (b *httpBackend) Put(ctx context.Context, sha, src string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
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
|
||||
return b.putFile(ctx, BlobKey(sha), src, sha)
|
||||
}
|
||||
|
||||
// Get fetches sha into dest via temp file + rename, re-hashing and deleting
|
||||
|
||||
Reference in New Issue
Block a user