kvcrdt

package
v1.7.1 Latest Latest
Warning

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

Go to latest
Published: Oct 8, 2026 License: MIT Imports: 9 Imported by: 0

Documentation

Overview

Package kvcrdt bridges grove/crdt types into KV storage, enabling distributed, eventually-consistent state without a relational database.

It provides CRDT-backed distributed counters, registers, sets, maps, lists, and documents that use a KV Store for persistence and can be synchronized across stores.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Counter

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

Counter is a distributed PNCounter backed by a KV store. It stores the full crdt.PNCounterState under the key so that merge is always possible across nodes.

func NewCounter

func NewCounter(store *kv.Store, key string, opts ...Option) *Counter

NewCounter creates a new CRDT counter backed by the given store and key.

func (*Counter) Decrement

func (c *Counter) Decrement(ctx context.Context, delta int64) error

Decrement adds delta to this node's decrement counter.

func (*Counter) Increment

func (c *Counter) Increment(ctx context.Context, delta int64) error

Increment adds delta to this node's increment counter.

func (*Counter) Merge

func (c *Counter) Merge(ctx context.Context, remote *crdt.PNCounterState) error

Merge merges a remote counter state into the local state.

func (*Counter) State

func (c *Counter) State(ctx context.Context) (*crdt.PNCounterState, error)

State returns the raw PNCounterState for sync purposes.

func (*Counter) Value

func (c *Counter) Value(ctx context.Context) (int64, error)

Value returns the current counter value (sum of all nodes).

type Document

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

Document is a distributed nested CRDT document backed by a KV store. Each field at a dot-separated path is independently mergeable with its own CRDT type (LWW, Counter, Set, List), enabling JSON-like nested structures.

func NewDocument

func NewDocument(store *kv.Store, key string, opts ...Option) *Document

NewDocument creates a new CRDT Document backed by the given store and key.

func (*Document) Delete

func (d *Document) Delete(ctx context.Context, path string) error

Delete removes a field and all its children at the given path.

func (*Document) Get

func (d *Document) Get(ctx context.Context, path string, dest any) error

Get returns the value at the given path, decoded into dest.

func (*Document) Merge

func (d *Document) Merge(ctx context.Context, remote *crdt.DocumentCRDTState) error

Merge merges a remote document state into the local state.

func (*Document) Resolve

func (d *Document) Resolve(ctx context.Context) (map[string]any, error)

Resolve returns the full document as a nested map by materializing all field paths into a tree structure.

func (*Document) Set

func (d *Document) Set(ctx context.Context, path string, value any) error

Set sets a field at the given dot-separated path with LWW semantics.

func (*Document) SetCounter

func (d *Document) SetCounter(ctx context.Context, path string, delta int64) error

SetCounter sets a counter field at the given path, incrementing by delta. If the path does not yet have a counter, a new PNCounterState is created.

func (*Document) SetFieldState

func (d *Document) SetFieldState(ctx context.Context, path string, fs *crdt.FieldState) error

SetFieldState sets a typed field state at a path. Use this for advanced use cases where you need to set a specific CRDT type (counter, set, list) at a document path.

func (*Document) State

func (d *Document) State(ctx context.Context) (*crdt.DocumentCRDTState, error)

State returns the raw DocumentCRDTState for sync purposes.

type List

type List[T any] struct {
	// contains filtered or unexported fields
}

List is a distributed RGA (Replicated Growable Array) backed by a KV store. It stores a crdt.RGAListState under the key, providing ordered sequence semantics with support for concurrent inserts, deletes, and moves.

func NewList

func NewList[T any](store *kv.Store, key string, opts ...Option) *List[T]

NewList creates a new CRDT RGA List backed by the given store and key.

func (*List[T]) Append

func (l *List[T]) Append(ctx context.Context, value T) error

Append adds an element to the end of the list. The element is inserted after the last visible element.

func (*List[T]) Delete

func (l *List[T]) Delete(ctx context.Context, id crdt.HLC) error

Delete removes the element at the given position ID by marking it as tombstoned.

func (*List[T]) Elements

func (l *List[T]) Elements(ctx context.Context) ([]T, error)

Elements returns all visible (non-tombstoned) elements in order.

func (*List[T]) InsertAfter

func (l *List[T]) InsertAfter(ctx context.Context, afterID crdt.HLC, value T) error

InsertAfter inserts an element after the given position ID. Use a zero HLC to insert at the beginning of the list.

func (*List[T]) Len

func (l *List[T]) Len(ctx context.Context) (int, error)

Len returns the number of visible (non-tombstoned) elements.

func (*List[T]) Merge

func (l *List[T]) Merge(ctx context.Context, remote *crdt.RGAListState) error

Merge merges a remote RGA list state into the local state.

func (*List[T]) NodeIDs

func (l *List[T]) NodeIDs(ctx context.Context) ([]crdt.HLC, error)

NodeIDs returns the HLC IDs of visible elements in order. Use these IDs for InsertAfter and Delete operations.

func (*List[T]) State

func (l *List[T]) State(ctx context.Context) (*crdt.RGAListState, error)

State returns the raw RGAListState for sync purposes.

type Map

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

