readstate

package
v0.4.3 Latest Latest
Warning

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

Go to latest
Published: Sep 23, 2026 License: MIT Imports: 9 Imported by: 0

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

View Source
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 New

func New(db anystore.DB, spaceId string, seq SeqFunc, resolve Resolver) *Engine

func (*Engine) ChangedSince

func (e *Engine) ChangedSince(ctx context.Context, since uint64, limit int) ([]ObjectState, error)

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) Counts

func (e *Engine) Counts(ctx context.Context, objectId string) (map[string]int, error)

Counts returns the per-tag unread counters for an object.

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

func (e *Engine) MarkReadUpTo(ctx context.Context, objectId, upTo string) (MarkResult, error)

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

func (e *Engine) MarkSeeded(ctx context.Context, objectId string) error

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

func (e *Engine) NotifyState(objectId string, stateSeq uint64)

NotifyState fires the subscribed pings. Call AFTER the tx that produced stateSeq committed.

func (*Engine) SeedFrontier

func (e *Engine) SeedFrontier(ctx context.Context, objectId string, heads []string) (bool, error)

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

func (e *Engine) Seeded(ctx context.Context, objectId string) (bool, error)

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

func (e *Engine) SubscribeState(cb func(objectId string, stateSeq uint64)) (cancel func())

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

func (e *Engine) TrackChange(ctx context.Context, t Track) error

Track records one classified change from the apply path. Must run inside the apply's WriteTx (pass tx.Context()).

func (*Engine) UnreadEntries

func (e *Engine) UnreadEntries(ctx context.Context, objectId string) ([]Entry, uint64, error)

UnreadEntries returns the object's current unread set, ascending by versionId ("" upTo = all).

func (*Engine) WriteTx

func (e *Engine) WriteTx(ctx context.Context, fn func(txCtx context.Context) error) error

WriteTx runs fn inside a write transaction on the engine's DB — the tx opener for callers outside the apply path (marks, merges, reconcile). fn's error rolls the tx back and is returned.

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

type ObjectState struct {
	ObjectId string
	StateSeq uint64
}

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

type SeqFunc func(ctx context.Context) (uint64, error)

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.

Jump to

Keyboard shortcuts

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