diff --git a/cmd/stellar-rpc/internal/rpcv2/backfill/perf_test.go b/cmd/stellar-rpc/internal/rpcv2/backfill/perf_test.go index 1ed84ba52..f7febfd68 100644 --- a/cmd/stellar-rpc/internal/rpcv2/backfill/perf_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/backfill/perf_test.go @@ -105,10 +105,10 @@ func TestStreamingRebuild_ByteIdenticalToColdPath(t *testing.T) { // --------------------------------------------------------------------------- // TestStreamingBin_MatchesSpecFormat asserts the .bin a frozen chunk leaves on -// disk matches gettransaction §6.1: a uint64-LE entry-count header, then the -// 16-byte index secret the keys were blinded with, then 20-byte [16-byte key | -// 4-byte LE seq] entries. freezeChunkBin uses the real txhash.WriteColdBin, so -// this is the producer's actual on-disk contract. +// disk matches gettransaction §6.1: the "SBIN" magic-and-version prelude, a +// uint64-LE entry-count, the 16-byte index secret the keys were blinded with, +// then 20-byte [16-byte key | 4-byte LE seq] entries. freezeChunkBin uses the +// real txhash.WriteColdBin, so this is the producer's actual on-disk contract. func TestStreamingBin_MatchesSpecFormat(t *testing.T) { cat, _ := smallTxHashIndexCatalog(t, 4) @@ -119,12 +119,13 @@ func TestStreamingBin_MatchesSpecFormat(t *testing.T) { raw, err := os.ReadFile(cat.Layout().TxHashBinPath(0)) require.NoError(t, err) - // §6.1: 24-byte header (8-byte uint64-LE count + 16-byte index secret) + - // N * 20-byte entries. + // §6.1: 32-byte header (8-byte magic/version prelude + 8-byte uint64-LE + // count + 16-byte index secret) + N * 20-byte entries. const ( + preludeW = 8 // magic "SBIN" + version + reserved countW = 8 // uint64-LE entry count secretW = stores.SecretLen // index secret the keys were blinded with - hdrSize = countW + secretW + hdrSize = preludeW + countW + secretW keyW = 16 // streamhash.MinKeySize seqW = 4 entryW = keyW + seqW // 20 bytes exactly @@ -134,7 +135,9 @@ func TestStreamingBin_MatchesSpecFormat(t *testing.T) { require.Equal(t, streamhash.MinKeySize, keyW, "16-byte key == streamhash routing-key width") require.Len(t, raw, hdrSize+wantCount*entryW, "header + 20-byte entries") - count := binary.LittleEndian.Uint64(raw[:countW]) + require.Equal(t, []byte("SBIN"), raw[:4], "magic in on-disk byte order") + require.Equal(t, byte(1), raw[4], ".bin format version") + count := binary.LittleEndian.Uint64(raw[preludeW : preludeW+countW]) require.Equal(t, uint64(wantCount), count, "uint64-LE entry-count header") // Each entry: the 16-byte secret-keyed routing key (keyed at ingest — never diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/catalog.go b/cmd/stellar-rpc/internal/rpcv2/catalog/catalog.go index b70802ba5..a1142a5e3 100644 --- a/cmd/stellar-rpc/internal/rpcv2/catalog/catalog.go +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/catalog.go @@ -46,6 +46,15 @@ func Open( return nil, err } c := &Catalog{store: store, logger: logger, layout: layout, txhashIndex: txhashIndex} + // Census before the secret mint below: a catalog holding entries outside + // this binary's vocabulary (a newer binary's formats, or corruption) is + // refused here, so Open writes no catalog entry of its own into a tree it + // refuses. RocksDB itself may still create housekeeping files and flush a + // previous binary's recovered WAL on Close. + if err := c.census(); err != nil { + _ = c.Close() + return nil, err + } // Mint-or-load the cold-index secret up front (get-or-create is not atomic; // here it runs single-threaded) and cache it, so post-Open Secret() reads are // lock-free and cannot fail. diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/census.go b/cmd/stellar-rpc/internal/rpcv2/catalog/census.go new file mode 100644 index 000000000..2cbcb3467 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/census.go @@ -0,0 +1,129 @@ +package catalog + +import ( + "errors" + "fmt" + "strconv" + "strings" + + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/geometry" +) + +// ErrForeignCatalog marks a catalog holding entries outside this binary's +// vocabulary. The daemon refuses to start on it: the entries were either +// written by a newer stellar-rpc (whose formats this binary cannot read, and +// whose artifacts the resolver would otherwise overwrite) or corrupted. +var ErrForeignCatalog = errors.New("catalog: entries unknown to this version of stellar-rpc") + +// censusMaxDetailed caps how many offending entries the refusal error spells +// out; the total count is always reported. +const censusMaxDetailed = 10 + +// census validates every key and value in the store against the exact +// vocabulary this binary writes. It runs once inside Open, before the secret +// mint, so a foreign catalog is refused before anything writes into it. The +// scan is read-only and touches only the catalog (never a hot DB or a file). +// +// Everything this binary durably writes must parse here; the whitelist spans +// the geometry key families AND the meta/catalog-secret key, and the value +// checks are the census's own exact-token comparisons (State/HotState reads +// elsewhere are raw casts, not validators). +func (c *Catalog) census() error { + var ( + offenders []string + total int + ) + flag := func(key, detail string) { + total++ + if len(offenders) < censusMaxDetailed { + offenders = append(offenders, fmt.Sprintf("%q: %s", key, detail)) + } + } + + for e, err := range c.prefixScan("") { + if err != nil { + return fmt.Errorf("catalog: census scan: %w", err) + } + switch e.Key { + case geometry.ConfigEarliestLedger: + if !isCanonicalUint32(e.Value) { + flag(e.Key, fmt.Sprintf("value %q is not a canonical decimal uint32", e.Value)) + } + case catalogSecretStoreKey: + // Never print the secret: it is the live index-blinding key. The + // width comes from the field's type, valid before the mint. + if len(e.Value) != len(c.secret) { + flag(e.Key, fmt.Sprintf("value is %d bytes, want %d (value redacted)", len(e.Value), len(c.secret))) + } + default: + if detail, ok := c.censusArtifactEntry(e.Key, e.Value); !ok { + flag(e.Key, detail) + } + } + } + + if total == 0 { + return nil + } + return fmt.Errorf("%w — either written by a newer stellar-rpc (deploy that version or newer) "+ + "or corrupted; %d offending entr%s: %s", + ErrForeignCatalog, total, plural(total), strings.Join(offenders, "; ")) +} + +// censusArtifactEntry validates one non-config entry against the three state +// key families. ok=false returns the reason. +func (c *Catalog) censusArtifactEntry(key, value string) (string, bool) { + switch { + case strings.HasPrefix(key, geometry.HotChunkPrefix): + if _, ok := geometry.ParseHotChunkKey(key); !ok { + return "malformed hot-chunk key", false + } + if !geometry.IsKnownHotState(geometry.HotState(value)) { + return fmt.Sprintf("unknown hot state %q", value), false + } + case strings.HasPrefix(key, geometry.ChunkPrefix): + if _, _, ok := geometry.ParseChunkKey(key); !ok { + return "malformed per-chunk artifact key", false + } + if !geometry.IsKnownState(geometry.State(value)) { + return fmt.Sprintf("unknown artifact state %q", value), false + } + case strings.HasPrefix(key, geometry.TxHashIndexPrefix): + cov, ok := geometry.ParseTxHashIndexKey(key) + if !ok { + return "malformed index coverage key", false + } + // The builder only ever covers chunks of the key's own index, so a + // cross-window coverage is something no binary wrote. Left accepted, + // a corrupt frozen one would suppress legitimate index rebuilds. + if c.txhashIndex.TxHashIndexID(cov.Lo) != cov.Index || + c.txhashIndex.TxHashIndexID(cov.Hi) != cov.Index { + return "index coverage endpoints outside the key's own index", false + } + if !geometry.IsKnownState(geometry.State(value)) { + return fmt.Sprintf("unknown artifact state %q", value), false + } + default: + // Never print an unknown key's value: a newer binary may store key + // material under a name this binary does not recognize. + return fmt.Sprintf("unknown key (value redacted, %d bytes)", len(value)), false + } + return "", true +} + +// isCanonicalUint32 reports whether v round-trips through ParseUint and +// FormatUint byte-identically — the exact form PinEarliestLedger writes. +func isCanonicalUint32(v string) bool { + n, err := strconv.ParseUint(v, 10, 32) + if err != nil { + return false + } + return strconv.FormatUint(n, 10) == v +} + +func plural(n int) string { + if n == 1 { + return "y" + } + return "ies" +} diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/census_test.go b/cmd/stellar-rpc/internal/rpcv2/catalog/census_test.go new file mode 100644 index 000000000..55b3ac385 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/census_test.go @@ -0,0 +1,133 @@ +package catalog + +import ( + "fmt" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/chunk" + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/geometry" + "github.com/stellar/stellar-rpc/cmd/stellar-rpc/internal/rpcv2/rocksdb" +) + +// reopenAfter opens a catalog at path, applies mutate, closes it, and returns +// the error from reopening — the census's verdict on the mutated store. +func reopenAfter(t *testing.T, mutate func(c *Catalog)) error { + t.Helper() + path := t.TempDir() + cat, err := openKVAt(t, path) + require.NoError(t, err) + mutate(cat) + require.NoError(t, cat.Close()) + reopened, err := openKVAt(t, path) + if reopened != nil { + t.Cleanup(func() { _ = reopened.Close() }) + } + return err +} + +// TestCensus_AcceptsOwnVocabulary pins that every key and value this binary +// writes — all three state families in every state, the pin, and the minted +// secret — reopens cleanly. This is the whole release-1 write surface. +func TestCensus_AcceptsOwnVocabulary(t *testing.T) { + err := reopenAfter(t, func(c *Catalog) { + for i, s := range geometry.AllStates() { + id := chunk.ID(i) + for _, kind := range geometry.AllKinds() { + require.NoError(t, c.put(geometry.ChunkKey(id, kind), string(s))) + } + // Coverage endpoints must lie in the key's own index window. + first := chunk.ID(i) * chunk.ID(geometry.ChunksPerTxhashIndex) + idxKey := geometry.TxHashIndexKey(geometry.TxHashIndexID(i), first, first+5) + require.NoError(t, c.put(idxKey, string(s))) + } + for i, s := range geometry.AllHotStates() { + require.NoError(t, c.put(geometry.HotChunkKey(chunk.ID(7+i)), string(s))) + } + require.NoError(t, c.PinEarliestLedger(2)) + }) + require.NoError(t, err) +} + +// TestCensus_AcceptsFirstStartResidues pins the crash-shaped first-start +// catalogs: secret-only (crash after Open, before the pin) and secret+pin. +func TestCensus_AcceptsFirstStartResidues(t *testing.T) { + require.NoError(t, reopenAfter(t, func(*Catalog) {}), "secret-only catalog") + require.NoError(t, reopenAfter(t, func(c *Catalog) { + require.NoError(t, c.PinEarliestLedger(10002)) + }), "secret+pin catalog") +} + +// TestCensus_RefusesForeignEntries walks the refusal matrix: every entry shape +// a newer binary (or corruption) could leave must fail reopen with +// ErrForeignCatalog and an actionable message. +func TestCensus_RefusesForeignEntries(t *testing.T) { + cases := []struct { + name, key, value string + }{ + {"novel prefix", "format:events", "2"}, + {"novel meta key", "meta/other", "WOULD-BE-SECRET-BYTES"}, + {"unknown kind under chunk prefix", "chunk:00000001:bogus", string(geometry.StateFrozen)}, + {"version suffix on chunk state", geometry.ChunkKey(1, geometry.KindEvents), "frozen@2"}, + {"version suffix on hot state", geometry.HotChunkKey(3), "ready@2"}, + {"version suffix on index state", geometry.TxHashIndexKey(0, 0, 9), "frozen@2"}, + {"unknown hot state", geometry.HotChunkKey(4), "warm"}, + {"unpadded chunk id", "chunk:123:ledgers", string(geometry.StateFrozen)}, + {"index lo above hi", "index:00000000:00000005:00000002", string(geometry.StateFrozen)}, + {"cross-window index coverage", "index:00000001:00000000:00003000", string(geometry.StateFrozen)}, + {"non-canonical pin", geometry.ConfigEarliestLedger, "007"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + err := reopenAfter(t, func(c *Catalog) { + require.NoError(t, c.put(tc.key, tc.value)) + }) + require.ErrorIs(t, err, ErrForeignCatalog) + require.ErrorContains(t, err, "deploy that version or newer") + require.ErrorContains(t, err, tc.key) + require.NotContains(t, err.Error(), "WOULD-BE-SECRET-BYTES", + "an unknown key's value must never be printed; it may be a newer binary's key material") + }) + } +} + +// TestCensus_CapsDetailButCountsAll pins the refusal shape on a store with +// more offenders than the detail cap: every offender is counted, only the +// first censusMaxDetailed are spelled out. +func TestCensus_CapsDetailButCountsAll(t *testing.T) { + const n = censusMaxDetailed + 2 + err := reopenAfter(t, func(c *Catalog) { + for i := range n { + require.NoError(t, c.put(fmt.Sprintf("future:%02d", i), "x")) + } + }) + require.ErrorIs(t, err, ErrForeignCatalog) + require.ErrorContains(t, err, fmt.Sprintf("%d offending entries", n)) + require.ErrorContains(t, err, fmt.Sprintf("future:%02d", censusMaxDetailed-1)) + require.NotContains(t, err.Error(), fmt.Sprintf("future:%02d", censusMaxDetailed)) +} + +// TestCensus_RefusalIsWriteFree pins that a refused Open leaves the store +// untouched: in particular it must NOT mint a secret into a foreign tree (the +// census runs before ensureSecret). Simulates a tree whose secret a newer +// binary relocated: a foreign key present, the secret key absent. +func TestCensus_RefusalIsWriteFree(t *testing.T) { + path := t.TempDir() + cat, err := openKVAt(t, path) + require.NoError(t, err) + require.NoError(t, cat.put("future:key", "x")) + require.NoError(t, cat.del(catalogSecretStoreKey)) + require.NoError(t, cat.Close()) + + _, err = openKVAt(t, path) + require.ErrorIs(t, err, ErrForeignCatalog) + + // Inspect the raw store: the secret key must still be absent. + store, err := rocksdb.New(rocksdb.Config{Path: path, Logger: silentLogger()}) + require.NoError(t, err) + t.Cleanup(func() { _ = store.Close() }) + _, found, err := store.Get("", []byte(catalogSecretStoreKey)) + require.NoError(t, err) + require.False(t, found, "the refused Open must not have minted a secret") +} diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/kv_test.go b/cmd/stellar-rpc/internal/rpcv2/catalog/kv_test.go index 81263ef98..78fdcffa4 100644 --- a/cmd/stellar-rpc/internal/rpcv2/catalog/kv_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/kv_test.go @@ -293,18 +293,19 @@ func TestKV_GracefulCloseAndReopen(t *testing.T) { first, err := openKVAt(t, path) require.NoError(t, err) - require.NoError(t, first.put("streaming:last_committed_ledger", "999")) - require.NoError(t, first.put("config:ledgers_per_tx_index", "100000")) - require.NoError(t, first.put("chunk:00000001:lfs", "1")) + // Keys must stay inside the census vocabulary — reopen refuses anything else. + require.NoError(t, first.put(geometry.ConfigEarliestLedger, "999")) + require.NoError(t, first.put(geometry.HotChunkKey(3), string(geometry.HotReady))) + require.NoError(t, first.put(geometry.ChunkKey(1, geometry.KindLedgers), string(geometry.StateFrozen))) require.NoError(t, first.Close()) second, err := openKVAt(t, path) require.NoError(t, err) t.Cleanup(func() { _ = second.Close() }) - assert.Equal(t, "999", kvGetHit(t, second, "streaming:last_committed_ledger")) - assert.Equal(t, "100000", kvGetHit(t, second, "config:ledgers_per_tx_index")) - assert.Equal(t, "1", kvGetHit(t, second, "chunk:00000001:lfs")) + assert.Equal(t, "999", kvGetHit(t, second, geometry.ConfigEarliestLedger)) + assert.Equal(t, string(geometry.HotReady), kvGetHit(t, second, geometry.HotChunkKey(3))) + assert.Equal(t, string(geometry.StateFrozen), kvGetHit(t, second, geometry.ChunkKey(1, geometry.KindLedgers))) } func TestKV_ConcurrentOpsAndCloseRaceFree(t *testing.T) { diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/secret.go b/cmd/stellar-rpc/internal/rpcv2/catalog/secret.go index 63f8c2175..ae825c722 100644 --- a/cmd/stellar-rpc/internal/rpcv2/catalog/secret.go +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/secret.go @@ -2,7 +2,6 @@ package catalog import ( "crypto/rand" - "fmt" ) // catalogSecretStoreKey holds the deployment's cold-index secret. @@ -17,7 +16,8 @@ const catalogSecretStoreKey = "meta/catalog-secret" func (c *Catalog) Secret() [32]byte { return c.secret } // ensureSecret loads the persisted cold-index secret, minting and persisting a -// fresh random one on first call. Open runs it single-threaded and caches the +// fresh random one on first call. Open runs it single-threaded, after the +// census has already validated any persisted value's width, and caches the // result; nothing else should call it (get-or-create is not atomic). func (c *Catalog) ensureSecret() ([32]byte, error) { var s [32]byte @@ -26,9 +26,6 @@ func (c *Catalog) ensureSecret() ([32]byte, error) { return s, err } if found { - if len(v) != len(s) { - return s, fmt.Errorf("persisted cold-index secret is %d bytes, want %d", len(v), len(s)) - } copy(s[:], v) return s, nil } diff --git a/cmd/stellar-rpc/internal/rpcv2/catalog/secret_test.go b/cmd/stellar-rpc/internal/rpcv2/catalog/secret_test.go index 50d8f9bd0..e338b4b2e 100644 --- a/cmd/stellar-rpc/internal/rpcv2/catalog/secret_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/catalog/secret_test.go @@ -29,7 +29,8 @@ func TestSecret_MintedOnceAndStable(t *testing.T) { // TestSecret_RejectsCorruptedPersisted pins that a wrong-length persisted // secret fails Open loudly instead of silently reminting: a truncated write or // downgraded encoding must not swap in a fresh secret, which would orphan every -// txhash .bin already keyed under the old one. +// txhash .bin already keyed under the old one. The census catches it (before +// the mint), and its message must not print the secret's value. func TestSecret_RejectsCorruptedPersisted(t *testing.T) { path := t.TempDir() @@ -40,5 +41,7 @@ func TestSecret_RejectsCorruptedPersisted(t *testing.T) { _, err = openKVAt(t, path) require.Error(t, err, "a corrupted persisted secret must fail Open, not remint") - require.ErrorContains(t, err, "persisted cold-index secret is 5 bytes") + require.ErrorIs(t, err, ErrForeignCatalog) + require.ErrorContains(t, err, "value is 5 bytes, want 32") + require.NotContains(t, err.Error(), "short", "the secret's value must never appear in the error") } diff --git a/cmd/stellar-rpc/internal/rpcv2/geometry/keys.go b/cmd/stellar-rpc/internal/rpcv2/geometry/keys.go index 2b67c0005..9e2fb1b8f 100644 --- a/cmd/stellar-rpc/internal/rpcv2/geometry/keys.go +++ b/cmd/stellar-rpc/internal/rpcv2/geometry/keys.go @@ -62,6 +62,43 @@ var allKinds = []Kind{KindLedgers, KindEvents, KindTxHash} // AllKinds returns the per-chunk artifact kinds in canonical order. func AllKinds() []Kind { return append([]Kind(nil), allKinds...) } +// allStates / allHotStates are the state-token registries, the single source +// of truth the census derives its vocabulary from — a token added here is +// automatically accepted, so the daemon can never refuse its own writes. +// +//nolint:gochecknoglobals // immutable state registries, single source of truth +var ( + allStates = []State{StateFreezing, StateFrozen, StatePruning} + allHotStates = []HotState{HotTransient, HotReady} +) + +// AllStates returns the artifact lifecycle states. +func AllStates() []State { return append([]State(nil), allStates...) } + +// AllHotStates returns the hot-DB lifecycle states. +func AllHotStates() []HotState { return append([]HotState(nil), allHotStates...) } + +// IsKnownState reports whether s is a registered artifact lifecycle state. +// The switch names every constant so a state added to the const block but not +// here is caught by the exhaustive linter, not by the census refusing the +// daemon's own catalog on the next restart. +func IsKnownState(s State) bool { + switch s { + case StateFreezing, StateFrozen, StatePruning: + return true + } + return false +} + +// IsKnownHotState reports whether s is a registered hot-DB lifecycle state. +func IsKnownHotState(s HotState) bool { + switch s { + case HotTransient, HotReady: + return true + } + return false +} + // TxHashIndexID identifies a tx-hash index: a contiguous run of // chunks_per_txhash_index chunks. Distinct type from chunk.ID (both uint32) so // index ids and chunk ids never silently interchange. diff --git a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go index c3334b807..edb3dfde4 100644 --- a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go +++ b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb.go @@ -56,6 +56,17 @@ func OpenSnapshots() int64 { return openSnapshots.Load() } const ( dirPerm os.FileMode = 0o700 defaultCFName = "default" + + // pinnedTableFormatVersion pins the block-based table format written to + // disk, so a grocksdb/librocksdb upgrade cannot silently change the on-disk + // format. 6 is the default in librocksdb 10.9.1 (the version + // scripts/install-rocksdb.sh builds), so the pin changes no byte + // today. RocksDB's own header advises leaving format_version at the + // default so improvements arrive automatically; we choose the opposite on + // purpose. Raising it is a format-touching change: an older binary cannot + // open the newer tables, so it must ship as a declared storage-format bump, + // never as a side effect of a dependency bump. + pinnedTableFormatVersion = 6 ) // Config — per-store knobs. Each Layer-2 facade owns one, built in @@ -121,10 +132,10 @@ type Store struct { // cache is the block cache shared across every CF in this store, // created in applyTuning when BlockCacheMB is set. bbtos are the - // per-CF block-based-table options (one per CF that has a cache, - // bloom filter, or block-size override); each may own a moved-in - // bloom filter. Both are destroyed in Close after opts/cfOpts, - // which hold C-side refs we must drop first. + // per-CF block-based-table options (one per CF, carrying the pinned + // table format version); each may own a moved-in bloom filter. Both + // are destroyed in Close after opts/cfOpts, which hold C-side refs + // we must drop first. cache *grocksdb.Cache bbtos []*grocksdb.BlockBasedTableOptions @@ -839,10 +850,10 @@ func applyDBTuning(opts *grocksdb.Options, t Tuning) { // override (when set). The Store retains every BBTO and the cache; // Close destroys them after opts/cfOpts. // -// A BBTO is installed on a CF iff the shared cache, that CF's bloom -// filter, or that CF's BlockSize override is configured — preserving -// the previous behavior of leaving RocksDB's default BBTO untouched -// when no table-level knob is set. +// Every CF gets an explicit BBTO, tuned or not, so the on-disk table +// format stays pinned (pinnedTableFormatVersion) instead of riding +// grocksdb's default; restoring a skip-when-untuned path would +// silently unpin it. // // The bloom filter is built per CF because SetFilterPolicy MOVES the // policy into the BBTO (it nils the source pointer), so a single @@ -854,10 +865,8 @@ func (s *Store) applySharedTableOptions(cfNames []string, cfOpts []*grocksdb.Opt } for i, o := range cfOpts { override := s.cfg.PerCFOptions[cfNames[i]] - if s.cache == nil && override.BloomFilterBitsPerKey == 0 && override.BlockSize == 0 { - continue - } bbto := grocksdb.NewDefaultBlockBasedTableOptions() + bbto.SetFormatVersion(pinnedTableFormatVersion) if s.cache != nil { bbto.SetBlockCache(s.cache) } diff --git a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb_test.go b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb_test.go index 750ae9c20..11a9162a8 100644 --- a/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/rocksdb/rocksdb_test.go @@ -762,7 +762,7 @@ func TestStore_TuningZeroValue(t *testing.T) { t.Cleanup(func() { _ = s.Close() }) assert.Nil(t, s.cache) - assert.Nil(t, s.bbtos) + assert.Len(t, s.bbtos, 1, "every CF carries a BBTO even with zero tuning") require.NoError(t, s.Put(defaultCFName, []byte("k"), []byte("v"))) v, found, err := s.Get(defaultCFName, []byte("k")) diff --git a/cmd/stellar-rpc/internal/rpcv2/rpcv2test/rpcv2test.go b/cmd/stellar-rpc/internal/rpcv2/rpcv2test/rpcv2test.go index 5ec828499..3499995a6 100644 --- a/cmd/stellar-rpc/internal/rpcv2/rpcv2test/rpcv2test.go +++ b/cmd/stellar-rpc/internal/rpcv2/rpcv2test/rpcv2test.go @@ -333,13 +333,17 @@ func ReadColdBin(t *testing.T, path string) []txhash.ColdEntry { require.NoError(t, err) defer f.Close() - const headerSize = 8 + stores.SecretLen // uint64-LE count + index secret + // Prelude (magic "SBIN", version, 3 reserved) + uint64-LE count + secret. + const preludeSize = 8 + const headerSize = preludeSize + 8 + stores.SecretLen entrySize := txhash.ColdKeySize + 4 var header [headerSize]byte _, err = io.ReadFull(f, header[:]) require.NoError(t, err) - count := binary.LittleEndian.Uint64(header[:8]) + require.Equal(t, []byte("SBIN"), header[:4], "cold .bin magic") + require.Equal(t, byte(1), header[4], "cold .bin version") + count := binary.LittleEndian.Uint64(header[preludeSize : preludeSize+8]) info, err := f.Stat() require.NoError(t, err) diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/blobversion.go b/cmd/stellar-rpc/internal/rpcv2/stores/blobversion.go new file mode 100644 index 000000000..5be6b064b --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/blobversion.go @@ -0,0 +1,22 @@ +package stores + +import ( + "errors" + "fmt" +) + +// CheckBlobVersion gates a blob that leads with its own version byte: an +// empty blob and an unrecognized version each refuse, the version checked +// before any length so a differently-sized newer-format blob reports as an +// upgrade problem rather than a corruption-shaped size mismatch. Callers +// keep their own length checks on the remainder. +func CheckBlobVersion(blob []byte, want byte) error { + if len(blob) == 0 { + return errors.New("empty blob") + } + if blob[0] != want { + return fmt.Errorf("unsupported version 0x%02x, want 0x%02x (written by a newer stellar-rpc?)", + blob[0], want) + } + return nil +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go index 2c6bffecb..da09a6b99 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format.go @@ -95,6 +95,71 @@ const indexPackChecksum = packfile.ChecksumCRC32C // false positives before deserializing the bitmap. const IndexRecordFingerprintLen = 4 +// ────────────────────────────────────────────────────────────────── +// index.pack build stamp. +// +// Embedded in index.pack's app-data slot: +// +// offset size field +// 0 1 version (0x01) +// 1 2 term schema version (uint16 BE) +// 3 8 indexed-field bitmask (uint64 BE) +// +// The stamp records which term-derivation scheme and field set the +// index was built under, making the artifact self-describing: an +// index missing a term family becomes distinguishable from one that +// simply matched nothing. Freeze and walk write identical stamps +// (all three values are compile-time constants), so freeze-vs-walk +// byte identity is unaffected. Decoding ignores trailing bytes so a +// future version can extend the blob without moving these fields. +// ────────────────────────────────────────────────────────────────── + +const ( + indexStampVersion byte = 0x01 + indexStampLen = 1 + 2 + 8 +) + +func encodeIndexBuildStamp() []byte { + buf := make([]byte, indexStampLen) + buf[0] = indexStampVersion + binary.BigEndian.PutUint16(buf[1:3], TermSchemaVersion) + binary.BigEndian.PutUint64(buf[3:11], IndexedFieldMask()) + return buf +} + +// checkIndexBuildStamp refuses an index.pack whose build stamp names a term +// schema or field set other than this binary's own. +func checkIndexBuildStamp(indexPackPath string, r *stores.PackReader) error { + ad, err := r.AppData() + if err != nil { + return fmt.Errorf("events: read build stamp of %s: %w", indexPackPath, err) + } + schema, mask, err := decodeIndexBuildStamp(ad) + if err != nil { + return fmt.Errorf("events: %s: %w", indexPackPath, err) + } + if schema != TermSchemaVersion || mask != IndexedFieldMask() { + return fmt.Errorf( + "events: %s was built under term schema %d with field mask %#x; this binary expects "+ + "schema %d with mask %#x (rebuilt index required, or a binary matching the artifact)", + indexPackPath, schema, mask, TermSchemaVersion, IndexedFieldMask()) + } + return nil +} + +// decodeIndexBuildStamp recovers (termSchema, fieldMask) from an index.pack +// app-data blob, rejecting a short blob or an unknown stamp version. Bytes +// past the stamp are ignored (future extension room). +func decodeIndexBuildStamp(data []byte) (uint16, uint64, error) { + if err := stores.CheckBlobVersion(data, indexStampVersion); err != nil { + return 0, 0, fmt.Errorf("events: index.pack build stamp: %w", err) + } + if len(data) < indexStampLen { + return 0, 0, fmt.Errorf("events: index.pack build stamp is %d bytes, want at least %d", len(data), indexStampLen) + } + return binary.BigEndian.Uint16(data[1:3]), binary.BigEndian.Uint64(data[3:11]), nil +} + // ────────────────────────────────────────────────────────────────── // events.pack record codec. // ────────────────────────────────────────────────────────────────── @@ -160,12 +225,9 @@ const LedgerOffsetsFormatVersion byte = 0x01 const ledgerOffsetsHeaderLen = 1 + 4 + 4 -// ErrUnknownLedgerOffsetsVersion is returned when decoding app data -// whose leading version byte isn't recognized by this binary. -var ErrUnknownLedgerOffsetsVersion = errors.New("events: unknown LedgerOffsets format version") - -// ErrShortLedgerOffsets is returned when the app data buffer is -// shorter than the declared header or trailing cumulative array. +// ErrShortLedgerOffsets is returned when the app data buffer is empty, +// carries an unknown version byte, or is shorter than the declared header +// or trailing cumulative array. var ErrShortLedgerOffsets = errors.New("events: LedgerOffsets app data too short") // encodeLedgerOffsets serializes o for packfile app-data embedding. @@ -190,12 +252,12 @@ func encodeLedgerOffsets(o *LedgerOffsets) ([]byte, error) { // encodeLedgerOffsets back into a *LedgerOffsets. Used by the cold // reader (PR-3a). func DecodeLedgerOffsets(data []byte) (*LedgerOffsets, error) { + if err := stores.CheckBlobVersion(data, LedgerOffsetsFormatVersion); err != nil { + return nil, fmt.Errorf("%w: %w", ErrShortLedgerOffsets, err) + } if len(data) < ledgerOffsetsHeaderLen { return nil, ErrShortLedgerOffsets } - if data[0] != LedgerOffsetsFormatVersion { - return nil, fmt.Errorf("%w: 0x%02x", ErrUnknownLedgerOffsetsVersion, data[0]) - } startLedger := binary.BigEndian.Uint32(data[1:5]) n := binary.BigEndian.Uint32(data[5:9]) expected := ledgerOffsetsHeaderLen + int(n)*4 @@ -269,12 +331,12 @@ func encodeEventsMeta(secret [stores.SecretLen]byte) []byte { func decodeEventsMeta(data []byte) ([stores.SecretLen]byte, error) { var secret [stores.SecretLen]byte + if err := stores.CheckBlobVersion(data, eventsMetaVersion); err != nil { + return secret, fmt.Errorf("%w: %w", errBadIndexMetadata, err) + } if len(data) != eventsMetaLen { return secret, fmt.Errorf("%w: %d bytes, want %d", errBadIndexMetadata, len(data), eventsMetaLen) } - if data[0] != eventsMetaVersion { - return secret, fmt.Errorf("%w: unknown version 0x%02x", errBadIndexMetadata, data[0]) - } copy(secret[:], data[1:]) return secret, nil } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format_test.go index 9ccf37225..fbd335040 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_format_test.go @@ -65,7 +65,9 @@ func TestLedgerOffsets_DecodeRejectsShortBuffer(t *testing.T) { _, err := DecodeLedgerOffsets(nil) require.ErrorIs(t, err, ErrShortLedgerOffsets) - _, err = DecodeLedgerOffsets(make([]byte, ledgerOffsetsHeaderLen-1)) + short := make([]byte, ledgerOffsetsHeaderLen-1) + short[0] = LedgerOffsetsFormatVersion // valid version, so the length check fires + _, err = DecodeLedgerOffsets(short) assert.ErrorIs(t, err, ErrShortLedgerOffsets) } @@ -73,7 +75,7 @@ func TestLedgerOffsets_DecodeRejectsUnknownVersion(t *testing.T) { buf := make([]byte, ledgerOffsetsHeaderLen) buf[0] = 0xff // not LedgerOffsetsFormatVersion _, err := DecodeLedgerOffsets(buf) - assert.ErrorIs(t, err, ErrUnknownLedgerOffsetsVersion) + assert.ErrorContains(t, err, "written by a newer stellar-rpc") } func TestLedgerOffsets_DecodeRejectsTruncatedArray(t *testing.T) { diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index.go index 6f4c959fa..c3a38b5cb 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index.go @@ -163,6 +163,7 @@ func WriteColdIndex( pw, err := packfile.Create(indexPackPath, packfile.WriterOptions{ Format: indexPackFormat, ItemsPerRecord: indexPackItemsPerRecord, + ContentHash: true, Overwrite: true, RecordChecksum: indexPackChecksum, }) @@ -208,5 +209,5 @@ func writeIndexPackEntries(pw *packfile.Writer, entries []indexEntry) error { return fmt.Errorf("events: write slot %d to index.pack: %w", e.slot, err) } } - return pw.Finish(nil) + return pw.Finish(encodeIndexBuildStamp()) } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index_test.go index a0d5eee0f..cc7c034c5 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_index_test.go @@ -336,3 +336,35 @@ func TestWriteIndex_RecordEncoding(t *testing.T) { // just to lock the endianness contract. _ = binary.LittleEndian.Uint32(record[:IndexRecordFingerprintLen]) } + +// TestWriteColdIndex_StampAndContentHash pins index.pack's app-data build +// stamp (schema, field mask, trailing-bytes-ignored, newer-version refusal) +// and its content hash. +func TestWriteColdIndex_StampAndContentHash(t *testing.T) { + dir := t.TempDir() + require.NoError(t, WriteColdIndex(context.Background(), indexTestChunkID, indexFixture(t, 4), dir, testIndexSecret)) + r := packfile.Open(filepath.Join(dir, IndexPackName(indexTestChunkID)), packfile.ReaderOptions{}) + t.Cleanup(func() { _ = r.Close() }) + + ad, err := r.AppData() + require.NoError(t, err) + schema, mask, err := decodeIndexBuildStamp(ad) + require.NoError(t, err) + assert.Equal(t, TermSchemaVersion, schema) + assert.Equal(t, IndexedFieldMask(), mask) + + // Bytes past the stamp are extension room: the decoder ignores them. + _, _, err = decodeIndexBuildStamp(append(append([]byte(nil), ad...), 0xAB, 0xCD)) + require.NoError(t, err) + + // An unknown stamp version refuses with the newer-binary hint. + newer := append([]byte(nil), ad...) + newer[0] = indexStampVersion + 1 + _, _, err = decodeIndexBuildStamp(newer) + require.ErrorContains(t, err, "written by a newer stellar-rpc") + + _, hashed, err := r.ContentHash() + require.NoError(t, err) + assert.True(t, hashed, "index.pack carries a content hash") + require.NoError(t, r.Verify(context.Background())) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go index f24760141..9c0a3efa9 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader.go @@ -175,9 +175,11 @@ func OpenColdReader(chunkID chunk.ID, bucketDir string, opts ColdReaderOptions) if err != nil { return err } - // index.pack's own integrity is checked for every chunk, eventless or - // not: a foreign or unchecked pack is not servable whatever its term - // count, and the open is already in flight either way. + // Format, record-checksum, and build-stamp checks run for eventless + // chunks too, so a foreign, unchecked, or mis-schemed pack refuses on the + // first indexed lookup regardless of chunk content. The guarantee is + // lookup-path-only by design: payload reads (FetchEvents, All) never + // consult the index pair and stay valid against events.pack's own checks. tr, terr := c.index.Trailer() if terr != nil { return fmt.Errorf("events: open %s: %w", indexPackPath, terr) @@ -194,6 +196,9 @@ func OpenColdReader(chunkID chunk.ID, bucketDir string, opts ColdReaderOptions) return fmt.Errorf("%w: %s: built without a record checksum (stale build)", stores.ErrCorrupt, indexPackPath) } + if err := checkIndexBuildStamp(indexPackPath, c.index); err != nil { + return err + } if idx.isEmpty() { // A zero-term index is only valid for an eventless chunk: cross-check // events.pack's count so a mispaired empty index fails loudly instead @@ -212,10 +217,11 @@ func OpenColdReader(chunkID chunk.ID, bucketDir string, opts ColdReaderOptions) // Non-empty index: bind the pair to this chunk before serving from // it — index.pack/index.hash carry no chunk ID of their own, so a // mispaired index would silently return an incomplete subset of - // matches. Two cheap checks on top of the pack's own validation - // above: index.hash keys == index.pack records (halves of one - // build), and non-empty index ⇒ non-empty events.pack (converse of - // the empty-index check below). + // matches. Two cheap checks beyond the shared format, checksum, and + // stamp gates above: + // index.hash keys == index.pack records (halves of one build), and + // non-empty index ⇒ non-empty events.pack (converse of the + // empty-index check above). if uint64(tr.TotalItems) != idx.numKeys() { return fmt.Errorf( "events: index pair mismatch for chunk %s: index.hash holds %d keys "+ diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_test.go index a53f29fe8..5e504f631 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_reader_test.go @@ -942,3 +942,53 @@ func TestColdReader_UncheckedIndexPackIsCorrupt(t *testing.T) { // so a bare ErrCorrupt assertion would survive the guard's removal. require.Contains(t, err.Error(), "built without a record checksum") } + +// rewriteIndexPackStamp replaces the chunk's index.pack with a zero-record +// pack carrying the given build stamp. The stamp gate runs before every +// count cross-check, so this exercises the schema/mask refusal on both +// eventful and eventless chunks. +func rewriteIndexPackStamp(t *testing.T, dir string, chunkID chunk.ID, schema uint16, mask uint64) { + t.Helper() + stamp := make([]byte, indexStampLen) + stamp[0] = indexStampVersion + binary.BigEndian.PutUint16(stamp[1:3], schema) + binary.BigEndian.PutUint64(stamp[3:11], mask) + pw, err := packfile.Create(filepath.Join(dir, IndexPackName(chunkID)), packfile.WriterOptions{ + Format: indexPackFormat, + ItemsPerRecord: indexPackItemsPerRecord, + RecordChecksum: indexPackChecksum, + Overwrite: true, + }) + require.NoError(t, err) + require.NoError(t, pw.Finish(stamp)) +} + +// TestColdReader_RejectsMismatchedBuildStamp pins the capability gate: an +// index built under a different term schema or field set must refuse at the +// lookup path, on eventful and eventless chunks alike, rather than serve +// as if it could answer this binary's filters. +func TestColdReader_RejectsMismatchedBuildStamp(t *testing.T) { + cases := []struct { + name string + schema uint16 + mask uint64 + }{ + {"schema", TermSchemaVersion + 1, IndexedFieldMask()}, + {"mask", TermSchemaVersion, IndexedFieldMask() | 1<<63}, + } + for _, eventsPerLedger := range []int{2, 0} { + for _, tc := range cases { + t.Run(fmt.Sprintf("eventsPerLedger=%d/%s", eventsPerLedger, tc.name), func(t *testing.T) { + const chunkID = chunk.ID(0) + dir, _ := buildColdFixture(t, chunkID, eventsPerLedger, 2) + rewriteIndexPackStamp(t, dir, chunkID, tc.schema, tc.mask) + + cr, err := OpenColdReader(chunkID, dir, ColdReaderOptions{}) + require.NoError(t, err) + t.Cleanup(func() { _ = cr.Close() }) + _, err = cr.LookupKeys(context.Background(), []TermKey{{1}}) + require.ErrorContains(t, err, "was built under term schema") + }) + } + } +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer.go index 8d1491f04..fa540be7e 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer.go @@ -95,9 +95,12 @@ func NewColdWriter(chunkID chunk.ID, bucketDir string, opts ColdWriterOptions) ( Format: eventsPackFormat, ItemsPerRecord: eventsPackItemsPerRecord, NewRecordEncoder: newEventsPackEncoder, - Concurrency: opts.Concurrency, - BytesPerSync: opts.BytesPerSync, - Overwrite: true, + // Items reach AppendItem in canonical (uncompressed) payload form, so + // the content hash is independent of the zstd encoder version. + ContentHash: true, + Concurrency: opts.Concurrency, + BytesPerSync: opts.BytesPerSync, + Overwrite: true, }) if err != nil { return nil, fmt.Errorf("events: create events.pack at %s: %w", path, err) diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer_test.go index 8ad80fca9..662b45254 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/cold_writer_test.go @@ -279,3 +279,19 @@ func TestEventsPack_TrailerPinsFormatAndRecordSize(t *testing.T) { assert.Equal(t, uint32(eventsPackItemsPerRecord), tr.ItemsPerRecord, "events.pack ItemsPerRecord must match eventsPackItemsPerRecord constant") } + +// TestColdWriter_EventsPackContentHash pins that events.pack carries a +// content hash over the canonical payload bytes and that it verifies through +// the production decoder, matching the ledger and index.pack assertions. +func TestColdWriter_EventsPackContentHash(t *testing.T) { + const chunkID = chunk.ID(0) + dir, _ := buildColdFixture(t, chunkID, 3, 1) + + r := packfile.Open(filepath.Join(dir, EventsPackName(chunkID)), + packfile.ReaderOptions{RecordDecoder: eventsPackDecoder}) + t.Cleanup(func() { _ = r.Close() }) + _, hashed, err := r.ContentHash() + require.NoError(t, err) + require.True(t, hashed, "events.pack carries a content hash") + require.NoError(t, r.Verify(context.Background())) +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/format_golden_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/format_golden_test.go new file mode 100644 index 000000000..0549426a5 --- /dev/null +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/format_golden_test.go @@ -0,0 +1,119 @@ +package event + +// Byte-level goldens for the term-key derivation. The term keys are on-disk +// format: they are the identities every frozen index.pack/index.hash is built +// over, so any change to the hash function, the field-byte prefix, or a +// field's value encoding silently invalidates every existing index. These +// fixtures make such a change fail CI with "bytes changed: bump +// TermSchemaVersion (and the artifact format) or revert" instead of relying +// on review to notice. (The shared blinding primitives are pinned at their +// own level, in stores/blind_test.go.) + +import ( + "encoding/hex" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/stellar/go-stellar-sdk/xdr" +) + +func hexKey(t *testing.T, k TermKey) string { + t.Helper() + return hex.EncodeToString(k[:]) +} + +func TestComputeTermKey_Golden(t *testing.T) { + val := []byte("stellar-rpc-term-golden") + want := map[Field]string{ + FieldContractID: "68828a4d63340ec7a182c7bacec54f0a", + FieldTopic0: "fe2646f2447419a4c3b5899d18f19517", + FieldTopic1: "5d3c6660724f4267f75f0314a8bd99bc", + FieldTopic2: "924913135dfd660d50dd75352d8b4bb8", + FieldTopic3: "9a981d29f52a36df9c3d9fa2d19ffc91", + FieldEventType: "b8cba0b1d3c82ec9ba18a11d780d0d23", + FieldTopicCount: "8188dbfd6145c54e6cb442dc02b3c324", + } + require.Len(t, want, len(allFields), "extend the golden map when adding a Field") + for _, f := range allFields { + w, pinned := want[f] + require.True(t, pinned, "field %d has no golden pin; add one when adding a Field", f) + assert.Equal(t, w, hexKey(t, ComputeTermKey(val, f)), "field %d", f) + } +} + +func TestFieldTermKeys_Golden(t *testing.T) { + assert.Equal(t, "092967d65ec25b057c6a54e627dec4e6", + hexKey(t, EventTypeTermKey(xdr.ContractEventTypeContract))) + + wantCounts := []string{ + "5a21d2e3bd688e548ac939e4babaeb9b", + "8e389795ce9351bb0abbe6c00fff6052", + "a68396e0abe112008816c89412c9113b", + "291796ea36f7a80a1d7515faf2a22ec7", + "db3d5df76fa2f6cd611dde2b96f16401", + "70c2b104db4b2ba115eeb4954daac675", // the overflow bucket + } + for n, w := range wantCounts { + assert.Equal(t, w, hexKey(t, TopicCountTermKey(n)), "topic count %d", n) + } + assert.Equal(t, TopicCountTermKey(topicCountOverflowBucket), TopicCountTermKey(99), + "counts past the overflow bucket clamp into it") +} + +// goldenContractEventBytes builds the fixed marshaled ContractEvent the +// TermsForBytes golden runs over: a contract ID of bytes 0..31 and four +// topics exercising every indexed topic position. +func goldenContractEventBytes(t *testing.T) []byte { + t.Helper() + var cid xdr.ContractId + for i := range cid { + cid[i] = byte(i) + } + sym0, sym2 := xdr.ScSymbol("transfer"), xdr.ScSymbol("to") + u1, u3 := xdr.Uint32(7), xdr.Uint32(9) + ev := xdr.ContractEvent{ + ContractId: &cid, + Type: xdr.ContractEventTypeContract, + Body: xdr.ContractEventBody{ + V: 0, + V0: &xdr.ContractEventV0{ + Topics: []xdr.ScVal{ + {Type: xdr.ScValTypeScvSymbol, Sym: &sym0}, + {Type: xdr.ScValTypeScvU32, U32: &u1}, + {Type: xdr.ScValTypeScvSymbol, Sym: &sym2}, + {Type: xdr.ScValTypeScvU32, U32: &u3}, + }, + Data: xdr.ScVal{Type: xdr.ScValTypeScvU32, U32: &u1}, + }, + }, + } + b, err := ev.MarshalBinary() + require.NoError(t, err) + return b +} + +// TestTermsForBytes_Golden pins the full derivation from a marshaled event: +// not just the hash and field prefix (TestComputeTermKey_Golden's job) but +// the value encoding of every field TermsForBytes extracts, so a change to +// how a contract ID or topic becomes hash input cannot pass unnoticed. +// Cross-anchors: key 0 equals the eventType(contract) golden and key 2 the +// topicCount(4) golden above. +func TestTermsForBytes_Golden(t *testing.T) { + keys, err := TermsForBytes(goldenContractEventBytes(t)) + require.NoError(t, err) + want := []string{ + "092967d65ec25b057c6a54e627dec4e6", // event type (contract) + "f96521b095e9249b0af1546bfbd9a80f", // contract ID 0..31 + "db3d5df76fa2f6cd611dde2b96f16401", // topic count 4 + "8b9390866f0062b43022a0896fe6b606", // topic 0: symbol "transfer" + "de4348720f97c94df5621c5cc5304a3d", // topic 1: u32 7 + "9ea0933498a23b8f6c3d867631e945c8", // topic 2: symbol "to" + "dacc6ffd833f6932052ef2cd55f75844", // topic 3: u32 9 + } + require.Len(t, keys, len(want)) + for i, k := range keys { + assert.Equal(t, want[i], hex.EncodeToString(k[:]), "key %d", i) + } +} diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/event/index.go b/cmd/stellar-rpc/internal/rpcv2/stores/event/index.go index a897d1c13..189f8423c 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/event/index.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/event/index.go @@ -26,6 +26,42 @@ const ( FieldTopicCount Field = 6 ) +// TermSchemaVersion names the term-derivation scheme: the hash function and +// byte encoding behind ComputeTermKey plus each field's value encoding. Bump +// it whenever any of those change. Adding a whole field sets a new bit in +// IndexedFieldMask instead and is a storage-format change in its own right; +// see the design doc's "Format versioning and upgrades" section. Both values +// are recorded in every index.pack's build stamp, and the reader accepts +// exactly its own pair. +const TermSchemaVersion uint16 = 1 + +// allFields is the field registry, the single source of truth the mask and +// the golden tests derive from. Extend it in the same change that adds a +// Field constant. +// +//nolint:gochecknoglobals // immutable field registry, single source of truth +var allFields = []Field{ + FieldContractID, FieldTopic0, FieldTopic1, FieldTopic2, FieldTopic3, + FieldEventType, FieldTopicCount, +} + +// indexedFieldMask is derived from allFields once at init. It is unexported so +// no importer can reassign the value every index.pack build stamp records. +// +//nolint:gochecknoglobals // derived from the immutable field registry +var indexedFieldMask = func() uint64 { + var m uint64 + for _, f := range allFields { + m |= 1 << f + } + return m +}() + +// IndexedFieldMask is the set of indexed fields as a bitmask (bit i set means +// Field i is indexed). The build stamp records it per artifact and the +// term-key golden tests iterate the same registry. +func IndexedFieldMask() uint64 { return indexedFieldMask } + // ComputeTermKey computes a 16-byte term key by hashing the field byte // followed by the value bytes: xxh3_128(field || value), encoded as // two little-endian uint64s. diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader.go b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader.go index 8a5648146..7c7020d55 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader.go @@ -32,10 +32,14 @@ func MissingPackOpens() uint64 { return missingPackOpens.Load() } // store. Shared by the reader and the writer (same package). const formatLedgerCold packfile.Format = 1 -// appDataSize — firstSeq (4 BE). lastSeq is derived from -// trailer.TotalItems at open. Shared by the reader and the writer -// (same package). -const appDataSize = 4 +// AppData layout: a leading version byte, then firstSeq (4 BE). +// lastSeq is derived from trailer.TotalItems at open. Shared by the +// reader and the writer (same package). Every app-data blob leads +// with its own version byte so it is self-describing on its own, +// independent of the trailer Format that names the whole encoding. +const coldAppDataVersion byte = 0x01 + +const appDataSize = 1 + 4 // version byte + firstSeq (uint32 BE) // coldPackDecoder is the process-wide zstd decoder for cold ledger // pack records. packfile.RecordDecoder must be concurrent-safe and @@ -107,10 +111,13 @@ func (c *ColdReader) loadHeader() (coldHeader, error) { if err != nil { return coldHeader{}, fmt.Errorf("cold: read AppData %q: %w", c.path, err) } + if err := stores.CheckBlobVersion(ad, coldAppDataVersion); err != nil { + return coldHeader{}, fmt.Errorf("cold %q: AppData: %w", c.path, err) + } if len(ad) != appDataSize { return coldHeader{}, fmt.Errorf("cold %q: expected %d-byte AppData, got %d", c.path, appDataSize, len(ad)) } - first := binary.BigEndian.Uint32(ad) + first := binary.BigEndian.Uint32(ad[1:]) if uint64(first)+uint64(tr.TotalItems)-1 > math.MaxUint32 { return coldHeader{}, fmt.Errorf( "cold %q: lastSeq overflows uint32 (firstSeq=%d, items=%d)", diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader_test.go index dba5642be..6b08e3d35 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_reader_test.go @@ -175,22 +175,55 @@ func TestColdReader_LazyOpen(t *testing.T) { } func TestColdReader_RejectsWrongAppDataSize(t *testing.T) { - path := filepath.Join(t.TempDir(), "bad-appdata.pack") + // Every blob carries a valid version byte, so the size check (which runs + // after the version check) is what fires. appDataSize is 5: too long and + // too short both refuse, and the 1-byte case pins that the firstSeq read + // never runs on a blob that only holds the version byte. + for _, tc := range []struct { + name string + ad []byte + }{ + {"too long", []byte{coldAppDataVersion, 's', 'e', 'v', 'e', 'n', 'b'}}, + {"version byte only", []byte{coldAppDataVersion}}, + } { + t.Run(tc.name, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "bad-appdata.pack") + pw, err := packfile.Create(path, packfile.WriterOptions{ + ItemsPerRecord: 1, + Format: formatLedgerCold, + }) + require.NoError(t, err) + require.NoError(t, pw.AppendItem([]byte("v"))) + require.NoError(t, pw.Finish(tc.ad)) + + c, err := OpenColdReader(path) + require.NoError(t, err) + t.Cleanup(func() { _ = c.Close() }) + _, err = c.LastSeq() + require.Error(t, err) + assert.Contains(t, err.Error(), "AppData") + }) + } +} + +func TestColdReader_RejectsNewerAppDataVersion(t *testing.T) { + path := filepath.Join(t.TempDir(), "newer-appdata.pack") pw, err := packfile.Create(path, packfile.WriterOptions{ ItemsPerRecord: 1, Format: formatLedgerCold, }) require.NoError(t, err) require.NoError(t, pw.AppendItem([]byte("v"))) - // 7-byte payload — appDataSize is 4. - require.NoError(t, pw.Finish([]byte("seven-b"))) + // A longer blob under an unknown version byte must report as a version + // problem, not a size mismatch. + require.NoError(t, pw.Finish([]byte{coldAppDataVersion + 1, 0, 0, 0, 2, 9})) c, err := OpenColdReader(path) require.NoError(t, err) t.Cleanup(func() { _ = c.Close() }) _, err = c.LastSeq() require.Error(t, err) - assert.Contains(t, err.Error(), "AppData") + assert.Contains(t, err.Error(), "written by a newer stellar-rpc") } func TestColdReader_RejectsWrongFormat(t *testing.T) { @@ -235,7 +268,8 @@ func TestColdReader_LastSeqOverflowRejected(t *testing.T) { require.NoError(t, pw.AppendItem([]byte{byte(i)})) } var ad [appDataSize]byte - binary.BigEndian.PutUint32(ad[:], math.MaxUint32-1) // firstSeq = MaxUint32-1, items=3 → overflow + ad[0] = coldAppDataVersion + binary.BigEndian.PutUint32(ad[1:], math.MaxUint32-1) // firstSeq = MaxUint32-1, items=3 → overflow require.NoError(t, pw.Finish(ad[:])) c, err := OpenColdReader(path) diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer.go b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer.go index 9907374e2..99b0950ac 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer.go @@ -69,8 +69,11 @@ func NewColdWriter(path string, firstSeq uint32, opts ColdWriterOptions) (*ColdW Format: formatLedgerCold, Overwrite: true, NewRecordEncoder: newColdPackEncoder, - Concurrency: opts.Concurrency, - BytesPerSync: opts.BytesPerSync, + // Items reach AppendItem as raw LCM bytes, so the content hash is + // independent of the zstd encoder version. + ContentHash: true, + Concurrency: opts.Concurrency, + BytesPerSync: opts.BytesPerSync, }) if err != nil { return nil, fmt.Errorf("cold: create packfile %q: %w", path, err) @@ -105,7 +108,8 @@ func (w *ColdWriter) Commit() error { return fmt.Errorf("cold %q: commit with no appends", w.path) } var ad [appDataSize]byte - binary.BigEndian.PutUint32(ad[:], w.firstSeq) + ad[0] = coldAppDataVersion + binary.BigEndian.PutUint32(ad[1:], w.firstSeq) if err := w.pw.Finish(ad[:]); err != nil { return translateWriterErr(err) } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer_test.go index 1d5df0dc1..34245aa88 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/ledger/cold_writer_test.go @@ -1,6 +1,7 @@ package ledger import ( + "context" "encoding/binary" "fmt" "os" @@ -93,7 +94,9 @@ func TestColdWriter_CommitEmitsTrailerAndAppData(t *testing.T) { } require.NoError(t, w.Commit()) - r := packfile.Open(path, packfile.ReaderOptions{}) + // The record decoder matters here: Verify recomputes the content hash over + // decoded (raw) items, the form the writer hashed. + r := packfile.Open(path, packfile.ReaderOptions{RecordDecoder: coldPackDecoder}) t.Cleanup(func() { _ = r.Close() }) total, err := r.TotalItems() @@ -103,7 +106,15 @@ func TestColdWriter_CommitEmitsTrailerAndAppData(t *testing.T) { ad, err := r.AppData() require.NoError(t, err) require.Len(t, ad, appDataSize) - assert.Equal(t, firstSeq, binary.BigEndian.Uint32(ad)) + assert.Equal(t, coldAppDataVersion, ad[0]) + assert.Equal(t, firstSeq, binary.BigEndian.Uint32(ad[1:])) + + // The pack carries a content hash over the raw (pre-compression) items, + // and it verifies. + _, hashed, err := r.ContentHash() + require.NoError(t, err) + assert.True(t, hashed, "cold packs carry a content hash") + require.NoError(t, r.Verify(context.Background())) } func TestNewColdWriter_TruncatesPreexistingFile(t *testing.T) { diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin.go index 85fac628e..9827b1ac3 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin.go @@ -10,7 +10,10 @@ package txhash // // File layout: // -// header uint64 LE entry count +// header uint32 LE magic ("SBIN" in on-disk byte order) +// uint8 version (1) +// 3 B reserved (zero) +// uint64 LE entry count // stores.SecretLen B index secret the keys were blinded with // entry ColdKeySize B blinded txhash[:ColdKeySize] // uint32 LE absolute ledger seq @@ -45,15 +48,41 @@ const ( // coldBinEntrySize is the per-entry width in the cold .bin file: // ColdKeySize bytes of blinded key + the ledger seq. coldBinEntrySize = ColdKeySize + coldBinSeqSize - // coldBinCountSize is the leading uint64 LE entry count. + // coldBinMagic identifies a cold txhash .bin; a mis-pointed or foreign + // file fails the header scan on it. + coldBinMagic = "SBIN" + // coldBinVersion is the .bin format version; the header scan rejects + // files written by a newer binary rather than misreading them. + coldBinVersion byte = 1 + // coldBinPreludeSize is the magic + version + 3 reserved zero bytes. + coldBinPreludeSize = 8 + // coldBinCountSize is the uint64 LE entry count after the prelude. coldBinCountSize = 8 - // coldBinHeaderSize is the count followed by the index secret the keys were - // blinded with (stores.SecretLen). The build reads the secret back and - // adopts it, so an index can never be built under a secret that disagrees - // with the one its .bin keys were keyed with (see BuildColdIndex). - coldBinHeaderSize = coldBinCountSize + stores.SecretLen + // coldBinHeaderSize is the prelude, the count, and the index secret the + // keys were blinded with (stores.SecretLen). The build reads the secret + // back and adopts it, so an index can never be built under a secret that + // disagrees with the one its .bin keys were keyed with (see BuildColdIndex). + coldBinHeaderSize = coldBinPreludeSize + coldBinCountSize + stores.SecretLen ) +// checkBinPrelude validates the magic, version, and reserved bytes at the +// start of a .bin header. Both readers call it, so a format bump cannot leave +// one of them reading a newer header as entry bytes. +func checkBinPrelude(path string, hdr []byte) error { + if magic := string(hdr[:4]); magic != coldBinMagic { + return fmt.Errorf("txhash: %s is not a cold txhash .bin (magic %q, want %q)", path, magic, coldBinMagic) + } + if hdr[4] != coldBinVersion { + return fmt.Errorf("txhash: %s has .bin version %d unsupported by this binary "+ + "(written by a newer stellar-rpc?)", path, hdr[4]) + } + if hdr[5]|hdr[6]|hdr[7] != 0 { + return fmt.Errorf("txhash: %s has reserved header bytes set "+ + "(written by a newer stellar-rpc, or corrupted)", path) + } + return nil +} + // ColdEntry is one (blinded key, ledger seq) tuple in a cold .bin file. type ColdEntry struct { Key [ColdKeySize]byte @@ -102,8 +131,10 @@ func WriteColdBin(path string, secret [stores.SecretLen]byte, entries []ColdEntr bw := bufio.NewWriterSize(f, 1<<20) var header [coldBinHeaderSize]byte - binary.LittleEndian.PutUint64(header[:coldBinCountSize], uint64(len(entries))) - copy(header[coldBinCountSize:], secret[:]) + copy(header[:4], coldBinMagic) + header[4] = coldBinVersion + binary.LittleEndian.PutUint64(header[coldBinPreludeSize:], uint64(len(entries))) + copy(header[coldBinPreludeSize+coldBinCountSize:], secret[:]) if _, werr := bw.Write(header[:]); werr != nil { return fmt.Errorf("txhash: write header: %w", werr) } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin_test.go index 15ce40409..c3a71a715 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_bin_test.go @@ -38,7 +38,7 @@ func readColdBin(path string) ([]ColdEntry, error) { if _, err := io.ReadFull(br, header[:]); err != nil { return nil, fmt.Errorf("txhash: read header of %s: %w", path, err) } - count := binary.LittleEndian.Uint64(header[:coldBinCountSize]) + count := binary.LittleEndian.Uint64(header[coldBinPreludeSize : coldBinPreludeSize+coldBinCountSize]) info, err := f.Stat() if err != nil { @@ -77,8 +77,9 @@ func TestColdBin_RoundTrip(t *testing.T) { assert.Equal(t, entries, got) } -// TestColdBin_HeaderAndLayout pins the raw on-disk layout: uint64 LE count -// header followed by fixed-width (key, uint32 LE seq) entries. +// TestColdBin_HeaderAndLayout pins the raw on-disk layout: the "SBIN" magic, +// the version byte, three reserved zero bytes, the uint64 LE count, the +// secret, then fixed-width (key, uint32 LE seq) entries. func TestColdBin_HeaderAndLayout(t *testing.T) { dir := t.TempDir() path := filepath.Join(dir, "out.bin") @@ -91,13 +92,52 @@ func TestColdBin_HeaderAndLayout(t *testing.T) { data, err := os.ReadFile(path) require.NoError(t, err) require.Len(t, data, coldBinHeaderSize+2*coldBinEntrySize) - assert.Equal(t, uint64(2), binary.LittleEndian.Uint64(data[:coldBinCountSize])) - assert.Equal(t, testBinSecret[:], data[coldBinCountSize:coldBinHeaderSize], "secret recorded after the count") + assert.Equal(t, []byte("SBIN"), data[:4], "magic in on-disk byte order") + assert.Equal(t, coldBinVersion, data[4]) + assert.Equal(t, []byte{0, 0, 0}, data[5:coldBinPreludeSize], "reserved bytes zero") + assert.Equal(t, uint64(2), + binary.LittleEndian.Uint64(data[coldBinPreludeSize:coldBinPreludeSize+coldBinCountSize])) + assert.Equal(t, testBinSecret[:], data[coldBinPreludeSize+coldBinCountSize:coldBinHeaderSize], + "secret recorded after the count") assert.Equal(t, byte(0xaa), data[coldBinHeaderSize]) assert.Equal(t, uint32(7), binary.LittleEndian.Uint32(data[coldBinHeaderSize+ColdKeySize:coldBinHeaderSize+coldBinEntrySize])) } +// TestColdBin_ScanRejectsForeignHeader pins the prelude refusals in BOTH .bin +// consumers, the pre-scan and the self-defending merge reader: a foreign magic, +// a newer version byte, and a set reserved byte each fail loudly instead of +// being misread as entry data. +func TestColdBin_ScanRejectsForeignHeader(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "out.bin") + require.NoError(t, WriteColdBin(path, testBinSecret, []ColdEntry{{Key: [ColdKeySize]byte{1}, Seq: 2}})) + data, err := os.ReadFile(path) + require.NoError(t, err) + + cases := []struct { + name string + mutate func([]byte) + want string + }{ + {"bad magic", func(b []byte) { copy(b, "JUNK") }, "not a cold txhash .bin"}, + {"newer version", func(b []byte) { b[4] = coldBinVersion + 1 }, "written by a newer stellar-rpc"}, + {"reserved byte set", func(b []byte) { b[6] = 0x01 }, "reserved header bytes set"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + bad := append([]byte(nil), data...) + tc.mutate(bad) + badPath := filepath.Join(dir, tc.name+".bin") + require.NoError(t, os.WriteFile(badPath, bad, 0o600)) + _, _, err := scanBinHeader(badPath) + require.ErrorContains(t, err, tc.want) + _, err = newFileReader(badPath, 0) + require.ErrorContains(t, err, tc.want) + }) + } +} + // TestColdBin_CreateFails forces os.Create on the destination to fail by // pre-creating the final path as a DIRECTORY (so create returns EISDIR). The // error must propagate; the pre-existing directory is untouched. diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format.go index 3f5793712..4e1fa305a 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format.go @@ -37,9 +37,14 @@ const ColdFingerprintSize = 1 // coldPayloadMax is the largest offset that fits ColdPayloadSize bytes. const coldPayloadMax = uint64(1)<<(ColdPayloadSize*8) - 1 +// coldMetadataVersion is the metadata blob's leading version byte. Every +// app-metadata blob leads with its own version byte so it is self-describing +// on its own, independent of the streamhash container version. +const coldMetadataVersion byte = 0x01 + // coldMetadataSize is the metadata blob width: -// [MinLedger:4 LE][MaxLedger:4 LE][routing secret:16]. -const coldMetadataSize = 8 + stores.SecretLen +// [version:1][MinLedger:4 LE][MaxLedger:4 LE][routing secret:16]. +const coldMetadataSize = 1 + 8 + stores.SecretLen // coldRoutingDomain is the DeriveIndexSecret domain for txhash cold indexes. const coldRoutingDomain = "txhash" @@ -55,28 +60,34 @@ func ColdIndexSecret(catalogSecret []byte, indexID uint32) [stores.SecretLen]byt return stores.DeriveIndexSecret(catalogSecret, coldRoutingDomain, indexID) } -// EncodeColdMetadata packs [minLedger, maxLedger, secret] into the metadata blob. +// EncodeColdMetadata packs [version, minLedger, maxLedger, secret] into the +// metadata blob. func EncodeColdMetadata(minLedger, maxLedger uint32, secret [stores.SecretLen]byte) []byte { buf := make([]byte, coldMetadataSize) - binary.LittleEndian.PutUint32(buf[:4], minLedger) - binary.LittleEndian.PutUint32(buf[4:8], maxLedger) - copy(buf[8:], secret[:]) + buf[0] = coldMetadataVersion + binary.LittleEndian.PutUint32(buf[1:5], minLedger) + binary.LittleEndian.PutUint32(buf[5:9], maxLedger) + copy(buf[9:], secret[:]) return buf } // ParseColdMetadata recovers [minLedger, maxLedger, secret] from the metadata -// blob, rejecting a wrong size or maxLedger < minLedger with ErrInvalidMetadata. +// blob. Every refusal (empty, unknown version, wrong size, maxLedger < +// minLedger) wraps ErrInvalidMetadata. func ParseColdMetadata(metadata []byte) (uint32, uint32, [stores.SecretLen]byte, error) { var secret [stores.SecretLen]byte + if err := stores.CheckBlobVersion(metadata, coldMetadataVersion); err != nil { + return 0, 0, secret, fmt.Errorf("%w: %w", ErrInvalidMetadata, err) + } if len(metadata) != coldMetadataSize { return 0, 0, secret, fmt.Errorf("%w: got %d bytes, want %d", ErrInvalidMetadata, len(metadata), coldMetadataSize) } - minLedger := binary.LittleEndian.Uint32(metadata[:4]) - maxLedger := binary.LittleEndian.Uint32(metadata[4:8]) + minLedger := binary.LittleEndian.Uint32(metadata[1:5]) + maxLedger := binary.LittleEndian.Uint32(metadata[5:9]) if maxLedger < minLedger { return 0, 0, secret, fmt.Errorf("%w: maxLedger %d < minLedger %d", ErrInvalidMetadata, maxLedger, minLedger) } - copy(secret[:], metadata[8:]) + copy(secret[:], metadata[9:]) return minLedger, maxLedger, secret, nil } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format_test.go index 4601edc98..a0b03c6ba 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_format_test.go @@ -26,11 +26,25 @@ func TestEncodeParseColdMetadata_RoundTrip(t *testing.T) { } func TestParseColdMetadata_WrongSizeErrors(t *testing.T) { - // 8 is the pre-secret layout; 24 is the only valid width. - for _, sz := range []int{0, 1, 4, 7, 8, 9, 16, 23, 25} { - _, _, _, err := ParseColdMetadata(make([]byte, sz)) + // 8 is the pre-secret layout, 24 the pre-version one; 25 is the only + // valid width. Blobs carry a valid version byte so the size check is what + // fires (the version check runs first by design). + for _, sz := range []int{1, 4, 7, 8, 9, 16, 24, 26} { + blob := make([]byte, sz) + blob[0] = coldMetadataVersion + _, _, _, err := ParseColdMetadata(blob) assert.ErrorIs(t, err, ErrInvalidMetadata, "size %d should error", sz) } + _, _, _, err := ParseColdMetadata(nil) + require.ErrorIs(t, err, ErrInvalidMetadata, "empty blob should error") + require.ErrorContains(t, err, "empty") +} + +func TestParseColdMetadata_UnknownVersionErrors(t *testing.T) { + blob := EncodeColdMetadata(5, 9, testSecret()) + blob[0] = coldMetadataVersion + 1 + _, _, _, err := ParseColdMetadata(blob) + require.ErrorContains(t, err, "written by a newer stellar-rpc") } func TestParseColdMetadata_MaxBelowMinErrors(t *testing.T) { @@ -93,7 +107,8 @@ func TestOpenColdReader_BadMetadataErrors(t *testing.T) { require.NoError(t, sb.Close()) _, err = OpenColdReader(path) - assert.ErrorIs(t, err, ErrInvalidMetadata) + require.ErrorIs(t, err, ErrInvalidMetadata, "a metadata-less index must fail open") + require.ErrorContains(t, err, "empty") } func TestOpenColdReader_WrongPayloadSizeErrors(t *testing.T) { diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index.go index e983f3c60..e1b811e9c 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index.go @@ -169,7 +169,11 @@ func scanBinHeader(path string) (uint64, [stores.SecretLen]byte, error) { if _, err := io.ReadFull(f, hdr[:]); err != nil { return 0, secret, fmt.Errorf("txhash: read header of %s: %w", path, err) } - copy(secret[:], hdr[coldBinCountSize:]) - count, err := coldBinCount(path, fi.Size(), binary.LittleEndian.Uint64(hdr[:coldBinCountSize])) + if err := checkBinPrelude(path, hdr[:]); err != nil { + return 0, secret, err + } + copy(secret[:], hdr[coldBinPreludeSize+coldBinCountSize:]) + count, err := coldBinCount(path, fi.Size(), + binary.LittleEndian.Uint64(hdr[coldBinPreludeSize:coldBinPreludeSize+coldBinCountSize])) return count, secret, err } diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index_test.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index_test.go index d792f378c..a96f19537 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index_test.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_index_test.go @@ -125,19 +125,17 @@ func writeBinFile(t *testing.T, path string, entries []fixtureEntry, secret [sto sort.Slice(keyed, func(i, j int) bool { return bytes.Compare(keyed[i].Key[:], keyed[j].Key[:]) < 0 }) + require.NoError(t, WriteColdBin(path, secret, keyed)) +} - var buf bytes.Buffer +// rawBinHeader hand-builds a syntactically valid .bin header (magic, version, +// zero secret) claiming count entries, for corrupt-file fixtures. +func rawBinHeader(count uint64) [coldBinHeaderSize]byte { var hdr [coldBinHeaderSize]byte - binary.LittleEndian.PutUint64(hdr[:coldBinCountSize], uint64(len(keyed))) - copy(hdr[coldBinCountSize:], secret[:]) - buf.Write(hdr[:]) - var seqBuf [coldBinSeqSize]byte - for _, e := range keyed { - buf.Write(e.Key[:]) - binary.LittleEndian.PutUint32(seqBuf[:], e.Seq) - buf.Write(seqBuf[:]) - } - require.NoError(t, os.WriteFile(path, buf.Bytes(), 0o600)) + copy(hdr[:4], coldBinMagic) + hdr[4] = coldBinVersion + binary.LittleEndian.PutUint64(hdr[coldBinPreludeSize:], count) + return hdr } // writeFixtureBins partitions entries into per-chunk .bin files under @@ -387,8 +385,7 @@ func TestBuildColdIndex_TruncatedFileErrors(t *testing.T) { // holds): the open-time size cross-check rejects it; no index. dir := t.TempDir() var buf bytes.Buffer - var hdr [coldBinHeaderSize]byte - binary.LittleEndian.PutUint64(hdr[:], 5) // claim 5 + hdr := rawBinHeader(5) // claim 5 buf.Write(hdr[:]) // ...but write only one entry. buf.Write(make([]byte, coldBinEntrySize)) @@ -408,8 +405,7 @@ func TestBuildColdIndex_HeaderUndercountErrors(t *testing.T) { // silently drop the trailing entries; it must error instead. dir := t.TempDir() var buf bytes.Buffer - var hdr [coldBinHeaderSize]byte - binary.LittleEndian.PutUint64(hdr[:], 1) // declare 1... + hdr := rawBinHeader(1) // declare 1... buf.Write(hdr[:]) buf.Write(make([]byte, coldBinEntrySize*3)) // ...but write 3 entries p := filepath.Join(dir, "00000005.bin") @@ -431,8 +427,7 @@ func TestBuildColdIndex_HeaderOverflowRejected(t *testing.T) { // reject this cleanly (no allocation, no panic). dir := t.TempDir() var buf bytes.Buffer - var hdr [coldBinHeaderSize]byte - binary.LittleEndian.PutUint64(hdr[:], math.MaxUint64) // wildly overstated count + hdr := rawBinHeader(math.MaxUint64) // wildly overstated count buf.Write(hdr[:]) buf.Write(make([]byte, coldBinEntrySize)) // one real entry p := filepath.Join(dir, "00000005.bin") diff --git a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_merge.go b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_merge.go index 26bfb795e..a1c96ddfc 100644 --- a/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_merge.go +++ b/cmd/stellar-rpc/internal/rpcv2/stores/txhash/cold_merge.go @@ -255,6 +255,12 @@ func newFileReader(path string, bufBytes int) (*fileReader, error) { _ = f.Close() return nil, fmt.Errorf("txhash: %s too short for header (%d bytes)", path, n) } + // Self-defending: the merge consumes raw entry bytes, so it verifies the + // prelude itself rather than trusting that scanAndValidate ran upstream. + if err := checkBinPrelude(path, buf); err != nil { + _ = f.Close() + return nil, err + } return &fileReader{ path: path, f: f, diff --git a/design-docs/full-history-streaming-workflow.md b/design-docs/full-history-streaming-workflow.md index 55f960073..b2fccf167 100644 --- a/design-docs/full-history-streaming-workflow.md +++ b/design-docs/full-history-streaming-workflow.md @@ -977,7 +977,7 @@ Each invariant has a distinct audit. INV-1 you check by issuing reads or by re-d ### Convergence -**Startup converges from any on-disk state.** Whatever a partial-completion crash, an operator action, or surgical recovery leaves behind, startup drives the system to a settled state satisfying INV-1 ∧ INV-2 ∧ INV-3 ∧ INV-4. Startup here is the backfill pass followed by the first lifecycle run — fired by the startup seed when a complete chunk exists; a young store with none is settled by construction — and it reaches a settled state within that first run, typically seconds after serving opens, bounded by the run's freeze, rebuild, and prune workload. From any state reachable *during* a run, the lifecycle run alone converges, within a bounded number of runs. And since a runtime op failure fails the run and hands recovery back to startup, every state a run can leave behind is one startup is built to converge. +**Startup converges from any on-disk state this binary's own protocol can produce.** Whatever a partial-completion crash, an operator action, or surgical recovery leaves behind, startup drives the system to a settled state satisfying INV-1 ∧ INV-2 ∧ INV-3 ∧ INV-4. Catalog content outside the binary's compiled-in vocabulary is deliberately not converged: the open-time census refuses startup instead ([Format versioning and upgrades](#format-versioning-and-upgrades)), because "converging" a newer binary's state would mean overwriting artifacts this binary cannot read. Startup here is the backfill pass followed by the first lifecycle run — fired by the startup seed when a complete chunk exists; a young store with none is settled by construction — and it reaches a settled state within that first run, typically seconds after serving opens, bounded by the run's freeze, rebuild, and prune workload. From any state reachable *during* a run, the lifecycle run alone converges, within a bounded number of runs. And since a runtime op failure fails the run and hands recovery back to startup, every state a run can leave behind is one startup is built to converge. Two accepted limits. Convergence needs disk headroom: the prune that frees space runs after the backfill that needs it (the peak-disk note under [Startup](#startup)), so a disk sized below the transient peak fails the run on every pass until space is added. And a same-window 16-byte hash collision — the ~10⁻²⁰-per-window event the transactions design accepts (§8.2) — fails that window's index build permanently: the accepted risk presents as a run that never completes, never as wrong data. @@ -1043,6 +1043,22 @@ A storage walk against the invariants is enough to find these without inspecting --- +## Format versioning and upgrades + +Every durable byte the daemon writes belongs to a storage format, and the policy for changing one is deliberately asymmetric: **forward is seamless, backward refuses.** + +**The census.** At every catalog open, before the first-start pin can be written and before any hot DB is touched, the daemon scans every catalog key and value against the exact vocabulary this binary writes: the three state key families with their exact state tokens, the `config:earliest_ledger` pin as a canonical decimal, and the 32-byte `meta/catalog-secret`. Anything else refuses startup with one message: the entry is either from a newer stellar-rpc (deploy that version or newer) or corrupted. The refusal writes no catalog entry of its own (in particular, no secret is minted into a refused tree) and never prints the secret's value; RocksDB itself may still create housekeeping files and flush a previous binary's recovered WAL on Close. This is what makes rollback safe by construction: without it, an older binary would treat a newer binary's artifacts as missing work and overwrite them. + +**Roll-forward-only.** A release that changes any storage format, adds a catalog key family, or changes a value form cannot be rolled back past: the older binary refuses the tree, and the remedy is redeploying the newer binary. Release notes for a format-touching release must say so. + +**Where format identity lives.** The catalog value that makes an artifact visible also names its format: the planned grammar widens `"frozen"` to `"frozen@id"`, with a bare value meaning format 1 forever. Release 1 writes only bare values, so the grammar ships with the first bump. The grammar applies to the three state families only, never to `config:*` or `meta/*` values (the secret's random bytes can contain `@`). The files themselves stay self-describing as witnesses: the packfile trailer carries a container version and a per-store `Format` id, the streamhash files carry their own header version, and every app-data blob leads with its own version byte. Two extension rules coexist and are both deliberate: the `index.pack` build stamp ignores trailing bytes, so a later version can append fields without a version bump, while the ledger app data and the txhash `.idx` metadata require their exact length, so any growth is a version bump. The packfile content hash is a witness for offline verification and repair tooling; the serving path never recomputes it. A change to the term schema or the indexed-field set ships as a new events format id, never as a stamp-only bump: nothing ties the stamp to the trailer's format id, and a stamp-only bump would leave every chunk `frozen` in the catalog with an index that fails on every filtered lookup. + +**The bump razor.** Any reader-visible change to a family's bytes or semantics is a new format id. Three non-obvious corollaries. A new indexed term family bumps the hot format too, or the live chunk's already-ingested half would silently lack the new terms. A dependency bump that changes RocksDB's on-disk behavior is format-touching, which is why the block-based-table `format_version` is pinned in code. And a change to what the writers hash or emit (such as enabling the content hash) must switch the backfill walk and the freeze in the same release, or freeze-vs-walk byte identity breaks inside one unbumped id. The term-key derivation and blinding are pinned by byte-level golden fixtures, so an unbumped change to them fails CI rather than silently forking a format id. + +**Upgrades, per tier.** Cold artifacts are write-new, read-old: a newer binary reads every id in its read set forever, writes only its current id, and mixed-id trees are permanent and normal. Cold ids are never demoted automatically in either direction, and mass rebuilds are never an upgrade path. The hot tier is the one automatic reaction: a newer binary that finds hot state at a known old hot id discards it (demote to transient, wipe, recreate) and re-ingests at most one chunk from the cold boundary, minutes of work. A capability the artifacts of an old id cannot answer (a term family they never indexed) is a loud, actionable error naming the range and the remedy, never a silent empty result; an optional background reindex can converge old chunks when wanted. + +--- + ## Related documents - The transactions design ([gettransaction-full-history-design.md](./gettransaction-full-history-design.md)) — the tx-by-hash subsystem end to end: the hot `txhash` CF, the `.bin`/`.idx` formats, the rolling window index rebuild — its streamhash merge internals and safety argument — the `getTransaction` read path, and the capacity numbers. Canonical for the streamhash `.bin`/`.idx` formats, the index merge internals, and the index-key coverage semantics this doc summarizes. diff --git a/design-docs/gettransaction-full-history-design.md b/design-docs/gettransaction-full-history-design.md index 5c92a4f89..a4414caaa 100644 --- a/design-docs/gettransaction-full-history-design.md +++ b/design-docs/gettransaction-full-history-design.md @@ -102,11 +102,15 @@ The `.bin` lives at `txhash/raw/{bucket:05d}/{chunk:08d}.bin`, with catalog key **Format** (the streamhash merge format): ``` +uint32 LE magic ("SBIN" in on-disk byte order) +uint8 version (1) +3 bytes reserved (zero) uint64 LE entry count +16 bytes index secret the keys were blinded with entry × count 20 bytes each: [key: 16][seq: 4 LE] ``` -- `key` is the **first 16 bytes of the transaction hash**. The index uses only these 16 bytes to place and find a transaction; what happens when two hashes share a 16-byte prefix is in §8.2. +- `key` is the **secret-blinded first 16 bytes of the transaction hash** (blinded with the index secret recorded in the header, so the deferred build always adopts the secret its inputs were keyed with). The index uses only these 16 bytes to place and find a transaction; what happens when two hashes share a 16-byte prefix is in §8.2. - Entries are sorted ascending by `key`, **bytewise over all 16 bytes** — a total order, so the same entries always produce byte-identical files (rebuilds are deterministic). The `.bin` is a pre-sorted file, and a lookup never reads it directly. It is sorted because streamhash builds an index **much faster, and with much less memory, when its keys arrive already sorted** — its *sorted-builder mode*. @@ -117,7 +121,7 @@ A `.bin` is kept while it is still a rebuild input — every rebuild re-merges t The `.idx` lives at `txhash/index/{window:08d}/{lo:08d}-{hi:08d}.idx`, tracked by the catalog key `index:{window:08d}:{lo:08d}:{hi:08d}`. There is one minimal-perfect-hash file per **coverage** — a coverage being the chunk range `[lo, hi]` the file actually hashes. Streamhash's `SortedBuilder` builds it from the k-way merge of `.bin[lo..hi]`. The index carries two per-entry fields: -- **Payload (3 bytes): the answer the hash maps to — a ledger seq.** It is stored as an offset from the coverage's first ledger (`MinLedger = chunkFirstLedger(lo)`) rather than as a full seq, to save bytes. A window spans 10,000,000 ledgers, so the largest offset (`10_000_000 - 1`) fits in a 24-bit field. Streamhash writes the payload width into the index file's header; the coverage's ledger range `[MinLedger, MaxLedger]`, which streamhash does not model itself, rides in the file's user-metadata slot as two 4-byte little-endian values. Both are read back when the index is opened, so there is no separate sidecar file. +- **Payload (3 bytes): the answer the hash maps to — a ledger seq.** It is stored as an offset from the coverage's first ledger (`MinLedger = chunkFirstLedger(lo)`) rather than as a full seq, to save bytes. A window spans 10,000,000 ledgers, so the largest offset (`10_000_000 - 1`) fits in a 24-bit field. Streamhash writes the payload width into the index file's header; the coverage's ledger range `[MinLedger, MaxLedger]`, which streamhash does not model itself, rides in the file's user-metadata slot behind a leading version byte, as two 4-byte little-endian values followed by the 16-byte routing secret. All of it is read back when the index is opened, so there is no separate sidecar file. - **Fingerprint (1 byte, fixed by the format): screens out wrong hashes** before the expensive fetch-and-verify. One byte rejects ~255/256 of foreign hashes per window probe (§8.2), costing one byte per transaction. All-in, the index costs ≈4.2 bytes per transaction (MPHF structure + payload + fingerprint) — ≈12.5 GB for a dense full window, versus the ≈60 GB of `.bin` runs it consumes.