p2p

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: 38 Imported by: 0

Documentation

Index

Constants

View Source
const (

	// DefaultServiceType is the DNS-SD service type SDK peers announce.
	// Deliberately NOT anytype-heart's "_anytype._tcp" — the protocols
	// are separate networks.
	DefaultServiceType = "_any._tcp"
)
View Source
const PickTimeout = 50 * time.Millisecond

PickTimeout bounds a pool.Pick on a global peer. The pool's Pick waits on an in-flight load for the same id (a LAN+global peer can have a LAN dial in flight), so an unbounded ctx would turn "never dials" into "waits for somebody else's dial".

View Source
const RecordKey = "p2p/iroh"

RecordKey is the key-value key of a device's global p2p record. The store keeps one row per (key, peer), so every device of a space owns exactly one row: its endpoint ticket, re-set as a heartbeat.

Variables

This section is empty.

Functions

func CurrentOwnAddresses

func CurrentOwnAddresses(port int) sdkp2p.OwnAddresses

CurrentOwnAddresses is this device's announce right now: LAN IPv4s plus the given listen port. Used by the exchange's proactive re-handshake.

func PickLive

func PickLive(ctx context.Context, p Picker, id string) (peer.Peer, error)

PickLive returns the live pool connection to a peer, or an error within PickTimeout. It never dials.

Types

type AddrBook

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

AddrBook is the only writer of peer addresses into any-sync's peer service. A peer is registered with either its LAN addresses (from the space exchange) or its iroh ticket (from key-value records), never both: a LAN dial that fails must answer in one RTT, not fall through into a relay dial that can take the whole dial timeout. LAN wins while present; the ticket takes over once the LAN entry is cleared (mDNS lost, dial strikes). The one address the book does not own is the push node's (config.Push), registered by the SDK directly: it is neither a LAN nor a global peer and never enters the peer store.

func NewAddrBook

func NewAddrBook() *AddrBook

func (*AddrBook) ClearLAN

func (b *AddrBook) ClearLAN(peerId string)

ClearLAN forgets a peer's LAN addresses; its ticket, if any, takes over.

func (*AddrBook) ClearTicket

func (b *AddrBook) ClearTicket(peerId string)

ClearTicket forgets a peer's ticket.

func (*AddrBook) HasLAN

func (b *AddrBook) HasLAN(peerId string) bool

HasLAN reports whether the peer currently has LAN addresses.

func (*AddrBook) Init

func (b *AddrBook) Init(a *app.App) error

func (*AddrBook) Name

func (b *AddrBook) Name() string

func (*AddrBook) SetLAN

func (b *AddrBook) SetLAN(peerId string, addrs []string)

SetLAN registers a peer's LAN addresses (already scheme-prefixed). Empty clears them.

func (*AddrBook) SetTicket

func (b *AddrBook) SetTicket(peerId, ticket string)

SetTicket registers a peer's iroh ticket. Empty clears it.

func (*AddrBook) Ticket

func (b *AddrBook) Ticket(peerId string) string

Ticket returns a peer's registered ticket, empty when none.

type Discovery

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

Discovery announces this device on the LAN and browses for other SDK peers, feeding each sighting to the Notifier (the SpaceExchangeV2 handshake).

Concurrency model, deliberately simpler than anytype-heart's:

  • one supervisor goroutine runs announce+browse "sessions", restarting them (scoped child context) when the interface set changes or the driver dies — never tearing down the component;
  • driver callbacks only enqueue onto a bounded channel;
  • one consumer goroutine owns the known-peer map and performs all notifier calls serially.

Nothing is ever reassigned after Run; Close cancels one context and waits (bounded) on one WaitGroup.

func NewDiscovery

func NewDiscovery(cfg config.P2P, peerId string, portFn func() (int, bool), notifier Notifier) *Discovery

func (*Discovery) Close

func (d *Discovery) Close(_ context.Context) error

func (*Discovery) Enabled

func (d *Discovery) Enabled() bool

