staterecord

package
v0.18.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 19, 2026 License: MPL-2.0 Imports: 26 Imported by: 0

Documentation

Overview

Package staterecord is a small, versioned key/value Store with first-class conditional writes: Get, PutIfVersion, PutIfAbsent, Delete, List — nothing else. It backs the micro-state records issue #73's charter describes (record-less residue: null_resource, terraform_data, time_*, non-sensitive random_* run through the stock provider lifecycle against an in-memory state hydrated from and CAS-persisted to one small record per resource), but this package itself knows nothing about that. It has no notion of an estate, a resource, redaction, or anything else choudoufu-specific — keys are opaque strings, payloads are opaque bytes, and every choudoufu concept (what a key names, what goes in a payload, which resources get one) lives entirely in the caller.

Why that separation is the point

This package is meant to be upstream-adoptable verbatim: proposable to OpenTofu as a lightweight state backend on its own merits, independent of choudoufu ever existing. Concretely, that shapes three decisions:

  • The Store interface follows upstream's own backend conventions — a clean Get/Put-with-condition/Delete surface, no fork-specific types anywhere in its signatures, workspace-agnostic naming (a "key", not a "workspace" or a "resource address").
  • Conditional-write/CAS is a first-class interface concept, not something bolted onto a plain Put as an optional flag. Upstream's own s3-locking-with-conditional-writes RFC (20250211) already shows appetite for exactly this primitive as a first-class one.
  • The package directory holds only store implementations and their tests — nothing that imports estate configuration, redaction rules, or resource-selection logic. A third store (issue #73's ruling: "design the interface so a third store is a new file, not a refactor") is one new file implementing Store, never a change to this one.

The interface contract, precisely

  • Keys are opaque strings. Every implementation accepts a reasonably portable subset — this package itself only rejects the empty string, a NUL byte, and a ".." path segment (see validateKey) — but each store's own backend (a filesystem, an S3 object key) may reject a key its own naming rules forbid; that surfaces as an ordinary error, not a Store-defined one.
  • Payloads are opaque []byte. No implementation inspects, parses, or redacts a payload's content; that is the caller's job, every time, before a payload reaches this package and after one leaves it.
  • Versions are opaque strings with exactly one universal meaning: "" denotes "no record exists here." No implementation ever assigns "" as a live record's version, so a caller can treat it as a stable sentinel without inspecting which store it is talking to. Beyond that, a version's shape is entirely implementation-defined — a content hash, an S3 ETag — and Store callers are expected to hold it opaque too: compare it for equality, pass it to PutIfVersion/Delete, never parse it.
  • Every conditional operation that fails on a version mismatch reports exactly one error type: *VersionConflictError, naming both the version the caller expected and the version the store actually found (or "" for "no record"). A caller never has to distinguish "conflict" from "some other failure" by parsing prose.
  • "Conditional" means real compare-and-swap with no read-compare-write race window, on every store: LocalStore and S3Store both give it. That is a requirement of the interface and not a property two implementations happen to share. See "The store that was retired".

The two implementations

LocalStore (a directory of files, the zero-configuration default — solo development, tests, air-gapped runs, mirroring plain local state's own "just works" shape) and S3Store (S3 conditional writes, for anything more than one operator shares). Both implement the identical Store interface; a caller choosing between them is choosing an operational tradeoff, never a different programming model.

The store that was retired

Until GitHub issue #1346 there was a third, on AWS Systems Manager Parameter Store, and it was the default recommendation for a team. It was retired as a RECORD store for three reasons. Standard parameters cap at 10,000 per account and region, against the customer's own quota. Past that, every parameter bills monthly on the advanced tier. And it has no general conditional write: it could create-if-absent and nothing else, so every update and delete was a read-compare-write with a race window, where this package's whole consistency story is a per-key conditional write.

No migration was written, because no estate was on it when it was retired. That is the reason, and it is recorded so nobody later assumes a migration path was designed and lost.

This says nothing about Parameter Store for SECRET values. Keeping secret material out of the bucket, in SSM, is planned (#1244 section 3) and not built; nothing in this package does it today.

Index

Constants

View Source
const (
	KMSDeniedByKeyPolicy      = "key-policy"
	KMSDeniedByIdentityPolicy = "identity-policy"
	KMSDeniedExplicitly       = "explicit-deny"
)

What AWS blames a KMS denial on.

View Source
const DefaultS3GetAllParallelism = 8

DefaultS3GetAllParallelism is how many GetObject calls S3Store.GetAll has in flight at once unless S3Config.GetAllParallelism says otherwise.

Eight is chosen to be unremarkable, not tuned, and it has not been measured at scale. The namespace read here holds one record per managed instance - internal/live/projection's write-back records every instance, an identity envelope for an ordinary taggable resource as well as the whole value of a record-backed one - so the read is N GetObject calls for an estate of N instances and this bound sets how long that takes. An earlier version of this comment called the namespace "the record-backed slice only, a small fraction of an estate" and the bound "not load-bearing". That was the design text's claim and the code never matched it. The bound is configurable because the estate that needs otherwise will know why and should not have to patch the binary to find out. GitHub issue #1336.

Variables

BucketSettings is every asserted setting, in the order findings are reported.

Functions

func BucketContractRefusal added in v0.18.0

func BucketContractRefusal(bucket string, f BucketFinding) (summary, detail string)

BucketContractRefusal is the headline and the paragraph for one failed finding, in internal/command's statelessCommandRefusals shape: what was refused, then what it protects against and what to do instead. Empty for a finding that passed.

func BucketWaiverCost added in v0.18.0

func BucketWaiverCost(setting BucketSetting) string

BucketWaiverCost says what an estate gives up by waiving setting, as a clause that completes "is waived, so ...". GitHub issue #1340: the warning names the setting and its cost in the same sentence, never a generic "running with reduced checks", because a cost the reader has to look up is a cost they have already decided not to read.

func NamespacePrefix added in v0.18.0

func NamespacePrefix(prefix string) string

NamespacePrefix returns prefix with exactly one trailing "/", and "" for "".

Store.List and BulkReader.GetAll match an ordinary string prefix, and so does S3's ListObjectsV2. A namespace handed to either without its trailing delimiter therefore matches every sibling whose name merely starts the same way: "tofu-records/prod" lists "tofu-records/prod-eu/..." too. GitHub issue #1335 measured that for two estates sharing one store, which the bucket backend (#1332) makes the recommended arrangement. Under that backend's IAM model the listing is defended by the s3:prefix condition ALONE - an object tag cannot condition a LIST, which touches no object - so the delimiter is what the isolation rests on, not tidiness.

Every layer that turns a namespace into a List or GetAll prefix goes through this one function, so the delimiter cannot be present in the key builder and missing from the listing, or the other way round.

func ObjectTags added in v0.18.0

func ObjectTags(ctx context.Context) map[string]string

ObjectTags is what WithObjectTags put in ctx, nil when nothing did.

func ResetRunCacheForTest added in v0.13.0

func ResetRunCacheForTest(t testing.TB)

ResetRunCacheForTest clears the process-wide "something has been written" switch that RunCache uses to decide whether it may still trust its snapshot (see RunCache's doc comment, "Why it cannot serve a stale value", for why that switch is one atomic per process rather than one per cache).

The switch is deliberately sticky for the product - one process is one run, and a run that has written once must never serve a remembered value again - but that same stickiness means it is sticky across every test in a `go test -count=N` process: iteration 1's write turns every cache off for iterations 2..N, so a test whose premise is "the cache is currently serving" has no way to establish that premise from inside the test.

Call it before any assertion that depends on the switch's state. t.Cleanup restores whatever the switch held before the call, so this cannot leak a fixed state into whatever else runs in the same process afterward.

func SortedByCount added in v0.5.0

func SortedByCount(m map[string]int) []string

SortedByCount renders a bucket map as lines ordered by descending count, ties broken by name, for a report a human reads.

func WithObjectTags added in v0.18.0

func WithObjectTags(ctx context.Context, tags map[string]string) context.Context

WithObjectTags returns ctx carrying tags for the next write made with it. A backend with nowhere to put tags ignores them. Tags already in ctx are kept, with the new ones winning on a shared key.

Types

type BucketContractAPI added in v0.18.0

type BucketContractAPI interface {
	GetBucketVersioning(ctx context.Context, in *s3.GetBucketVersioningInput, optFns ...func(*s3.Options)) (*s3.GetBucketVersioningOutput, error)
	GetBucketLifecycleConfiguration(ctx context.Context, in *s3.GetBucketLifecycleConfigurationInput, optFns ...func(*s3.Options)) (*s3.GetBucketLifecycleConfigurationOutput, error)
	GetPublicAccessBlock(ctx context.Context, in *s3.GetPublicAccessBlockInput, optFns ...func(*s3.Options)) (*s3.GetPublicAccessBlockOutput, error)
}

BucketContractAPI is the three reads the contract needs, and the permissions they cost: s3:GetBucketVersioning, s3:GetLifecycleConfiguration and s3:GetBucketPublicAccessBlock. *s3.Client satisfies it.

type BucketContractChecker added in v0.18.0

type BucketContractChecker interface {
	CheckBucketContract(ctx context.Context, namespaces []string) ([]BucketFinding, error)
}

BucketContractChecker is implemented by a store that lives in a bucket. The local store does not implement it: a directory has no such settings, and a store with nothing to assert is not a store that failed.

func AsBucketContractChecker added in v0.18.0

func AsBucketContractChecker(s Store) (BucketContractChecker, bool)

AsBucketContractChecker finds the bucket-backed store under s, looking through this package's own wrappers (RunCache, CountingStore). False means there is nothing to assert - a local store - which is a different answer from a bucket that failed.

type BucketFinding added in v0.18.0

type BucketFinding struct {
	Setting BucketSetting

	// OK is true when the bucket satisfies the assertion.
	OK bool

	// Unreadable is true when the setting could not be read at all - the
	// role lacks the Get* permission, typically. It is never true together
	// with OK. From the caller's side it is the same refusal a wrong setting
	// gets, because a bucket nobody could check is not a bucket that passed.
	Unreadable bool

	// Found says what the bucket actually has, in one clause, for the
	// refusal to quote: "versioning is Suspended", "no lifecycle
	// configuration", "s3:GetBucketVersioning was denied".
	Found string
}

BucketFinding is what one setting turned out to be.

func CheckBucketContract added in v0.18.0

func CheckBucketContract(ctx context.Context, api BucketContractAPI, bucket string, namespaces []string) ([]BucketFinding, error)

CheckBucketContract reads the three settings of bucket and reports one finding per setting, always all three and always in BucketSettings order: a caller that refuses on the first bad one would make an operator fix them one run at a time.

namespaces are the object-key prefixes the caller writes under - an estate's records, hint and outputs. The lifecycle assertion needs them: a rule scoped to some other prefix expires nothing of ours, and a rule that covers the records but not the outputs leaves the outputs growing. With none given, only a rule with no prefix filter counts.

The error return is for a failure that is not about the bucket's settings at all - a cancelled context, an unreachable endpoint. A denied read is NOT an error: it is a finding with Unreadable set.

func SplitWaived added in v0.18.0

func SplitWaived(findings []BucketFinding, waived []string) (refused, waivedFailing []BucketFinding)

SplitWaived sorts the findings that did not pass into the ones the run must refuse on and the ones waived names, leaving passing findings out of both. A waiver reaches exactly the settings it names: waiving one leaves a failure of either of the others in refused.

An unreadable setting is waived by the same name as a wrong one. From the caller's side they are one refusal - the run cannot rely on the setting - and an operator whose role cannot read the bucket's configuration has no other way to proceed.

type BucketSetting added in v0.18.0

type BucketSetting string

BucketSetting names one asserted setting. The values are the names an operator writes in configuration (#1340's allow_insecure), so they are part of the configuration language and do not change casually.

const (
	BucketVersioning        BucketSetting = "versioning"
	BucketLifecycle         BucketSetting = "lifecycle"
	BucketPublicAccessBlock BucketSetting = "public_access_block"
)

type BulkReader added in v0.5.0

type BulkReader interface {
	GetAll(ctx context.Context, keyPrefix string) (map[string]Record, error)
}

BulkReader is the optional half of Store that loads a whole namespace in one call. It exists because a plan needs the entire estate's records and stock OpenTofu gets its whole equivalent — the state file — in one read. Without it, a converged plan's cheapest possible shape is still one call per instance for information no per-instance decision needed separately.

GetAll returns every record whose key begins with keyPrefix, keyed by the same key Store.Get and Store.List use. An empty keyPrefix means every key. The returned map is the CALLER's; implementations must not retain or reuse it.

The result is complete for that prefix: a key absent from the map holds no record. That is what makes it usable as a snapshot rather than a warm cache — a reader can answer "there is nothing recorded for this address" from it without going back to the store.

It is optional rather than part of Store because it is an optimization and not a semantic: a backend with no way to enumerate values simply does not implement it, and every caller keeps working through Store.Get.

type CountingStore added in v0.5.0

type CountingStore struct {
	// contains filtered or unexported fields
}

CountingStore wraps a Store and counts every operation that reaches it, which is the record-store half of what internal/live/flocitest's CountingProxy counts for provider traffic. The proxy stands in front of the AWS endpoint and is therefore blind to these: a record read goes to a local directory or an S3 object, neither of which the provider's endpoint ever sees. Until this existed, no instrument in this repository counted them at all, so "a plan costs N calls" was a partial number by construction.

A "trip" is one call to the wrapped store, because that is the unit that costs something: against LocalStore a stat plus a read, against S3Store a network round trip. Stock OpenTofu makes zero of them — it reads its whole state once, from one file.

Each trip records more than a method name, because a bare per-method total cannot answer either question a reduction needs answered:

  • Site is the first frame outside this package and outside projection.(*RecordStore) — the code that actually wanted the record. This is the per-site breakdown; without it, "158 Gets" names no line to fix.
  • Via is the outermost projection.(*RecordStore) method the call came through (GetIdentity, GetResidue, getProvisioned, ...), so several sites reading the same underlying envelope through different accessors stay distinguishable.
  • Key is the store key, so [CountingStore.RepeatTrips] can report how many trips re-read a key some earlier trip already read. That number is the size of the prize a cache can win, measured rather than assumed.

It is safe for concurrent use; the sweep and the projection both run goroutines.

func NewCountingStore added in v0.5.0

func NewCountingStore(inner Store, log io.Writer) *CountingStore

NewCountingStore wraps inner. log may be nil; when it is not, every trip is written to it as a single Trip.String line followed by a newline, under this store's own lock. A caller sharing one writer between several counting stores is responsible for that writer being safe to call from several goroutines.

func (*CountingStore) Counts added in v0.5.0

func (c *CountingStore) Counts() TripCounts

Counts summarizes this store's own trips.

func (*CountingStore) Delete added in v0.5.0

func (c *CountingStore) Delete(ctx context.Context, key string, expectedVersion string) error

func (*CountingStore) Get added in v0.5.0

func (c *CountingStore) Get(ctx context.Context, key string) ([]byte, string, bool, error)

func (*CountingStore) GetAll added in v0.5.0

func (c *CountingStore) GetAll(ctx context.Context, keyPrefix string) (map[string]Record, error)

GetAll forwards the optional bulk read, counting it as the one trip it is. Forwarding matters as much as counting: without it a RunCache stacked above this counter could not see that the store beneath can bulk-read, and the measurement would report the per-key cost of a stack that only has it because it is being measured.

func (*CountingStore) List added in v0.5.0

func (c *CountingStore) List(ctx context.Context, keyPrefix string) ([]string, error)

func (*CountingStore) PutIfAbsent added in v0.5.0

func (c *CountingStore) PutIfAbsent(ctx context.Context, key string, payload []byte) (string, error)

func (*CountingStore) PutIfVersion added in v0.5.0

func (c *CountingStore) PutIfVersion(ctx context.Context, key string, payload []byte, expectedVersion string) (string, error)

func (*CountingStore) Reset added in v0.5.0

func (c *CountingStore) Reset()

Reset discards every trip recorded so far, so one process can measure several phases separately. It does not touch the log.

func (*CountingStore) Total added in v0.5.0

func (c *CountingStore) Total() int

Total is how many trips reached the wrapped store.

func (*CountingStore) Trips added in v0.5.0

func (c *CountingStore) Trips() []Trip

Trips returns every trip so far, in order, as a copy.

func (*CountingStore) Unwrap added in v0.18.0

func (c *CountingStore) Unwrap() Store

Unwrap returns the wrapped store, for AsBucketContractChecker.

type KMSDeniedError added in v0.18.0

type KMSDeniedError struct {
	// Action is the KMS action refused, for example "kms:Decrypt".
	Action string
	// KeyARN is the key that refused it. Empty if AWS did not say.
	KeyARN string
	// Principal is who was refused, as AWS names them. Empty if AWS did not say.
	Principal string
	// Where is which policy AWS blamed: KMSDeniedByKeyPolicy,
	// KMSDeniedByIdentityPolicy, KMSDeniedExplicitly, or "" when the message
	// did not say.
	Where string
	// Err is the S3 error as it arrived.
	Err error
}

KMSDeniedError is an S3 request refused because of the bucket's KMS key, and not because of anything about S3.

A bucket whose default encryption is a customer managed key makes every GetObject a kms:Decrypt and every PutObject a kms:GenerateDataKey, made by S3 with the caller's identity. When KMS refuses, S3 reports it as a plain 403 AccessDenied on the S3 operation, and the only sign that the cause is the key is inside the message text. An operator who reads "AccessDenied ... PutObject" goes to the bucket policy and the role's S3 statements, which are correct, and the mistake is almost always somewhere they did not look: a key policy that does not name the role. A customer managed key is usable only by the principals its key policy allows, and an IAM policy alone never grants it unless the key policy delegates to IAM.

The wording below was written against what real AWS says (GitHub issue #1345, us-east-2, 2026-09-18), which is, on one line:

User: arn:aws:sts::<acct>:assumed-role/<role>/<session> is not authorized
to perform: kms:GenerateDataKey on resource: arn:aws:kms:<region>:<acct>:key/<id>
because no resource-based policy allows the kms:GenerateDataKey action

func (*KMSDeniedError) Error added in v0.18.0

func (e *KMSDeniedError) Error() string

func (*KMSDeniedError) Headline added in v0.18.0

func (e *KMSDeniedError) Headline() string

Headline is the one sentence of what happened: which key refused what to whom.

func (*KMSDeniedError) Remedy added in v0.18.0

func (e *KMSDeniedError) Remedy() string

Remedy says where to look, which depends on the policy AWS blamed.

func (*KMSDeniedError) Unwrap added in v0.18.0

func (e *KMSDeniedError) Unwrap() error

type LocalStore

type LocalStore struct {
	// contains filtered or unexported fields
}

LocalStore is a Store backed by a directory of files: one file per key, nested directories mirroring any "/" the key contains. It is the zero-configuration default — solo development, tests, air-gapped runs — mirroring plain local state's own "just works, no backend to configure" shape.

Version

A record's version is its content hash ("sha256:<hex>"), not a timestamp or a counter: two writes of byte-identical payloads carry the same version, and the version never depends on the clock or on how many times the key has been written.

Atomicity and its limit: single-operator only

LocalStore.PutIfAbsent is a single O_CREATE|O_EXCL open — atomic on its own, no locking needed, exactly like a real filesystem's create primitive already guarantees. LocalStore.PutIfVersion and LocalStore.Delete are a read, a compare, and a write (or removal); nothing in POSIX makes that sequence atomic by itself, so each one holds a sidecar "<file>.lock" file — itself an O_CREATE|O_EXCL create — for the duration, and the write half lands via a temp file plus os.Rename so a reader never observes a half-written file.

That gives real compare-and-swap for every writer on one machine: two goroutines in one process, or two separate `tofu` invocations racing on the same directory, serialize through the lockfile and the loser gets a *VersionConflictError rather than a silently clobbered write. It gives nothing across machines — there is no network protocol here, only local filesystem primitives — which is the store's fundamental limit rather than an oversight: LocalStore is for a single operator (or a single machine's worth of concurrent processes), never for a team sharing state across laptops. Reaching for S3Store is what "more than one operator" means in this package.

What this store does not manage

The directory's location, its presence or absence in version control, and its backup story are the caller's to decide — this store only reads and writes files under the directory it is given. That is an acceptable hands-off position specifically because a micro-state record's blast radius is small (an effect re-runs, a random id regenerates) in a way a full Terraform state file's loss never was.

func NewLocalStore

func NewLocalStore(dir string) (*LocalStore, error)

NewLocalStore builds a LocalStore rooted at dir, creating dir (and any missing parents) if it does not exist yet.

func (*LocalStore) Delete

func (s *LocalStore) Delete(ctx context.Context, key string, expectedVersion string) error

Delete implements Store, under the same lockfile discipline as LocalStore.PutIfVersion.

func (*LocalStore) Get

func (s *LocalStore) Get(_ context.Context, key string) ([]byte, string, bool, error)

Get implements Store.

func (*LocalStore) GetAll added in v0.5.0

func (s *LocalStore) GetAll(_ context.Context, keyPrefix string) (map[string]Record, error)

GetAll reads every record under keyPrefix from the store directory in one walk. Costs no network at all, so the whole namespace is one traversal plus one file read each — the local backend's equivalent of stock reading its state file.

The exclusions are LocalStore.List's exactly: a lockfile is not a record, and a temp file is a write in progress that no reader may observe.

func (*LocalStore) List

func (s *LocalStore) List(_ context.Context, keyPrefix string) ([]string, error)

List implements Store by walking the whole directory tree and filtering by a plain string prefix — the store's own layout already mirrors key hierarchy in directories, but List's contract is the interface's ordinary string-prefix match, not a path-boundary match, so this walks everything under s.dir rather than trying to shortcut to a subdirectory.

func (*LocalStore) PutIfAbsent

func (s *LocalStore) PutIfAbsent(_ context.Context, key string, payload []byte) (string, error)

PutIfAbsent implements Store. It is a single O_CREATE|O_EXCL open, so unlike LocalStore.PutIfVersion it needs no lockfile of its own: the filesystem's own create primitive is already the atomicity.

func (*LocalStore) PutIfVersion

func (s *LocalStore) PutIfVersion(ctx context.Context, key string, payload []byte, expectedVersion string) (string, error)

PutIfVersion implements Store. expectedVersion == "" delegates to LocalStore.PutIfAbsent, which needs no lockfile; any other value takes this key's lockfile for a read-compare-write critical section — see the type doc's "Atomicity and its limit" section for exactly what that does and does not protect against.

type Record added in v0.5.0

type Record struct {
	Payload []byte
	Version string
}

Record is one stored record's content and version, as BulkReader.GetAll returns it — the same pair Store.Get returns for one key, for a caller that asked for many.

type RunCache added in v0.5.0

type RunCache struct {
	// contains filtered or unexported fields
}

RunCache is the estate's records loaded the way stock OpenTofu loads its state file: once, in bulk, at the moment something first asks for any of it, held in memory for the rest of the read phase, and never written back through.

What it replaces

A migrated plan asks the store about the same instance three to eight times. Every accessor on projection's RecordStore - identity, residue, the provisioner bit, the deposed set, the envelope kind - decodes the same physical key, and each went to the store for itself. Measured at scale 1 over 78 instances: 377 trips over 80 distinct keys. Caching the repeats takes that to 80, one per instance. Loading the namespace in bulk takes it to 1, which is what stock pays, and is the only figure that makes the two state models comparable at all.

The bulk load is lazy - it happens on the first read under the namespace, not at construction - because that is exactly when stock reads its state file, and because a command that touches no record should pay for none. A store that does not implement BulkReader silently keeps the per-key behaviour: still one trip per key instead of one per accessor.

Why it cannot serve a stale value

One rule, and it is not a balance of risks: **the cache is switched off permanently, process-wide, by the first write through any RunCache.** Nothing it serves was ever read after something was written. A plan writes nothing, so a plan is served entirely from the snapshot; the instant a write-back, a migration or a seeder writes anything, every later read in that process goes to the store, for good.

That is deliberately blunter than invalidating the key that was written. Invalidation has to be right about which reads a write can affect, and the three reads it must never be wrong about - projection's mergeEnvelope read-modify-write, its currentVersion, and the seeders' read-before-write halves - are the ones where being wrong means a lost update rather than a slow run. A switch that is simply off after the first write cannot be wrong about any of them. Those three call sites additionally bypass this cache explicitly, through RunCache.Uncached, so they are correct even before the first write has happened.

Nothing conditional is ever decided here either: PutIfVersion, PutIfAbsent and Delete always go to the wrapped store, so the compare-and-swap that guards against a writer OUTSIDE this process is performed by the store on the store's own current version, exactly as before.

Lifetime

The process, and never longer. A record is live state; whether one has gone stale between runs is the charter's business, and a cache must never be what answers it. There is no expiry, no file, no shared daemon.

func (*RunCache) Delete added in v0.5.0

func (c *RunCache) Delete(ctx context.Context, key string, expectedVersion string) error

func (*RunCache) Get added in v0.5.0

func (c *RunCache) Get(ctx context.Context, key string) ([]byte, string, bool, error)

func (*RunCache) GetAll added in v0.5.0

func (c *RunCache) GetAll(ctx context.Context, keyPrefix string) (map[string]Record, error)

GetAll passes a bulk read through to the wrapped store, so a caller that wants the whole namespace still gets it in one call. It is deliberately NOT served from the snapshot: the only caller of a bulk read that is not this cache is one that wants the store's own current answer.

func (*RunCache) List added in v0.5.0

func (c *RunCache) List(ctx context.Context, keyPrefix string) ([]string, error)

func (*RunCache) PutIfAbsent added in v0.5.0

func (c *RunCache) PutIfAbsent(ctx context.Context, key string, payload []byte) (string, error)

func (*RunCache) PutIfVersion added in v0.5.0

func (c *RunCache) PutIfVersion(ctx context.Context, key string, payload []byte, expectedVersion string) (string, error)

func (*RunCache) Uncached added in v0.5.0

func (c *RunCache) Uncached() Store

Uncached returns the wrapped store, for the reads that must never be served from a snapshot: a read-modify-write's own read, a version observed so a compare-and-swap can catch an outside writer, and a seeder's read-before-write. See this type's doc comment.

It is a method rather than a field so the intent is stated at every call site that needs it, and so a store that is not a RunCache needs no special case: see Fresh.

func (*RunCache) Unwrap added in v0.18.0

func (c *RunCache) Unwrap() Store

Unwrap returns the wrapped store, for AsBucketContractChecker.

type S3Config

type S3Config struct {
	// Client is the S3 client every call goes through. The caller builds
	// and authenticates it — region, credentials, any endpoint override
	// for a local emulator — this package has no opinion on any of that.
	Client *s3.Client

	// Bucket is the S3 bucket every key lives in.
	Bucket string

	// KeyPrefix is joined ahead of every key this store is asked for, so
	// one bucket can host more than one caller's keyspace without either
	// seeing the other's keys in [S3Store.List]. Empty means keys map
	// directly to object keys. This package does not interpret
	// KeyPrefix's structure at all — it is an opaque string, the same as
	// every key passed to the [Store] interface.
	KeyPrefix string

	// GetAllParallelism bounds how many GetObject calls [S3Store.GetAll] has
	// in flight at once. Zero or negative takes
	// [DefaultS3GetAllParallelism]; 1 is a sequential read.
	GetAllParallelism int

	// BaseTags go on every object this store writes. The record store sets
	// tofu-estate here, because every object in an estate's namespaces -
	// its records, its sentinel, its hint, its outputs - is the estate's,
	// whether or not it records a resource. See [WithObjectTags] for the
	// per-write half. GitHub issue #1337.
	BaseTags map[string]string
}

S3Config configures an S3Store.

type S3Store

type S3Store struct {
	// contains filtered or unexported fields
}

S3Store is a Store backed by S3 object versions via conditional writes: If-Match and If-None-Match, the ETag-based compare-and-swap primitive S3 added for general-purpose buckets. This is the store's strongest offering — a real, server-enforced CAS, not a read-compare-write approximation — and a version here is exactly an object's ETag, unmodified.

What is genuinely atomic

Every conditional operation is a single S3 request carrying the condition; there is no read-compare-write window for this store to document a caveat about:

  • S3Store.PutIfAbsent and a "" S3Store.PutIfVersion call send If-None-Match: * — S3 rejects the write with HTTP 412 if any object already exists at the key.
  • A non-"" S3Store.PutIfVersion call sends If-Match: <version> — S3 rejects the write with HTTP 412 if the object's current ETag does not match.
  • S3Store.Delete sends If-Match: <version> on DeleteObject, which S3 honors for general-purpose buckets, not only the directory-bucket case the S3 API docs otherwise reserve conditional deletes for.

On a 412, this store issues one extra read (Get) purely to populate VersionConflictError.ActualVersion with an accurate answer; that read is not part of the conditional guarantee itself; the conditional write already failed atomically before it.

What this store does not manage

Bucket creation, lifecycle policy, and encryption configuration are the caller's concern — S3Store only issues GetObject/PutObject/DeleteObject/ ListObjectsV2 against a bucket and (optional) key prefix it is given. It does not build or authenticate the s3.Client itself; the caller supplies one already configured for the target account, region and endpoint, which is what keeps this store's own surface free of anything AWS-credential-shaped.

func NewS3Store

func NewS3Store(cfg S3Config) (*S3Store, error)

NewS3Store builds an S3Store from cfg.

func (*S3Store) CheckBucketContract added in v0.18.0

func (s *S3Store) CheckBucketContract(ctx context.Context, namespaces []string) ([]BucketFinding, error)

CheckBucketContract implements BucketContractChecker. namespaces are store-relative, like every key this store is handed; S3Config.KeyPrefix is joined ahead of each.

func (*S3Store) Delete

func (s *S3Store) Delete(ctx context.Context, key string, expectedVersion string) error

Delete implements Store. expectedVersion == "" against an absent key is a no-op (checked with a HeadObject first, since DeleteObject's If-Match has no "only if absent" form); any other value sends If-Match: <expectedVersion> on DeleteObject itself, S3's real conditional delete.

func (*S3Store) Get

func (s *S3Store) Get(ctx context.Context, key string) ([]byte, string, bool, error)

Get implements Store.

func (*S3Store) GetAll added in v0.5.0

func (s *S3Store) GetAll(ctx context.Context, keyPrefix string) (map[string]Record, error)

GetAll reads every record under keyPrefix: one ListObjectsV2 pagination, then one GetObject per key, at most S3Config.GetAllParallelism of them in flight at once.

S3 is the backend that genuinely cannot bulk-fetch. There is no batch-read operation in the S3 API - ListObjectsV2 returns each object's key and ETag but never its body - so N objects cost N GetObject calls whatever this function does. Overlapping them is the only saving there is.

A bulk read is complete or it fails

That is the constraint, and it matters more than the speedup. BulkReader promises the result is complete for its prefix: a key absent from the map holds no record. A plan reads a record key with no configuration behind it as an instruction to destroy, and reads a declared instance with no record as something to create. So a map that silently lacks a key is the worst thing this backend can produce, and parallelism is exactly where it would come from: the sequential loop this replaced got completeness for free by returning on its first error, and a fan-out has to be written to keep it.

So: results are written by index into a slice, never into a shared map. The first failure wins, cancels the rest, and is the error returned, naming its key. The map is built only after every worker has stopped, and only if nothing failed AND every key was handed to a worker - a cancelled context that stopped the feed with no GET in flight would otherwise leave no error at all and a short map behind it.

One omission is legitimate and is kept: a key that 404s between the LIST and its GET was deleted in between, and leaving it out of the map is the correct way to say so. That is a different thing from a GET that failed.

func (*S3Store) List

func (s *S3Store) List(ctx context.Context, keyPrefix string) ([]string, error)

List implements Store by paginating ListObjectsV2 with Prefix set to this store's own key prefix plus keyPrefix — S3's list primitive is already an ordinary string prefix, the same contract Store.List promises, so no client-side filtering beyond stripping s.keyPrefix back off is needed.

func (*S3Store) ObjectKey added in v0.14.0

func (s *S3Store) ObjectKey(key string) string

ObjectKey is the S3 object key this store will read and write key at: S3Config.KeyPrefix joined ahead of it. Exported so a caller can say where a record actually is, in words an operator can paste into the AWS CLI - issue #916.

func (*S3Store) PutIfAbsent

func (s *S3Store) PutIfAbsent(ctx context.Context, key string, payload []byte) (string, error)

PutIfAbsent implements Store.

func (*S3Store) PutIfVersion

func (s *S3Store) PutIfVersion(ctx context.Context, key string, payload []byte, expectedVersion string) (string, error)

PutIfVersion implements Store. expectedVersion == "" sends If-None-Match: *; any other value sends If-Match: <expectedVersion> — see the type doc for why both are a single atomic S3 request rather than a read-compare-write.

type Store

type Store interface {
	// Get reads the current record at key. exists is false when no record
	// is there; payload and version are then the zero value and err is nil
	// — a missing key is not itself an error. version is "" if and only if
	// exists is false: no implementation ever assigns "" as a live record's
	// version, so a caller may treat it as a stable "absent" sentinel.
	Get(ctx context.Context, key string) (payload []byte, version string, exists bool, err error)

	// PutIfVersion writes payload to key, but only if the record's current
	// version equals expectedVersion. expectedVersion == "" asserts that no
	// record exists yet at key (the same assertion PutIfAbsent makes,
	// reachable here for callers that hold a uniform "expected version"
	// value rather than branching on whether they have seen the key
	// before). On success it returns the record's new version. On a
	// mismatch, it returns a *VersionConflictError naming both
	// expectedVersion and the version the store actually found — never a
	// bare error a caller has to parse.
	PutIfVersion(ctx context.Context, key string, payload []byte, expectedVersion string) (newVersion string, err error)

	// PutIfAbsent creates key with payload, but only if no record exists
	// there yet. It is PutIfVersion(ctx, key, payload, "") under a name
	// that does not require the caller to know the empty-string
	// convention. On success it returns the new record's version; on
	// conflict it returns a *VersionConflictError with ExpectedVersion ""
	// and ActualVersion set to whatever is already there.
	PutIfAbsent(ctx context.Context, key string, payload []byte) (version string, err error)

	// Delete removes key, but only if its current version equals
	// expectedVersion — the same conditional discipline PutIfVersion
	// applies to writes, applied here to removal. Deleting an
	// already-absent key with expectedVersion == "" succeeds silently
	// (idempotent); deleting an already-absent key with a non-empty
	// expectedVersion, or a present key whose version does not match, both
	// return a *VersionConflictError.
	Delete(ctx context.Context, key string, expectedVersion string) error

	// List returns every key currently stored whose name begins with
	// keyPrefix, as an ordinary Go string prefix (not a path-hierarchy
	// match), sorted lexically. keyPrefix == "" lists every key. See each
	// implementation's own doc comment for how closely its underlying
	// primitive matches this; the returned set is exactly this contract
	// regardless.
	List(ctx context.Context, keyPrefix string) ([]string, error)
}

Store is a conditional-write key/value backend: a name for a small blob, versioned so a writer can prove it is updating what it last read rather than clobbering someone else's change. It has no notion of what a key names or what a payload contains — that belongs entirely to the caller. See doc.go for the full contract every implementation must honor.

func Fresh added in v0.5.0

func Fresh(s Store) Store

Fresh returns the store beneath any read cache in s, or s itself when there is none. A caller that must not read a remembered value asks for this rather than testing for a cache type it should not have to know about.

func NewRunCache added in v0.5.0

func NewRunCache(inner Store, prefix string) Store

NewRunCache wraps inner, snapshotting the namespace at prefix. A nil inner returns nil, so a caller that already treats "no store configured" as nil keeps doing so. An empty prefix disables the bulk load and leaves ordinary per-key caching.

type Trip added in v0.5.0

type Trip struct {
	Method string
	Key    string
	Via    string
	Site   string
}

Trip is one operation that reached the wrapped store. See CountingStore for what each field is for.

func ParseTripLog added in v0.5.0

func ParseTripLog(data []byte) ([]Trip, error)

ParseTripLog reads back what CountingStore's log writer wrote. A line that does not have the four tab-separated fields is an error rather than a skipped line: a partially readable log would under-report, and an under-reported cost is exactly the failure this instrument exists to end.

func (Trip) String added in v0.5.0

func (t Trip) String() string

String renders a trip as the one TSV line CountingStore's log writes and ParseTripLog reads back.

type TripCounts added in v0.5.0

type TripCounts struct {
	Total int

	ByMethod map[string]int
	BySite   map[string]int
	ByVia    map[string]int

	// DistinctKeys is how many different keys the trips touched, and
	// RepeatTrips is Total minus DistinctKeys: the trips that re-read
	// something an earlier trip had already read. A cache can remove at
	// most RepeatTrips of them and never fewer than zero of the rest,
	// which is the whole reason this is reported beside the total.
	DistinctKeys int
	RepeatTrips  int
}

TripCounts is a set of trips summarized several ways at once. Every field is a total over the same trips, so a reader can check them against each other rather than having to trust one.

func Summarize added in v0.5.0

func Summarize(trips []Trip) TripCounts

Summarize buckets trips every way TripCounts names.

type VersionConflictError

type VersionConflictError struct {
	Key             string
	ExpectedVersion string
	ActualVersion   string
}

VersionConflictError reports that a conditional operation's expected version did not match what the store actually holds for Key. It names both versions so a caller can decide how to react — reread and retry, surface a merge conflict, give up loudly — without parsing an error string.

ActualVersion is "" when the store holds no record at Key at all (including the case where ExpectedVersion was itself "" and the key turned out to already exist would instead set ActualVersion to that existing version — "" only ever means "no record").

func (*VersionConflictError) Error

func (e *VersionConflictError) Error() string

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL