Documentation
¶
Overview ¶
Package crdt implements the version-gated record store CRDT defined in docs/crdt.md and docs/crdt-spec.md.
It owns the apply algorithm, _ver (per-field order map) handling and operation semantics over anyenc values, persisted to any-store inside one WriteTx per change.
Index ¶
- Constants
- Variables
- func BackfillApplySeq(ctx context.Context, coll anystore.Collection, spaceId string) error
- func CompareVersion(a, b VersionId) int
- func ComposeVersion(sdkVersion, consumerVersion int) int
- func DeriveRecordId(changeId string) string
- func LoadMeta(ctx context.Context, coll anystore.Collection, objectId string) (maxAddSeq, maxApplySeq uint64, handlerVersions map[string]int, err error)
- func LoadOrInitGeneration(ctx context.Context, coll anystore.Collection, spaceId string) (string, error)
- func LoadSpaceDeletedGate(ctx context.Context, coll anystore.Collection, spaceId string) (head string, ver int, err error)
- func LoadSpaceMaxAddSeq(ctx context.Context, coll anystore.Collection, spaceId string) (uint64, error)
- func MaxObjectApplySeq(ctx context.Context, coll anystore.Collection, spaceId string) (uint64, error)
- func MetaExists(ctx context.Context, coll anystore.Collection, objectId, spaceId string) (bool, error)
- func MetaOwnership(v *anyenc.Value) (spaceId string, deleted bool)
- func NormalizedVersion(v int) int
- func PersistDeletionMark(ctx context.Context, coll anystore.Collection, objectId, spaceId string, ...) error
- func PersistMeta(ctx context.Context, coll anystore.Collection, objectId string, ...) error
- func PersistSpaceDeletedGate(ctx context.Context, coll anystore.Collection, spaceId, head string, ver int) error
- func PersistSpaceMaxAddSeq(ctx context.Context, coll anystore.Collection, spaceId string, ...) error
- func PurgeSpaceMeta(ctx context.Context, coll anystore.Collection, spaceId string, ...) error
- func ResolveRecordIds(ch Change) ([]string, error)
- func SetRecordVersion(arena *anyenc.Arena, record *anyenc.Value, version VersionId, path ...string)
- func SpaceMetaKey(spaceId string) string
- func StaleObjects(ctx context.Context, coll anystore.Collection, spaceId string, ...) ([]string, error)
- type ApplyHook
- type ApplyResult
- type ApplySeqAllocator
- type Change
- type ChangeCtx
- type Controller
- func (c *Controller) ApplyChange(ctx context.Context, ch Change) error
- func (c *Controller) ApplyChangeWithResult(ctx context.Context, ch Change) (ApplyResult, error)
- func (c *Controller) CloseOwnedCollections() error
- func (c *Controller) Collection(ctx context.Context, dataset string) anystore.Collection
- func (c *Controller) DatasetSchemaRev(dataset string) string
- func (c *Controller) Get(ctx context.Context, dataset, id string) *anyenc.Value
- func (c *Controller) HandlerVersions() map[string]int
- func (c *Controller) HasDataset(dataset string) bool
- func (c *Controller) IsShared(dataset string) bool
- func (c *Controller) LoadAndSeedMeta(ctx context.Context, metaColl anystore.Collection) (handlerVersions map[string]int, err error)
- func (c *Controller) LocalFields(dataset string) []string
- func (c *Controller) MaxAddSeq() uint64
- func (c *Controller) NextLocalVersion(ctx context.Context, ch *Change) VersionId
- func (c *Controller) ObjectId() string
- func (c *Controller) PersistMeta(ctx context.Context, metaColl anystore.Collection) error
- func (c *Controller) PersistVersions(ctx context.Context) error
- func (c *Controller) PreValidateLocal(ctx context.Context, ch *Change) error
- func (c *Controller) RecordGetter(ctx context.Context) func(dataset, id string) *anyenc.Value
- func (c *Controller) Records(ctx context.Context, dataset string) []*anyenc.Value
- func (c *Controller) RegisterHandler(ctx context.Context, reg HandlerReg) error
- func (c *Controller) RegisteredDatasets() []string
- func (c *Controller) ReindexLocalLeaves() []LocalLeaf
- func (c *Controller) ReindexPending() bool
- func (c *Controller) ResetForReindex(ctx context.Context, leaves []LocalLeaf) error
- func (c *Controller) SetApplyHook(h ApplyHook)
- func (c *Controller) SetApplySeqAllocator(a *ApplySeqAllocator)
- func (c *Controller) SetMaxAddSeq(seq uint64)
- func (c *Controller) SetSpaceId(spaceId string)
- func (c *Controller) StaleDatasets() []string
- func (c *Controller) ValidateChange(ch Change) error
- type DefaultHandler
- func (DefaultHandler) BeforeCreate(_ *ChangeCtx, _ *RecordChange, _ *Sink) error
- func (DefaultHandler) BeforeDelete(_ *ChangeCtx, _ *RecordChange, _ *Sink) error
- func (DefaultHandler) BeforeModify(_ *ChangeCtx, _ *RecordChange, _ *Op, _ *Sink) error
- func (DefaultHandler) Init(_ context.Context) error
- type Handler
- type HandlerReg
- type LocalLeaf
- type LocalPreValidator
- type LocalPreValidatorMulti
- type ObjectSeq
- type ObjectStamp
- type ObjectStamper
- type Op
- type OpRejection
- type OpType
- type ReadClassification
- type ReadClassifier
- type ReadSeedMode
- type ReadTracking
- type RecordChange
- type RecordGetter
- type SchemaHandler
- func (h *SchemaHandler) BeforeCreate(ctx *ChangeCtx, rec *RecordChange, sink *Sink) error
- func (h *SchemaHandler) BeforeDelete(ctx *ChangeCtx, _ *RecordChange, _ *Sink) error
- func (h *SchemaHandler) BeforeModify(ctx *ChangeCtx, _ *RecordChange, op *Op, sink *Sink) error
- func (h *SchemaHandler) Init(_ context.Context) error
- func (h *SchemaHandler) PreValidateMulti(ch *Change, get RecordGetter) error
- type SharedCollections
- type Sibling
- type Sink
- type VersionId
Constants ¶
const ( IdField = "id" DeletedAtField = "_deletedAt" // AddSeqField records any-sync's per-space AddSeq watermark for the // change that last touched this record. Stamped by the apply path // (not a handler), reserved by the `_` prefix so user writes can't // collide. Backs the consumer-side change-index "records changed // since N" scans. AddSeqField = "_addSeq" )
Reserved field names.
const ApplySeqField = "_applySeq"
ApplySeqField is the reserved record field carrying the per-space apply sequence — the consumer-feed watermark. Unlike AddSeqField (any-sync's DAG delivery counter, which only DAG-borne changes have), _applySeq is minted by the SDK on EVERY apply that mutates a record: DAG changes, the tech-space account mirror's injected applies, and device-local LocalSet writes. It answers "is a consumer (indexer) caught up with any-store", while AddSeq keeps answering "is any-store caught up with any-sync" — two reconciliation links, each keyed by its upstream's coordinate. See docs/scoped-properties-proposal.md § applySeq.
Strictly local: never on the wire, never compared across peers. Gaps are normal (a rolled-back tx skips its allocated seqs); only monotonicity matters.
const MetaCollectionName = "_meta"
MetaCollectionName is the shared collection for per-object metadata. Not prefixed by objectId — one collection holds metadata for all objects in the DB, keyed by objectId.
const SchemaHandlerVersion = 2
SchemaHandlerVersion is the generic schema handler's LOCAL logic version (HandlerReg.Version). Every dataset the SDK applies through this handler carries it, so a change here rebuilds their rows from the DAG (docs/versioning.md). v2: the createTime / modifyTime stamps are TypeDateTime instants, not epoch numbers.
const TracesKey = "_traces"
TracesKey is the reserved field that holds per-versionId trace lists on a record. Shape: `{ "<versionId>": ["traceA", "traceB"], ... }`. Keyed by versionId purely for compaction — multiple fields stamped by the same change share one entry.
const VersionsKey = "_ver"
VersionsKey is the reserved field on every record holding the per-field version map (`_ver` in the spec).
Variables ¶
var ( ErrMissingRecordId = errors.New("crdt: RecordChange.Id is empty and Change.ChangeId is also empty") ErrEmptyIdRequiresUpsert = errors.New("crdt: RecordChange.Id is empty but Upsert is false") ErrMissingDataVersion = errors.New("crdt: Change.DataVersion is empty") // ErrStrictSkipAbsent surfaces the documented strict-mode behaviour // (Upsert=false on an absent record = silent no-op) as a per-record // rejection. Without this, callers who happen to send a // non-existent record id with Upsert=false get a fully successful // WriteResult (RecordIds populated, Rejections empty) but the // record never lands in the projection — and the same change still // commits to the tree, so peers see the same skip. The rejection // lets writers detect the case via ApplyResult.Rejections without // changing the strict-mode semantics that other paths rely on. ErrStrictSkipAbsent = errors.New("crdt: strict modify (Upsert=false) skipped because record id does not exist") )
Sentinel errors.
var ErrInvalidPath = errors.New("crdt: invalid path")
ErrInvalidPath is returned when an op references a malformed or reserved field path. Joined with ErrValidation so callers can detect both classes with errors.Is(err, ErrValidation).
var ErrRecordDeleted = errors.New("crdt: record is deleted; the id cannot be reused")
ErrRecordDeleted is the rejection sentinel for a modify (including upsert) that landed on a tombstoned record. Delete-wins absorbs the write — the tombstone is sticky and none of the ops land — and without this rejection the absorbed write is indistinguishable from a successful create to a local caller (ModifyResult with recordIds and no rejections, nothing stored). The absorption itself is the convergence rule and stays; this sentinel only makes it visible. Deleting an already-deleted record stays silent (idempotent).
var ErrUnknownDataset = errors.New("crdt: no handler for dataset")
ErrUnknownDataset is returned when a change targets a dataset that has no registered handler. Per spec §8.2 the SDK still persists such changes via any-sync; this in-memory state just skips them.
var ErrValidation = errors.New("crdt: validation rejected")
ErrValidation is the sentinel returned (via errors.Join) when a handler rejects an operation. Per-op rejections drop just the offending op from the apply step; whole-change validation failures (path syntax) abort the Change entirely.
Functions ¶
func BackfillApplySeq ¶
BackfillApplySeq seeds applySeq := addSeq on this space's legacy per-object _meta rows (rows written before the applySeq watermark existed). One-off per space, guarded by a flag on the space's own _meta row; idempotent. Keeps consumer cursors valid across the re-key: every historical position N (AddSeq units) means the same thing on the applySeq axis, and the allocator seeds past it.
func CompareVersion ¶
CompareVersion returns -1, 0, or +1 for a < b, a == b, and a > b respectively. Provided as a documented entry point so the comparison contract is discoverable; internally just delegates to cmp.Compare. Callers inside the package can use raw `<` / `>` / `==` on VersionId directly since it's a string type.
The empty string is the "no version" sentinel and sorts less than any non-empty version.
func ComposeVersion ¶
ComposeVersion packs an SDK-side and a consumer-side handler version into the single integer persisted per dataset. Only equality matters downstream, so any injective encoding works — this one keeps both halves legible (2_000_003 = SDK 2, consumer 3) and, unlike a sum, cannot let one side's bump cancel the other's.
func DeriveRecordId ¶
DeriveRecordId produces the default record id for an empty-id RecordChange.
The derivation is base58(xxh3-64(changeId)). changeId is a content- addressable CID, so its bytes are uniformly random; xxh3 preserves that uniformity, and base58 encodes the resulting 8 bytes in up to 11 chars.
Collision math (2^64 space): P(any collision) < 1e-6 up to ~6M derived ids within a single storage namespace; comfortable for any realistic scale.
func LoadMeta ¶
func LoadMeta(ctx context.Context, coll anystore.Collection, objectId string) (maxAddSeq, maxApplySeq uint64, handlerVersions map[string]int, err error)
LoadMeta reads the per-object metadata from the _meta collection. Returns zero values if the document doesn't exist yet.
func LoadOrInitGeneration ¶
func LoadOrInitGeneration(ctx context.Context, coll anystore.Collection, spaceId string) (string, error)
LoadOrInitGeneration returns the per-space rebuild epoch on the space:<id> row, minting a fresh UUID when the row (or the gen field) is absent — an sdk.db wipe drops the row, so the epoch changes and signals a renumbered applySeq axis to consumers. Idempotent; call once at store load (single-threaded) to mint eagerly, and on every consumer read.
func LoadSpaceDeletedGate ¶
func LoadSpaceDeletedGate(ctx context.Context, coll anystore.Collection, spaceId string) (head string, ver int, err error)
LoadSpaceDeletedGate reads the persisted deletion-reconcile gate for a space — the settings-tree head last reconciled (dh) and the reconcile-logic version last applied (dv). Returns ("", 0) when no row exists (fresh/rebuilt sdk.db), a guaranteed mismatch that forces exactly one reconcile sweep.
func LoadSpaceMaxAddSeq ¶
func LoadSpaceMaxAddSeq(ctx context.Context, coll anystore.Collection, spaceId string) (uint64, error)
LoadSpaceMaxAddSeq reads the persisted space-level head-store watermark — the lower bound of "we've already replayed any-sync trees up to this LastAddSeq for this space". Returns 0 when no row exists yet (first boot or never caught up).
func MaxObjectApplySeq ¶
func MaxObjectApplySeq(ctx context.Context, coll anystore.Collection, spaceId string) (uint64, error)
MaxObjectApplySeq returns the highest persisted per-object applySeq in spaceId — the current upper bound a change-index cursor can reach, and the ApplySeqAllocator's seed. Returns 0 when the space has no scoped object rows yet. Run BackfillApplySeq first so legacy rows participate.
func MetaExists ¶
func MetaExists(ctx context.Context, coll anystore.Collection, objectId, spaceId string) (bool, error)
MetaExists reports whether a space-scoped per-object _meta row exists — i.e. the object was ever materialized/indexed in this space. Gates the del-stamp for objects that were fed but never got a shared `objects` row (base-dataset-only). Read with the caller's (tx) ctx.
func MetaOwnership ¶
MetaOwnership reads the ownership fields off a per-object _meta row: the scoping spaceId (empty on pre-scoping rows) and the sticky purge marker. The orphan-collection GC uses it to decide whether a per-object collection's owner is live — a row claiming a live space without the purge marker proves liveness even when the object never got a shared `objects` row (base-dataset-only objects).
func NormalizedVersion ¶
NormalizedVersion resolves a HandlerReg.Version to the value that is actually persisted and compared: an unset version means 1.
func PersistDeletionMark ¶
func PersistDeletionMark(ctx context.Context, coll anystore.Collection, objectId, spaceId string, applySeq uint64) error
PersistDeletionMark stamps the object's kept _meta row as deleted, in place: del=true, as=applySeq (a fresh seq > the object's last content applySeq), and re-asserts sp so a sparse-history object is guaranteed visible to the change-index query. Leaves q/hv intact. Call inside the same WriteTx as the shared-`objects` row removal so the deletion is announced atomically with the purge.
func PersistMeta ¶
func PersistMeta(ctx context.Context, coll anystore.Collection, objectId string, maxAddSeq, maxApplySeq uint64, handlerVersions map[string]int, spaceId string) error
PersistMeta writes per-object metadata to the _meta collection. Call inside the same WriteTx as the record mutations for atomicity.
spaceId scopes the row for the change-index query; pass "" to leave it unset (unit tests without a space). An unset row is invisible to QueryChangedObjects, which is the accepted lazy-backfill behaviour.
func PersistSpaceDeletedGate ¶
func PersistSpaceDeletedGate(ctx context.Context, coll anystore.Collection, spaceId, head string, ver int) error
PersistSpaceDeletedGate writes the deletion-reconcile gate. Sets only dh/dv via ModifyFunc so the forward-catchup watermark "q" and the generation "gen" on the same space:<id> row are preserved. Called ONLY after a fully successful reconcile sweep.
func PersistSpaceMaxAddSeq ¶
func PersistSpaceMaxAddSeq(ctx context.Context, coll anystore.Collection, spaceId string, maxAddSeq uint64) error
PersistSpaceMaxAddSeq writes the space-level head-store watermark. Written once at the end of a successful catch-up pass — per-change progress is already captured by the per-object _meta rows, so this value is a coarse-grained "no work needed at startup" hint, not a per-change atomic counter.
func PurgeSpaceMeta ¶
func PurgeSpaceMeta(ctx context.Context, coll anystore.Collection, spaceId string, objectIds []string) error
PurgeSpaceMeta removes every _meta row belonging to spaceId: the space-scoped per-object watermark rows (matched on sp), any explicit objectIds (pre-scoping rows without sp), and the `space:<id>` row — dropping it rotates the generation on the next rebuild, the signal change-index consumers full-reindex on. Offload calls this: a stale MaxAddSeq watermark is NOT inert once a space can re-materialize (guest rejoin) — the rebuilt controllers would trust it and skip the entire cold-restore replay, leaving synced trees with zero rows.
func ResolveRecordIds ¶
ResolveRecordIds replaces empty RecordChange.Ids with ChangeId- derived values and returns the resolved id per record. First empty-id gets `base58(xxh3-64(ChangeId))`; subsequent empty ids get that seed with `:<index>` appended. Caller-supplied ids pass through unchanged.
The suffix separator is `:` (not `/`) so resolved ids are safe to embed in REST URL path segments without extra encoding.
Exposed so the write path can return the resolved record ids alongside VersionId/ChangeId — useful for property creates where the propId is the resolved (auto-derived) record id.
func SetRecordVersion ¶
SetRecordVersion writes `version` to record._ver at `path` (creating _ver if missing) and applies the bottom-up collapse rule. Use after a successful gated mutation.
Setting at the leaf overwrites whatever was at that exact path (string or subtree). Intermediate string nodes on the path are expanded into objects with a defaultKey holding the previous string, preserving inherited versions for sibling fields.
func SpaceMetaKey ¶
SpaceMetaKey returns the _meta document id for a space's row.
func StaleObjects ¶
func StaleObjects(ctx context.Context, coll anystore.Collection, spaceId string, registered map[string]int) ([]string, error)
StaleObjects returns the ids of objects in spaceId whose persisted handler versions differ from `registered`, plus those whose rebuild was interrupted (captured leaves still on the row) — the sweep's work list. Purged objects are skipped: their rows are already gone and a rebuild must never resurrect them.
One scan of the space's _meta rows, bounded by its object count. The hv map is small (one entry per dataset the object wrote), so the per-row compare is cheap.
Types ¶
type ApplyHook ¶
ApplyHook runs inside the apply WriteTx after every record of a change has applied, before the watermark persist and commit — a returned error rolls the whole change back. recordIds are the resolved per-record ids (ChangeId sugar applied, shared-dataset collapse done). The read-tracking classification closure installs here so unread entries are atomic with the change.
type ApplyResult ¶
type ApplyResult struct {
Rejections []OpRejection
DerivedOps [][]Op
// ApplySeq is the per-space apply sequence allocated for this
// change (0 when the Controller has no allocator — unit tests).
// Forwarded to the change-index feed so consumers cursor on it.
ApplySeq uint64
// ObjectStamps are the writes ObjectStampers landed on the
// object's rows in shared datasets alongside this change (see
// ObjectStamper). Empty when the change targets the shared dataset
// itself, arrived on the local/account route, wrote nothing, the
// row is absent/tombstoned, or every stamp op lost its LWW gate;
// Ops holds only the ops that took it. The dispatcher projects
// each as an update of that row so live queries on the shared
// dataset see it.
ObjectStamps []ObjectStamp
}
ApplyResult is the auxiliary return from ApplyChangeWithResult — the change still committed (with a fresh VersionId), but specific ops or records were dropped. Empty Rejections means everything landed.
DerivedOps surfaces, per ch.Records index, the extra ops the apply path stamped onto the record beyond the input ops: handler-emitted sink.derived ops (author, createdAt, …) and the modifier's own auto-stamped fields (_ver.id on create). The dispatcher merges these into the EventRecord so a viewer can reconstruct a fresh record without having to know which fields are "auto" vs "user". Nil when no record had extras (the common update path); per-record entry is nil when that record had none.
type ApplySeqAllocator ¶
type ApplySeqAllocator struct {
// contains filtered or unexported fields
}
ApplySeqAllocator mints the per-space apply sequence. One allocator per space, shared by every object Controller in it (the same sharing pattern as the per-space VersionAllocator).
Seeding is lazy: the first Next runs seedFn to learn the highest already-persisted applySeq (after the one-off backfill that seeds legacy rows from their AddSeq, this is ≥ every historical AddSeq, so pre-existing consumer cursors stay valid on one continuous axis).
Crash safety needs no separate counter row: every allocated seq that matters is persisted via some object's _meta row in the same WriteTx as the record stamp, and re-seeding reads the max of those — the counter can never regress below a persisted stamp. Callers must allocate AFTER acquiring the WriteTx (any-store's single writer then serializes allocation order = commit order, so an ascending-cursor consumer can't skip a seq that commits late).
func NewApplySeqAllocator ¶
func NewApplySeqAllocator(seedFn func(ctx context.Context) (uint64, error)) *ApplySeqAllocator
NewApplySeqAllocator builds an allocator that lazily seeds from seedFn. A nil seedFn seeds at 0 (fresh space, unit tests).
type Change ¶
type Change struct {
SpaceId string
ObjectId string
Dataset string
ChangeId string
VersionId VersionId
AddSeq uint64
// PrevIds are the DAG parents of this change (the any-sync change's
// PreviousIds), stamped by the apply pipeline next to Creator. Empty
// on Local/Injected changes (no DAG) and on hand-built test changes.
// Consumed by read tracking: unread entries persist them so the
// mark-read ancestor closure never needs the tree.
PrevIds []string
// ApplySeq is the per-space apply sequence the Controller allocates
// inside the apply transaction (never set by callers, never on the
// wire). Stamped on every written record as _applySeq and persisted
// as the per-object maxApplySeq — the consumer-feed watermark that,
// unlike AddSeq, also covers non-DAG applies (account mirror,
// device-local writes). Zero when the Controller has no allocator.
ApplySeq uint64
Records []RecordChange
// DataVersion pins the change to the schema or handler state its writer
// used: `typeId:shortId` pairs for schema state, or a handler version
// string. The per-dataset mapping is in
// docs/types-properties-proposal.md § "Change-level DataVersion".
//
// The string is opaque to the CRDT layer — only equality matters, not
// ordering. Empty is invalid and causes ApplyChange to reject the change.
DataVersion string
// Timestamp carried alongside the change for tombstones (delete writes a
// `deletedAt` field). The CRDT layer treats it as advisory metadata; it is
// not used for conflict resolution.
Timestamp int64
// ObjectAuthor is the StrKey-encoded identity (PubKey.Account()) of
// the peer that created the OBJECT — i.e. the signer of the tree's
// root change. Constant across every change in a tree; the apply
// pipeline stamps it before calling Controller.ApplyChange. Used by
// SystemPropertiesHandler to auto-stamp `author` at row root on
// first record creation. Empty on hand-built changes (tests) where
// no tree is wired.
//
// Distinct from Creator below: ObjectAuthor is the constant root
// signer (object creator); Creator is the per-change signer (who
// wrote THIS change). They coincide for the root change but diverge
// for subsequent changes in shared / multi-author objects.
ObjectAuthor string
// Creator is the StrKey-encoded identity (PubKey.Account()) of the
// peer that signed THIS change — i.e. the per-change author. Stamped
// by Object.stampObjectMeta from the underlying any-sync change's
// Identity (looked up via tree.GetChange by ChangeId). Empty on
// hand-built changes (tests) where no tree is wired.
Creator string
// ObjectCreatedAt is the Unix-seconds timestamp of the tree's root
// change — the moment the object was created. Constant across every
// change in a tree. Used by SystemPropertiesHandler to auto-stamp
// `createdAt` at row root on first record creation.
ObjectCreatedAt int64
// TraceIds are opaque correlation tokens attached by the originating
// caller (e.g. a UI session id, an AI operation id). They travel with the
// change through any-sync (stored in a dedicated, non-encrypted field on
// the wire) so the whole space can be queried by trace. On the record
// side, the CRDT stamps them into `_traces[versionId]` for any fields
// this change actually landed, so per-record queries can map a field's
// `_ver` entry back to the traces that produced it. An empty slice means
// "no trace for this change" — any prior `_traces[versionId]` entry for
// this versionId is cleared.
TraceIds []string
// Local marks a device-local materialization that does NOT flow
// through the any-sync DAG. Set by Object.LocalSet; never by a
// synced write or replay. When true the apply path: (1) skips the
// dataset handler (local fields are handler-exclusive), (2) only
// admits local-class fields (classifyFieldWrite), and the VersionId
// is locally allocated (NextVersion of the field's current version)
// rather than an any-sync OrderId. When false, the apply path
// rejects any op targeting a local-class field. The field classes
// are disjoint by construction — see version.go.
Local bool
// Injected marks the account mirror's materialization: like Local
// it does NOT flow through this object's DAG and skips the dataset
// handler, but the VersionId is CALLER-SUPPLIED — the tech-space
// carrier tree's orderId for the change that produced the value —
// rather than locally minted. Set by Object.InjectedSet; never by
// a synced write, replay, or LocalSet (Local and Injected are
// mutually exclusive). Only account-class fields (and undeclared
// heads of DynamicScopeByKey datasets, whose per-key scopes the
// mirror resolves itself) are writable on this route. Safe because
// account paths are written by exactly one writer per device (the
// mirror) sourcing one tech tree — versions stay monotonic per
// path. See docs/scoped-properties-proposal.md § Account transport.
Injected bool
}
Change is one batch — exactly one any-sync DAG change. Every op in the batch shares the same VersionId and targets the same Dataset inside the same (SpaceId, ObjectId) tree.
Three identifiers travel with each change, playing three different roles:
- VersionId (= any-sync orderId) — lexicographically sortable DAG order used for gating. Scoped to one object tree. See VersionId docs.
- ChangeId — content-addressable, immutable DAG change hash. Used for observability and as the seed for the "empty record id → ChangeId" sugar on RecordChange. Distinct from VersionId.
- AddSeq — monotonic local delivery sequence assigned by any-sync's spacestorage as changes are received. VersionId says "where in the DAG did this happen"; AddSeq says "when did we receive it locally". The Controller tracks the max AddSeq it has observed so the space layer can efficiently ask any-sync "give me everything with seq > max" on restore or reindex.
type ChangeCtx ¶
type ChangeCtx struct {
Change *Change
Before *anyenc.Value
// SelfIdentity is this replica's account identity
// (PubKey.Account() encoding). Populated ONLY on the read-tracking
// classify path — read state is device-local, so identity-relative
// verdicts (ReadClassification.Audience) are sound there. It stays
// zero in the Before* handler hooks on purpose: handler validation
// runs on every replica for the same change and must be
// replica-independent, or CRDT replicas diverge.
SelfIdentity string
// Get resolves the CURRENT value of a record in one of this
// object's datasets by explicit id — nil when the record is absent,
// the id is empty, or the invoker wired no reader (hand-built test
// contexts). Reads happen inside the apply WriteTx, so a record
// created earlier in the same change is visible. Tombstones are
// returned as-is (check _deletedAt); their content fields are wiped.
//
// Determinism contract for Before* hooks: the hook runs on every
// replica for the same change and its output must not diverge, so a
// handler may consult ONLY fields that are immutable post-create
// (derived creation stamps such as creator/createdAt qualify;
// freely-edited fields do not) and only on records that are causal
// ancestors of the triggering change — the writer saw them, so
// every replica applies them first. Known accepted edge: a record
// deleted CONCURRENTLY with the triggering change loses its fields
// on the tombstone, so a replica that applied the delete first
// reads nil-equivalent state; derived output disagrees between
// replicas on that race, which is tolerable only because derived
// ops are local re-derivation, never synced payload.
Get func(dataset, id string) *anyenc.Value
// RecordId is the resolved id of the record being classified.
// Populated ONLY on the read-tracking classify path (like
// SelfIdentity) — a create with an auto-derived id carries an empty
// RecordChange.Id, and a classifier that wants to point-read the
// just-applied record via Get needs the real id. Zero in the
// Before* handler hooks.
RecordId string
}
ChangeCtx is the per-callback context handed to a Handler. It carries the originating Change (read-only metadata: VersionId, ChangeId, Timestamp, Creator, ObjectAuthor, …) and the record's pre-op state.
ObjectAuthor is the constant root-signer (object creator); Creator is the per-change signer (who wrote THIS change). For the root change they coincide; for shared spaces / multi-author objects they diverge — handlers that gate "only the author of this message can edit it" read Creator; handlers stamping the object's creation facts (author, createdAt) read ObjectAuthor, and per-change provenance (modifiedBy) reads Creator.
Before evolves across ops in the same RecordChange: op[1]'s Before is op[0]'s after — but only counting ops that actually landed (rejected ops don't mutate state). Nil during BeforeCreate, since the record didn't exist before this Change.
Pointer-passed; do not retain the pointer past the callback.
type Controller ¶
type Controller struct {
// contains filtered or unexported fields
}
Controller is the CRDT entrypoint for one any-sync object tree. It applies version-gated changes to an any-store database, one collection per dataset. The Controller owns no arena — any-store's DocBuffer pool provides arenas inside UpsertId/UpdateId modifiers.
Transaction model: ApplyChange wraps its mutations in a WriteTx. When the caller already holds a WriteTx (batch apply for restore), pass its context — ApplyChange will create a savepoint inside the existing tx.
Not safe for concurrent use.
func NewController ¶
func NewController(ctx context.Context, objectId string, db anystore.DB, regs ...HandlerReg) (*Controller, error)
NewController opens or creates the per-dataset collections and returns a ready Controller. The caller provides the any-store DB (which may be shared across multiple Controllers for different objects). Collection names are `{objectId}/{datasetName}` unless a shared collection is supplied for that dataset, in which case the shared one is used and row ids in writes default to the change's ObjectId.
func NewControllerWithShared ¶
func NewControllerWithShared(ctx context.Context, objectId string, db anystore.DB, shared SharedCollections, regs ...HandlerReg) (*Controller, error)
NewControllerWithShared opens the per-dataset collections, accepting shared overrides for specific datasets (typically the per-space `objects` values collection). Datasets not in `shared` get the usual per-object collection.
func (*Controller) ApplyChange ¶
func (c *Controller) ApplyChange(ctx context.Context, ch Change) error
ApplyChange is the simple entry point — wraps ApplyChangeWithResult and discards the per-op rejection list. Existing callers that don't need rejection info (tests, replay paths) keep the original signature.
func (*Controller) ApplyChangeWithResult ¶
func (c *Controller) ApplyChangeWithResult(ctx context.Context, ch Change) (ApplyResult, error)
ApplyChangeWithResult applies a single change per spec §7 and returns the list of per-op rejections. The change still commits when ops are rejected; the rejections let the writer detect "committed with a hole" without scanning post-state.
Wraps all mutations in a WriteTx for atomicity. When the ctx already carries a WriteTx (batch mode), a savepoint is used instead — lightweight and correct.
Content validation (op-path syntax, field class) and handler validation drop only the offending op or key, recorded into ApplyResult.Rejections. Local writers fail fast in ValidateChange before the change enters the DAG.
func (*Controller) CloseOwnedCollections ¶
func (c *Controller) CloseOwnedCollections() error
CloseOwnedCollections closes and forgets every cached PER-OBJECT collection handle (`<objectId>_<dataset>`). Shared (per-space) handles and the `_meta` collection stay untouched — the Store owns those, and they are shared across every controller in the space.
Called on ocache eviction (Object.Close / TryClose): each open any-store handle pins its query-planner sketch and value caches for the process lifetime, so with a collection-per-object schema an account's heap grows linearly with objects ever touched unless idle objects release their handles.
Recovery is open-by-name: the lazy paths (collectionForWrite / collectionForRead) re-open on the next touch, and after this call they stop caching (collsReleased), so a stale Controller retained past eviction can never hold a closed — or later-closed — handle.
Serialization: the caller (Object close path) holds — or held — tree.Lock with Object.closed set, so no apply is in flight on this controller. Unlocked readers racing the close get a one-shot ErrCollectionClosed and re-open by name on retry.
func (*Controller) Collection ¶
func (c *Controller) Collection(ctx context.Context, dataset string) anystore.Collection
Collection returns the any-store collection backing the named dataset. Exposed so query / iteration paths can hit any-store directly without going through Controller's own helpers. Returns nil when the dataset has no registered handler or no writer has materialised the per-object collection on disk yet.
func (*Controller) DatasetSchemaRev ¶
func (c *Controller) DatasetSchemaRev(dataset string) string
DatasetSchemaRev returns the SchemaRev the dataset registered with ("" for static regs / unknown datasets). Immutable after construction; concurrency-safe like HasDataset.
func (*Controller) Get ¶
Get returns the record by id from the dataset, or nil if absent.
The returned *anyenc.Value is cloned off any-store's pooled doc buffer — safe to retain past the FindId call. Without the clone the value would alias a buffer that's released back to the pool before this function returns (any-store's FindId releases via defer).
func (*Controller) HandlerVersions ¶
func (c *Controller) HandlerVersions() map[string]int
HandlerVersions returns a map of dataset→version from the Controller's registered handlers.
func (*Controller) HasDataset ¶
func (c *Controller) HasDataset(dataset string) bool
HasDataset reports whether a handler is registered for the dataset. The handler map is immutable after construction, so this is safe to call concurrently (same contract as ValidateChange). Used by the store to detect controllers built before a runtime dataset appeared.
func (*Controller) IsShared ¶
func (c *Controller) IsShared(dataset string) bool
IsShared reports whether the dataset uses a per-space shared collection (write-side rule: row id = ObjectId, not RecordChange.Id). Subscribers projecting changes back to a path-based wire need this to look up the right post-apply row.
func (*Controller) LoadAndSeedMeta ¶
func (c *Controller) LoadAndSeedMeta(ctx context.Context, metaColl anystore.Collection) (handlerVersions map[string]int, err error)
LoadAndSeedMeta reads metadata from the _meta collection, seeds the Controller's maxAddSeq, and returns the stored handler versions so the caller can compare them with current versions for re-indexing decisions.
func (*Controller) LocalFields ¶
func (c *Controller) LocalFields(dataset string) []string
LocalFields returns the dataset's declared local-scope top-level field names. Local values are written by Object.LocalSet and never enter the DAG, so a replay cannot reproduce them — the re-index path snapshots these fields before the wipe and re-applies them after.
func (*Controller) MaxAddSeq ¶
func (c *Controller) MaxAddSeq() uint64
func (*Controller) NextLocalVersion ¶
func (c *Controller) NextLocalVersion(ctx context.Context, ch *Change) VersionId
NextLocalVersion returns the VersionId to stamp on a device-local (Change.Local) write: one version for the whole change, strictly greater than every targeted field's current version on its record. Because local fields are handler-exclusive and never written by a synced change, a single bump above their current versions is monotonic and cannot be overshadowed by a competing synced write. Reads the current records by their explicit ids (local writes never use empty-id resolution). Safe to call under the apply lock.
func (*Controller) ObjectId ¶
func (c *Controller) ObjectId() string
func (*Controller) PersistMeta ¶
func (c *Controller) PersistMeta(ctx context.Context, metaColl anystore.Collection) error
PersistMeta writes the Controller's watermarks and handler versions to the _meta collection. The caller should pass a context carrying the same WriteTx as the record mutations for atomicity.
func (*Controller) PersistVersions ¶
func (c *Controller) PersistVersions(ctx context.Context) error
PersistVersions stamps the current watermarks and handler versions and clears the in-flight mark (the persisted local leaves), ending the rebuild. The replay path stamps versions from inside each apply; this also covers the object whose replay re-applied nothing (an empty tree, or one whose every change belongs to a dataset that is no longer registered), which would otherwise be found stale again on every load.
func (*Controller) PreValidateLocal ¶
func (c *Controller) PreValidateLocal(ctx context.Context, ch *Change) error
PreValidateLocal runs the dataset handler's optional LocalPreValidator against a LOCAL change before it enters the DAG. No-op when the dataset has no handler or the handler doesn't implement the interface. Reads the current value of the change's target record and hands it to PreValidate; a non-nil error rejects the whole write.
Unlike ValidateChange (pure, structural), this does a DB read — but it runs only on the local-write path, never on the inbound hot path. ch.ChangeId is not yet known here, so empty record ids can't be resolved; for the shared `objects` dataset the row id is the controller's own objectId, which is always known.
func (*Controller) RecordGetter ¶
RecordGetter returns a ChangeCtx.Get implementation bound to ctx — the one wrapper both the handler-hook and read-tracking classify seams share. When ctx carries a WriteTx (the apply path), reads reuse that tx (any-store's getReadTx), so calling it from inside a Modify callback or an ApplyHook is safe. Empty ids resolve to nil without a store round-trip.
func (*Controller) Records ¶
Records returns all live (non-tombstone) records in the dataset.
Each value is cloned off the iterator's reusable buffer onto its own parser arena so caller-held pointers stay valid past iteration (any-store's iter.Doc reuses one DocBuffer across Next() calls). For streaming over many rows, hit Collection(dataset).Find(...) + Iter directly and clone selectively — Records always materialises.
func (*Controller) RegisterHandler ¶
func (c *Controller) RegisterHandler(ctx context.Context, reg HandlerReg) error
RegisterHandler adds a handler at runtime (for late-bound datasets).
func (*Controller) RegisteredDatasets ¶
func (c *Controller) RegisteredDatasets() []string
RegisteredDatasets returns every dataset name this controller handles, sorted. The re-index wipe walks them to find what to clear.
func (*Controller) ReindexLocalLeaves ¶
func (c *Controller) ReindexLocalLeaves() []LocalLeaf
ReindexLocalLeaves returns the leaves an in-flight rebuild persisted. Empty unless ReindexPending.
func (*Controller) ReindexPending ¶
func (c *Controller) ReindexPending() bool
ReindexPending reports whether a rebuild of this object was started and never finished restoring its local-scope leaves — a crash, or the sweep cancelled at shutdown, between ResetForReindex and PersistVersions. The next load resumes: the replay continues from the persisted watermark and the leaves come back from the _meta row.
func (*Controller) ResetForReindex ¶
func (c *Controller) ResetForReindex(ctx context.Context, leaves []LocalLeaf) error
ResetForReindex rewinds the controller to "never applied" so the next ColdRestore replays the whole tree, and persists the rewound watermark together with the captured local-scope leaves BEFORE the caller wipes the materialized rows.
Order matters: a crash between the wipe and the first re-applied change must leave a watermark of 0 over the missing rows, so the next boot replays them. Persisting after the wipe would leave a window where the stored watermark claims changes are applied whose rows are gone. The leaves ride the same upsert for the same reason — once the rows are gone they exist nowhere else, and a replay interrupted at any point (including the sweep's own shutdown cancel) must still find them.
The stored handler versions are deliberately left alone — they keep the stale verdict (and therefore the wipe) armed until a replay actually re-applies a change, at which point its PersistMeta stamps the current versions. The leaves outlive that stamp: only PersistVersions clears them.
func (*Controller) SetApplyHook ¶
func (c *Controller) SetApplyHook(h ApplyHook)
SetApplyHook wires the in-tx apply hook. Call once right after construction, before the first apply. No-op on a nil controller.
func (*Controller) SetApplySeqAllocator ¶
func (c *Controller) SetApplySeqAllocator(a *ApplySeqAllocator)
SetApplySeqAllocator wires the per-space apply-sequence allocator. Call once right after construction, before the first apply. No-op on a nil controller.
func (*Controller) SetMaxAddSeq ¶
func (c *Controller) SetMaxAddSeq(seq uint64)
SetMaxAddSeq seeds the watermark from persisted storage on restore.
func (*Controller) SetSpaceId ¶
func (c *Controller) SetSpaceId(spaceId string)
SetSpaceId scopes this controller's persisted _meta row to a space. Call once right after construction, before the first apply, so the change-index query can filter object rows by space. No-op on a nil controller.
func (*Controller) StaleDatasets ¶
func (c *Controller) StaleDatasets() []string
StaleDatasets returns the datasets whose materialized rows were built by a different version of their handler. Empty for a fresh object (no stored versions) and for the common case where nothing changed.
func (*Controller) ValidateChange ¶
func (c *Controller) ValidateChange(ch Change) error
ValidateChange runs the cheap up-front checks ApplyChange would otherwise do — without needing the change's final ChangeId or any DB reads:
- DataVersion non-empty.
- Dataset has a registered handler + collection.
- ObjectId, when set, matches the controller.
- Every op's path syntax legal.
- Every op honors its dataset's field-class scope (a synced change may not write a local/derived field, etc.) — so a scope-violating change never reaches AddContent and the DAG.
- Records that submit empty Id flag Upsert=true (the only legal way to get an auto-derived id).
Used by the local-write path to fail fast BEFORE handing the payload to any-sync's AddContent — a structurally invalid change stays out of the DAG entirely. Handler-level rejection (kind mismatch, terminal-status, etc.) still fires at apply time because it depends on the existing record state and needs a DB read.
Safe to call concurrently; no state mutation.
type DefaultHandler ¶
type DefaultHandler struct{}
DefaultHandler is a no-op handler accepting every op. Useful as a base for tests and as a convenient embed; the dataset name / version / indexes live on the HandlerReg, not here.
func (DefaultHandler) BeforeCreate ¶
func (DefaultHandler) BeforeCreate(_ *ChangeCtx, _ *RecordChange, _ *Sink) error
func (DefaultHandler) BeforeDelete ¶
func (DefaultHandler) BeforeDelete(_ *ChangeCtx, _ *RecordChange, _ *Sink) error
func (DefaultHandler) BeforeModify ¶
func (DefaultHandler) BeforeModify(_ *ChangeCtx, _ *RecordChange, _ *Op, _ *Sink) error
type Handler ¶
type Handler interface {
// Init runs once at registration (NewController / RegisterHandler).
// Use it to acquire dependencies; nil for stateless handlers.
Init(ctx context.Context) error
BeforeCreate(ctx *ChangeCtx, rec *RecordChange, sink *Sink) error
BeforeModify(ctx *ChangeCtx, rec *RecordChange, op *Op, sink *Sink) error
BeforeDelete(ctx *ChangeCtx, rec *RecordChange, sink *Sink) error
}
Handler is the lifecycle behavior for one dataset. Hooks fire inside any-store's Modify callback so handlers can read pre-state from ctx.Before and emit derived or sibling writes via sink without an extra DB round-trip. Cross-record reads go through ctx.Get, subject to its determinism contract (immutable fields of causal ancestors only) — hooks run on every replica and must not diverge.
Per-op error returns from BeforeModify drop just the offending op; other ops in the same RecordChange still apply. BeforeCreate / BeforeDelete are per-record, and an error there drops the whole RecordChange.
A handler is pure behavior: its dataset name, wire DataVersion, handler version, and indexes are declared alongside it in a HandlerReg at registration (see NewController / HandlerReg), not via methods on the handler. All three Before* methods may be no-ops (see DefaultHandler).
type HandlerReg ¶
type HandlerReg struct {
Name string
Version int
Handler Handler
Indexes []anystore.IndexInfo
// Schema is the dataset's required, JSON-Schema-compatible field
// declaration. Each field carries a class (Scope: synced/derived/
// local/account) the apply path enforces: derived fields are
// handler-only, local fields never sync, and a dataset that isn't
// Dynamic rejects undeclared fields. Free-form datasets (shortIds,
// the per-type `objects` namespace) set Schema.Dynamic.
Schema schema.Dataset
// SchemaRev is an opaque fingerprint of the registered schema for
// runtime-defined datasets. The store compares a resident
// controller's rev against the current catalog rev to detect stale
// registrations (a field added/removed after construction) and
// evict lazily. Empty for static registrations.
SchemaRev string
// DynamicScopeByKey marks a Dynamic dataset whose UNDECLARED field
// heads carry per-key scopes owned by the dataset's own layer (the
// per-space `objects` dataset: scope lives on the property
// definition, resolved by the writer's routing and the handler's
// registry validation — the controller can't see second path
// segments). For such datasets the head-level local/synced
// direction check is skipped for undeclared heads; DECLARED fields
// (e.g. the derived author/createdAt/spaceId) stay fully enforced.
DynamicScopeByKey bool
// ReadTracking opts the dataset into read/unread tracking; nil =
// untracked. See readtracking.go.
ReadTracking *ReadTracking
// DisableFilteredReplay opts the dataset out of the record-filtered
// history fast path (docs/version-history-proposal.md §4.2). The
// fast path is sound only while apply hooks stay record-local —
// a handler that reads OTHER records during apply must set this,
// forcing per-record history through the full-object slow path.
DisableFilteredReplay bool
// SkipHistory keeps this dataset out of the persistent history
// index (docs/version-history-proposal.md §4.4): no index rows are
// written and the dataset is invisible in history listings. For
// chatty machine-written datasets (presence-like state) whose
// permanent index would leak disk for history nobody asks for.
// DAG changes still retain everything — flipping the flag later
// just requires a backfill.
SkipHistory bool
}
HandlerReg binds a dataset name to its handler behavior and the metadata the Controller needs: the handler Version (persisted in _meta for re-index decisions; defaults to 1) and any indexes to ensure on the dataset's collection. The wire DataVersion peers gate against is a write-time concern owned by the caller (it stamps Change.DataVersion), not the Controller — so it lives on the public handler.Dataset, not here.
type LocalLeaf ¶
LocalLeaf is one local-scope value captured before a rebuild's wipe. Local values never entered the DAG, so no replay reproduces them; the re-index path re-applies these after the replay. Value is the leaf's anyenc encoding — the document it came from dies with its read.
type LocalPreValidator ¶
LocalPreValidator is an optional interface a Handler may implement to validate a LOCAL change before it enters the DAG. The Controller calls PreValidate from the local-write path (not on inbound apply); a non-nil error rejects the whole write, keeping a malformed change out of any-sync entirely. Used for strict writer-side schema validation that returns an agent-readable error, while inbound apply stays read-tolerant.
`before` is the current value of the change's target record (nil when the change creates it). The Controller resolves it once and hands it to PreValidate, which inspects ch.Records itself. For the shared `objects` dataset every RecordChange collapses onto the object's single row, so one `before` suffices.
type LocalPreValidatorMulti ¶
type LocalPreValidatorMulti interface {
PreValidateMulti(ch *Change, get RecordGetter) error
}
LocalPreValidatorMulti is the batch-friendly variant of LocalPreValidator: instead of one pre-resolved `before`, the handler receives a getter and resolves pre-state per record — required for changes carrying N explicit-id records (e.g. recording N networkSigns in one change). When a handler implements both interfaces the Controller prefers this one.
type ObjectSeq ¶
ObjectSeq pairs an object id with its persisted max applySeq. Returned by QueryChangedObjects for the consumer-side change-index feed. Deleted is true when the row is a purged-object marker (the consumer evicts it) rather than a content change.
func QueryChangedObjects ¶
func QueryChangedObjects(ctx context.Context, coll anystore.Collection, spaceId string, since uint64, limit int) ([]ObjectSeq, error)
QueryChangedObjects returns the objects in spaceId whose persisted max applySeq exceeds `since`, ordered ascending so the caller can page by passing the last returned ApplySeq as the next `since`. limit<=0 means no cap. Backs the change-index "what changed" pull path.
Keyed on applySeq (not AddSeq) so non-DAG mutations — the account mirror's injected applies and device-local writes — surface to consumers; legacy rows are seeded applySeq := addSeq by BackfillApplySeq, keeping pre-existing cursors on one axis.
The spaceId filter naturally excludes the `space:<id>` watermark rows (no `sp` field) and per-object rows written before space-scoping.
type ObjectStamp ¶
ObjectStamp is one ObjectStamper write that landed: the shared dataset, the row (the change's ObjectId) and the derived ops applied to it with the change's VersionId.
type ObjectStamper ¶
ObjectStamper is an optional interface for the handler of a SHARED dataset (row id = ObjectId, see SharedCollections). StampObject runs once per applied synced change on any OTHER dataset of the same object — inside the apply tx, after the change's records landed, and only when the change wrote something: a record materialized, an op applied, or a tombstone was written (a fully rejected change stamps nothing, matching the handler's own BeforeModify gate). Ops queued on sink via Derive/DeriveOnce are applied to the object's row in the stamper's dataset with the triggering change's VersionId, so the row carries the "object changed" marks (modifiedAt / modifiedBy) that per-dataset handlers cannot see. Strict: an absent or tombstoned row is left untouched — creation stamps it through BeforeCreate. Sink.Project is not honored here.
SDK-internal: not part of the consumer handler API. A stamper on a dataset that is not shared never runs.
ctx carries the change only (Before is nil); the determinism contract of Handler hooks applies — the output must not depend on replica-local state. Op payloads must live on memory owned by the call (a fresh arena): they are retained past the store write, into the event dispatched to subscribers.
type Op ¶
Op is one operation inside a RecordChange.
Path is the dotted field path the op applies to. For $set and $unset, an empty Path activates the multi-field form: Payload is interpreted as an object whose keys are dot-separated field paths and whose values are applied as parallel $set/$unset ops sharing the change's versionId.
Payload is op-specific (see spec §5):
- $set: value to assign at Path, OR a multi-field object when Path is empty
- $unset: unused for single-path; multi-field form uses an object whose values are ignored
- $addToSet: the element to add to the set at Path
- $pull: the element to remove from the set at Path
- $inc: a number value, the delta
- $incGated: same as $inc
- delete: unused (nil); the record id comes from RecordChange.Id
type OpRejection ¶
type OpRejection struct {
RecordIndex int
OpIndex int // -1 for whole-record drops (BeforeCreate / BeforeDelete)
RecordId string
Err error
}
OpRejection records one per-op handler rejection. The op was validated, its prerequisites checked, and the handler returned an error — so the apply path skipped this op while letting siblings in the same RecordChange land. Surfaces through ApplyResult so the writer can tell that the change committed but with a hole.
Whole-record drops (handler.BeforeCreate / BeforeDelete returning an error) also surface here with OpIndex = -1.
type OpType ¶
type OpType string
OpType is one of the operation kinds defined in spec §5.
There is no `insert` op. Records are auto-created by the first modify that targets a non-existent id (the record starts empty with `_ver = {}` and the op then runs as normal). This eliminates the dual-insert convergence problem and the "skip if exists" footgun. See the spec divergence memory note and docs/crdt-spec.md §5 for rationale.
type ReadClassification ¶
type ReadClassification struct {
Track bool
// Tags label the unread entry and bucket the per-object counters
// (chat: "message", "mention", "reaction"). Ignored when
// Track=false.
Tags []string
// Key is the optional supersede key: a later classification with
// the same key replaces (Track=true) or clears (Track=false) the
// prior unread entry.
Key string
// Audience optionally restricts WHO counts this entry unread: the
// entry tracks only on replicas where the target record (post-
// apply) matches this typed any-store filter. On every other
// account the change applies as untracked — a Key still
// supersedes/clears. nil = everyone (except the change author, who
// is always born read). Build identity-relative conditions from
// ChangeCtx.SelfIdentity, e.g. "count a reaction only for the
// reacted-to message's author":
//
// query.Key{Path: []string{"creator"},
// Filter: query.NewCompValue(query.CompOpEq, arena.NewString(ctx.SelfIdentity))}
//
// This is the sanctioned way for a verdict to depend on record
// state the classifier cannot see (ChangeCtx.Before is nil there):
// the engine resolves it with one in-tx point read of the record
// after the change applied. A missing record tracks for nobody.
// The filter must only consult fields immutable post-create
// (derived-scope creation stamps qualify; freely-edited fields do
// not) so the verdict is replay-deterministic. Ignored when
// Track=false.
Audience query.Filter
}
ReadClassification is a ReadClassifier's verdict for one applied record change. Track=false records no unread entry, but a non-empty Key still clears a previously-tracked same-key entry (a reaction toggle: react tracks with a key, un-react clears it).
type ReadClassifier ¶
type ReadClassifier func(ctx *ChangeCtx, rec *RecordChange) ReadClassification
ReadClassifier classifies one record change for read tracking. It runs on the apply path for every change on an opted-in dataset — keep it pure and cheap (no locks). It receives the change envelope: classification runs after the record loop, so ChangeCtx.Before is ALWAYS nil here (unlike the Before* hooks) — classify from the ops and the envelope (Upsert, Creator, paths), never from prior record state. When the verdict must depend on the stored record, two tools exist: return an Audience filter and the engine resolves it against the post-apply record (it gates the WHOLE entry — every tag), or point-read the post-apply record yourself via ChangeCtx.Get + ChangeCtx.RecordId (both populated on this path) when only PART of the verdict depends on record state — e.g. adding a "mention" tag next to an unconditional "message" tag by inspecting a handler-derived field. Verdicts are device-local, so combining SelfIdentity with post-apply state is sound; incremental replays reapply from the persisted watermark, so post-apply state at re-classification matches the original run (first restore is seed-skipped). Self-authored changes are born read regardless of the verdict; the classifier still runs for them so a supersede Key can clear entries (own un-react clears the unread reaction another device tracked).
type ReadSeedMode ¶
type ReadSeedMode int
ReadSeedMode picks the initial frontier for an object that has no stored read state yet.
const ( // ReadSeedAtFirstSight (default) marks everything present at the // object's first tracked load as read — a fresh joiner starts // clean and only subsequent changes count as unread. ReadSeedAtFirstSight ReadSeedMode = iota // ReadSeedAllUnread starts with an empty frontier: the object's // whole tracked history is unread. ReadSeedAllUnread )
type ReadTracking ¶
type ReadTracking struct {
// Classify is required: the per-change verdict (track / tags /
// supersede key).
Classify ReadClassifier
// Seed picks the initial frontier for objects with no stored read
// state.
Seed ReadSeedMode
// CounterFields materializes per-tag unread counters as
// local-scope properties on the object's row in the shared
// `objects` dataset (tag → property id), e.g.
// "message" → "unreadCount". Optional.
CounterFields map[string]string
// RecordFlags materializes per-record unread booleans as
// local-scope fields on this dataset's records (tag → field id),
// e.g. "message" → "unread". A flag is true iff at least one
// unread entry with that tag references the record. The fields
// must be declared local-scope in the dataset Schema. Optional.
RecordFlags map[string]string
}
ReadTracking opts a dataset into read/unread tracking (see docs/read-tracking-proposal.md). Attached to the dataset's registration; nil = untracked.
type RecordChange ¶
RecordChange groups all ops applied to one record inside a single change.
Upsert toggles the create-on-absent behavior for the modify ops in this batch:
false (default, strict): if the target id has no record, every modify in Ops is skipped and reported as an ErrStrictSkipAbsent rejection. Use this for "update if exists" semantics — the safe default that prevents accidental record creation from typos or stale ids.
true: an empty record `{id, _ver:{}}` is created when the target id is absent, and the ops then run normally. This is the explicit "create or update" mode and is how callers create records (there is no separate insert op). The first creating change is typically a multi-field $set with Upsert=true.
Tombstones are sticky in both modes: a delete'd record never resurrects regardless of the flag.
`delete` ops ignore the flag — they always produce a tombstone (creating one on absent records too, so deletes that race ahead of creates still give "delete wins absolutely").
type RecordGetter ¶
RecordGetter resolves a record's CURRENT value by explicit id (nil when absent or the id is empty). Handed to LocalPreValidatorMulti so a batch validator can read per-record pre-state.
type SchemaHandler ¶
type SchemaHandler struct {
// contains filtered or unexported fields
}
SchemaHandler is the generic dataset handler: it enforces a schema.Dataset's behavioral declaration (required fields, write-once / author-gated mutability, apply-time stamps, id rules, delete gates, declared value shapes) with no bespoke code. A dataset registration with a declared Schema and no Handler gets one automatically.
Every verdict is arrival-order-independent, per the replica-determinism contract (see ChangeCtx): id and required checks are intrinsic to the change; write-once fields are writable ONLY in the record's creating change (never a presence probe — concurrent late fills would diverge); author gates compare the per-change Creator against the record's creator stamp, a derived-at-create immutable fact; author-gated deletes of never-created records are rejected, which converges because an author's own create→delete is self-causally ordered, so a delete arriving before its record's create can only be a non-author's.
The declaration is compiled once at construction into a flat rule table; the accept path of BeforeModify is one map lookup with zero allocations. BeforeModify writes to the sink only on acceptance (the modifyTime bump via DeriveOnce), so the controller's multi-field per-key salvage may safely re-probe it.
func NewSchemaHandler ¶
func NewSchemaHandler(ds schema.Dataset) (*SchemaHandler, error)
NewSchemaHandler compiles a dataset declaration into a generic handler. The declaration is validated first; the returned handler is read-only after construction.
func (*SchemaHandler) BeforeCreate ¶
func (h *SchemaHandler) BeforeCreate(ctx *ChangeCtx, rec *RecordChange, sink *Sink) error
BeforeCreate validates the creating change (id rule, required fields, declared value shapes) and derives the declared stamps from the change's metadata — no DB reads. An error drops the whole RecordChange.
func (*SchemaHandler) BeforeDelete ¶
func (h *SchemaHandler) BeforeDelete(ctx *ChangeCtx, _ *RecordChange, _ *Sink) error
BeforeDelete enforces the dataset's delete gate. Author-gated deletes require the record's creator stamp to match the deleting change's Creator; a never-created record (no creation marker) is rejected — convergent because an author's own create→delete is causally ordered, so a delete racing ahead of its create is never the author's.
func (*SchemaHandler) BeforeModify ¶
func (h *SchemaHandler) BeforeModify(ctx *ChangeCtx, _ *RecordChange, op *Op, sink *Sink) error
BeforeModify is the per-op gate on existing records. An error drops just this op (the controller salvages multi-field ops per key by re-probing, so combined rejections keep their innocent keys).
func (*SchemaHandler) PreValidateMulti ¶
func (h *SchemaHandler) PreValidateMulti(ch *Change, get RecordGetter) error
PreValidateMulti is the strict local gate: it re-runs the structural checks (id rules, required fields, value shapes, write-once / unknown-field rules) with agent-readable errors before the change enters the DAG. Author-identity gates are left to apply time — the local change's Creator isn't stamped yet at pre-validation, and local apply is synchronous, surfacing those rejections in the ModifyResult.
type SharedCollections ¶
type SharedCollections map[string]anystore.Collection
SharedCollections lets callers supply pre-opened collections for specific datasets, overriding the default per-object naming. Used for the per-space `objects` values collection.
type Sibling ¶
type Sibling struct {
Dataset string
Record RecordChange
}
Sibling describes a write the handler wants applied to a different dataset on the same Controller, atomically with the triggering change. The Record's `_ver` entries are stamped by the apply loop with the triggering Change's VersionId — handlers must not set them.
type Sink ¶
type Sink struct {
// contains filtered or unexported fields
}
Sink collects same-record derived ops and cross-dataset sibling writes emitted during a single RecordChange's Modify callback. Pooled at the Controller, reset between RecordChanges, drained by the apply loop.
Derived ops are folded into the same any-store Modify call (zero extra reads/writes); sibling writes execute as separate UpsertIds in the same transaction after the outer Modify returns.
func (*Sink) Derive ¶
Derive queues an op to apply on the same record, in the same Modify callback as the triggering op. Inherits the change's VersionId.
func (*Sink) DeriveOnce ¶
DeriveOnce queues op unless a derived op with the same Path is already queued. For per-change stamps (e.g. modifiedAt) emitted from the per-op BeforeModify hook, which may fire several times for one RecordChange — without the guard each op would queue a duplicate stamp, and every duplicate is projected onto the wire as a separate derived op.
func (*Sink) Project ¶
func (s *Sink) Project(dataset string, rec RecordChange)
Project queues a write to a different dataset on the same Controller. Applied in the same WriteTx as the triggering change.
type VersionId ¶
type VersionId string
VersionId is a lexicographically-sortable local ordering key maintained by any-sync's tree storage (= its `orderId`). It is **local to one peer's view of one object tree**, not a globally-consistent identity:
- Scoped to one any-sync object tree — versionIds from different trees are not comparable and must never be mixed.
- Peer-local — two peers holding the same logical DAG may encode the "same" logical change with different versionId strings. The literal value in one peer's `_ver` map is meaningless to another peer.
- Used in this package for two things: (1) gating in the CRDT apply algorithm, which is always a LOCAL comparison within one peer's Controller, and (2) as an ordered index key for queries against the local store (e.g. `_ver.id` creation order).
Peer-local convergence: each Controller applies the changes delivered by its local any-sync, gating against its own (locally-consistent) versionIds. Cross-peer convergence happens because both peers eventually process the same DAG changes and reach the same logical record content — they do NOT happen to hold the same `_ver` strings, and nothing outside this peer's Controller should compare them to another peer's versionIds.
Callers (middleware, and the clients middleware serves) DO see the `_ver` tree attached to query results and carried in subscription events — clients need per-field version info to reconcile their own optimistic in-memory state with the SDK's any-store. The tree shape documented in the spec (§3.1, §3.2) is the contract; clients walk it with the same lookup rules the SDK uses internally (see GetRecordVersion).
The empty string is the "no version" sentinel and compares less than any real version inside the same tree.
The CRDT layer treats VersionId as opaque and never invents versions — every value passed through ApplyChange comes from any-sync's local orderId for that peer.
func GetRecordVersion ¶
GetRecordVersion returns the version stored in record._ver for the given path. Returns "" when no information exists for the path.
func IsCollapsible ¶
IsCollapsible reports whether `node` is an object that can be replaced with a single string version without changing any lookup. Required: the defaultKey is present and every entry is a string holding the same version. Only a node that already claims authority over unenumerated siblings (via `*`) may become a collapsed string, which claims the same authority — collapsing sibling enumerations WITHOUT a default would invent a version for never-written fields and gate out concurrent older writes to them, diverging from peers that applied those writes before the collapse.
func NextVersion ¶
NextVersion returns a VersionId strictly greater than v in the shared lexid order. Empty v yields the smallest version. Used by the device-local write path to advance a local field past its current version (read from the record's _ver), which is correct because no synced change ever writes a Local-class field to compete with it.