Map is a distributed CRDT Map where each field is an independent LWW register. It stores a crdt.State with per-field FieldState entries.

func NewMap

func NewMap(store *kv.Store, key string, opts ...Option) *Map

NewMap creates a new CRDT Map backed by the given store and key.

func (*Map) All

func (m *Map) All(ctx context.Context) (map[string]json.RawMessage, error)

All returns all fields as a map of field name to raw JSON values.

func (*Map) Delete

func (m *Map) Delete(ctx context.Context, field string) error

Delete removes a field by setting a tombstone.

func (*Map) Get

func (m *Map) Get(ctx context.Context, field string, dest any) error

Get reads a field value and decodes it into dest.

func (*Map) Keys

func (m *Map) Keys(ctx context.Context) ([]string, error)

Keys returns all field names in the map.

func (*Map) Merge

func (m *Map) Merge(ctx context.Context, remote *crdt.State) error

Merge merges a remote CRDT State into the local state using per-field LWW merge.

func (*Map) Set

func (m *Map) Set(ctx context.Context, field string, value any) error

Set sets a field to a value using LWW semantics.

func (*Map) State

func (m *Map) State(ctx context.Context) (*crdt.State, error)

State returns the raw crdt.State for sync purposes.

type Option

type Option func(*crdtConfig)

Option configures a CRDT KV type.

func WithClock

func WithClock(clock crdt.Clock) Option

WithClock sets the Clock implementation for CRDT timestamps. The clock must implement crdt.Clock (e.g., crdt.NewHybridClock).

func WithCodec

func WithCodec(cc codec.Codec) Option

WithCodec sets the codec for serializing CRDT state.

func WithNodeID

func WithNodeID(id string) Option

WithNodeID sets the node identifier for CRDT operations. Each node in a distributed system should have a unique ID.

type Register

type Register[T any] struct {
	// contains filtered or unexported fields
}

Register is a distributed LWW-Register backed by a KV store. The value with the highest HLC timestamp wins.

func NewRegister

func NewRegister[T any](store *kv.Store, key string, opts ...Option) *Register[T]

NewRegister creates a new CRDT LWW-Register backed by the given store and key.

func (*Register[T]) Get

func (r *Register[T]) Get(ctx context.Context) (T, error)

Get reads the current register value.

func (*Register[T]) Merge

func (r *Register[T]) Merge(ctx context.Context, remote *crdt.LWWRegister) error

Merge merges a remote register, keeping the winner per MergeLWW.

func (*Register[T]) Set

func (r *Register[T]) Set(ctx context.Context, value T) error

Set writes a new value, timestamped with the local HLC.

func (*Register[T]) State

func (r *Register[T]) State(ctx context.Context) (*crdt.LWWRegister, error)

State returns the raw LWWRegister for sync purposes.

type Set

type Set[T any] struct {
	// contains filtered or unexported fields
}

Set is a distributed OR-Set (Observed-Remove Set) backed by a KV store. Concurrent add and remove of the same element results in the element being present (add-wins semantics).

func NewSet

func NewSet[T any](store *kv.Store, key string, opts ...Option) *Set[T]

NewSet creates a new CRDT OR-Set backed by the given store and key.

func (*Set[T]) Add

func (s *Set[T]) Add(ctx context.Context, element T) error

Add inserts an element into the set.

func (*Set[T]) Contains

func (s *Set[T]) Contains(ctx context.Context, element T) (bool, error)

Contains returns true if the element is in the set.

func (*Set[T]) Members

func (s *Set[T]) Members(ctx context.Context) ([]T, error)

Members returns all elements currently in the set.

func (*Set[T]) Merge

func (s *Set[T]) Merge(ctx context.Context, remote *crdt.ORSetState) error

Merge merges a remote OR-Set state into the local state.

func (*Set[T]) Remove

func (s *Set[T]) Remove(ctx context.Context, element T) error

Remove removes an element from the set by marking all its current tags as removed.

func (*Set[T]) State

func (s *Set[T]) State(ctx context.Context) (*crdt.ORSetState, error)

State returns the raw ORSetState for sync purposes.

type Syncer

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

Syncer synchronizes CRDT state between two KV stores. It scans for CRDT keys in the primary store and merges them bidirectionally with the replica store.

func NewSyncer

func NewSyncer(primary, replica *kv.Store, opts ...SyncerOption) *Syncer

NewSyncer creates a new CRDT syncer between two KV stores.

func (*Syncer) Start

func (s *Syncer) Start(ctx context.Context)

Start begins a background sync loop.

func (*Syncer) Stop

func (s *Syncer) Stop()

Stop stops the background sync loop.

func (*Syncer) Sync

func (s *Syncer) Sync(ctx context.Context) (*crdt.SyncReport, error)

Sync performs a single round of bidirectional CRDT merge.

type SyncerOption

type SyncerOption func(*syncerConfig)

SyncerOption configures the CRDT Syncer.

func WithKeyPattern

func WithKeyPattern(pattern string) SyncerOption

WithKeyPattern sets the key pattern for scanning CRDT keys during sync.

func WithSyncInterval

func WithSyncInterval(d time.Duration) SyncerOption

WithSyncInterval sets the interval between sync rounds.

Jump to

Keyboard shortcuts

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