strata

package module
v0.5.0 Latest Latest
Warning

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

Go to latest
Published: Mar 27, 2026 License: MIT Imports: 24 Imported by: 0

README

Strata

CI Go Reference Go Report Card

An embeddable, S3-durable key-value store for Go.

  • Embedded-firststrata.Open(cfg) is the entire API. No sidecar, no daemon.
  • S3-durable — WAL segments and periodic snapshots are uploaded to S3. A node that loses its disk recovers automatically.
  • Multi-node — Leader elected via an S3 lock. Followers stream the WAL in real time and forward writes transparently.
  • etcd v3 compatible — The standalone binary speaks the etcd v3 gRPC protocol. Any etcd client works against it unchanged.

Embedded usage

import "github.com/makhov/strata"

node, err := strata.Open(strata.Config{
    DataDir: "/var/lib/myapp/strata",
})
defer node.Close()

rev, err := node.Put(ctx, "/config/timeout", []byte("30s"), 0)

kv, err := node.Get("/config/timeout")
fmt.Println(string(kv.Value)) // 30s

events, _ := node.Watch(ctx, "/config/", 0)
for e := range events {
    fmt.Printf("%s %s=%s\n", e.Type, e.KV.Key, e.KV.Value)
}
With S3 durability
import (
    awsconfig "github.com/aws/aws-sdk-go-v2/config"
    "github.com/aws/aws-sdk-go-v2/service/s3"
    "github.com/makhov/strata/pkg/object"
)

awsCfg, _ := awsconfig.LoadDefaultConfig(ctx)
store := object.NewS3Store(s3.NewFromConfig(awsCfg), "my-bucket", "strata/")

node, err := strata.Open(strata.Config{
    DataDir:     "/var/lib/myapp/strata",
    ObjectStore: store,
})

Standalone binary

The strata binary exposes the etcd v3 gRPC protocol. Use etcdctl, the official Go client, or any other etcd v3 compatible tool.

go install github.com/makhov/strata/cmd/strata@latest

# Single node, local only
strata --data-dir /var/lib/strata --listen 0.0.0.0:3379

# Single node with S3
strata --data-dir /var/lib/strata --listen 0.0.0.0:3379 \
       --s3-bucket my-bucket --s3-prefix strata/

# Verify
etcdctl --endpoints=localhost:3379 put /hello world
etcdctl --endpoints=localhost:3379 get /hello

Multi-node and production setup: see docs/operations.md.


Documentation

Document Contents
docs/api.md Full Go API reference — methods, types, errors
docs/configuration.md All config fields and CLI flags
docs/operations.md Multi-node, S3, mTLS, observability, point-in-time restore
docs/architecture.md Internals — WAL, checkpoints, election, replication

Documentation

Overview

Package strata provides an embeddable, S3-durable key-value store.

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrKeyExists = errors.New("strata: key already exists")
	ErrNotLeader = errors.New("strata: this node is not the leader; writes are rejected")
)

Sentinel errors.

Functions

This section is empty.

Types

type Config

