Documentation
¶
Index ¶
- Constants
- Variables
- func DecodeSortableFloat64(b [8]byte) float64
- func DecodeWriteConflictLabel(label string) (kind, keyClass string, ok bool)
- func EncodeSortableFloat64(f float64) [8]byte
- func EncodeStreamID(ms, seq uint64) []byte
- func EncodeWriteConflictLabel(kind, keyClass string) string
- func ExtractHashFieldName(key, userKey []byte) []byte
- func ExtractHashUserKeyFromDelta(key []byte) []byte
- func ExtractHashUserKeyFromDeltaScanPrefix(key []byte) []byte
- func ExtractHashUserKeyFromField(key []byte) []byte
- func ExtractHashUserKeyFromMeta(key []byte) []byte
- func ExtractLegacyListUserKeyFromDelta(key []byte) []byte
- func ExtractLegacyListUserKeyFromDeltaScanPrefix(key []byte) []byte
- func ExtractListItemSeq(itemKey, userKey []byte) (int64, bool)
- func ExtractListUserKey(key []byte) []byte
- func ExtractListUserKeyFromClaim(key []byte) []byte
- func ExtractListUserKeyFromClaimScanKey(key []byte) []byte
- func ExtractListUserKeyFromClaimScanPrefix(key []byte) []byte
- func ExtractListUserKeyFromDelta(key []byte) []byte
- func ExtractListUserKeyFromDeltaScanKey(key []byte) []byte
- func ExtractListUserKeyFromDeltaScanPrefix(key []byte) []byte
- func ExtractSetMemberName(key, userKey []byte) []byte
- func ExtractSetUserKeyFromDelta(key []byte) []byte
- func ExtractSetUserKeyFromDeltaScanPrefix(key []byte) []byte
- func ExtractSetUserKeyFromMember(key []byte) []byte
- func ExtractSetUserKeyFromMeta(key []byte) []byte
- func ExtractStreamUserKeyFromEntry(key []byte) []byte
- func ExtractStreamUserKeyFromEntryScanPrefix(key []byte) []byte
- func ExtractStreamUserKeyFromMeta(key []byte) []byte
- func ExtractZSetMemberName(key, userKey []byte) []byte
- func ExtractZSetScoreAndMember(key, userKey []byte) (score float64, member []byte, ok bool)
- func ExtractZSetUserKeyFromDelta(key []byte) []byte
- func ExtractZSetUserKeyFromDeltaScanPrefix(key []byte) []byte
- func ExtractZSetUserKeyFromMember(key []byte) []byte
- func ExtractZSetUserKeyFromMeta(key []byte) []byte
- func ExtractZSetUserKeyFromScore(key []byte) []byte
- func ExtractZSetUserKeyFromScoreScanPrefix(key []byte) []byte
- func HashFieldKey(userKey, fieldName []byte) []byte
- func HashFieldScanPrefix(userKey []byte) []byte
- func HashMetaDeltaKey(userKey []byte, commitTS uint64, seqInTxn uint32) []byte
- func HashMetaDeltaScanPrefix(userKey []byte) []byte
- func HashMetaKey(userKey []byte) []byte
- func IsHashFieldKey(key []byte) bool
- func IsHashMetaDeltaKey(key []byte) bool
- func IsHashMetaKey(key []byte) bool
- func IsListClaimKey(key []byte) bool
- func IsListItemKey(key []byte) bool
- func IsListMetaDeltaKey(key []byte) bool
- func IsListMetaDeltaValue(value []byte) bool
- func IsListMetaKey(key []byte) bool
- func IsSetMemberKey(key []byte) bool
- func IsSetMetaDeltaKey(key []byte) bool
- func IsSetMetaKey(key []byte) bool
- func IsStreamEntryKey(key []byte) bool
- func IsStreamMetaKey(key []byte) bool
- func IsZSetMemberKey(key []byte) bool
- func IsZSetMetaDeltaKey(key []byte) bool
- func IsZSetMetaKey(key []byte) bool
- func IsZSetScoreKey(key []byte) bool
- func LegacyListMetaDeltaScanPrefix(userKey []byte) []byte
- func ListClaimKey(userKey []byte, seq int64) []byte
- func ListClaimScanPrefix(userKey []byte) []byte
- func ListItemKey(userKey []byte, seq int64) []byte
- func ListMetaDeltaKey(userKey []byte, commitTS uint64, seqInTxn uint32) []byte
- func ListMetaDeltaScanPrefix(userKey []byte) []byte
- func ListMetaDeltaScanPrefixes(userKey []byte) [][]byte
- func ListMetaKey(userKey []byte) []byte
- func MarshalHashMeta(m HashMeta) []byte
- func MarshalHashMetaDelta(d HashMetaDelta) []byte
- func MarshalListMeta(meta ListMeta) ([]byte, error)
- func MarshalListMetaDelta(d ListMetaDelta) []byte
- func MarshalSetMeta(m SetMeta) []byte
- func MarshalSetMetaDelta(d SetMetaDelta) []byte
- func MarshalStreamMeta(m StreamMeta) ([]byte, error)
- func MarshalZSetMeta(m ZSetMeta) []byte
- func MarshalZSetMetaDelta(d ZSetMetaDelta) []byte
- func MarshalZSetScore(score float64) []byte
- func NewWriteConflictError(key []byte) error
- func PrefixScanEnd(prefix []byte) []byte
- func SetMemberKey(userKey, member []byte) []byte
- func SetMemberScanPrefix(userKey []byte) []byte
- func SetMetaDeltaKey(userKey []byte, commitTS uint64, seqInTxn uint32) []byte
- func SetMetaDeltaScanPrefix(userKey []byte) []byte
- func SetMetaKey(userKey []byte) []byte
- func StreamEntryKey(userKey []byte, ms, seq uint64) []byte
- func StreamEntryScanPrefix(userKey []byte) []byte
- func StreamEntryScanStart(userKey []byte, meta StreamMeta) []byte
- func StreamMetaKey(userKey []byte) []byte
- func UnmarshalZSetScore(b []byte) (float64, error)
- func WriteConflictKey(err error) ([]byte, bool)
- func WriterRegistryFor(s MVCCStore) (encryption.WriterRegistryStore, error)
- func ZSetMemberKey(userKey, member []byte) []byte
- func ZSetMemberScanPrefix(userKey []byte) []byte
- func ZSetMetaDeltaKey(userKey []byte, commitTS uint64, seqInTxn uint32) []byte
- func ZSetMetaDeltaScanPrefix(userKey []byte) []byte
- func ZSetMetaKey(userKey []byte) []byte
- func ZSetScoreKey(userKey []byte, score float64, member []byte) []byte
- func ZSetScoreRangeScanPrefix(userKey []byte, score float64) []byte
- func ZSetScoreScanPrefix(userKey []byte) []byte
- type ActiveStorageKeyID
- type CounterNonceFactory
- type ExportVersionsOptions
- type ExportVersionsResult
- type HashMeta
- type HashMetaDelta
- type HybridClock
- type ImportVersionsOptions
- type ImportVersionsResult
- type KVPair
- type KVPairMutation
- type ListMeta
- type ListMetaDelta
- type MVCCStore
- type MVCCStoreOption
- type MVCCVersion
- type NonceFactory
- type OpType
- type PebbleStoreOption
- func WithEncryption(cipher *encryption.Cipher, nf NonceFactory, activeKeyID ActiveStorageKeyID) PebbleStoreOption
- func WithPebbleLogger(l *slog.Logger) PebbleStoreOption
- func WithSSTIngestSnapshots(enabled bool) PebbleStoreOption
- func WithStorageEnvelopeGate(active StorageEnvelopeActive) PebbleStoreOption
- func WithStorageEnvelopeV2Gate(active StorageEnvelopeV2Active) PebbleStoreOption
- func WithStorageRegistrationGate(registered StorageRegistered) PebbleStoreOption
- type RetentionController
- type SetMeta
- type SetMetaDelta
- type Snapshot
- type StorageEnvelopeActive
- type StorageEnvelopeV2Active
- type StorageRegistered
- type StreamMeta
- type VersionPresenceReader
- type VersionedValue
- type WriteConflictError
- type ZSetMeta
- type ZSetMetaDelta
Constants ¶
const ( HashMetaPrefix = "!hs|meta|" HashFieldPrefix = "!hs|fld|" HashMetaDeltaPrefix = "!hs|meta|d|" )
Hash wide-column key layout:
Base Metadata: !hs|meta|<userKeyLen(4)><userKey> → [Len(8)][ExpireAtMs(8)] Field Key: !hs|fld|<userKeyLen(4)><userKey><fieldName> → field value bytes Delta Key: !hs|meta|d|<userKeyLen(4)><userKey><commitTS(8)><seqInTxn(4)> → [LenDelta(8)]
const ( // ListMetaDeltaPrefix is the prefix for all list metadata delta keys. // Layout: !lst|delta|<userKeyLen(4)><userKey><commitTS(8)><seqInTxn(4)> ListMetaDeltaPrefix = "!lst|delta|" // LegacyListMetaDeltaPrefix is the pre-upgrade list metadata delta prefix. // Writers use ListMetaDeltaPrefix, but readers and compactors keep scanning // this prefix until old uncompacted deltas have been drained. LegacyListMetaDeltaPrefix = "!lst|meta|d|" // ListClaimPrefix is the prefix for list claim keys used by POP operations. // Layout: !lst|claim|<userKeyLen(4)><userKey><seq(8-byte sortable)> ListClaimPrefix = "!lst|claim|" // MaxDeltaScanLimit is the hard limit on delta scan results per read. // If a key accumulates more than this many uncompacted deltas, reads // return ErrDeltaScanTruncated until background compaction catches up. // Set high enough to tolerate burst traffic between compaction ticks // (default 30 s) without sacrificing read availability on hot keys. MaxDeltaScanLimit = 1024 )
Delta/Claim key constants.
const ( ListMetaPrefix = "!lst|meta|" ListItemPrefix = "!lst|itm|" )
const ( SetMetaPrefix = "!st|meta|" SetMemberPrefix = "!st|mem|" SetMetaDeltaPrefix = "!st|meta|d|" )
Set wide-column key layout:
Base Metadata: !st|meta|<userKeyLen(4)><userKey> → [Len(8)][ExpireAtMs(8)] Member Key: !st|mem|<userKeyLen(4)><userKey><member> → (empty value) Delta Key: !st|meta|d|<userKeyLen(4)><userKey><commitTS(8)><seqInTxn(4)> → [LenDelta(8)]
const ( StreamMetaPrefix = "!stream|meta|" StreamEntryPrefix = "!stream|entry|" // StreamMetaTrimBinarySize is the current fixed binary size for StreamMeta: // Len(8) || LastMs(8) || LastSeq(8) || ExpireAtMs(8) || TrimmedMs(8) || TrimmedSeq(8). StreamMetaTrimBinarySize = streamMetaTrimBinarySize // StreamIDBytes is the fixed size of the binary StreamID suffix on an entry key: // 8 bytes big-endian ms || 8 bytes big-endian seq. Big-endian so lex order // over the raw key bytes matches the (ms, seq) numeric order used by XADD / XRANGE. StreamIDBytes = 16 )
Stream wide-column key layout:
Meta: !stream|meta|<userKeyLen(4)><userKey> → [Len(8)][LastMs(8)][LastSeq(8)][ExpireAtMs(8)] Entry: !stream|entry|<userKeyLen(4)><userKey><StreamID(16)>
const ( WriteConflictKindWrite = "write" WriteConflictKindRead = "read" )
Write-conflict classification labels. These are emitted as the "key_prefix" Prometheus label on elastickv_store_write_conflict_total so operators can tell WHICH namespace is producing OCC conflicts. The label set is deliberately bounded and stable — adding a new bucket is a conscious schema change, not a cardinality surprise.
Buckets are aligned with the internal key prefixes declared across the codebase:
!txn|lock|, !txn|int|, !txn|cmt|, !txn|rb| (kv/txn_keys.go) !redis|str|, !redis|ttl|, !redis|hll| (adapter/redis_compat_types.go) !hs|, !st|, !zs|, !lst| (store/*_helpers.go) !ddb| (kv/shard_key.go)
const ( ZSetMetaPrefix = "!zs|meta|" ZSetMemberPrefix = "!zs|mem|" ZSetScorePrefix = "!zs|scr|" ZSetMetaDeltaPrefix = "!zs|meta|d|" )
ZSet wide-column key layout:
Base Metadata: !zs|meta|<userKeyLen(4)><userKey> → [Len(8)][ExpireAtMs(8)] Member Key: !zs|mem|<userKeyLen(4)><userKey><member> → [Score(8)] IEEE 754 Score Index: !zs|scr|<userKeyLen(4)><userKey><sortableScore(8)><member> → (empty) Delta Key: !zs|meta|d|<userKeyLen(4)><userKey><commitTS(8)><seqInTxn(4)> → [LenDelta(8)]
Variables ¶
var ErrEncryptedReadCompression = errors.New("store: authenticated encrypted value contains an invalid Snappy payload")
ErrEncryptedReadCompression is returned when an authenticated envelope is marked as Snappy-compressed but its decrypted body is not a valid Snappy stream. Disk corruption normally fails GCM first; reaching this sentinel means a writer produced a malformed authenticated payload, so reads fail closed instead of surfacing the compressed bytes as user data.
var ErrEncryptedReadIntegrity = errors.New("store: encrypted value failed integrity check (GCM tag mismatch); refusing to surface plaintext")
ErrEncryptedReadIntegrity wraps encryption.ErrIntegrity for storage-layer callers (Get / scan / iterator). Per design §4.1, callers MUST treat this as a typed read error and never silently zero the value or skip the row.
Callers can disambiguate it from any other read error with errors.Is.
var ErrEncryptedValueReservedState = errors.New("store: value header carries reserved encryption_state; binary too old to read this entry")
ErrEncryptedValueReservedState indicates decodeValue saw an encryption_state value (0b10 or 0b11) that the current build does not know how to interpret. Per design §7.1, this is a fail-closed trip-wire so an old binary cannot silently treat a future-version encrypted entry as cleartext bytes.
var ErrExpired = errors.New("expired")
var ErrImportBatchGap = errors.New("migration import batch gap")
var ErrInvalidChecksum = errors.New("invalid checksum")
var ErrInvalidExportBudget = errors.New("migration export requires a positive version budget")
ErrInvalidExportBudget rejects an export asked for with no version budget. Answering Done for a zero budget lets a migration driver record a non-empty bracket as fully copied and move on to cutover, so the completion signal must never be the answer to "you gave me nothing to fill".
var ErrInvalidExportCursor = errors.New("invalid export cursor")
var ErrKeyNotFound = errors.New("not found")
var ErrNotSupported = errors.New("not supported")
var ErrReadTSCompacted = errors.New("read timestamp has been compacted")
var ErrSnapshotKeyTooLarge = errors.New("mvcc snapshot key too large")
var ErrSnapshotVersionCountTooLarge = errors.New("mvcc snapshot version count too large")
var ErrUnknownOp = errors.New("unknown op")
var ErrUnsupportedStoreForWriterRegistry = errors.New("store: WriterRegistryFor requires a Pebble-backed MVCCStore")
ErrUnsupportedStoreForWriterRegistry is returned by WriterRegistryFor when the supplied MVCCStore is not backed by the in-tree Pebble implementation. Test fakes and alternate backends are detected at construction time so the failure mode is "binary refuses to start" rather than "FSM apply panics inside Raft loop".
var ErrValueTooLarge = errors.New("value too large")
var ErrWriteConflict = errors.New("write conflict")
var ErrWriterNotRegistered = errors.New("store: storage envelope active but writer not yet registered; refusing to emit nonce before registration")
ErrWriterNotRegistered is returned by the direct write path when the §7.1 storage envelope is active but this process load has not yet confirmed its §4.1 writer registration for the active storage DEK (Stage 7a-2). It is a transient, retryable condition: the caller should retry once registration commits (the barrier closes). Callers disambiguate it with errors.Is so a fail-closed pre-registration write is never mistaken for a permanent storage fault.
var Tombstone = []byte{0x00}
Functions ¶
func DecodeSortableFloat64 ¶
DecodeSortableFloat64 decodes a sortable 8-byte representation back to float64.
func DecodeWriteConflictLabel ¶
DecodeWriteConflictLabel splits a flat label back into (kind, keyClass). Returns ok=false for malformed input so collector code can skip rather than emit bogus series.
func EncodeSortableFloat64 ¶
EncodeSortableFloat64 encodes a float64 into a sortable 8-byte representation. For positive floats: XOR the sign bit to make them sort above negative. For negative floats: XOR all bits to reverse the order. This produces a byte sequence that sorts correctly with standard byte comparison.
func EncodeStreamID ¶
EncodeStreamID writes a StreamID as 16 big-endian bytes: ms || seq.
func EncodeWriteConflictLabel ¶
EncodeWriteConflictLabel joins a (kind, keyClass) tuple into the flat map key used by WriteConflictCountsByPrefix snapshots. Exposed so the monitoring collector can build the same string without re-implementing the convention.
func ExtractHashFieldName ¶
ExtractHashFieldName extracts the field name from a hash field key.
func ExtractHashUserKeyFromDelta ¶
ExtractHashUserKeyFromDelta extracts the logical user key from a hash delta key.
func ExtractHashUserKeyFromDeltaScanPrefix ¶
ExtractHashUserKeyFromDeltaScanPrefix extracts the user key from a hash metadata delta scan start/prefix.
func ExtractHashUserKeyFromField ¶
ExtractHashUserKeyFromField extracts the logical user key from a hash field key.
func ExtractHashUserKeyFromMeta ¶
ExtractHashUserKeyFromMeta extracts the logical user key from a hash meta key.
func ExtractLegacyListUserKeyFromDelta ¶
ExtractLegacyListUserKeyFromDelta extracts the user key from an old !lst|meta|d| delta key. Callers that only have a key must be careful: this legacy layout overlaps with base !lst|meta| keys whose user key begins with d|. Prefer using it when the associated value is known to be a delta value.
func ExtractLegacyListUserKeyFromDeltaScanPrefix ¶
ExtractLegacyListUserKeyFromDeltaScanPrefix extracts the user key from a legacy-list delta scan start/prefix.
func ExtractListItemSeq ¶
ExtractListItemSeq extracts the sequence number encoded in a list item key. Returns (seq, true) on success; (0, false) if key is not a valid item key for userKey.
func ExtractListUserKey ¶
ExtractListUserKey returns the logical user key from a list meta or item key. If the key is not a list key, it returns nil.
func ExtractListUserKeyFromClaim ¶
ExtractListUserKeyFromClaim extracts the logical user key from a list claim key.
func ExtractListUserKeyFromClaimScanKey ¶
ExtractListUserKeyFromClaimScanKey extracts the logical user key from a claim scan prefix, a full claim key, or a scan cursor within that prefix.
func ExtractListUserKeyFromClaimScanPrefix ¶
ExtractListUserKeyFromClaimScanPrefix extracts the logical user key from the exact prefix produced by ListClaimScanPrefix.
func ExtractListUserKeyFromDelta ¶
ExtractListUserKeyFromDelta extracts the logical user key from a list delta key.
func ExtractListUserKeyFromDeltaScanKey ¶
ExtractListUserKeyFromDeltaScanKey extracts the logical user key from a delta scan prefix, a full delta key, or a scan cursor within that prefix. Unlike ExtractListUserKeyFromDeltaScanPrefix it tolerates trailing bytes, because a resumed scan starts from a cursor inside the prefix rather than from the prefix itself.
func ExtractListUserKeyFromDeltaScanPrefix ¶
ExtractListUserKeyFromDeltaScanPrefix extracts the user key from a new-list delta scan start/prefix.
func ExtractSetMemberName ¶
ExtractSetMemberName extracts the member name from a set member key.
func ExtractSetUserKeyFromDelta ¶
ExtractSetUserKeyFromDelta extracts the logical user key from a set delta key.
func ExtractSetUserKeyFromDeltaScanPrefix ¶
ExtractSetUserKeyFromDeltaScanPrefix extracts the user key from a set metadata delta scan start/prefix.
func ExtractSetUserKeyFromMember ¶
ExtractSetUserKeyFromMember extracts the logical user key from a set member key.
func ExtractSetUserKeyFromMeta ¶
ExtractSetUserKeyFromMeta extracts the logical user key from a set meta key.
func ExtractStreamUserKeyFromEntry ¶
ExtractStreamUserKeyFromEntry extracts the logical user key from a stream entry key.
See ExtractStreamUserKeyFromMeta for the rationale of the uint64 bounds check; the entry variant additionally has to account for the trailing StreamIDBytes (16 bytes) suffix.
func ExtractStreamUserKeyFromEntryScanPrefix ¶
ExtractStreamUserKeyFromEntryScanPrefix extracts the logical user key from the prefix used to scan all entries for a stream.
func ExtractStreamUserKeyFromMeta ¶
ExtractStreamUserKeyFromMeta extracts the logical user key from a stream meta key.
The bounds check is done in uint64 so a corrupted length prefix near math.MaxUint32 cannot wrap (uint32(wideColKeyLenSize)+ukLen) and pass a false negative, which would then panic on the trimmed[lo:hi] slice below.
func ExtractZSetMemberName ¶
ExtractZSetMemberName extracts the member name from a zset member key.
func ExtractZSetScoreAndMember ¶
ExtractZSetScoreAndMember extracts the score and member name from a zset score index key.
func ExtractZSetUserKeyFromDelta ¶
ExtractZSetUserKeyFromDelta extracts the logical user key from a zset delta key.
func ExtractZSetUserKeyFromDeltaScanPrefix ¶
ExtractZSetUserKeyFromDeltaScanPrefix extracts the user key from a zset metadata delta scan start/prefix.
func ExtractZSetUserKeyFromMember ¶
ExtractZSetUserKeyFromMember extracts the logical user key from a zset member key.
func ExtractZSetUserKeyFromMeta ¶
ExtractZSetUserKeyFromMeta extracts the logical user key from a zset meta key.
func ExtractZSetUserKeyFromScore ¶
ExtractZSetUserKeyFromScore extracts the logical user key from a zset score index key.
func ExtractZSetUserKeyFromScoreScanPrefix ¶
ExtractZSetUserKeyFromScoreScanPrefix extracts the user key from a zset score-index scan start/prefix that does not yet include the sortable score.
func HashFieldKey ¶
HashFieldKey builds the per-field key for a hash field.
func HashFieldScanPrefix ¶
HashFieldScanPrefix returns the prefix to scan all fields of a hash.
func HashMetaDeltaKey ¶
HashMetaDeltaKey builds the delta key for a hash metadata change.
func HashMetaDeltaScanPrefix ¶
HashMetaDeltaScanPrefix returns the prefix to scan all delta keys for a hash.
func HashMetaKey ¶
HashMetaKey builds the metadata key for a hash.
func IsHashFieldKey ¶
IsHashFieldKey reports whether the key is a hash field key.
func IsHashMetaDeltaKey ¶
IsHashMetaDeltaKey reports whether the key is a hash metadata delta key.
func IsHashMetaKey ¶
IsHashMetaKey reports whether the key is a hash metadata key.
func IsListClaimKey ¶
IsListClaimKey reports whether the key is a list claim key.
func IsListItemKey ¶
func IsListMetaDeltaKey ¶
IsListMetaDeltaKey reports whether the key is a list metadata delta key.
func IsListMetaDeltaValue ¶
IsListMetaDeltaValue reports whether value has the fixed delta encoding.
func IsListMetaKey ¶
IsListMetaKey Exported helpers for other packages (e.g., Redis adapter).
func IsSetMemberKey ¶
IsSetMemberKey reports whether the key is a set member key.
func IsSetMetaDeltaKey ¶
IsSetMetaDeltaKey reports whether the key is a set metadata delta key.
func IsSetMetaKey ¶
IsSetMetaKey reports whether the key is a set metadata key.
func IsStreamEntryKey ¶
IsStreamEntryKey reports whether the key is a stream entry key.
func IsStreamMetaKey ¶
IsStreamMetaKey reports whether the key is a stream metadata key.
func IsZSetMemberKey ¶
IsZSetMemberKey reports whether the key is a sorted set member key.
func IsZSetMetaDeltaKey ¶
IsZSetMetaDeltaKey reports whether the key is a sorted set metadata delta key.
func IsZSetMetaKey ¶
IsZSetMetaKey reports whether the key is a sorted set metadata key.
func IsZSetScoreKey ¶
IsZSetScoreKey reports whether the key is a sorted set score index key.
func LegacyListMetaDeltaScanPrefix ¶
LegacyListMetaDeltaScanPrefix returns the pre-upgrade delta scan prefix for a userKey.
func ListClaimKey ¶
ListClaimKey builds the claim key for a list item at the given sequence number.
func ListClaimScanPrefix ¶
ListClaimScanPrefix returns the prefix used to scan all claim keys for a userKey.
func ListItemKey ¶
ListItemKey builds the item key for a user key and sequence number.
func ListMetaDeltaKey ¶
ListMetaDeltaKey builds the delta key for a list: prefix + 4-byte userKeyLen + userKey + 8-byte commitTS + 4-byte seqInTxn.
func ListMetaDeltaScanPrefix ¶
ListMetaDeltaScanPrefix returns the prefix used to scan all delta keys for a userKey.
func ListMetaDeltaScanPrefixes ¶
ListMetaDeltaScanPrefixes returns all prefixes that may contain visible list metadata deltas for userKey. New writers emit only ListMetaDeltaPrefix, but upgrade reads must include LegacyListMetaDeltaPrefix until compaction drains it.
func ListMetaKey ¶
ListMetaKey builds the metadata key for a user key.
func MarshalHashMeta ¶
MarshalHashMeta encodes HashMeta into a fixed 16-byte binary format.
func MarshalHashMetaDelta ¶
func MarshalHashMetaDelta(d HashMetaDelta) []byte
MarshalHashMetaDelta encodes HashMetaDelta into a fixed 8-byte binary format.
func MarshalListMeta ¶
MarshalListMeta encodes ListMeta into a fixed 32-byte binary format.
func MarshalListMetaDelta ¶
func MarshalListMetaDelta(d ListMetaDelta) []byte
MarshalListMetaDelta encodes a ListMetaDelta into a fixed 16-byte binary format.
func MarshalSetMeta ¶
MarshalSetMeta encodes SetMeta into a fixed 16-byte binary format.
func MarshalSetMetaDelta ¶
func MarshalSetMetaDelta(d SetMetaDelta) []byte
MarshalSetMetaDelta encodes SetMetaDelta into a fixed 8-byte binary format.
func MarshalStreamMeta ¶
func MarshalStreamMeta(m StreamMeta) ([]byte, error)
MarshalStreamMeta encodes StreamMeta into a fixed binary format.
func MarshalZSetMeta ¶
MarshalZSetMeta encodes ZSetMeta into a fixed 16-byte binary format.
func MarshalZSetMetaDelta ¶
func MarshalZSetMetaDelta(d ZSetMetaDelta) []byte
MarshalZSetMetaDelta encodes ZSetMetaDelta into a fixed 8-byte binary format.
func MarshalZSetScore ¶
MarshalZSetScore encodes a float64 score in IEEE 754 big-endian format.
func NewWriteConflictError ¶
func PrefixScanEnd ¶
PrefixScanEnd returns the exclusive end key for a prefix scan. It increments the last byte of the prefix; if overflow occurs (all 0xFF), it returns a nil slice which callers must interpret as "scan to end of keyspace".
func SetMemberKey ¶
SetMemberKey builds the per-member key for a set member.
func SetMemberScanPrefix ¶
SetMemberScanPrefix returns the prefix to scan all members of a set.
func SetMetaDeltaKey ¶
SetMetaDeltaKey builds the delta key for a set metadata change.
func SetMetaDeltaScanPrefix ¶
SetMetaDeltaScanPrefix returns the prefix to scan all delta keys for a set.
func SetMetaKey ¶
SetMetaKey builds the metadata key for a set.
func StreamEntryKey ¶
StreamEntryKey builds the per-entry key for a stream.
func StreamEntryScanPrefix ¶
StreamEntryScanPrefix returns the common prefix for every entry belonging to userKey, used as the [start, end) prefix for XRANGE / XREAD scans.
func StreamEntryScanStart ¶
func StreamEntryScanStart(userKey []byte, meta StreamMeta) []byte
StreamEntryScanStart returns the inclusive start key for scans over a stream's live entries. For streams with a persisted trim cursor, the start skips IDs known to have been removed from the logical head.
func StreamMetaKey ¶
StreamMetaKey builds the per-stream metadata key.
func UnmarshalZSetScore ¶
UnmarshalZSetScore decodes a float64 score from IEEE 754 big-endian format.
func WriteConflictKey ¶
func WriterRegistryFor ¶
func WriterRegistryFor(s MVCCStore) (encryption.WriterRegistryStore, error)
WriterRegistryFor returns an encryption.WriterRegistryStore backed by the Pebble DB underlying the supplied MVCCStore. Returns ErrUnsupportedStoreForWriterRegistry if the store is not a Pebble-backed MVCCStore — callers should treat this as a startup-fatal misconfiguration (the only in-tree non-Pebble MVCCStore is the in-memory test fake, which has no encryption requirements).
func ZSetMemberKey ¶
ZSetMemberKey builds the per-member key storing the score for a sorted set.
func ZSetMemberScanPrefix ¶
ZSetMemberScanPrefix returns the prefix to scan all members of a sorted set.
func ZSetMetaDeltaKey ¶
ZSetMetaDeltaKey builds the delta key for a sorted set metadata change.
func ZSetMetaDeltaScanPrefix ¶
ZSetMetaDeltaScanPrefix returns the prefix to scan all delta keys for a sorted set.
func ZSetMetaKey ¶
ZSetMetaKey builds the metadata key for a sorted set.
func ZSetScoreKey ¶
ZSetScoreKey builds the score index key for a sorted set entry. Layout: !zs|scr|<userKeyLen(4)><userKey><sortableScore(8)><member>
func ZSetScoreRangeScanPrefix ¶
ZSetScoreRangeScanPrefix returns the prefix for scanning scores in [minScore, maxScore].
func ZSetScoreScanPrefix ¶
ZSetScoreScanPrefix returns the prefix to scan all score index keys for a sorted set.
Types ¶
type ActiveStorageKeyID ¶
ActiveStorageKeyID reports the currently-active storage DEK identifier. The bool is false when no storage DEK is active (i.e. the cluster has not run Phase 1 of the §7.1 rollout yet) — in that case the storage layer writes cleartext as if no cipher were configured. Stage 5/6 wires this from the sidecar's Active.Storage slot; Stage 2 takes it as a closure so test code can flip it independently.
type CounterNonceFactory ¶
type CounterNonceFactory struct {
// contains filtered or unexported fields
}
CounterNonceFactory is a test-only NonceFactory that produces the design §4.1 deterministic nonce shape (`node_id ‖ local_epoch ‖ write_count`) without the writer-registry round-trip Stage 7 brings. Production wiring uses the registry-backed factory; this implementation is only safe for tests where the caller controls every node_id / local_epoch combination.
Exposed (vs. living in a *_test.go file) so the encryption integration tests in other packages can build on the same implementation without re-deriving the byte layout. It is nevertheless test-grade — the doc comment on NonceFactory emphasises that production callers MUST guarantee (node_id, local_epoch, write_count) uniqueness.
func NewCounterNonceFactory ¶
func NewCounterNonceFactory(nodeID, localEpoch uint16) *CounterNonceFactory
NewCounterNonceFactory constructs a CounterNonceFactory pinned to the given (nodeID, localEpoch). write_count starts at 0 and monotonically increments on every Next().
func (*CounterNonceFactory) Next ¶
func (f *CounterNonceFactory) Next() ([encryption.NonceSize]byte, error)
Next produces the next 12-byte nonce. Layout matches design §4.1:
bytes 0-1 node_id (big-endian uint16) bytes 2-3 local_epoch (big-endian uint16) bytes 4-11 write_count (big-endian uint64)
Big-endian is chosen so a hex dump of consecutive nonces is human-readable as a counter; the AAD does NOT include the nonce bytes (the cipher composes the nonce into AES-GCM directly), so the byte order is internal to the factory.
type ExportVersionsOptions ¶
type ExportVersionsOptions struct {
StartKey []byte
EndKey []byte
MinCommitTSExclusive uint64
MaxCommitTSInclusive uint64
Cursor []byte
MaxVersions int
MaxBytes uint64
MaxScannedBytes uint64
KeyFamily uint32
AcceptKey func([]byte) bool
AcceptVersion func(key []byte, value []byte) bool
}
ExportVersionsOptions selects a raw MVCC-version export window.
type ExportVersionsResult ¶
type ExportVersionsResult struct {
Versions []MVCCVersion
NextCursor []byte
Done bool
ScannedBytes uint64
ExportedBytes uint64
AcceptedRows uint64
}
ExportVersionsResult is one resumable chunk of raw MVCC versions.
type HashMeta ¶
HashMeta is the base metadata for a hash collection.
func UnmarshalHashMeta ¶
UnmarshalHashMeta decodes HashMeta from the legacy 8-byte or current 16-byte binary format.
type HashMetaDelta ¶
type HashMetaDelta struct {
LenDelta int64
}
HashMetaDelta holds a signed change in field count.
func UnmarshalHashMetaDelta ¶
func UnmarshalHashMetaDelta(b []byte) (HashMetaDelta, error)
UnmarshalHashMetaDelta decodes HashMetaDelta from the fixed 8-byte binary format.
type HybridClock ¶
type HybridClock interface {
Now() uint64
}
HybridClock provides monotonically increasing timestamps (HLC).
type ImportVersionsOptions ¶
type ImportVersionsOptions struct {
JobID uint64
BracketID uint64
BatchSeq uint64
Versions []MVCCVersion
Cursor []byte
}
ImportVersionsOptions applies one idempotent migration-import batch.
type ImportVersionsResult ¶
ImportVersionsResult reports the cursor durably acknowledged by the target.
type KVPairMutation ¶
type KVPairMutation struct {
Op OpType
Key []byte
Value []byte
// ExpireAt is an HLC timestamp; 0 means no TTL.
ExpireAt uint64
}
KVPairMutation is a small helper struct for MVCC mutation application.
type ListMeta ¶
type ListMeta struct {
Head int64 `json:"h"`
Tail int64 `json:"t"`
Len int64 `json:"l"`
ExpireAt uint64 `json:"e,omitempty"`
}
func UnmarshalListMeta ¶
UnmarshalListMeta decodes ListMeta from the legacy 24-byte or current 32-byte binary format.
type ListMetaDelta ¶
ListMetaDelta holds the signed deltas applied by a single PUSH/POP operation.
func UnmarshalListMetaDelta ¶
func UnmarshalListMetaDelta(b []byte) (ListMetaDelta, error)
UnmarshalListMetaDelta decodes a ListMetaDelta from the fixed 16-byte binary format.
type MVCCStore ¶
type MVCCStore interface {
// GetAt returns the newest version whose commit timestamp is <= ts.
GetAt(ctx context.Context, key []byte, ts uint64) ([]byte, error)
// ExistsAt reports whether a visible, non-tombstone version exists at ts.
ExistsAt(ctx context.Context, key []byte, ts uint64) (bool, error)
// ScanAt returns versions visible at the given timestamp.
ScanAt(ctx context.Context, start []byte, end []byte, limit int, ts uint64) ([]*KVPair, error)
// ScanKeysAt returns keys with visible versions at the given timestamp.
ScanKeysAt(ctx context.Context, start []byte, end []byte, limit int, ts uint64) ([][]byte, error)
// ReverseScanAt returns visible versions in descending key order for keys in [start, end).
ReverseScanAt(ctx context.Context, start []byte, end []byte, limit int, ts uint64) ([]*KVPair, error)
// PutAt commits a value at the provided commit timestamp and optional expireAt.
PutAt(ctx context.Context, key []byte, value []byte, commitTS uint64, expireAt uint64) error
// DeleteAt commits a tombstone at the provided commit timestamp.
DeleteAt(ctx context.Context, key []byte, commitTS uint64) error
// PutWithTTLAt stores a value with a precomputed expireAt (HLC) at the given commit timestamp.
PutWithTTLAt(ctx context.Context, key []byte, value []byte, commitTS uint64, expireAt uint64) error
// ExpireAt sets/renews TTL using a precomputed expireAt (HLC) at the given commit timestamp.
ExpireAt(ctx context.Context, key []byte, expireAt uint64, commitTS uint64) error
// LatestCommitTS returns the commit timestamp of the newest version.
// The boolean reports whether the key has any version.
LatestCommitTS(ctx context.Context, key []byte) (uint64, bool, error)
// CommittedVersionAt reports whether a committed version stamped
// EXACTLY commitTS exists for key. Unlike GetAt (newest version
// <= ts) this is an exact-timestamp existence check, used by the
// one-phase transaction idempotency probe to ask "did the previous
// attempt — which committed at this exact commit_ts — land?". Because
// commit timestamps are issued by the strictly-monotonic, unique HLC
// (Clock().Next()), a version at an exact commitTS on a given key can
// only have come from the transaction that was assigned that timestamp,
// so an exact hit unambiguously identifies that attempt. A tombstone
// counts as a landed version (the attempt committed a delete). See
// docs/design/2026_05_21_proposed_txn_secondary_idempotency.md.
CommittedVersionAt(ctx context.Context, key []byte, commitTS uint64) (bool, error)
// ApplyMutations atomically validates and appends the provided mutations.
// It must return ErrWriteConflict if any mutation key or any read key has
// a newer commit timestamp than startTS. readKeys carries the transaction's
// read set for read-write conflict detection; pass nil when no read set
// validation is needed.
//
// Isolation guarantees vary by transaction topology:
//
// Single-shard transactions: readKeys are included in the Raft log entry
// and validated atomically under the FSM's applyMu lock alongside
// write-write conflict detection. The adapter's pre-Raft validateReadSet
// call is kept as a fast-fail optimization but the FSM check is
// authoritative. No TOCTOU window; full SSI.
//
// Multi-shard (2PC) write shards: readKeys are included in the
// PREPARE Raft entry and validated atomically under the FSM's applyMu
// lock. No TOCTOU window; full SSI.
//
// Multi-shard (2PC) read-only shards: validated via a linearizable
// read barrier followed by LatestCommitTS outside the FSM lock. A
// small TOCTOU window exists between the barrier and the check.
ApplyMutations(ctx context.Context, mutations []*KVPairMutation, readKeys [][]byte, startTS, commitTS uint64) error
// ApplyMutationsRaft is the raft-apply variant of ApplyMutations. It
// carries identical MVCC semantics but is governed by the FSM-commit
// sync-mode knob (ELASTICKV_FSM_SYNC_MODE). Callers MUST only use this
// when the write is part of a raft-log apply — the raft WAL is the
// durability backstop that makes an un-fsynced Pebble write safe.
//
// Direct (non-raft) callers (catalog bootstrap, admin snapshots,
// migrations, tests) must use ApplyMutations, which is always
// pebble.Sync and therefore safe without raft-log replay.
ApplyMutationsRaft(ctx context.Context, mutations []*KVPairMutation, readKeys [][]byte, startTS, commitTS uint64) error
// ApplyMutationsRaftAt is ApplyMutationsRaft with the raft entry
// index threaded through. The leaf bundles metaAppliedIndex in the
// same pebble.Batch as the data mutation so a successful Apply
// implies LastAppliedIndex >= appliedIndex; the cold-start
// snapshot-restore skip gate uses this invariant (PR #910 / B2).
// appliedIndex==0 is treated as "no index, do not bump the meta
// key", matching ApplyMutationsRaft semantics for callers that have
// not yet been wired to the raftengine.ApplyIndexAware seam.
ApplyMutationsRaftAt(ctx context.Context, mutations []*KVPairMutation, readKeys [][]byte, startTS, commitTS, appliedIndex uint64) error
// DeletePrefixAt atomically deletes all visible (non-tombstone, non-expired)
// keys matching prefix at commitTS by writing tombstone versions. An empty
// prefix means "all keys". Keys matching excludePrefix are preserved.
// No conflict checking is performed; this is intended for bulk operations
// such as FLUSHALL where the caller knows no conflict check is needed.
DeletePrefixAt(ctx context.Context, prefix []byte, excludePrefix []byte, commitTS uint64) error
// DeletePrefixAtRaft is the raft-apply variant of DeletePrefixAt with
// the same durability contract as ApplyMutationsRaft.
DeletePrefixAtRaft(ctx context.Context, prefix []byte, excludePrefix []byte, commitTS uint64) error
// DeletePrefixAtRaftAt is DeletePrefixAtRaft with the raft entry
// index threaded through. handleDelPrefix builds an independent
// pebble.Batch separate from applyMutationsWithOpts; this overload
// bundles metaAppliedIndex in that batch so DEL_PREFIX entries
// also advance the meta key. PR #910 design §2 "why both leaves".
DeletePrefixAtRaftAt(ctx context.Context, prefix []byte, excludePrefix []byte, commitTS, appliedIndex uint64) error
// LastCommitTS returns the highest commit timestamp applied on this node.
LastCommitTS() uint64
// WriteConflictCountsByPrefix returns a snapshot of the MVCC
// write-conflict counters keyed by "<kind>|<key_prefix>" where
// kind is "read" or "write" and key_prefix is a bounded
// classification of the conflicting key. The map is a copy; the
// caller may mutate it freely. Implementations that do not track
// conflicts may return an empty (non-nil) map.
WriteConflictCountsByPrefix() map[string]uint64
// Compact removes versions older than minTS that are no longer needed.
Compact(ctx context.Context, minTS uint64) error
// ExportVersions exports raw committed MVCC versions for range migration.
ExportVersions(ctx context.Context, opts ExportVersionsOptions) (ExportVersionsResult, error)
// ImportVersions applies a migration import batch idempotently by
// (jobID, bracketID, batchSeq), preserving tombstones and expireAt.
ImportVersions(ctx context.Context, opts ImportVersionsOptions) (ImportVersionsResult, error)
// MigrationHLCFloor returns the full-HLC target-local migration floor
// persisted by ImportVersions for jobID.
MigrationHLCFloor(ctx context.Context, jobID uint64) (uint64, error)
// RetireMigration removes target-local import progress metadata for a
// completed migration job. Data imported by the job is left intact.
RetireMigration(ctx context.Context, jobID uint64) error
Snapshot() (Snapshot, error)
Restore(buf io.Reader) error
Close() error
}
MVCCStore extends Store with multi-version concurrency control helpers. The interface is timestamp-explicit; callers must supply the snapshot or commit timestamp for every operation.
func NewMVCCStore ¶
func NewMVCCStore(opts ...MVCCStoreOption) MVCCStore
NewMVCCStore creates a new MVCC-enabled in-memory store.
func NewPebbleStore ¶
func NewPebbleStore(dir string, opts ...PebbleStoreOption) (MVCCStore, error)
NewPebbleStore creates a new Pebble-backed MVCC store.
type MVCCStoreOption ¶
type MVCCStoreOption func(*mvccStore)
MVCCStoreOption configures the MVCCStore.
func WithLogger ¶
func WithLogger(l *slog.Logger) MVCCStoreOption
WithLogger sets a custom logger for the store.
type MVCCVersion ¶
type MVCCVersion struct {
Key []byte
CommitTS uint64
Tombstone bool
Value []byte
KeyFamily uint32
ExpireAt uint64
}
MVCCVersion is a raw committed MVCC version for range migration. Unlike scan results, it preserves tombstones and TTL expiry metadata.
type NonceFactory ¶
type NonceFactory interface {
Next() ([encryption.NonceSize]byte, error)
}
NonceFactory produces unique 12-byte AES-GCM nonces for the storage envelope (§4.1). The factory is responsible for the cluster-wide uniqueness invariant across `(node_id, local_epoch, write_count)` — the storage layer just calls Next() and uses what comes back.
Stage 7 of the encryption rollout will replace the in-tree reference implementation (CounterNonceFactory) with a writer-registry-backed factory that guarantees uniqueness across voters, learners, and historical replicas. The interface stays the same; only the construction changes. Implementations MUST NOT return the same nonce twice under the same DEK — AES-GCM nonce reuse is catastrophic (see encryption.Cipher doc).
type PebbleStoreOption ¶
type PebbleStoreOption func(*pebbleStore)
PebbleStoreOption configures the PebbleStore.
func WithEncryption ¶
func WithEncryption(cipher *encryption.Cipher, nf NonceFactory, activeKeyID ActiveStorageKeyID) PebbleStoreOption
WithEncryption configures the pebble-backed store to wrap every committed value in the §4.1 storage envelope.
All three arguments must be non-nil. activeKeyID is called on every Put — when it returns ok=false the store writes cleartext (encryption_state = 0b00) even though a cipher is wired, matching the §7.1 Phase 0 / Phase 1 split where capability is provisioned before activation. Reads that observe encryption_state = 0b01 always go through the cipher regardless of activeKeyID, so a cluster mid-cutover stays readable.
Calling WithEncryption with any nil argument is a no-op (the store stays in legacy cleartext-only mode). This keeps the option backwards-compatible with every existing NewPebbleStore caller and keeps the Stage 2 wiring trivially reversible.
func WithPebbleLogger ¶
func WithPebbleLogger(l *slog.Logger) PebbleStoreOption
WithPebbleLogger sets a custom logger.
func WithSSTIngestSnapshots ¶
func WithSSTIngestSnapshots(enabled bool) PebbleStoreOption
WithSSTIngestSnapshots controls whether Snapshot emits the SST-ingest format. It is primarily useful for tests and embedded callers; production resolves ELASTICKV_PEBBLE_SST_INGEST_SNAPSHOT at store construction time.
func WithStorageEnvelopeGate ¶
func WithStorageEnvelopeGate(active StorageEnvelopeActive) PebbleStoreOption
WithStorageEnvelopeGate wires the Stage 6D-5 §6.2 cutover gate in front of the existing envelope-emit path. After this option, encryptForKey writes cleartext when active() returns false even if a cipher and active DEK are present — exactly the §7.1 Phase 0 / Phase 1 split described in the parent encryption design doc.
Passing nil is a no-op: the store stays in the pre-6D-5 "encrypt whenever activeKeyID returns a DEK" posture used by existing test fixtures and by production deployments that have not yet run the EnableStorageEnvelope RPC (lands in 6D-6). Operators must thread this option in alongside WithEncryption to actually opt into the cutover semantics.
The gate is consulted on every Put. Reads never consult it — on-disk versions with `encryption_state == 0b01` always go through the cipher regardless of the gate, so a cluster mid- cutover (some versions cleartext, others encrypted) stays readable.
func WithStorageEnvelopeV2Gate ¶
func WithStorageEnvelopeV2Gate(active StorageEnvelopeV2Active) PebbleStoreOption
WithStorageEnvelopeV2Gate prevents V2 writes until the cluster-wide capability fan-out has observed support from every voter and learner. Reads always accept both V1 and V2. A nil gate preserves the standalone and test-fixture behavior from before capability negotiation was wired.
func WithStorageRegistrationGate ¶
func WithStorageRegistrationGate(registered StorageRegistered) PebbleStoreOption
WithStorageRegistrationGate wires the Stage 7a-2 §4.1 registration gate in front of the envelope-emit path on the DIRECT write path only. After this option, when the envelope would encrypt (cipher + active DEK + StorageEnvelopeActive all true) but registered() reports false, the direct path (PutAt / ExpireAt / ApplyMutations) returns ErrWriterNotRegistered instead of emitting a nonce. The FSM-apply path (ApplyMutationsRaft) passes gateRegistration=false into encryptForKey and is never gated — see design §1: replicated apply must stay deterministic and may run before this node's own registration entry commits, so fail-closing it would halt the apply loop (and deadlock a node whose storage entry is ordered before its registration entry).
Passing nil is a no-op: the store keeps the pre-7a-2 posture where the direct path emits envelopes as soon as the cutover gate is open, used by existing fixtures and by deployments that have not wired the writer registry. Operators thread cache.Registered in via main.go.
type RetentionController ¶
RetentionController exposes the minimum timestamp still retained by a store after MVCC compaction. Reads older than this watermark may fail with ErrReadTSCompacted.
type SetMeta ¶
SetMeta is the base metadata for a set collection.
func UnmarshalSetMeta ¶
UnmarshalSetMeta decodes SetMeta from the legacy 8-byte or current 16-byte binary format.
type SetMetaDelta ¶
type SetMetaDelta struct {
LenDelta int64
}
SetMetaDelta holds a signed change in member count.
func UnmarshalSetMetaDelta ¶
func UnmarshalSetMetaDelta(b []byte) (SetMetaDelta, error)
UnmarshalSetMetaDelta decodes SetMetaDelta from the fixed 8-byte binary format.
type Snapshot ¶
Snapshot streams a consistent point-in-time store image to a writer. Implementations may back this with a temp file or an engine-native snapshot.
type StorageEnvelopeActive ¶
type StorageEnvelopeActive func() bool
StorageEnvelopeActive reports whether the §7.1 Phase 1 cutover has fired on this node (the Stage 6D-4 sidecar field of the same name, surfaced to the storage layer via WithStorageEnvelopeGate). A return value of false forces every Put to write cleartext even when a cipher + active DEK are wired; true allows the existing activeStorageKeyID-driven envelope emit to proceed.
The contract intentionally separates "encryption is provisioned" (ActiveStorageKeyID has a DEK to use) from "encryption is activated for new writes" (the cutover has fired). The two signals diverge during the §7.1 Phase 0 → Phase 1 window: every node in the cluster has the DEK installed, but the cluster has not yet quiesced the cutover Raft entry that flips the bit on every replica simultaneously.
Reads ignore this gate — once an on-disk version's encryption_state bit is 0b01, the read path always invokes the cipher to unwrap regardless of the gate's current value. A cluster that flipped from active back to inactive (not a path the design supports) would still decrypt old envelopes correctly.
type StorageEnvelopeV2Active ¶
type StorageEnvelopeV2Active func() bool
StorageEnvelopeV2Active reports whether every current voter and learner can read V2 storage envelopes. When wired and false, encrypted writes remain V1 and uncompressed so a rolling-upgrade peer can still read snapshots.
type StorageRegistered ¶
type StorageRegistered func() bool
StorageRegistered reports whether this process load's §4.1 writer registration has committed for the currently-active storage DEK (Stage 7a-2). It gates only the DIRECT write path: when it returns false while the envelope would otherwise encrypt, encryptForKey refuses to emit a nonce and returns ErrWriterNotRegistered so the self-originated caller (catalog bootstrap Save, admin snapshot, migration) retries until registration lands rather than emitting an envelope ahead of its writer-registry row.
The FSM-apply path never consults it — see WithStorageRegistrationGate.
type StreamMeta ¶
type StreamMeta struct {
Length int64
LastMs uint64
LastSeq uint64
ExpireAt uint64
TrimmedMs uint64
TrimmedSeq uint64
}
StreamMeta is the per-stream metadata. Length is authoritative for XLEN; LastMs/LastSeq track the highest ID ever appended so XADD '*' stays strictly monotonic even after XTRIM removes the current tail. TrimmedMs/TrimmedSeq track the largest ID removed by a head trim. They are a seek cursor, not part of Redis-visible stream state: when non-zero, scans for the logical head can start strictly after that ID and skip old tombstoned MVCC history.
func UnmarshalStreamMeta ¶
func UnmarshalStreamMeta(b []byte) (StreamMeta, error)
UnmarshalStreamMeta decodes StreamMeta from the legacy 24-byte, 32-byte, or current 48-byte binary format.
type VersionPresenceReader ¶
type VersionPresenceReader interface {
VersionExistsAtOrBefore(ctx context.Context, key []byte, ts uint64) (bool, error)
}
VersionPresenceReader reports whether a key has any committed version at or before a read timestamp, including tombstones and expired values.
type VersionedValue ¶
type VersionedValue struct {
TS uint64
Value []byte
Tombstone bool
ExpireAt uint64 // HLC timestamp; 0 means no TTL
}
VersionedValue represents a single committed version in MVCC storage.
type WriteConflictError ¶
type WriteConflictError struct {
// contains filtered or unexported fields
}
func (*WriteConflictError) Error ¶
func (e *WriteConflictError) Error() string
func (*WriteConflictError) Unwrap ¶
func (e *WriteConflictError) Unwrap() error
type ZSetMeta ¶
ZSetMeta is the base metadata for a sorted set collection.
func UnmarshalZSetMeta ¶
UnmarshalZSetMeta decodes ZSetMeta from the legacy 8-byte or current 16-byte binary format.
type ZSetMetaDelta ¶
type ZSetMetaDelta struct {
LenDelta int64
}
ZSetMetaDelta holds a signed change in member count.
func UnmarshalZSetMetaDelta ¶
func UnmarshalZSetMetaDelta(b []byte) (ZSetMetaDelta, error)
UnmarshalZSetMetaDelta decodes ZSetMetaDelta from the fixed 8-byte binary format.
Source Files
¶
- encryption_glue.go
- encryption_test_helpers.go
- hash_helpers.go
- list_helpers.go
- lsm_migration.go
- lsm_store.go
- migration_versions.go
- mvcc_store.go
- pebble_memory_budget.go
- pebble_memory_linux.go
- set_helpers.go
- snapshot_file.go
- snapshot_pebble.go
- snapshot_pebble_sst.go
- store.go
- stream_helpers.go
- write_conflict_counter.go
- zset_helpers.go