Enabled is the local-discovery switch state: config p2p.localDiscovery at start, SetEnabled afterwards. Off as well while p2p is disabled in config, since discovery never runs then whatever the switch says.

func (*Discovery) Init

func (d *Discovery) Init(_ *app.App) error

func (*Discovery) Name

func (d *Discovery) Name() string

func (*Discovery) Port

func (d *Discovery) Port() int

Port is the announced QUIC listen port (zero before Run).

func (*Discovery) Possibility

func (d *Discovery) Possibility() sdkp2p.Possibility

Possibility is the current discovery-possibility state.

func (*Discovery) RegisterPossibilityHook

func (d *Discovery) RegisterPossibilityHook(fn func(sdkp2p.Possibility))

RegisterPossibilityHook adds a callback fired (outside locks) on every possibility change. Used by sync status.

func (*Discovery) Run

func (d *Discovery) Run(_ context.Context) error

func (*Discovery) SetEnabled

func (d *Discovery) SetEnabled(enabled bool)

SetEnabled switches mDNS announce and browse on or off at runtime (SDK.SetLocalDiscoveryEnabled has the contract). Off ends a live session at once and keeps the supervisor from starting another; on starts a session as soon as the probe allows, cutting short any retry backoff. A restatement is a no-op.

type Exchange

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

Exchange runs the SpaceExchangeV2 handshake with discovered local peers: both sides learn each other's dialable addresses and which spaces they SHARE, recorded in the PeerStore. Reuses any-sync's clientspaceproto wire shape.

Peers exchange per-space HMAC tokens keyed by a member-only discovery key and learn only the INTERSECTION of their space sets. A matching token proves the sender's membership, so strangers on the LAN learn nothing, can't track a device across sessions, and can't poison the peer store with spaces they don't hold. Spaces whose discovery key isn't derivable yet (ACL not synced) are skipped until it is.

The ACL-derived key alone dead-locks the offline cold restore: a fresh device of the SAME account knows a space's id (from the synced tech-space index) but can't derive its discovery key before pulling the space — and can't pull over LAN without the key. PROBE tokens break the cycle: for known-but-keyless spaces the caller sends a token keyed by an account-derived key (HKDF of the account signing key — only the account's own devices hold it). The responder answers a probe with a membership proof but records NOTHING: a probe claims interest, not possession, so it must not put the caller into the responder's per-space peer set. The caller records the responder as holding the space and pulls from it; the post-pull re-handshake (storage set change → Broadcast) then advertises the space normally.

The legacy plaintext SpaceExchange v1 is NOT supported: the SDK never calls it, and the inbound handler refuses it — full space-id lists must never leave this device, and there is no fallback an attacker could downgrade to. Pre-v2 peers simply don't pair over LAN.

func NewExchange

func NewExchange(selfPeerId string, store *PeerStore, allSpaceIds func() []string, discoveryKeys func(ctx context.Context, spaceIds []string) map[string][]byte, onPeerUpdated func(peerId string, spaceIds []string)) *Exchange

func (*Exchange) Broadcast

func (e *Exchange) Broadcast(ctx context.Context)

Broadcast re-runs the handshake with every known local peer. Called when this device's own space set changes (space created, pulled, or deleted) so peers learn the new set promptly instead of on the next discovery resweep.

func (*Exchange) Init

func (e *Exchange) Init(a *app.App) error

func (*Exchange) Name

func (e *Exchange) Name() string

func (*Exchange) PeerDiscovered

func (e *Exchange) PeerDiscovered(ctx context.Context, discovered sdkp2p.DiscoveredPeer, own sdkp2p.OwnAddresses)

PeerDiscovered is the discovery notifier: register the peer's addresses, dial, and run the handshake. Errors are logged, not returned — discovery re-announces periodically, so a failed attempt retries on the next sighting.

func (*Exchange) PeerLost

func (e *Exchange) PeerLost(peerId string)

PeerLost is the discovery notifier for a peer that left the LAN: its LAN addresses and presence go, so its iroh ticket, if any, takes over. The next sighting re-adds it through PeerDiscovered.