type Config struct {

	// DataDir is the directory used for local Pebble data and WAL segments.
	// Required.
	DataDir string

	// ObjectStore is used to archive WAL segments and checkpoints and to run
	// leader election. If nil the node runs in single-node mode.
	ObjectStore object.Store

	// RestorePoint, if set, causes the node to bootstrap from a specific
	// point in time on first boot rather than reading the latest checkpoint
	// from ObjectStore. See RestorePoint for details.
	RestorePoint *RestorePoint

	// SegmentMaxSize is the byte threshold that triggers WAL segment rotation.
	// Default: 50 MB.
	SegmentMaxSize int64

	// SegmentMaxAge is the time threshold that triggers WAL segment rotation.
	// Default: 10 s.
	SegmentMaxAge time.Duration

	// CheckpointInterval controls how often the leader writes a checkpoint.
	// Default: 15 minutes.
	CheckpointInterval time.Duration

	// CheckpointEntries triggers a checkpoint after this many WAL entries
	// regardless of time. 0 means disabled.
	CheckpointEntries int64

	// NodeID is a stable, unique identifier for this node.
	// Defaults to the machine hostname.
	NodeID string

	// PeerListenAddr is the address on which the peer WAL-streaming gRPC
	// server listens (e.g. "0.0.0.0:3380"). Empty → single-node mode.
	PeerListenAddr string

	// AdvertisePeerAddr is the address followers use to reach this node's peer
	// server. Defaults to PeerListenAddr.
	AdvertisePeerAddr string

	// LeaderWatchInterval is how often the leader reads the lock from S3 to
	// detect if it has been superseded. Read-only; no renewals.
	// Default: 5 minutes.
	LeaderWatchInterval time.Duration

	// FollowerMaxRetries is the number of consecutive stream failures a follower
	// tolerates before attempting a TakeOver election.
	// Default: 5.
	FollowerMaxRetries int

	// PeerBufferSize is the number of WAL entries the leader buffers for
	// follower catch-up. Default: 10 000.
	PeerBufferSize int

	// PeerServerTLS is the transport credentials used by the leader's peer
	// gRPC server. Nil means plaintext (only safe inside a trusted network).
	PeerServerTLS credentials.TransportCredentials

	// PeerClientTLS is the transport credentials used by a follower's peer
	// gRPC client. Must be set when PeerServerTLS is set on the leader.
	PeerClientTLS credentials.TransportCredentials

	// MetricsAddr is the TCP address for the Prometheus /metrics, /healthz,
	// and /readyz HTTP endpoints (e.g. "0.0.0.0:9090"). Empty means disabled.
	MetricsAddr string
}

Config holds all configuration for a Node.

type Event

type Event struct {
	Type   EventType
	KV     *KeyValue
	PrevKV *KeyValue // nil for creates
}

Event is a single watch notification.

type EventType

type EventType int

EventType classifies a watch event.

const (
	EventPut    EventType = iota // create or update
	EventDelete                  // deletion
)

type KeyValue

type KeyValue struct {
	Key            string
	Value          []byte
	Revision       int64
	CreateRevision int64
	PrevRevision   int64
	Lease          int64
}

KeyValue is a versioned key-value pair.

type Node

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

Node is the top-level Strata instance.

Single-node mode (PeerListenAddr == ""):

Writes: WAL.Append (fsync) → store.Apply → notify watchers
Background: WAL segments uploaded to S3, periodic checkpoints

Leader mode:

Same as single-node, plus fan-out to followers via peer gRPC stream.
Holds the S3 leader lock; watches it infrequently for supersession.

Follower mode:

Reads are served locally. Writes are forwarded to the leader via the
peer gRPC channel and the response is returned transparently to the caller.
After persistent stream failure, attempts a TakeOver election.

func Open

func Open(cfg Config) (*Node, error)

Open creates and starts a Node.

func (*Node) Close

func (n *Node) Close() error

Close shuts down the node cleanly.

func (*Node) Compact

func (n *Node) Compact(ctx context.Context, revision int64) error

Compact removes log entries at or below revision.

func (*Node) CompactRevision

func (n *Node) CompactRevision() int64

func (*Node) Config

func (n *Node) Config() Config

func (*Node) Count

func (n *Node) Count(prefix string) (int64, error)

func (*Node) Create

func (n *Node) Create(ctx context.Context, key string, value []byte, lease int64) (int64, error)

Create creates key only if it does not already exist.

func (*Node) CurrentRevision

func (n *Node) CurrentRevision() int64

func (*Node) Delete

func (n *Node) Delete(ctx context.Context, key string) (int64, error)

Delete removes key unconditionally.

func (*Node) DeleteIfRevision

func (n *Node) DeleteIfRevision(ctx context.Context, key string, revision int64) (int64, *KeyValue, bool, error)

