anysyncx

package
v0.4.7 Latest Latest
Warning

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

Go to latest
Published: Oct 3, 2026 License: MIT Imports: 84 Imported by: 0

Documentation

Overview

Package anysyncx is the sole importer of github.com/anyproto/any-sync. All other packages in this SDK reach any-sync through this adapter, so any-sync upgrades have a single blast radius.

Wraps:

  • SpaceService — Create / Join / Derive spaces
  • ObjectTree — content-addressable DAG + AddContent
  • ocache — TTL-based live-instance cache for spaces/objects
  • ACL client — invite, accept/decline, change permissions, ownership transfer, self-remove
  • KeyValue service — per-space KV (used by techspace chat-read tracking)
  • Identity repo — encrypted user metadata

Re-exports the any-sync types the rest of the SDK needs so downstream packages never import any-sync directly.

Client wiring for any-sync's commonspace/pubsub: the Deps adapters (peers / crypto / membership), the RPC-registration component, and the App-level pass-throughs the space layer calls. The engine itself (pubsub.New) is registered as a regular app component in app.go; the adapters hold a late-bound *App because their methods only run after app.Start.

Index

Constants

This section is empty.

Variables

View Source
var ErrPubSubGuestSpace = errors.New("anysyncx: pubsub unavailable in guest-mode spaces")

ErrPubSubGuestSpace rejects pubsub on guest-mode (public access) spaces: the engine signs and handshakes as the real account, which a guest space's ACL does not contain, so peers would silently reject every publish and subscribe as non-member. Fail fast instead.

View Source
var ErrPubSubNoKey = errors.New("anysyncx: no pubsub key")

ErrPubSubNoKey is returned by the pubsub crypto adapter when the local identity has no read access to the space (keyless reader) or the referenced key-record id is unknown. Publishing without a read key fails rather than silently sending plaintext into an encrypted space.

View Source
var ErrSpaceRegistryUnset = errors.New("anysyncx: space registry not set")

ErrSpaceRegistryUnset is returned when GetTree fires before the space registry has been wired up. Indicates init order bug, not a runtime miss.

View Source
var ErrTreeTypeSkipped = errors.New("anysyncx: tree type not selected for sync")

ErrTreeTypeSkipped is a tree fetch declined by selective sync before any tree-storage write. The registry records a heads-only stub in its place (or finds the entry already deleted), so the diff converges with nothing to materialize and the syncer never parks it.

Functions

This section is empty.

Types

type App

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

App is the running any-sync app plus the components downstream packages need to reach. Constructed by New, torn down by Close.

func New

func New(ctx context.Context, cfg config.Config, provider auth.Provider) (*App, error)

New brings up the any-sync app. Order matters: keys first (provider may block on user input), then nodeconf parsing, then component registration mirroring the legacy bootstrap order.

Logger setup is the caller's job — call (logger.Config{}).ApplyGlobal once before Open. We intentionally don't touch the global logger here because ApplyGlobal mutates named-logger structs in place, and a second Open in the same process would race goroutines from the first SDK that are already using those loggers.

func (*App) AccountKeys

func (a *App) AccountKeys() *accountdata.AccountKeys

AccountKeys holds the decoded peer/sign keys.

func (*App) AclSnapshot

func (a *App) AclSnapshot(ctx context.Context, spaceId string) (list.AclList, error)

AclSnapshot builds an in-memory, verified ACL for spaceId from the records the nodes serve — the read the joining client and the ACL waiter each make on their own. It needs no local storage and materializes nothing, so it is safe on a space this account is not a member of: the join controller reads membership / the pending request off it for a row whose request was posted elsewhere.

func (*App) BroadcastP2P

func (a *App) BroadcastP2P()

BroadcastP2P re-runs the LAN handshake with every known local peer — called when the tech-space index changes (new row synced, a joining row flipped active) so a fresh device starts probing for the space right away instead of on the next discovery resweep. The index change is also a "key may be derivable now" signal, so the discovery-key negative cache is reset first — without that, a fresh join's failed request-time derivation would suppress the space from the handshake for up to negativeRetryAfter even though its ACL has synced in.

