Documentation
¶
Overview ¶
Package s3 implements a backend.Backend over an S3-compatible object store. The store-specific calls are isolated behind the small ObjectStore interface: the AWS SDK adapter (NewAWS) implements it for real S3, and tests implement it with an in-memory fake, so all of the Backend's mapping logic — key prefixing, sorted listing, 404 → backend.ErrNotExist translation, conditional-put, and existence-checked delete — is exercised by the shared backend conformance suite without a live bucket.
The object store is the durable, stateless tier: nodes hold no authoritative state, so the read path is reconstructed entirely from objects (DESIGN.md §3, §11). Whole-object Get/Put is sufficient because a part maps to one key prefix with one object per column/marks/manifest.
Index ¶
- Variables
- type AWSAPI
- type Backend
- func (b *Backend) CompareAndSwap(ctx context.Context, key string, expected backend.Version, data []byte) (backend.Version, bool, error)
- func (b *Backend) CreateObject(ctx context.Context, key string) (backend.ObjectWriter, error)
- func (b *Backend) Delete(ctx context.Context, key string) error
- func (*Backend) IsEphemeral() bool
- func (b *Backend) List(ctx context.Context, prefix string) ([]string, error)
- func (b *Backend) PutIfAbsent(ctx context.Context, key string, data []byte) (bool, error)
- func (b *Backend) Read(ctx context.Context, key string) ([]byte, error)
- func (b *Backend) ReadAt(ctx context.Context, key string, off, n int64) ([]byte, error)
- func (b *Backend) ReadVersioned(ctx context.Context, key string) ([]byte, backend.Version, error)
- func (b *Backend) StreamsWrites() bool
- func (b *Backend) Write(ctx context.Context, key string, data []byte) error
- type MultipartObjectStore
- type ObjectStore
- type Option
- type RangeObjectStore
Constants ¶
This section is empty.
Variables ¶
var ErrObjectNotFound = errors.New("s3: object not found")
ErrObjectNotFound is returned (wrapped) by ObjectStore.GetObject for an absent key. The Backend translates it to backend.ErrNotExist.
Functions ¶
This section is empty.
Types ¶
type AWSAPI ¶
type AWSAPI interface {
GetObject(ctx context.Context, in *awss3.GetObjectInput, optFns ...func(*awss3.Options)) (*awss3.GetObjectOutput, error)
PutObject(ctx context.Context, in *awss3.PutObjectInput, optFns ...func(*awss3.Options)) (*awss3.PutObjectOutput, error)
HeadObject(ctx context.Context, in *awss3.HeadObjectInput, optFns ...func(*awss3.Options)) (*awss3.HeadObjectOutput, error)
DeleteObject(ctx context.Context, in *awss3.DeleteObjectInput, optFns ...func(*awss3.Options)) (*awss3.DeleteObjectOutput, error)
ListObjectsV2(ctx context.Context, in *awss3.ListObjectsV2Input, optFns ...func(*awss3.Options)) (*awss3.ListObjectsV2Output, error)
CreateMultipartUpload(
ctx context.Context, in *awss3.CreateMultipartUploadInput, optFns ...func(*awss3.Options),
) (*awss3.CreateMultipartUploadOutput, error)
UploadPart(ctx context.Context, in *awss3.UploadPartInput, optFns ...func(*awss3.Options)) (*awss3.UploadPartOutput, error)
CompleteMultipartUpload(
ctx context.Context, in *awss3.CompleteMultipartUploadInput, optFns ...func(*awss3.Options),
) (*awss3.CompleteMultipartUploadOutput, error)
AbortMultipartUpload(
ctx context.Context, in *awss3.AbortMultipartUploadInput, optFns ...func(*awss3.Options),
) (*awss3.AbortMultipartUploadOutput, error)
}
AWSAPI is the subset of the aws-sdk-go-v2 *s3.Client the NewAWS adapter uses. The real client satisfies it; tests can fake it. It also satisfies s3.ListObjectsV2APIClient (used by the paginator).
type Backend ¶
type Backend struct {
// contains filtered or unexported fields
}
Backend is a backend.Backend over an ObjectStore. Keys are stored under an optional root prefix so several datasets can share one bucket.
The store's optional capabilities are resolved **once**, at construction, into rng and mp. That keeps an interface assertion off the ranged-read path — which a merge takes per frame — and it means a wrapper never has to be re-interrogated: whatever [newRetryStore] hands back is what the Backend holds.
func New ¶
func New(store ObjectStore, keyPrefix string, opts ...Option) *Backend
New returns a Backend over store, rooting all keys under keyPrefix (which may be empty). Pass WithRetry to make it resilient to a lossy/slow endpoint (per-attempt timeouts, bounded retries, hedged GETs).
func (*Backend) CompareAndSwap ¶ added in v0.40.0
func (b *Backend) CompareAndSwap( ctx context.Context, key string, expected backend.Version, data []byte, ) (backend.Version, bool, error)
CompareAndSwap stores data under key only if the object still carries the ETag the committer read, using S3's own conditional PUT (If-Match, or If-None-Match: * for the create). The store evaluates the condition, so this is a genuine CAS across processes and regions — not a read-then-write this backend could lose a race inside. Implements backend.Backend.
func (*Backend) CreateObject ¶ added in v0.48.0
CreateObject returns a writer that uploads key's object in parts once it exceeds [uploadPartBytes], and issues a single PutObject below that — which is also what it does for the whole object when the store cannot upload in parts at all. Implements backend.ObjectCreator.
func (*Backend) Delete ¶
Delete removes key, returning an backend.ErrNotExist-wrapping error if it is absent. Because S3 DeleteObject is idempotent, existence is checked first to honor the contract.
func (*Backend) IsEphemeral ¶
IsEphemeral reports false: objects persist in the store.
func (*Backend) PutIfAbsent ¶
PutIfAbsent stores data under key only if it does not already exist.
func (*Backend) Read ¶
Read returns the value stored under key, or an backend.ErrNotExist-wrapping error.
func (*Backend) ReadAt ¶ added in v0.37.0
ReadAt returns the object's [off, off+n) bytes with a ranged GET, so a reader taking one granule of a part column does not transfer the column. A store that does not implement RangeObjectStore falls back to fetching the object and slicing — correct, just not cheaper. Implements backend.ReaderAt.
func (*Backend) ReadVersioned ¶ added in v0.40.0
ReadVersioned returns the object and its ETag, which is the version a later Backend.CompareAndSwap conditions on. Implements backend.Backend.
func (*Backend) StreamsWrites ¶ added in v0.48.0
StreamsWrites reports whether the store can actually upload in parts. Without it the writer above still works — it just holds the object, which is what a caller sizing its output against memory needs to know. Implements backend.ObjectCreator.
type MultipartObjectStore ¶ added in v0.48.0
type MultipartObjectStore interface {
// CreateMultipartUpload starts an upload for key and returns its id. Nothing is visible under
// key until CompleteMultipartUpload.
CreateMultipartUpload(ctx context.Context, key string) (uploadID string, err error)
// UploadPart stores data as the upload's partNum-th part (1-based) and returns its ETag.
// Uploading the same part number again replaces it, so a retry by part number is safe.
UploadPart(ctx context.Context, key, uploadID string, partNum int32, data []byte) (etag string, err error)
// CompleteMultipartUpload publishes the parts, in etags order, as the object under key.
CompleteMultipartUpload(ctx context.Context, key, uploadID string, etags []string) error
// AbortMultipartUpload discards the upload and its parts. It is idempotent.
AbortMultipartUpload(ctx context.Context, key, uploadID string) error
}
MultipartObjectStore is an optional ObjectStore capability: assemble one object from parts uploaded separately. It is what makes backend.ObjectCreator available over S3, so a merge producing a part column far larger than it wants resident hands finished bytes to the store as they are produced instead of holding the whole object.
It is optional for the same reason RangeObjectStore is: every real S3-compatible store has it, but an embedder's narrow fake need not. A store without it makes the Backend buffer whole objects, which backend.StreamsWrites reports.
Enabling it carries an operational precondition: a crashed writer leaves an incomplete upload, which is not an object — it appears in no listing, so no orphan sweep can ever reclaim it. The bucket needs an AbortIncompleteMultipartUpload lifecycle rule (see ADMIN.md).
type ObjectStore ¶
type ObjectStore interface {
// GetObject returns the object bytes, or an error wrapping [ErrObjectNotFound] if absent.
GetObject(ctx context.Context, key string) ([]byte, error)
// PutObject stores data under key, overwriting any existing object.
PutObject(ctx context.Context, key string, data []byte) error
// PutObjectIfAbsent stores data under key only if it does not exist, returning whether it
// was created (the conditional write — S3 If-None-Match: *).
PutObjectIfAbsent(ctx context.Context, key string, data []byte) (bool, error)
// PutObjectIfVersion stores data under key only if the object's current ETag is etag, and
// returns the ETag the object then has. An empty etag demands that the object not exist
// (If-None-Match: *); a non-empty one is S3's If-Match. ok is false — with a nil error —
// when the condition did not hold, including an If-Match against an absent key, which S3
// answers 404 rather than 412.
PutObjectIfVersion(ctx context.Context, key string, data []byte, etag string) (newETag string, ok bool, err error)
// GetObjectVersion returns the object bytes and its ETag, or an error wrapping
// [ErrObjectNotFound] if absent.
GetObjectVersion(ctx context.Context, key string) ([]byte, string, error)
// HeadObject reports whether key exists.
HeadObject(ctx context.Context, key string) (bool, error)
// DeleteObject removes key. It is idempotent (no error if the key is absent), mirroring
// S3 DeleteObject.
DeleteObject(ctx context.Context, key string) error
// ListObjects returns every key under prefix (the implementation paginates internally).
// Order is unspecified; the [Backend] sorts.
ListObjects(ctx context.Context, prefix string) ([]string, error)
}
ObjectStore is the minimal object-store API the Backend needs. It is intentionally thin and S3-shaped (delete is idempotent, list is by prefix); the Backend layers the backend.Backend contract on top. Implementations must be safe for concurrent use.
func NewAWS ¶
func NewAWS(api AWSAPI, bucket string) ObjectStore
NewAWS returns an ObjectStore backed by an aws-sdk-go-v2 S3 client over the given bucket. Compose it with New to get a backend.Backend:
store := s3.NewAWS(awss3.NewFromConfig(cfg), "my-bucket") b := s3.New(store, "oteldb/")
This adapter is verified by compilation and exercised against real/MinIO S3 in integration tests; the Backend's contract logic is covered by the conformance suite over the in-memory fake.
type Option ¶ added in v0.5.0
type Option func(*config)
Option configures a Backend at construction.
func WithRetry ¶ added in v0.5.0
func WithRetry(c reliability.RetryConfig) Option
WithRetry makes the backend survive an unreliable S3 endpoint: each call gets a per-attempt timeout (so a hung request is abandoned instead of stalling for the provider's full timeout) and bounded retries, and GETs are additionally *hedged* — a slow read is re-issued on a fresh connection and the first response wins. Use reliability.LossyEnvironment for noisy networks. The zero config leaves the bare store (the AWS SDK's own retryer still applies).
type RangeObjectStore ¶ added in v0.37.0
type RangeObjectStore interface {
// GetObjectRange returns the object's [off, off+n) bytes, clamped to its end. Absent keys error
// like [ObjectStore.GetObject].
GetObjectRange(ctx context.Context, key string, off, n int64) ([]byte, error)
}
RangeObjectStore is an optional ObjectStore capability: fetch a byte range of an object (S3's Range header). NewAWS's adapter implements it — every S3-compatible store does — so it is optional only so that a fake ObjectStore need not.