DeleteIfRevision deletes key only if its current revision matches (CAS).

func (*Node) Get

func (n *Node) Get(key string) (*KeyValue, error)

func (*Node) HandleForward

func (n *Node) HandleForward(ctx context.Context, req *peer.ForwardRequest) (*peer.ForwardResponse, error)

HandleForward implements peer.ForwardHandler. Called by the peer gRPC server when a follower forwards a write. Dispatches to the appropriate Node method. Since HandleForward runs on the leader, all write methods execute directly.

func (*Node) IsLeader

func (n *Node) IsLeader() bool

func (*Node) List

func (n *Node) List(prefix string) ([]*KeyValue, error)

func (*Node) Put

func (n *Node) Put(ctx context.Context, key string, value []byte, lease int64) (int64, error)

Put creates or updates key with value. Returns the new revision.

func (*Node) Update

func (n *Node) Update(ctx context.Context, key string, value []byte, revision, lease int64) (int64, *KeyValue, bool, error)

Update updates key only if its current revision matches (CAS).

func (*Node) WaitForRevision

func (n *Node) WaitForRevision(ctx context.Context, rev int64) error

func (*Node) Watch

func (n *Node) Watch(ctx context.Context, prefix string, startRev int64) (<-chan Event, error)

type PinnedObject added in v0.3.0

type PinnedObject struct {
	Key       string
	VersionID string
}

PinnedObject identifies a specific version of an object in object storage.

type RestorePoint added in v0.3.0

type RestorePoint struct {
	// Store is the versioned object store to read pinned objects from.
	// It may use a different prefix than Config.ObjectStore (e.g. to read
	// from the source branch while writing to a new branch prefix).
	Store object.VersionedStore

	// CheckpointArchive is the pinned checkpoint archive object.
	CheckpointArchive PinnedObject

	// WALSegments are the WAL segments to replay after the checkpoint,
	// in ascending sequence order.
	WALSegments []PinnedObject
}

RestorePoint describes a precise point in time from which a node should bootstrap. When set in Config, the node restores the checkpoint and replays the listed WAL segments using their pinned S3 version IDs, rather than reading the latest objects from its own prefix.

This enables point-in-time restore, blue/green deployments, and copy-free forking: the source data is read directly from S3 by version ID — no objects are copied to the new prefix.

The node's own ObjectStore prefix is used for all subsequent writes after startup. RestorePoint is only applied on first boot (when the local data directory does not yet exist); it is ignored on subsequent restarts.

S3 versioning must be enabled on the source bucket.

Directories

Path Synopsis
cmd
strata command
Command strata runs a Strata node and exposes it as an etcd v3 gRPC endpoint.
Command strata runs a Strata node and exposes it as an etcd v3 gRPC endpoint.
Package etcd exposes a strata Node as an etcd v3 gRPC server.
Package etcd exposes a strata Node as an etcd v3 gRPC server.
internal
checkpoint
Package checkpoint handles creating, writing, and restoring Pebble snapshots to/from object storage.
Package checkpoint handles creating, writing, and restoring Pebble snapshots to/from object storage.
election
Package election implements S3-based leader election.
Package election implements S3-based leader election.
metrics
Package metrics defines Prometheus metrics for a Strata node.
Package metrics defines Prometheus metrics for a Strata node.
peer
Package peer implements the leader→follower WAL streaming gRPC service, plus write forwarding (follower→leader).
Package peer implements the leader→follower WAL streaming gRPC service, plus write forwarding (follower→leader).
store
Package store implements the Pebble-backed key-value state machine.
Package store implements the Pebble-backed key-value state machine.
wal
Package wal implements the write-ahead log.
Package wal implements the write-ahead log.
pkg
object
Package object provides a small interface for object storage operations used by Strata (WAL archive, checkpoints, manifest).
Package object provides a small interface for object storage operations used by Strata (WAL archive, checkpoints, manifest).

Jump to

Keyboard shortcuts

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