func (*Exchange) SetAccountKeysFn

func (e *Exchange) SetAccountKeysFn(fn func(ctx context.Context, spaceIds []string) map[string][]byte)

SetAccountKeysFn wires the account-derived discovery key source for probe tokens. Set once during app assembly.

func (*Exchange) SetKnownSpaceIdsFn

func (e *Exchange) SetKnownSpaceIdsFn(fn func() []string)

SetKnownSpaceIdsFn wires the known-space-ids source (the tech-space index) for probe tokens. Set after the SDK layers are up — discovery handshakes may already be running concurrently, hence the atomic; handshakes that run before simply don't probe.

func (*Exchange) SetOwnAddressesFn

func (e *Exchange) SetOwnAddressesFn(fn func() sdkp2p.OwnAddresses)

SetOwnAddressesFn wires the announce source Broadcast embeds in proactive re-handshakes. Set once during app assembly.

func (*Exchange) SpaceExchange

SpaceExchange refuses the legacy plaintext v1 handshake: it would hand our full space-id list to any LAN peer, and serving it at all would give an active attacker a downgrade target. The method exists only because the DRPC service interface requires it.

func (*Exchange) SpaceExchangeV2

SpaceExchangeV2 is the inbound side of the token handshake: compute this device's expected request token for every space it holds a discovery key for, intersect with what the caller sent, and answer with membership proofs for the intersection only — keyed by the caller's nonce, so they can't be precomputed or replayed.

type Global

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

Global is the internet-wide p2p layer: it publishes this device's endpoint ticket into every loaded space's key-value store, learns the other members' tickets from the same rows, keeps the peer store / addr book / status book in sync with them, gates inbound connections to known members, and runs the connector that maintains a bounded set of global connections. Everything network-facing happens on the connector; the key-value side never dials.

func NewGlobal

func NewGlobal(cfg config.GlobalP2P, selfPeerId, selfIdentity string, store *PeerStore, status *StatusBook, book *AddrBook) *Global

NewGlobal builds the layer; zero budget fields take their defaults.

func (*Global) AccountEnabled

func (g *Global) AccountEnabled() bool

AccountEnabled reports whether the account layer is on.

func (*Global) Close

func (g *Global) Close(_ context.Context) error

Close stops the workers (bounded wait), then flushes and closes the status book.

func (*Global) Init

func (g *Global) Init(a *app.App) error

func (*Global) LoadedSpaceIds

func (g *Global) LoadedSpaceIds() []string

LoadedSpaceIds lists the registered spaces.

func (*Global) Name

func (g *Global) Name() string

func (*Global) PeerStatus

func (g *Global) PeerStatus(peerId string) sdkp2p.PeerStatus

PeerStatus fills the liveness fields of one peer.

func (*Global) Republish

func (g *Global) Republish(spaceId string)

Republish re-sets the own row of a space (advertising switched on).

func (*Global) RepublishAccount

func (g *Global) RepublishAccount()

RepublishAccount runs a record cycle soon (ticket change, own relay change).

func (*Global) Run

func (g *Global) Run(_ context.Context) error

func (*Global) SetAccount

func (g *Global) SetAccount(keys *account.Keys, client accountClient, insecure bool, path string)

SetAccount turns the account layer on. Set before Run. insecure admits http:// relays named by sibling entries, for test relays; path is where the last decoded record is kept across restarts (empty keeps it in memory only).

func (*Global) SetAdvertiseFn

func (g *Global) SetAdvertiseFn(fn func(spaceId string) bool)

SetAdvertiseFn gates the per-space row: spaces the fn declines get no row and no heartbeat. Loaded spaces are re-evaluated at once.

func (*Global) SetKVSubscriber

func (g *Global) SetKVSubscriber(fn KVSubscriber)

SetKVSubscriber wires the per-space key-value dispatcher. Set during app assembly, before any space loads.

func (*Global) SetOnLive

func (g *Global) SetOnLive(fn func(peerId string))