func (*App) ClaimSpaceStorage

func (a *App) ClaimSpaceStorage(ctx context.Context, spaceId string) (SpaceStorageClaim, error)

ClaimSpaceStorage claims a space's any-sync store for deletion. The offload path claims before EvictSpace, so no load can reopen the space while it is torn down.

func (*App) Close

func (a *App) Close(ctx context.Context) error

Close stops the any-sync app and releases resources. The space cache is closed first so per-space goroutines wind down before the app's components; the space stores close last.

func (*App) Coordinator

func (a *App) Coordinator() coordinatorclient.CoordinatorClient

Coordinator client — used for SpaceMakeShareable, NetworkConfiguration, space delete confirmations, account delete.

func (*App) DRPCServer

func (a *App) DRPCServer() server.DRPCServer

DRPCServer is the inbound DRPC mux. The files p2p server registers its read-only FileP2P handler here after the file store is built.

func (*App) EvictSpace

func (a *App) EvictSpace(ctx context.Context, id string) error

EvictSpace forces immediate close + remove from cache. Used by space-delete paths to drop the in-memory state right away rather than waiting for TTL.

func (*App) FileNetworkId

func (a *App) FileNetworkId() string

FileNetworkId is the identity of the fileV2 fleet's shared receipt-signing key (nodeconf fileNetworkId). Custody receipts (networkSign) verify against this one stable key; empty on networks without a fileV2 fleet.

func (*App) FileV2Peers

func (a *App) FileV2Peers() []string

FileV2Peers is the current fileV2 fleet (nodeconf.NodeTypeFileV2). The files broker routes its RPCs across these peers.

func (*App) GetSpace

func (a *App) GetSpace(ctx context.Context, id string) (SpaceHandle, error)

GetSpace returns a handle for the given space id, loading it through the cache if necessary. The handle is valid for the duration of the operation; the cache TTL evicts idle spaces in the background.

func (*App) GlobalP2PEnabled

func (a *App) GlobalP2PEnabled() bool

GlobalP2PEnabled reports whether the internet-wide layer is on (cfg.P2P.Global).

func (*App) HeadCache

func (a *App) HeadCache() *HeadCache

HeadCache exposes the per-space hash cache. Useful for tests and the eventual SyncStatus integration; the spaceSyncHandler already has a direct reference for its fast path.

func (*App) Headless

func (a *App) Headless() bool

Headless reports whether the SDK runs in embedded-backend mode (cfg.Headless). See config.Config.Headless for the contract.

func (*App) InboxClient

func (a *App) InboxClient() inboxclient.InboxClient

InboxClient is the coordinator inbox transport (fetch / add-message). Always non-nil in this build (the component is registered unconditionally), but callers should treat a nil return as "inbox unavailable" so a future coordinator-less deployment degrades to the out-of-band 1-1 path.

func (*App) IsLocalOnly

func (a *App) IsLocalOnly(spaceId string) bool

IsLocalOnly reports whether spaceId is pinned to this device. The p2p-facing sync handlers use it to refuse serving or advertising local-only spaces even if a peer names one directly.

func (*App) JoiningClient

func (a *App) JoiningClient() aclclient.AclJoiningClient

JoiningClient is the account-level ACL client used to send RequestJoin / CancelJoin RPCs to the coordinator+nodes for a space the caller is not yet a member of.

func (*App) KeyValueStore

func (a *App) KeyValueStore(ctx context.Context, spaceId string) (keyvaluestorage.Storage, error)

KeyValueStore returns a space's default key-value store (loads the space if needed).

func (*App) LocalDiscoveryEnabled

func (a *App) LocalDiscoveryEnabled() bool

LocalDiscoveryEnabled is the mDNS switch state.

func (*App) LocalPeerHasSpace

func (a *App) LocalPeerHasSpace(peerId, spaceId string) bool

LocalPeerHasSpace reports whether a LAN peer advertised sharing spaceId in the exchange. The SpacePush handler uses it to bound which spaces a peer may seed onto this device.

