Documentation
¶
Overview ¶
Package readstate owns the per-space read/unread engine: the unread entries, the per-object seen-heads frontier, per-tag counters, and the transitions log backing the read-state feed. See docs/read-tracking-proposal.md.
The engine persists three space-shared collections (one set per DB, like _meta):
- _read_unread — one row per unread change, deleted when read.
- _read_state — one row per tracked object: frontier heads, pending remote heads, per-tag counters, and the object's last stateSeq — the dirty-object feed (ChangedSince) reads THIS table; there is no transitions log. A consumer re-pulls the object's unread snapshot and diffs against what it holds.
All mutating methods take the caller's tx context: the apply path calls Track inside the same WriteTx as the record mutations, so an unread entry is atomic with its change; marks and merges open their own tx at the service layer. The engine holds no locks and no in-memory state — callers serialize per object (the apply path via the tree lock, merges/marks via the service's per-object mutex).
Read marking is forward-only: MarkRead covers the given changes and their causal ancestry. The ancestor walk runs over unread rows (each row carries its change's prevIds) and traverses gaps — read, self-authored, or untracked changes between unread ones — through the Resolver, pruning where a walked change's versionId drops below the object's minimum unread versionId (an ancestor's versionId is always smaller than its descendant's, so nothing unread can hide below that line). MarkReadUpTo never walks: range coverage by versionId is gap-immune.
Index ¶
- Constants
- type Engine
- func (e *Engine) ChangedSince(ctx context.Context, since uint64, limit int) ([]ObjectState, error)
- func (e *Engine) ClearRecords(ctx context.Context, objectId string, recordIds []string, stateSeq uint64) error
- func (e *Engine) Counts(ctx context.Context, objectId string) (map[string]int, error)
- func (e *Engine) Frontier(ctx context.Context, objectId string) (heads, pending []string, err error)
- func (e *Engine) MarkRead(ctx context.Context, objectId string, changeIds []string) (MarkResult, error)
- func (e *Engine) MarkReadUpTo(ctx context.Context, objectId, upTo string) (MarkResult, error)
- func (e *Engine) MarkReadUpToChunk(ctx context.Context, objectId, upTo string, maxEntries int) (MarkResult, bool, error)
- func (e *Engine) MarkSeeded(ctx context.Context, objectId string) error
- func (e *Engine) MergeHeads(ctx context.Context, objectId string, heads []string) (MarkResult, error)
- func (e *Engine) NotifyState(objectId string, stateSeq uint64)
- func (e *Engine) SeedFrontier(ctx context.Context, objectId string, heads []string) (bool, error)
- func (e *Engine) Seeded(ctx context.Context, objectId string) (bool, error)
- func (e *Engine) SubscribeState(cb func(objectId string, stateSeq uint64)) (cancel func())
- func (e *Engine) TrackChange(ctx context.Context, t Track) error
- func (e *Engine) UnreadEntries(ctx context.Context, objectId string) ([]Entry, uint64, error)
- func (e *Engine) WriteTx(ctx context.Context, fn func(txCtx context.Context) error) error
- type Entry
- type MarkResult
- type ObjectState
- type Resolver
- type SeqFunc
- type Track
Constants ¶
const ( UnreadCollectionName = "_read_unread" StateCollectionName = "_read_state" )
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Engine ¶
type Engine struct {
// contains filtered or unexported fields
}
func (*Engine) ChangedSince ¶
ChangedSince returns objects whose read state advanced past since, ascending by stateSeq, capped at limit (0 = no cap). One element per object (this reads the per-object state rows, not a log): the consumer re-pulls UnreadEntries for each dirty object and diffs against what it holds. Page by passing the last StateSeq back.
func (*Engine) ClearRecords ¶
func (e *Engine) ClearRecords(ctx context.Context, objectId string, recordIds []string, stateSeq uint64) error
ClearRecords drops deleted records from unread entries: an entry referencing only deleted records is removed (counters and feed updated); an entry also touching surviving records just sheds the deleted ids.
func (*Engine) Frontier ¶
func (e *Engine) Frontier(ctx context.Context, objectId string) (heads, pending []string, err error)
Frontier returns the object's current seen-heads frontier and pending remote heads.
func (*Engine) MarkRead ¶
func (e *Engine) MarkRead(ctx context.Context, objectId string, changeIds []string) (MarkResult, error)
MarkRead covers the given change ids and their causal ancestry. Ids not yet arrived on this device are persisted as pending heads. Runs in the caller's tx; the caller serializes per object.
func (*Engine) MarkReadUpTo ¶
MarkReadUpTo covers every unread change with versionId <= upTo ("" = everything). Range coverage needs no ancestry walk.
func (*Engine) MarkReadUpToChunk ¶
func (e *Engine) MarkReadUpToChunk(ctx context.Context, objectId, upTo string, maxEntries int) (MarkResult, bool, error)
MarkReadUpToChunk is MarkReadUpTo bounded to at most maxEntries covered entries (0 = unbounded); done=false means more remain below upTo. Covering a versionId prefix is a valid frontier advance (forward-only marking), so a caller can commit each chunk in its own tx and a crash mid-way just resumes — the huge-ReadAll path uses this to keep transactions bounded.
func (*Engine) MarkSeeded ¶
MarkSeeded records first-sight seeding as done WITHOUT touching entries or frontier — the consult-KV seed path uses it: the account's published frontiers, merged separately in the same tx, are the real seed. Idempotent.
func (*Engine) MergeHeads ¶
func (e *Engine) MergeHeads(ctx context.Context, objectId string, heads []string) (MarkResult, error)
MergeHeads applies another device's published frontier. Identical semantics to MarkRead: known ids cover their ancestry, unknown ids park as pending heads.
func (*Engine) NotifyState ¶
NotifyState fires the subscribed pings. Call AFTER the tx that produced stateSeq committed.
func (*Engine) SeedFrontier ¶
SeedFrontier performs first-sight seeding: sets the frontier to the given heads (everything at or behind them is read), drops any unread entries that slipped in around the restore (their became- read transitions surface on the feed), preserves pending remote heads, and records the seed durably — all in the caller's tx. Idempotent: a second call on a seeded object is a no-op. Returns whether the seed ran.
func (*Engine) Seeded ¶
Seeded reports whether first-sight seeding already ran for the object. Durable: only SeedFrontier sets it, in the same tx as the seed itself, so the answer survives crashes on either side.
func (*Engine) SubscribeState ¶
SubscribeState registers a best-effort ping fired after a committed read-state change (new unread, mark, merge). cb runs synchronously on the notifying path — keep it small. Cancel is idempotent. A dropped ping is recovered by pulling TransitionsSince from the consumer's cursor.
func (*Engine) TrackChange ¶
Track records one classified change from the apply path. Must run inside the apply's WriteTx (pass tx.Context()).
func (*Engine) UnreadEntries ¶
UnreadEntries returns the object's current unread set, ascending by versionId ("" upTo = all).
type Entry ¶
type Entry struct {
ObjectId string
Dataset string
ChangeId string
VersionId string
AddSeq uint64
ApplySeq uint64
RecordIds []string
Tags []string
Key string
PrevIds []string
StateSeq uint64
}
Entry is one unread change as persisted.
type MarkResult ¶
type MarkResult struct {
Removed []Entry
Frontier []string
// Pending are marked ids not yet arrived on this device; they are
// persisted and resolved when the change applies.
Pending []string
StateSeq uint64 // 0 when nothing changed
}
MarkResult reports what a mark/merge covered.
type ObjectState ¶
ObjectState is one dirty-object feed element: an object whose read state changed, at the stateSeq that change advanced it to.
type Resolver ¶
type Resolver func(ctx context.Context, objectId, changeId string) (prevIds []string, versionId string, ok bool, err error)
Resolver looks up a change the unread rows no longer (or never did) cover — the gap-traversal fallback backed by any-sync's change storage. ok=false means the change has not arrived on this device.
type SeqFunc ¶
SeqFunc allocates the next per-space stateSeq. Wired to the space's ApplySeqAllocator so marks, merges, and applies share one monotonic cursor axis. Must be called with the tx held (allocation order = commit order).
type Track ¶
type Track struct {
ObjectId string
Dataset string
ChangeId string
VersionId string
AddSeq uint64
ApplySeq uint64 // stateSeq for the unread transition
RecordIds []string
Tags []string
Key string // supersede key; "" = none
PrevIds []string
// Tracked=false records no unread entry but still clears a
// superseded row (Key) and resolves a pending head.
Tracked bool
// SelfAuthored changes are born read and advance the frontier
// in place when their whole causal past is covered.
SelfAuthored bool
}
Track is the apply-path input: one tracked change, classified.