SetOnLive registers a callback for every global peer that becomes live (dialed or accepted). Set during app assembly.

func (*Global) SpaceLoaded

func (g *Global) SpaceLoaded(spaceId string, kv SpaceKV)

SpaceLoaded registers a loaded space: its records are read, its applied writes followed, and the own record published (or re-set when stale). Local-only and guest spaces must not be registered.

func (*Global) SpaceUnloaded

func (g *Global) SpaceUnloaded(spaceId string)

SpaceUnloaded drops a space: its records stop contributing to the peer store, and peers known only through it disappear.

func (*Global) Status

func (g *Global) Status() sdkp2p.GlobalStatus

Status is the debug snapshot of the layer.

type KVHandler

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

KVHandler receives applied key-value writes of one space. Runs on any-sync's apply path — decode and enqueue only.

type KVSubscriber

type KVSubscriber func(spaceId string, h KVHandler) (cancel func())

KVSubscriber registers a KVHandler for a space; the app layer wires its per-space dispatcher here.

type Notifier

type Notifier interface {
	PeerDiscovered(ctx context.Context, peer sdkp2p.DiscoveredPeer, own sdkp2p.OwnAddresses)
	// PeerLost reports a peer that left the LAN (driver lost event).
	PeerLost(peerId string)
}

Notifier consumes discovery results; implemented by Exchange.

type Observer

type Observer func(peerId string, before, after []string, removed bool)

Observer is notified after a peer's space set changes (union over sources). before and after are the peer's space ids around the change; removed is true when the peer is gone from every source (after is then nil). Called outside the store's lock — observers may call back into the store.

type PeerRecord

type PeerRecord struct {
	// LastSeen is the newest evidence the peer is alive: a key-value
	// heartbeat (publisher clock, clamped to now) or a local
	// connection.
	LastSeen time.Time `json:"lastSeen"`
	// LastAttempt is the last global dial attempt.
	LastAttempt time.Time `json:"lastAttempt,omitempty"`
	// Failures counts consecutive failed global dials; a success or
	// fresh evidence resets it.
	Failures int `json:"failures,omitempty"`
}

PeerRecord is the persisted liveness record of one peer.

type PeerStore

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

PeerStore tracks the peers that SHARE spaces with this device and which spaces, per source: LAN peers from the space exchange, global peers from key-value records, account peers from the account's discovery record. The per-space peer manager, pubsub and the files p2p source read LAN and global peers separately — LAN peers are dialed inline, global peers are only used while already connected. Account peers are global peers of every space: they carry no space set, GlobalPeerIds lists them for any space asked. In-memory; the LAN side is rediscovered from scratch on restart, the others are rebuilt from the records.

func NewPeerStore

func NewPeerStore() *PeerStore

func (*PeerStore) AccountPeerIds

func (p *PeerStore) AccountPeerIds() []string

AccountPeerIds returns the devices of this account known through the record, best first.

func (*PeerStore) AddObserver

func (p *PeerStore) AddObserver(o Observer)

AddObserver registers a change callback. No removal — the observer set is fixed at wiring time and lives as long as the app.

func (*PeerStore) AddSourceObserver

func (p *PeerStore) AddSourceObserver(o SourceObserver)

AddSourceObserver registers a per-source presence callback.

func (*PeerStore) AllGlobalPeers

func (p *PeerStore) AllGlobalPeers() []string

AllGlobalPeers returns every peer known through records — space rows or the account record — most recently seen first, disabled tier excluded.

func (*PeerStore) AllLocalPeers

func (p *PeerStore) AllLocalPeers() []string

AllLocalPeers returns every known LAN peer id, sorted.

func (*PeerStore) GlobalPeerIds

func (p *PeerStore) GlobalPeerIds(spaceId string) []string

GlobalPeerIds returns the global peers known to have spaceId — the space's record peers plus every device of this account — most recently seen first, disabled tier excluded.

func (*PeerStore) GlobalSpaceIds

func (p *PeerStore) GlobalSpaceIds(peerId string) []string