func (*App) MarkSpaceLocalOnly

func (a *App) MarkSpaceLocalOnly(spaceId string)

MarkSpaceLocalOnly pins spaceId to this device: its peer manager resolves no peers (no node subscribe, no diff-sync, no push) and the credential provider refuses to request a coordinator receipt for it. Must be called before the space is first loaded — the peer manager is chosen at NewSpace time. Used by the tech space in headless mode.

func (*App) NetworkId

func (a *App) NetworkId() string

NetworkId is the id of the any-sync network this app is bound to. Needed to build the signed space-delete confirmation, which the coordinator verifies against its own network id.

func (*App) NewAclWaiter

func (a *App) NewAclWaiter(spaceId, aclHeadId string, onFinish, onReject func(list.AclList) error) (aclwaiter.AclWaiter, error)

NewAclWaiter builds an any-sync ACL waiter bound to the running app: a background poller over the network ACL (no local space storage needed) that fires onFinish once the account gains permissions on spaceId, or onReject once the join request is declined at/after aclHeadId. The caller drives its lifecycle via Run(ctx) / Close(ctx). Used by the joiner-side post-acceptance loader.

func (*App) OnInboxMessage

func (a *App) OnInboxMessage(h InboxMessageHandler)

OnInboxMessage installs the push callback invoked when the coordinator signals a new inbox message. Replaces any previous handler; pass nil to detach (the notifier does this on Close). The forwarder registered before app.Start reads this atomically.

func (*App) OnKeyValues

func (a *App) OnKeyValues(spaceId string, h KeyValueHandler) (cancel func())

OnKeyValues registers a handler for a space's applied key-value writes. The returned cancel is idempotent. Writes applied while no handler is registered are not replayed — consumers reconcile from the store on startup (Iterate / GetAll).

func (*App) P2PEnabled

func (a *App) P2PEnabled() bool

P2PEnabled reports whether the local-network layer is on (cfg.P2P).

func (*App) P2PStatus

func (a *App) P2PStatus() sdkp2p.Status

P2PStatus is the account-wide p2p snapshot for the debug surface: LAN listener state, discovery possibility, every LAN peer with its live-connection flag, and the global layer.

func (*App) ParkedTreeCount

func (a *App) ParkedTreeCount(spaceId string) int

ParkedTreeCount returns the number of trees parked for retry in spaceId's treesyncer adapter — fetched to storage but never materialized into the CRDT projection (see treeSyncerAdapter.pending). 0 when the space was never loaded this session. Consumed by the SDK's close-time watermark gate: a space with parked trees must keep its boot replay, so its watermark is not snapshotted at Close.

func (*App) PeerStore

func (a *App) PeerStore() *p2p.PeerStore

PeerStore exposes the p2p local-peer registry (which LAN peers share which spaces). Used by the files p2p source for peer selection.

func (*App) PeerSyncStats

func (a *App) PeerSyncStats(spaceId string) []PeerSyncSnapshot

PeerSyncStats returns the latest per-peer SyncAll snapshots for spaceId. Empty slice if the space has never had an outbound diff round (or was never loaded). In-memory only — cleared on SDK restart. Used by the debug API; not a stable surface.

func (*App) PickSpace

func (a *App) PickSpace(ctx context.Context, id string) (SpaceHandle, bool)

PickSpace returns a handle only if the space is already resident in the cache; it never triggers a load (ocache.Pick does not call the LoadFunc). ok is false when the space is not loaded. Cheap read paths such as Spaces().List must use this rather than GetSpace: a commonspace.Init can block indefinitely when the responsible sync-node is unreachable, and a per-row List load would hang the whole call.

func (*App) Pool

func (a *App) Pool() pool.Pool

Pool exposes the any-sync peer pool (dial by peerId, addresses resolved from the nodeconf). Lets embedders/e2e speak node-side protocols (e.g. fileprotov2) over a connection that carries this account's identity in the handshake.

func (*App) PubSubCloseSpace

func (a *App) PubSubCloseSpace(spaceId string)