GlobalSpaceIds returns the spaces a global peer is known to have.

func (*PeerStore) HasAccountPeer

func (p *PeerStore) HasAccountPeer(peerId string) bool

HasAccountPeer reports whether the peer is a device of this account.

func (*PeerStore) HasGlobalPeer

func (p *PeerStore) HasGlobalPeer(peerId string) bool

HasGlobalPeer reports whether the peer is known through records — a space row or the account record — in any tier.

func (*PeerStore) HasLocalPeer

func (p *PeerStore) HasLocalPeer(peerId string) bool

HasLocalPeer reports whether the peer is known on the LAN.

func (*PeerStore) HasSpace

func (p *PeerStore) HasSpace(peerId, spaceId string) bool

HasSpace reports whether the peer is known to hold spaceId through any source; a device of this account holds every space.

func (*PeerStore) Init

func (p *PeerStore) Init(_ *app.App) error

func (*PeerStore) LocalPeerIds

func (p *PeerStore) LocalPeerIds(spaceId string) []string

LocalPeerIds returns the LAN peers known to have spaceId.

func (*PeerStore) Name

func (p *PeerStore) Name() string

func (*PeerStore) RemoveAccountPeer

func (p *PeerStore) RemoveAccountPeer(peerId string)

RemoveAccountPeer forgets a device of this account.

func (*PeerStore) RemoveGlobalPeer

func (p *PeerStore) RemoveGlobalPeer(peerId string)

RemoveGlobalPeer forgets a peer's global presence.

func (*PeerStore) RemoveLocalPeer

func (p *PeerStore) RemoveLocalPeer(peerId string)

RemoveLocalPeer forgets a peer's LAN presence (dial failure, connection closed and gone). Unknown peers are a no-op.

func (*PeerStore) SetStatus

func (p *PeerStore) SetStatus(b *StatusBook)

SetStatus wires the liveness book that orders global peers and hides disabled ones. nil keeps insertion order and hides nobody.

func (*PeerStore) Sources

func (p *PeerStore) Sources(peerId string) []Source

Sources returns the sources that know the peer.

func (*PeerStore) SpaceIds

func (p *PeerStore) SpaceIds(peerId string) []string

SpaceIds returns the spaces a peer is known to have, over every source.

func (*PeerStore) UpdateAccountPeer

func (p *PeerStore) UpdateAccountPeer(peerId string)

UpdateAccountPeer records a device of this account. It needs no space set: it holds every space this device holds.

func (*PeerStore) UpdateGlobalPeer

func (p *PeerStore) UpdateGlobalPeer(peerId string, spaceIds []string)

UpdateGlobalPeer records the full global space set for a peer.

func (*PeerStore) UpdateLocalPeer

func (p *PeerStore) UpdateLocalPeer(peerId string, spaceIds []string)

UpdateLocalPeer records the full LAN space set for a peer; an empty set keeps the peer known (it stays in AllLocalPeers).

type Picker

type Picker interface {
	Pick(ctx context.Context, id string) (peer.Peer, error)
}

Picker is the pool slice PickLive needs.

type Source

type Source uint8

Source is how a peer became known.

const (
	// SourceLAN — the SpaceExchangeV2 handshake on the local network.
	SourceLAN Source = iota
	// SourceGlobal — a key-value record in a shared space.
	SourceGlobal
	// SourceAccount — the account's device-discovery record: another
	// device of this account, which holds every space this device holds.
	SourceAccount
)

func (Source) String

func (s Source) String() string

String returns a stable lowercase token for logging / status.

type SourceObserver

type SourceObserver func(source Source, peerId string, present bool)

SourceObserver is notified when a peer enters (present) or leaves a source. Called outside the store's lock.

type SpaceKV

type SpaceKV interface {
	Store() keyvaluestorage.Storage
	// CanWrite reports whether this device may Set rows (writer+).
	CanWrite() bool
	// IsMember reports whether identity (account address) still holds
	// any permission in the space.
	IsMember(identity string) bool
}