PubSubCloseSpace drops every local subscription and remote interest for spaceId. Called from the space layer's deliberate teardown funnel (evict/offload) — deliberately NOT from cache eviction, which must not kill live subscriptions.

func (*App) PubSubPublish

func (a *App) PubSubPublish(ctx context.Context, spaceId, topic string, payload []byte) error

PubSubPublish signs, encrypts and fans payload out on spaceId/topic. Fire-and-forget past local validation; own publishes deliver to local subscribers synchronously. The network send is detached from the caller's cancelation (WithoutCancel): the pool runs the send closure after Publish returns, and the public contract promises delivery survives the caller's ctx.

func (*App) PubSubRevalidate

func (a *App) PubSubRevalidate(spaceId string)

PubSubRevalidate re-checks the engine's per-space membership state against the current ACL, dropping serving-side interest of removed members. Called on ACL record apply; a no-op for non-resident spaces.

func (*App) PubSubSubscribe

func (a *App) PubSubSubscribe(spaceId, pattern string, h pubsub.Handler) (cancel func(), err error)

PubSubSubscribe registers h for topics matching pattern in spaceId and pushes the interest to the space's peers. The returned cancel is idempotent.

func (*App) RepublishGlobalRecord

func (a *App) RepublishGlobalRecord(spaceId string)

RepublishGlobalRecord re-sets this device's global p2p row in a space (advertising switched on).

func (*App) SelectiveTreeTypes

func (a *App) SelectiveTreeTypes() []string

SelectiveTreeTypes is the selective-sync tree-type allowlist (cfg.Sync.TreeTypes). Empty = sync and materialize everything.

func (*App) SetAdvertiseFn

func (a *App) SetAdvertiseFn(fn func(spaceId string) bool)