SpaceKV is the slice of a loaded space the global layer reads and writes: its default key-value store and two ACL answers.

type StatusBook

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

StatusBook keeps every known peer's liveness record, persisted as one JSON file under DataDir (same low-ceremony persistence as the p2p port file), written debounced. Records outlive restarts so a device that was offline for weeks resumes with the right tiers instead of dialing every stale peer at boot.

func NewStatusBook

func NewStatusBook(path string, th Thresholds) *StatusBook

NewStatusBook creates a book persisted at path; empty path keeps it in memory only.

func (*StatusBook) Attempt

func (b *StatusBook) Attempt(peerId string, ok bool)

Attempt records the outcome of a global dial: success is liveness evidence and clears the failure streak; failure extends it.

func (*StatusBook) Close

func (b *StatusBook) Close() error

Close flushes pending changes, then refuses further saves.

func (*StatusBook) Flush

func (b *StatusBook) Flush() error

Flush writes the book now (fsync, atomic replace). No-op when clean. The book stays dirty until the replace succeeded, so a failed write is retried by the next change.

func (*StatusBook) Forget

func (b *StatusBook) Forget(peerId string)

Forget drops a peer's record.

func (*StatusBook) Get

func (b *StatusBook) Get(peerId string) (PeerRecord, bool)

Get returns a peer's record.

func (*StatusBook) LastSeen

func (b *StatusBook) LastSeen(peerId string) time.Time

LastSeen returns a peer's LastSeen, zero when unknown.

func (*StatusBook) Load

func (b *StatusBook) Load() error

Load reads the persisted records; a missing file is an empty book. Records past the disable threshold are dropped: the peer is disabled anyway and a fresh row re-creates its record.

func (*StatusBook) Seen

func (b *StatusBook) Seen(peerId string, at time.Time) bool

Seen records liveness evidence at time at. Publisher clocks are untrusted: at is clamped to now. Reports whether LastSeen advanced.

func (*StatusBook) SetOnAdvance

func (b *StatusBook) SetOnAdvance(fn func(peerId string))

SetOnAdvance installs the LastSeen-advanced hook. Set at wiring time.

func (*StatusBook) Thresholds

func (b *StatusBook) Thresholds() Thresholds

Thresholds returns the configured tier boundaries.

func (*StatusBook) Tier

func (b *StatusBook) Tier(peerId string) Tier

Tier classifies a peer now. Unknown peers are disabled: every global peer gets a Seen call from its key-value record before it is used.

type Thresholds

type Thresholds struct {
	Stale, Dormant, Disable time.Duration
}

Thresholds are the tier boundaries on the age of a peer's LastSeen.

func ThresholdsFrom

func ThresholdsFrom(g config.GlobalP2P) Thresholds

ThresholdsFrom reads the tier boundaries from the global config (defaults already applied by WithDefaults).

type Tier

type Tier uint8

Tier is a peer's liveness class, derived from how long ago it was last seen. It sets how eagerly the connector dials the peer and whether the peer is offered to sync at all.

const (
	// TierActive — seen recently; kept connected with a short backoff.
	TierActive Tier = iota
	// TierStale — probed on a slow cadence, after active peers.
	TierStale
	// TierDormant — probed at startup and every few hours.
	TierDormant
	// TierDisabled — never dialed, dropped from the addr book and the
	// inbound allowlist until fresh evidence arrives.
	TierDisabled
)

func TierFor

func TierFor(age time.Duration, th Thresholds) Tier

TierFor classifies a LastSeen age.

func (Tier) String

func (t Tier) String() string

String returns a stable lowercase token for logging / status.

Directories

Path Synopsis
Package account is the account-level device-discovery record: every device of an account registers itself in one pkarr record addressed by a key derived from the identity key, and every device — a fresh restore included — resolves its siblings from it.
Package account is the account-level device-discovery record: every device of an account registers itself in one pkarr record addressed by a key derived from the identity key, and every device — a fresh restore included — resolves its siblings from it.

Jump to

Keyboard shortcuts

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