SetAdvertiseFn wires the per-space p2p advertising decision (the tech-space row's switch; the tech space itself never gets a row — own devices come from the account record). Set before the tech space loads; loaded spaces are re-evaluated at once.

func (*App) SetFileStore

func (a *App) SetFileStore(st *filestore.Store)

SetFileStore injects the file store into the p2p file server, enabling it to serve stored CAR objects to LAN peers. Called by sdk.Open once the store is built. No-op when p2p is disabled.

func (*App) SetGuestKeyFn

func (a *App) SetGuestKeyFn(fn func(spaceId string) crypto.PrivKey)

SetGuestKeyFn wires the guest-identity resolver for guest-mode (public-access) spaces: fn returns the shared guest identity's private key for spaceId, or nil for regular spaces. Called once by sdk.Open after the tech space is up (the resolver reads the tech-space index). Guest spaces load with an account-service override so the commonspace signs as the guest identity — see loadSpaceForCache.

func (*App) SetKnownSpaceIdsFn

func (a *App) SetKnownSpaceIdsFn(fn func() []string)

SetKnownSpaceIdsFn wires the p2p exchange's probe source: the space ids this account knows of (tech-space index), regardless of whether they are stored locally. Local-only spaces are filtered out here — they must never reach the exchange in any form. Called once by sdk.Open after the tech space is up.

func (*App) SetLocalDiscoveryEnabled

func (a *App) SetLocalDiscoveryEnabled(enabled bool)

SetLocalDiscoveryEnabled switches mDNS announce and browse at runtime.

func (*App) SetPeerAddrs

func (a *App) SetPeerAddrs(peerId string, addrs []string)

SetPeerAddrs registers dial addresses for a peer that is NOT in the nodeconf — a direct out-of-band peer like the push-notification node (config.Push). After registration Pool().Get(peerId) dials it over the same secure transports as any node. Calling again replaces the address list; there is no removal (the entry is process-lifetime).

func (*App) SetSpaceRegistry

func (a *App) SetSpaceRegistry(r SpaceRegistry)

SetSpaceRegistry wires the tree manager to a space-level registry. Called once by the space package after it builds its ocache.

func (*App) SpaceDelete

func (a *App) SpaceDelete(ctx context.Context, spaceId string) error

SpaceDelete tells the coordinator to delete spaceId network-wide. It builds the signed deletion confirmation from the account sign key, peerId, and networkId — the coordinator verifies the signature and the embedded identity before moving the space to PendingDeletion. Only the space owner can delete; the coordinator rejects others.

Idempotent from the caller's view: re-sending for a space already pending/deleted is the reconciler's normal retry, and the coordinator returns the appropriate error which the caller can treat as done.

func (*App) SpaceExists

func (a *App) SpaceExists(spaceId string) bool

SpaceExists reports whether any-sync has local storage for spaceId. Used by the space layer to decide between Open and Create paths.

func (*App) SpaceService

func (a *App) SpaceService() commonspace.SpaceService

SpaceService is any-sync's per-account space create/derive/join surface.

func (*App) SpaceStatuses

func (a *App) SpaceStatuses(ctx context.Context, spaceIds []string) ([]*coordinatorproto.SpaceStatusPayload, error)

SpaceStatuses fetches the coordinator's view of the given spaces in one round trip — status (Created / PendingDeletion / Deleted / …) and our permissions (Owner / …) per space. The returned slice is aligned with spaceIds. Used by the deletion reconciler to decide whether we still owe a SpaceDelete (Owner + Created) or must offload a space the network reports gone.

func (*App) StreamPool

func (a *App) StreamPool() streampool.StreamPool

StreamPool is the outbound DRPC stream pool. The space layer uses it to send SpaceSubscription messages when a new space loads.

func (*App) SyncHeads

func (a *App) SyncHeads(ctx context.Context, id string) error

SyncHeads loads the space (through the cache) and forces an immediate head-sync (diff) round on it. Used to converge on demand instead of waiting for the periodic headsync timer.

func (*App) SyncStatus

func (a *App) SyncStatus() *syncstatus.Service

SyncStatus exposes the per-account sync-status registry. Its Trackers feed commonspace.Deps.SyncStatus; the space layer reads it for snapshot + subscribe.

type HeadCache

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

HeadCache keeps the current per-space space-hash in memory so the HeadSync RPC can answer the common full-range probe without loading the underlying commonspace.Space.

Only the new (DiffType_V3) hash is tracked — V2/V1 are legacy. A V2 request that reaches us falls through to the deep-diff path where any-sync's headsync subsystem still computes the right answer; we just don't accelerate it.

The cache is populated by a statestorage.Observer wired during space load, and survives space TTL eviction. The hash is opaque hex from any-sync — we never compute or compare it ourselves.

func (*HeadCache) Delete

func (c *HeadCache) Delete(spaceId string)

Delete drops the entry — used when a space is deleted locally.

func (*HeadCache) Get

func (c *HeadCache) Get(spaceId string) (string, bool)

Get returns the cached hash for spaceId. The bool is false when the cache has nothing — callers fall back to the deep diff.

func (*HeadCache) Set

func (c *HeadCache) Set(spaceId, hash string)

Set replaces the cached hash for spaceId.

type InboxMessageHandler

type InboxMessageHandler = func(*coordinatorproto.NotifySubscribeEvent)

InboxMessageHandler is the push callback the inbox notifier installs. The body-less NotifySubscribeEvent only signals "new mail" — the notifier's handler reacts by fetching.

type KeyValueHandler

type KeyValueHandler func(decryptor keyvaluestorage.Decryptor, kvs []innerstorage.KeyValue)

KeyValueHandler receives applied key-value writes (local Set and remote SetRaw) for one space. Runs synchronously on any-sync's apply/broadcast path — hand real work off to your own goroutine.

type PeerRoundCallback

type PeerRoundCallback func(peerId string, newCount, changedCount int, err error)

treeSyncerAdapter is registered by anysyncx as the per-space commonspace.Deps.TreeSyncer. On SyncAll it builds (fetches) missing trees and pushes existing ones via SyncWithPeer. Trees are not closed here — they stay in ocache for the async exchange to land.

SyncAll routes through SpaceRegistry.GetTree (NOT through the per-space TreeBuilder directly): the registry's per-space implementation is the only place that wires our SDK-side UpdateListener onto the tree, so that inbound changes fan out to the CRDT controller. A bare TreeBuilder.BuildTree call sets no listener, which silently breaks cold sync — the tree storage receives changes but our controller never sees them. See docs/sync-listener-wiring or the cold-sync e2e for the trail. PeerRoundCallback fires after every SyncAll completes. Wired by the App to route 0/0 success rounds into the syncstatus tracker for the space-level "all synced" sweep. peerId is the responsible peer (caller filters); newCount / changedCount are the post- deletionState-filter sizes (what SyncAll actually acted on); err is the SyncAll error, nil on success.

type PeerSyncSnapshot

type PeerSyncSnapshot struct {
	PeerId     string
	LastSyncAt time.Time
	New        int
	Changed    int
	// Pending is the space-wide parked-tree count at the time of this
	// round (see treeSyncerAdapter.pending) — nonzero means some trees'
	// GetTree failed and is awaiting retry. Same value on every peer's
	// row of the same round; kept per-row so the debug surface needs no
	// second query.
	Pending int
	LastErr string
}

PeerSyncSnapshot is the latest per-peer headsync result the adapter has observed. Returned by Stats() — debug surface only.

type SpaceHandle

type SpaceHandle interface {
	Id() string
	Inner() commonspace.Space
	// SyncHeads forces an immediate head-sync (diff) round on this
	// space, rather than waiting for the periodic timer. Blocks until
	// the round completes and returns its error verbatim.
	SyncHeads(ctx context.Context) error
}

SpaceHandle is the handle returned by App.GetSpace. Callers operate on it for the duration of a logical operation; the underlying commonspace.Space is reachable via Inner. The cache TTL releases idle spaces in the background — callers don't refcount manually.

type SpaceRegistry

type SpaceRegistry interface {
	// GetTree returns the live ObjectTree for (spaceId, treeId), loading
	// the space and the object via ocache as needed. Errors propagate
	// from any-sync directly.
	GetTree(ctx context.Context, spaceId, treeId string) (objecttree.ObjectTree, error)

	// HasTree reports whether the tree already exists in local storage.
	HasTree(ctx context.Context, spaceId, treeId string) (bool, error)

	// PutTree stores a remote-built tree payload — used by sync when a
	// node delivers a tree we don't have yet.
	PutTree(ctx context.Context, spaceId string, payload treestorage.TreeStorageCreatePayload) error

	// MarkTreeDeleted is the soft-delete hook fired when the settings
	// tree announces a deletion. Implementations typically clean up
	// cached state; the actual storage delete happens via DeleteTree.
	MarkTreeDeleted(ctx context.Context, spaceId, treeId string) error

	// DeleteTree removes a tree's local state.
	DeleteTree(ctx context.Context, spaceId, treeId string) error

	// ShouldPullTree decides whether a locally-missing tree announced by
	// a head update should be fetched (selective sync by tree type —
	// SYN-18). root is the tree's raw root change carried by the update;
	// heads are the sender's current heads. Implementations returning
	// false are expected to record the heads so the sync diff converges
	// without the tree's change bodies. Full-sync deployments always
	// return true.
	ShouldPullTree(ctx context.Context, spaceId, treeId string, root *treechangeproto.RawTreeChangeWithId, heads []string) bool
}

SpaceRegistry is the contract treeManager uses to resolve trees. Set once by the space layer after it builds its ocache. We keep this indirection so anysyncx doesn't depend on the higher-level space package — matches the pattern from anytype-heart's treemanager.

type SpaceStorageClaim

type SpaceStorageClaim interface {
	// Delete closes the store and removes `<DataDir>/anysync/<spaceId>.db`.
	Delete(ctx context.Context) error
	// Release gives the store back without deleting it.
	Release()
}

SpaceStorageClaim holds a space's any-sync store for deletion: nothing can open it until Delete removes it or Release gives it back.

type TreeFetchedCallback

type TreeFetchedCallback func(peerId, treeId string, heads []string)

TreeFetchedCallback fires after a tree missing locally was fetched from a peer during SyncAll, with the heads it now holds: the tree is in sync with that peer by construction.

Jump to

Keyboard shortcuts

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