Documentation
¶
Overview ¶
Package collab carries a github.com/go-crdt/crdt document between the people editing it: a gRPC service, a server that hosts documents, and a client that joins one.
The service is thin on purpose. The document is a CRDT, so the server never transforms an operation and never decides an outcome — it applies what it is sent to its own replica and hands it to everyone else. Two consequences follow that a server-authoritative design cannot offer: a participant may edit while disconnected and reconcile later, and the server may be restarted or replaced without any client losing work.
Over what ¶
Two carriers, and which one to use is decided by where the code runs rather than by taste. WebSocket carries a session's own framing over a plain WebSocket; GRPC carries it over gRPC. One server serves both at once — Server.ServeWebSocket beside the registered service — and a participant on each edits the same document.
The reason there are two is measured. Everything a session carries is bytes some encoder in github.com/go-crdt/crdt produced and will check on arrival, so protobuf is describing fields nobody reads through it — and compiled to wasm its reflection and registry machinery cannot be linked away. The browser test client, gzipped, is 919 KB over the framing and 4 461 KB over gRPC, against 633 KB for the CRDT alone. Outside a browser none of that matters, and gRPC brings deadlines, interceptors and the tooling built around them.
The client builds for js/wasm either way, so a browser tab and a server run the same code down to the merge. Two browsers with no server between them carry a session over a WebRTC data channel ([DataChannel]), and two tabs of one browser over a BroadcastChannel ([JoinBroadcastChannel]) with nothing to configure at all.
A document holds named parts ¶
What an editor holds is not one structure: the text of a file, the comments anchored into it, the record of who changed what, the messages beside it, the cells of a sheet. A document here is a github.com/go-crdt/crdt.Composite, so they travel together — one snapshot, one version, one decision about who may open it, and no instant at which the set of them disagrees.
A caller reaches for a part by name and gets a handle: Client.Text, Client.List, Client.Map. A handle edits and publishes in one step, which is why it exists rather than the replicated structure itself — a caller editing that directly would produce operations nobody ever heard, and drift away from everyone else while its own screen looked right.
Shape of a session ¶
One bidirectional stream per participant per document. The client opens with a collabpb.Join; the server answers with a collabpb.Welcome holding either the whole document or, for a participant that says what it already has, only what it missed. After that, operations and presence flow both ways until either side hangs up.
What this protects, and what it does not ¶
Nothing here is encrypted or signed. The server reads every document it holds, a store holds them in the clear, and a participant is whoever the transport says it is. That is a choice with a reason, and the reason is structural rather than a matter of effort.
A server that merges cannot be blind, and a server that is blind cannot merge. This one merges: it applies operations, hands a joining participant a snapshot built from them, and collects tombstones once every participant has delivered. Every one of those reads the document. Encrypting it end to end would leave the server with bytes it cannot combine, which is a different design and not a setting.
The field splits along exactly that line, and it is worth naming where the neighbours stand:
- Automerge and Yjs encrypt nothing in the format. Automerge's chunk header carries four bytes of a truncated SHA-256, which detects a chunk that changed and authenticates nobody; Yjs's updates carry no digest at all.
- Jazz states authorization as row-level policy that its serving node applies before accepting a write and before shipping a row — which it can only do by reading the row. Its relay links are a separate kind of peer, defined as having no permission subject at all.
- Evolu does encrypt end to end, and pays for it exactly here: its server is called a Relay and holds encrypted changes it reconciles by fingerprints over ranges of timestamps. It never merges anything, because it cannot.
So what is actually underneath a deployment of this package:
- The transport. A session runs over WebSocket or gRPC, and under TLS both authenticate the server and give every message a MAC — which is stronger than a checksum and covers forgery, not merely rot. Running either without TLS puts the document on the wire in the clear.
- Whoever the caller lets in. This package does not authenticate: a server serves the sessions its host hands it, and Config.Authorize is where a host decides who those are.
- The store, against a medium that changes bytes rather than against somebody who writes them. PackSnapshot and CheckSnapshot carry a CRC32C, and anything that can write a stored document can recompute one.
What a session costs, since nothing here bounds how many there are ¶
Config has no capacity limit, on purpose: this package decides who may be in a document, not how many machines are worth buying. A deployment bounds that outside — a reverse proxy, a connection limit, a quota — and to do that it needs the numbers, which are these.
A participant is cheap in time and not free in memory. Measured over an in-memory connection on one document: about 2.5 µs a participant an edit, flat from a hundred upwards, so a thousand watching and five typing at ten keystrokes a second is roughly 12% of a core. Bringing one in costs about 27 KB at a thousand, more below that while the document's own cost amortises. See BenchmarkFanOut for the table and for what its bytes column does and does not answer.
One message is not cheap. A session may send up to a gigabyte in a single message, because a document's whole snapshot travels in one and that bound is the largest document a session may open — not the size of an edit. Nothing caps the sum of those across sessions, so the arithmetic a deployment has to do is per-session peak times sessions allowed, and the lever is the second factor.
And the bytes are not the whole cost. What an operations message decodes INTO is sixteen to twenty-four times larger than the message, that being the ratio of a record in memory to the smallest encoded one — so a gibibyte of operations asks this server to reserve sixteen to twenty-four gibibytes, and it asks that of an honest sender as much as a hostile one. Config.MaxOperations is the bound on that, in the unit that amplifies rather than in bytes, and it is off by default (TestAServerBoundsWhatOneMessageMayReserve, and TestALinkIsBoundedLikeAnybodyElse for the link, which is the point of it).
The unit is the point, and it is measured: a parse allocates 80 bytes an operation, flat from a thousand to a million of them. What that is a multiple of is the SENDER's choice — 6.2 times a realistic single-character insert, 20.0 times the four-byte floor an operation can encode in — so a bound written in bytes would be a different bound for every peer, and this one is not.
That bound is the server's. Every path where this package is the server goes through it, a federation link included: Server.Follow hands what it receives to the same apply as a participant's batch, so a followed server run by another actor is bounded like anybody else. ClientConfig.MaxOperations is the same bound in the other direction, for a participant.
A participant usually needs no such thing, because its peer is the server its own operator runs: the one that reads every document it holds, keeps the store in the clear and decides who may join. Bounding a message from it protects against nothing that is not already extended.
The peer-to-peer carriers are why the field exists. Over [DataChannel] the far end is another PERSON's browser, and over [JoinBroadcastChannel] it is whichever tab won the election. The page that hosts builds a Server and could always bound what it was sent; until ClientConfig.MaxOperations the page that joined could not, so one session had a bound in one direction only. go-tex's playground is the deployment that made that concrete, and it now sets both. Sizing it is arithmetic: multiply by 80 to 96 bytes a record and compare against what one message may spend.
And what a participant can do once it is in: everything a replica can do to a document it holds. It can write anywhere and delete anything — a CRDT converges on what it is told, and does not adjudicate — and it can hand over operations another site made, which is worse than it sounds and is not refused by default. OwnSiteOnly is the one-line policy that refuses it, and Config.AuthorizeOperations is where a federating deployment writes a narrower one. Neither is a stricter merge: the merge is not the lever.
A site identity is claimed, not proved ¶
This bounds who may federate with whom, so it is worth being exact about. A github.com/go-crdt/crdt.SiteID is a number a joining session states, and nothing here binds the number to whoever states it: TLS authenticates the server rather than the site, Config.Authorize decides whether a session may join and cannot check a number it did not issue, and github.com/go-crdt/crdt.DeriveSiteID is a pure function of a name, so anyone who knows the name computes the identity.
Within one operator that is a topology decision — both ends hand out site identities, so a site is as trustworthy as the deployment. Across operators it is not, and the cost is measured in TestAFederatedPeerCanSpeakAsAnotherServersUser: a followed server whose participant claims a site belonging to the follower's user makes the follower hold a document that existed on neither server, my user's genuine prefix grafted to the tail of a forged sequence, with every character attributed to that user — and both replicas then report the SAME version vector, so each believes it is completely caught up with the other and neither will ever ask for anything again. Nothing returns an error: the operations were well formed and the merge converged on what it was told.
OwnSiteOnly is not the answer, for a reason worth knowing before reaching for it: the side it breaks is the FOLLOWER, because what arrives over a link names sites that server never authorised (TestALinkCarryingOtherSitesMeetsOwnSiteOnly, with TestAServerWithOwnSiteOnlyRefusesForgedOperations, TestAFederatedPeerCanSpeakAsAnotherServersUser and TestAnOperatorIsToldWhichBatchesWereRefused on the rest of that path). A server that federates cannot install it, which is why nothing on the wire distinguishes a link from a participant.
Those names are not decoration. Each was established by taking the guarantee out and watching what went red, and all four of this arc's safety properties have now been through it:
- deleting the version-and-digest comparison in MergeSnapshots fails exactly one test and nothing else;
- disabling Config.AuthorizeOperations fails four — at least four, because with it disabled the package does not finish, several tests waiting out deadlines for a refusal that never arrives;
- passing no bound where Config.MaxOperations goes fails two;
- and the welcome digest fails two from the server's side and one from the link's, depending on which half is removed.
A guarantee whose removal breaks nothing is a sentence, not a property, and reading the tests cannot tell the two apart.
Federating with somebody else's server, then ¶
It takes two rules, and both are the operator's to state because only the operator knows who the other actors are. Neither needs a change to the format or to this package, and there is a worked example of both in gitstore's federation_example_test.go, which is written as a consumer of this package rather than part of it so that anything it needs and cannot reach is a gap.
- SCOPE the site identity, so two actors cannot mint the same one. The example derives a site from an eduGAIN identifier — "ada@paris.example.ac" rather than "ada" — because only the home organisation issues inside its own scope. A bare name is the one failure this design cannot merge its way out of: github.com/go-crdt/crdt.DeriveSiteID is a function, so "42" is the same replica on every instance in the world.
- Install SpeaksFor and answer its one question: may this session hand over work that site made? It is Config.AuthorizeOperations about the RELATION rather than about the sender, which is the difference that matters — a participant speaks for itself, a link speaks for a server, and what a link carries is not what it joined as.
Two details decide whether such a policy works, both learned the hard way, and SpeaksFor now gets both right so a deployment does not have to.
Every KIND of operation, because a document holds three. gitstore's example read a batch's text and nothing else, so an unfederated site was refused when it wrote a character and allowed when it wrote a map entry, until TestTheScopeCheckSeesEveryKindOfOperation. Sites is that walk, exported once.
And this server's OWN scope must never be among the scopes a link may carry. Listing it is how a link comes to be allowed to write as one of this server's own users, which is the attack rather than a refinement of it — and the example listed it. With SpeaksFor there is nothing to list: a session's own site is allowed without the predicate being consulted at all, so the predicate names only the scopes you federate with. Held to it by TestAScopedPolicyStopsALinkSpeakingForOurOwnUsers, whose third case is the measured attack: a peer claiming the very user who wrote here, with a longer history so the tail is actually sent. The link's session ends at once naming the site it may not speak for, and this replica keeps what it had.
One trap in testing this, because it turns a defect into something that looks like a defence. When a peer claims a site one of our participants used and writes LESS than that participant did, our link joins saying it holds that site up to a higher clock, so the followed server sends nothing at all. Nothing arrives and nothing is refused. That is the version-vector collision measured in TestAFederatedPeerCanSpeakAsAnotherServersUser, and a test that only watched the document would read it as the policy working.
Two limits remain, and they are properties of this shape rather than gaps in it. Trust is hop by hop: if A follows B and B follows C, B relays C's sites, so A grants B the union and thereby trusts B about C — which is how mail and Matrix federation trust. And an operation carries no signature, so a link is believed about the attribution of everything it relays; what a server can check is which sites a link may speak for, not that a site really said this.
That second one is a choice the field makes deliberately, and the nearest neighbour with a written specification makes it in the same shape. Jazz states it as invariants: a client "MUST reject any attempt to attribute a write to another author" (INV-API-29), which is OwnSiteOnly; a session link's author "MUST equal that identity or be rejected" while a trusted backend link may differ (INV-RLS-18), which is the participant-and-link distinction this package draws; and a write whose author differs from the authenticated subject "MUST be accepted only via a trusted serving node" (INV-RLS-17), which is trust hop by hop. Their relay links carry "no permission subject" at all, and they say plainly that theirs is "a local relay forwarding contract, not an end-to-end signature protocol". Two designs that share almost nothing else — theirs has a trusted core, this has links between operators — arrive at the same place: attribution is checked at the edge, by the identity of the link, and not by a signature on the operation.
Closing that second one is not a matter of signing anything, which is worth knowing before reaching for a signature. The failure it would be reached for is two replicas holding different text while reporting the SAME version vector, and that is equivocation: a site claiming one name for two different operations. Signatures do not rule it out — Kleppmann says so explicitly in "Making CRDTs Byzantine Fault Tolerant" (PaPoC '22, §2.3.1), whose Figure 1 is that failure — because a signature says who produced a batch, not that a name is unique. What rules it out is something the two replicas can COMPARE: a digest of the document they each hold, beside the version vector they already exchange, so that disagreeing replicas cannot both conclude they are finished. That is github.com/go-crdt/crdt.Composite.Digest, and this package compares it in the two places where a document arrives from somewhere else.
A digest of state rather than of history, and that is decided by this design rather than chosen. github.com/go-crdt/crdt.Doc keeps no operations: it holds the document, and Doc.OpsSince synthesises operations by walking it. So there is no history to hash-link, in the way Automerge hash-links stored changes. The happy half is that github.com/go-crdt/crdt.Doc.Purge discards only runs whose every character is already deleted, so a digest over the VISIBLE document survives a purge, and two replicas that purged differently still compare equal.
Where it is compared:
- MergeSnapshots, which is where two documents are both in hand. Two snapshots claiming the same history and holding different documents are ErrDiverged, and the merge refuses rather than grafting one onto the other (TestAMergeRefusesTwoSnapshotsThatClaimOneHistory). Checked BEFORE anything is carried, because the carry would hide it: OpsSince selects by name, so an operation wearing a name the base already holds is never sent, and the merge would return the base unchanged and call that agreement.
- A link, on the welcome. A server appends its digest for a peer that announced CapDigest, and a link that catches up to exactly that version compares. Equal versions and different digests ends the session and reaches Config.OnOperationsRefused, since a site identity claimed by two replicas is not something a session can discover from inside itself (TestALinkFindsTwoReplicasWearingOneName, and TestAPeerThatDidNotAskIsNotSentADigest for the half that keeps the block off a peer that never announced it).
ONLY when the two versions are equal, and that is the whole rule. Two replicas at different points are supposed to hold different documents, and every honest catch-up passes through that state; comparing there would cry on the normal case, which is how an alarm stops being read.
What it does not do: attribute, prevent or repair. It says two replicas differ, which is the half of the failure that otherwise has no remedy at all — a replica that believes it is finished never asks again. And it cannot see a side that is merely AHEAD: when one version covers the other, the operations they disagree about are exactly the ones OpsSince will not carry.
Config.OnOperationsRefused is how an operator hears any of it happen. A refusal otherwise goes to the offending session and nowhere else, which is the wrong room for both of the cases above: a peer carrying sites it may not speak for, and two replicas that chose the same site, are things only an operator can take up — and an accidental collision means the identities a deployment hands out are not unique, which nobody can discover from inside a session.
Two servers do not share a store ¶
A Store holds snapshots and Store.Save replaces. A server holds the document in memory while it serves it, so two servers holding the same document at the same time each save their own replica and the later save replaces the earlier — losing every operation the other held and this one never saw. Measured, in TestTwoServersOverOneStoreLoseTheEarlierSave: one writes AAAA, the other BBBB, and the store ends holding BBBB alone. There is no error anywhere, because each save did exactly what Save is documented to do.
A ConditionalStore makes that loss loud rather than silent: it saves only while the store still holds what this server last read or wrote, and answers ErrChanged otherwise, at which point the document goes back to unsaved and the error reaches Config.OnPersistError. It does not make a shared store supported — the replicas in memory are still not shared, so the refused server is still missing what the other applied — and it is optional, because a store can only refuse if it can compare before writing. What it buys is an operator finding out.
It is the sequential case that works, and it is the one that makes this tempting: a server that starts, loads, serves and stops hands the next server everything, so a rolling restart is fine and a failover to a cold standby is fine. What is not fine is two of them up at once.
There are two supported ways to run more than one server, and both replace the shared store rather than adding to it:
- Federation. Server.Follow and Server.FollowWithRetry make one server a participant in another's document, so the operations travel and each server's own store holds the union. The cost is that both servers are up and reachable — and that a link is trusted for every site it carries, so both ends must be yours. See "A site identity is claimed, not proved".
- A gitstore with a remote. Each server has its own repository and pulls the other's, merging snapshots rather than replacing them, with ErrUnmergeable when two replicas have each discarded what the other needs. The cost is a pull interval instead of a link, and the benefit is that neither server has to be reachable from the other.
A shared PostgreSQL looks like a third way and is not one: the database is shared, the document in memory is not.
Index ¶
- Constants
- Variables
- func CheckSnapshot(snapshot []byte) []byte
- func GRPCServerOptions() []grpc.ServerOption
- func MergeSnapshots(ours, theirs []byte) ([]byte, error)
- func OwnSiteOnly(_ context.Context, _ string, from crdt.SiteID, batches []crdt.PartOps) error
- func PackSnapshot(snapshot []byte) []byte
- func Pipe() (Transport, *PipeConn)
- func Sites(batches ...crdt.PartOps) iter.Seq[crdt.SiteID]
- func SpeaksFor(...) func(context.Context, string, crdt.SiteID, []crdt.PartOps) error
- func UnpackSnapshot(stored []byte) ([]byte, error)
- type Archivable
- type Capabilities
- type Capability
- type Client
- func (c *Client) Changes() <-chan struct{}
- func (c *Client) Close() error
- func (c *Client) Document() string
- func (c *Client) Done() <-chan struct{}
- func (c *Client) Err() error
- func (c *Client) List(name string) (*List, error)
- func (c *Client) Map(name string) (*Map, error)
- func (c *Client) Parts() []crdt.Part
- func (c *Client) Peers() []awareness.Peer
- func (c *Client) SetCursor(cursor awareness.Cursor, meta map[string]string) error
- func (c *Client) Site() crdt.SiteID
- func (c *Client) Snapshot() []byte
- func (c *Client) TakeChanges() []crdt.PartChange
- func (c *Client) Text(name string) (*Text, error)
- func (c *Client) Version() crdt.CompositeVersion
- type ClientConfig
- type ConditionalStore
- type Config
- type Dialer
- type DirStore
- func (s *DirStore) Documents() ([]string, error)
- func (s *DirStore) Idle(_ context.Context, d time.Duration) ([]string, error)
- func (s *DirStore) Load(_ context.Context, document string) ([]byte, error)
- func (s *DirStore) LoadSites(_ context.Context, document string) ([]byte, error)
- func (s *DirStore) Release(ctx context.Context, document string, want []byte) error
- func (s *DirStore) Save(_ context.Context, document string, snapshot []byte) error
- func (s *DirStore) SaveSites(_ context.Context, document string, sites []byte) error
- type GRPCServer
- type LinkStatus
- type List
- func (l *List) Append(values ...[]byte) error
- func (l *List) Delete(pos, count int) error
- func (l *List) Get(pos int) ([]byte, error)
- func (l *List) Insert(pos int, values ...[]byte) error
- func (l *List) Len() int
- func (l *List) Name() string
- func (l *List) Part() crdt.Part
- func (l *List) Values() [][]byte
- type Map
- func (m *Map) Delete(key string) error
- func (m *Map) Edit(fn func(*crdt.Map) ([]crdt.MapOp, error)) error
- func (m *Map) Get(key string) ([]byte, bool)
- func (m *Map) Keys() []string
- func (m *Map) Len() int
- func (m *Map) Name() string
- func (m *Map) Part() crdt.Part
- func (m *Map) Read(fn func(*crdt.Map))
- func (m *Map) Set(key string, value []byte) error
- type MemoryStore
- func (s *MemoryStore) Documents() []string
- func (s *MemoryStore) Idle(_ context.Context, d time.Duration) ([]string, error)
- func (s *MemoryStore) Load(_ context.Context, document string) ([]byte, error)
- func (m *MemoryStore) LoadSites(_ context.Context, document string) ([]byte, error)
- func (s *MemoryStore) Release(_ context.Context, document string, want []byte) error
- func (s *MemoryStore) Save(_ context.Context, document string, snapshot []byte) error
- func (m *MemoryStore) SaveSites(_ context.Context, document string, sites []byte) error
- type MultiStore
- type PipeConn
- type RetryPolicy
- type Role
- type Server
- func (s *Server) Close(ctx context.Context) error
- func (s *Server) Flush(ctx context.Context) error
- func (s *Server) Follow(ctx context.Context, peer Transport, document string, as crdt.SiteID) error
- func (s *Server) FollowWithRetry(ctx context.Context, dial Dialer, document string, as crdt.SiteID, ...) error
- func (s *Server) ServePipe(ctx context.Context, sc *PipeConn) error
- func (s *Server) ServeWebSocket(origins ...string) http.Handler
- func (s *Server) Stable(name string) (crdt.CompositeVersion, bool)
- type SiteStore
- type Store
- type Text
- func (t *Text) Anchor(pos int) (crdt.ID, error)
- func (t *Text) AnchorUTF16(pos int) (crdt.ID, error)
- func (t *Text) AuthorRuns() []crdt.AuthorRun
- func (t *Text) AuthorRunsUTF16() []crdt.AuthorRun
- func (t *Text) Delete(pos, length int) error
- func (t *Text) DeleteUTF16(pos, length int) error
- func (t *Text) Insert(pos int, text string) error
- func (t *Text) InsertUTF16(pos int, text string) error
- func (t *Text) Len() int
- func (t *Text) LenUTF16() int
- func (t *Text) Name() string
- func (t *Text) Part() crdt.Part
- func (t *Text) Position(anchor crdt.ID) (int, bool)
- func (t *Text) PositionUTF16(anchor crdt.ID) (pos int, ok bool)
- func (t *Text) String() string
- func (t *Text) Visible(anchor crdt.ID) bool
- type Tiered
- type Token
- type Transport
- type WebSocketOption
Constants ¶
const DefaultBacklog = 256
DefaultBacklog is how many messages may be queued for one participant before the server gives up on it. See Config.
It is also, and less obviously, a capacity: a document where P participants edit in the same breath sends P-1 messages to each of them, so a backlog below P is a backlog that will be exceeded. Measured on one server with 800 participants each making one edit at once — a backlog of 256 disconnected 99% of them, 512 disconnected 32%, and 1024 disconnected none.
So 256 means "about two hundred and fifty people editing at the same instant", not "a queue that is usually long enough". A document with more than that wants a larger Config.Backlog, and the cost of one is a queue slot per participant rather than anything per document.
const DefaultRetryCeiling = 30 * time.Second
DefaultRetryCeiling is the longest a link waits between attempts, when RetryPolicy.Ceiling does not say. It bounds two things at once: how much work a peer that has been down for a day is asked to do, and how stale a replica can be once that peer comes back, since nothing crosses the link until the next attempt. Half a minute is the compromise, and an operator who knows their outages can say better.
const DefaultRetryWait = 250 * time.Millisecond
DefaultRetryWait is how long a link waits before its first attempt at coming back, when RetryPolicy.Wait does not say. It is short because most drops are brief — a process restarted, a route reconverging — and a link that is back within a second has lost nothing anybody typed.
Variables ¶
var ErrBroadcastClosed = errors.New("collab: broadcast channel closed")
ErrBroadcastClosed is why a BroadcastChannel carrier's [Recv] or [Send] returned: this end was closed, or the bus it spoke over was. It is the shared-bus counterpart of the error a dropped socket reports.
var ErrChanged = errors.New("collab: the document changed since it was read")
ErrChanged reports a document that was written since it was read, so it was not released. It is not a failure: the next pass will archive the newer one.
var ErrClosed = errors.New("collab: session closed")
ErrClosed is why a session ended when this participant closed it, and what an edit made afterwards returns.
var ErrDiverged = errors.New("collab: two snapshots claim the same history and hold different documents")
ErrDiverged reports two snapshots that claim the SAME history and do not hold the same document.
A version vector counts operations per site, so two replicas whose vectors match have, by that account, applied the same operations — and must therefore hold the same document. When they do not, some site put its name on two different operations, and the two replicas each believe they are completely caught up with the other. Neither will ever ask for anything again, which is why this is worth a merge refusing rather than a note in a log: merging them grafts one history onto the other and produces a document that existed on neither side, with every character attributed to whoever the names say.
It is not a corruption and the checksums will not see it. Both snapshots are well formed, and both replicas applied what they were sent. What differs is what the operations SAID, and only comparing the documents can show it — crdt.Composite.Digest is that comparison, and this is the one place in this package where two documents are both in hand.
An operator who meets this has two replicas claiming one site identity. Within one deployment that means the identities it hands out are not unique; across two, it means one of them is speaking for the other's users, which is what SpeaksFor refuses at a link and nothing refuses in a shared repository.
What it does NOT catch: a side that is genuinely ahead. If one version covers the other, the operations they disagree about are exactly the ones crdt.Composite.OpsSince will not carry, because it selects by name and the names match. That case is silent here and is why a link is the better place to federate from.
var ErrHostSuperseded = errors.New("collab: another tab is hosting this room")
ErrHostSuperseded is why a host's serve loop returned: another tab with a lower identifier — the one the tie-break gives priority — announced itself as host, so this tab steps down to let that one hold the room. It is the self-healing half of the election: when two tabs both reach RoleHost, exactly one keeps the document and the other yields, rather than the two drifting apart as separate documents.
A caller MUST re-join on it, CARRYING WHAT IT HELD ¶
Not "may", and not eventually: on this error the room has one host and this tab is not it, so a tab that only tears down is a tab holding nobody. A consumer that did exactly that showed "Connected" to a user who was alone in a room of two, until they clicked again -- the same-browser split-brain the election exists to prevent, arrived at from the other side. Re-joining finds the surviving host answering, because it was answering before this error was sent.
And re-join with Client.Snapshot in ClientConfig.Resume, not empty-handed. This tab was the room: its buffer was the seed and somebody may have typed into it. Arriving with nothing means adopting the survivor's document and losing that -- silently, and decided by which identifier happened to be lower, which is a tie-break that exists for something else.
Nothing new travels to make that work. Resume is documented as keeping the work done while disconnected, and a tab just superseded is a tab that was disconnected; the union does the rest. See TestASupersededHostKeepsItsWorkByBringingIt.
It is not rare ¶
[electRole] holds only while frames posted are also read inside its window, and a loaded machine does not read them. A consumer's CI met this on EVERY run for a week while its comment called the race vanishingly unlikely. Treat it as a path, not a corner.
var ErrNoDocument = errors.New("collab: a document must have a name")
ErrNoDocument reports a document with no name. The server refuses one at the door — a join must name a document — and this refuses it too rather than let it name the directory itself, which is what the empty name encodes to.
var ErrPipeClosed = errors.New("collab: pipe closed")
ErrPipeClosed is why a Pipe carrier's [Recv] or [Send] returned: either end was closed, or the session's context was cancelled. It is the in-process counterpart of the error a dropped socket reports.
var ErrProtocol = errors.New("collab: unexpected message")
ErrProtocol reports a message that is not part of a session: a kind that cannot arrive at that moment — a second welcome, or a join halfway through — or bytes that are not a message at all.
var ErrTransport = errors.New("collab: transport")
ErrTransport reports a carrier that could not be opened or that failed.
var ErrUnmergeable = errors.New("collab: neither snapshot can serve the other")
ErrUnmergeable reports two snapshots neither of which can be brought up to the other, each having discarded what the other still needs.
It is the one divergence MergeSnapshots cannot resolve, and there is no operation that resolves it later either: what each side purged is in no operation, so neither can be told about the other's past. An operator who meets it has two documents and has to choose one; the merge will not choose for them, because choosing here is losing text somebody wrote.
Functions ¶
func CheckSnapshot ¶ added in v0.51.0
CheckSnapshot gives a snapshot a checksum and leaves it uncompressed, for a Store on a medium that stores differences between versions. UnpackSnapshot reverses this as well as PackSnapshot; the two are told apart by what they write at the front, so a store may change from one to the other and still read what it already holds.
Prefer PackSnapshot unless the medium is one of these. A snapshot is mostly columns of identities and offsets, and it compresses about 13.9× — but a medium that deltas versions loses more to compression than compression saves, and the numbers for git are on [checkedRawMagic].
func GRPCServerOptions ¶ added in v0.44.0
func GRPCServerOptions() []grpc.ServerOption
GRPCServerOptions are what a grpc.Server carrying GRPCService needs:
srv := grpc.NewServer(collab.GRPCServerOptions()...) collabpb.RegisterCollabServer(srv, collab.GRPCService(s))
gRPC's default receive limit is four mebibytes, meant for request-response traffic. A session's largest message is not that: a participant coming back from a spell offline sends everything it wrote while it was away, and a rejoin carries a version and the operations behind it. The client call has raised its own limit since this carrier was written; the server is built by whoever runs it, so the library can only hand over the option -- and a server missing it works for small documents and drops large ones with "received message larger than max", which names a resource and not the document that outgrew it.
func MergeSnapshots ¶ added in v0.24.0
MergeSnapshots combines two snapshots of the same document into one that holds everything either of them still holds.
It is the operation that makes a snapshot safe to keep in more than one place. Two copies of a document that were written separately have not disagreed about anything — a snapshot is a set of operations, and the union of two sets of operations is a document, which is the whole reason this project exists.
Either argument may be nil, which is how a store says it has never held the document; merging with nothing gives back the other side.
It chooses a base, and the argument order is not what chooses it ¶
The union is built by carrying one side's operations onto the other, and the side carried onto — the base — brings something the operations cannot: the floors crdt.Doc.Purge and crdt.Map.Collect leave behind. Those live in the snapshot and in no operation, so taking the first argument as the base makes the argument order an input. That is not untidiness, it is loss in both directions:
- A replica that purged text emits neither the insertions nor the deletions of a purged run, so making it the receiver instead of the base sends it a history with a hole in it and hands back the text it discarded. crdt.Doc.CanServe is the question to ask about that, and this used not to ask it.
- A replica stale across a crdt.Map.Collect, made the base over the collected side, advances its version past a deletion it never learns, and the deleted key is alive again. A map has no CanServe to catch that.
So the base is chosen from the pair rather than from the order:
- A side the other cannot serve is the base. That is the correctness condition and not a preference: a replica that cannot hand over its deletions has to receive rather than send.
- Otherwise the side that has given up more — floors at least the other's on every part — is the base, so that floors here only rise, as they do everywhere else in this system. A merge that lowered one would undo an operator's purge on every read and write the un-purged document back on the next save, which is a purge that can never be made to stick.
- Floors that cross, each side having collected further than the other on a different part, cannot both be kept by one snapshot, so one economy is declined and the other side's tombstones are simply kept. Which side is kept is settled by comparing the bytes: arbitrary, but a function of the pair and not of the order, which is the property being bought.
Merging is therefore still symmetric to the byte, and now for a reason it can state: the encoding is canonical, both results hold the same operations, and the argument order is not among the inputs.
One side empty is not a merge ¶
With either side empty the other is returned VERBATIM, unread. Nothing is decoded, so nothing is validated: a store holding bytes no reader accepts gets those bytes back with a nil error, and Tiered then writes them to the tier that was empty.
That is deliberate rather than an oversight. The common case for an empty side is a tier that has not been filled yet, and decoding a whole document to confirm what will be written back unchanged would put that cost on every read of a cold tier -- to catch a corruption the framing's checksum already refuses one layer down (UnpackSnapshot). What it means for a caller is that a nil error here is not a statement about the bytes unless BOTH sides were documents.
What it cannot carry, it names ¶
It returns ErrDiverged when the two sides claim the same history and do not hold the same document, ErrUnmergeable when neither side can serve the other, and passes on crdt.ErrStranded when an operation cannot be carried onto the base — a write at or below a collected floor naming a key the base does not hold, which is a key that would otherwise come back alive. Both of those used to be a wrong document returned with a nil error, and a wrong document is found by the person who wrote the paragraph rather than by the operator.
func OwnSiteOnly ¶ added in v0.54.0
func PackSnapshot ¶ added in v0.50.0
PackSnapshot compresses a snapshot and gives it a checksum, for a Store that keeps bytes on a medium which can change them. UnpackSnapshot reverses it. DirStore calls it on every save.
It is exported for the stores that are not in this package — pgstore is one, and so is whatever a consumer writes against S3, a key-value store or a blobstore. A Store is a public interface, so the answer to "who checksums this?" has to be available to whoever implements one.
Rot, not forgery. Anything that can write the bytes can recompute the checksum in them, so this authenticates nothing and does not pretend to; it catches a medium that changed a byte nobody asked it to change. Of 1856 single-bit flips in one composite snapshot, 1048 loaded with no error at all as a DIFFERENT document, which a server would then serve and, at its next save, make permanent. Framed, none of them do. Both halves are counted by TestEverySingleBitFlipIsRefusedOnceItIsFramed rather than quoted from a probe that no longer exists.
It returns no error, and the reason is that it cannot have one: the destination is a bytes.Buffer made here, and writing to one never fails. Checking anyway would put two branches in the package that no test can reach and a reader has to wonder about — and would put a third in Save, which would have to decide what to do about a failure that does not happen.
func Pipe ¶ added in v0.21.0
Pipe returns the two ends of an in-process session: a Transport to hand to Join, and a handle to hand to Server.ServePipe. Nothing leaves the process — messages are carried on Go channels rather than a socket — so it is how a peer that both serves a document and edits it locally binds its own editor to that document without opening a loopback connection to itself.
What it is for ¶
A browser tab holding a document for others over a data channel is also a place someone is typing. Its own editor is a participant like any other, and the honest way to say so is to Join the document it is serving. Without this that means a second, real carrier looping back to the same page — a WebRTC data channel to itself, dialled, answered and encrypted, to carry bytes that never cross a wire. Pipe is that participant with the carrier taken out: the same session, the same merge, over two channels and nothing else.
The same holds for a native server that wants a replica of a document it hosts — a headless editor, a linter, an exporter — without the cost of a second server or the round trip of a real link.
Shape ¶
The two ends are one session and share its fate: closing either, or cancelling the context Join or Server.ServePipe was given, ends both. A [Recv] blocked when that happens returns ErrPipeClosed, and a [Send] after it does the same, so neither end can be left waiting on a peer that has gone.
A participant here edits the same document as one arriving over WebSocket or GRPC; it is a first-class participant, not a shortcut with less to it.
func Sites ¶ added in v0.64.0
OwnSiteOnly is a Config.AuthorizeOperations that refuses a batch carrying any site but the one the session joined as.
srv := collab.NewServer(collab.Config{
Store: store,
AuthorizeOperations: collab.OwnSiteOnly,
})
Why a server would want it ¶
Without a policy, a session may hand over operations another site made, and the server applies them. That is not merely wrong attribution. A site identity is the half of an operation's name that makes it unique, so two writers using one site id produce different characters with the same ID -- and a CRDT converges on names.
Measured, and it is the guarantee this whole package rests on: a forger writing "FORGED" as site 2, and the real site 2 working offline writing "GENUINE", give two replicas that received both batches two different documents. The divergence is permanent.
Since crdt v0.48.0 it is no longer silent: the second batch each replica sees is refused with github.com/go-crdt/crdt.ErrCollidingID. That changes the announcement and not the outcome — refusing the second batch does not undo the first, so the replicas still disagree — and it does not reach a forger who reproduces what the replica already holds before diverging. A diagnostic that arrives after the document is wrong is why this policy is still the thing that prevents it. Measured in TestTwoWritersOnOneSiteIdentityDiverge.
Why it is not the default ¶
A link is a session too, and it is meant to carry other sites: Server.Follow joins as one site and relays the work of everyone on the server it follows, which is thousands of sites this server never authorised. Refusing them would refuse federation. Nothing on the wire says which kind of session is speaking, so the server cannot choose for a deployment -- which is why this is a policy to install rather than a rule.
The side that must not have it is the FOLLOWER. What arrives over its link names sites it never authorised, so this refuses them and the link ends. The server being FOLLOWED sees an ordinary session and is untroubled: it may run this policy while federating. Measured all three ways in TestALinkCarryingOtherSitesMeetsOwnSiteOnly.
When it does refuse a link, it refuses loudly -- the link's session ends at once and the error names both sites, "a session of site 9001 sent an operation made by site 1", rather than leaving a link that retries forever and a document that never fills.
A deployment that federates wants a policy about the relationship instead: which sites a given link may speak for. See Config.AuthorizeOperations.
It refuses a participant that resumes under a fresh site ¶
ClientConfig.Resume carries what a participant did while it was away, and those operations were written by the site it had THEN. Coming back under a new one makes them somebody else's as far as this policy can tell, so they are refused -- and the tab keeps showing them, because they are in its own replica and nowhere else. It looks like it worked.
So resume as the site that wrote the work. Measured, in TestResumingUnderAFreshSiteLosesTheWorkToOwnSiteOnly: with no policy a fresh site's resume arrives, with this policy it does not, and with the authoring site it does. Sites is every site named by the operations in batches, in the order they appear and with repeats: a batch of a hundred characters from one site yields it a hundred times.
It exists because the question "which sites does this batch name" had three answers in this repository and one of them was wrong. OwnSiteOnly walks all three kinds of operation; so does the helper in this package's own tests; and gitstore's federation example walked a batch's text and nothing else, so an institution it did not federate with was refused for a character and ALLOWED for a map entry. A policy that can be stepped around by writing to a different part of the document is not one. So the walk is exported, once, and a deployment writing its own policy has nothing to get wrong.
All three kinds, and a valid batch carries exactly one of them, so two of the three loops do nothing on any given call.
func SpeaksFor ¶ added in v0.64.0
func SpeaksFor(mayCarry func(ctx context.Context, document string, session, carried crdt.SiteID) bool) func(context.Context, string, crdt.SiteID, []crdt.PartOps) error
SpeaksFor builds a Config.AuthorizeOperations from one question: may this session hand over work that site made?
srv := collab.NewServer(collab.Config{
Store: store,
AuthorizeOperations: collab.SpeaksFor(
func(_ context.Context, _ string, _, carried crdt.SiteID) bool {
return lyon[carried] // the scopes this server federates with
}),
})
What it gets right so a deployment does not have to ¶
Every operation in the batch rather than the sender, because that difference is the whole point: a participant speaks for itself, a link speaks for a server, and what a link carries is not what it joined as. Every KIND of operation, for the reason Sites gives. And a session may ALWAYS speak for the site it joined as, without being asked — which is OwnSiteOnly's rule, and is the part that is easy to get wrong in the dangerous direction.
That last one is worth spelling out. A policy written as "the scopes I federate with" invites listing your own among them, and gitstore's example did: a link from Lyon was then allowed to write as one of Paris's own users, which is the attack in go-crdt/collab#175 and not a refinement of it. Here the question is never asked about the session's own site, so there is nothing to list.
What it deliberately does not know ¶
How a site maps to whoever may speak for it. That is the deployment's, and it is why this takes a predicate rather than a set of scopes: a real one asks its federation metadata, which may be per document and may want the context's deadline, so both are passed through.
OwnSiteOnly is the same rule as this with a predicate that always says no, and it is written out rather than built from this. Measured on a one-operation batch, which is what a keystroke is: 3.50 ns inlined against 53.1 ns through here, so fifteen times, with the inlined arm measured before and after the other and agreeing to three decimals. Two indirect calls an operation -- the range-over-func yield and the predicate -- is what that buys.
For a deployment that federates the same 50 ns is nothing: the fan-out of one edit costs 2.5 µs a participant, and what this saves is writing the walk by hand and getting one of its three kinds wrong, which is what happened twice in this repository. For the policy every server installs it would be a cost on every batch for a convenience nobody there needs, which is why OwnSiteOnly does not pay it.
func UnpackSnapshot ¶ added in v0.50.0
UnpackSnapshot reverses PackSnapshot and CheckSnapshot, and passes through anything that was written by neither — which is every document written before this existed. A store that starts calling this pair reads what it already holds, and its documents gain the checksum one save at a time, with nothing to migrate.
It refuses bytes whose checksum does not match them. Refusing is the point: served instead, they would be a document nobody wrote.
Types ¶
type Archivable ¶ added in v0.26.1
type Archivable interface {
Store
// Idle returns the documents this store has not been asked to write for
// longer than d.
//
// Not written is the closest thing to not used that a store can know, and
// it is close enough for the reason that matters: a server persists a
// document somebody is in every [Config.PersistEvery], so a document with
// anybody in it is written constantly and never looks idle. What looks idle
// is what nobody has opened since it was last put away.
Idle(ctx context.Context, d time.Duration) ([]string, error)
// Release forgets a document, but only if this store still holds exactly
// want. It reports [ErrChanged] if it holds something else, and nil if it
// held nothing.
//
// The condition is the whole of it. An archiver copies a document
// elsewhere and then asks for it to be released, and between those two
// moments somebody may have joined the document and saved a newer one. A
// store that released it anyway would delete a version that is nowhere
// else. The comparison and the removal therefore happen together, under
// whatever the store uses to keep its own writes apart.
Release(ctx context.Context, document string, want []byte) error
}
An Archivable store can say which of its documents have gone quiet and let one go, which is what it takes to be the hot half of a Tiered.
Both methods are here rather than in Store because most stores have no business implementing them: a store that is only ever written through is complete with Load and Save, and asking every one of them for a way to forget a document would be asking for a way to lose one.
type Capabilities ¶ added in v0.38.0
type Capabilities map[Capability][]byte
Capabilities is what a peer says it understands: for each capability, every version of it the peer accepts.
A version set rather than a highest version, because they are not always contiguous: crdt reserves version 7 of a text for work on another branch and refuses it, so a build reading 1 to 6 and 8 told a peer "up to 8" would be sent a 7 and refuse it.
func Mine ¶ added in v0.38.0
func Mine() Capabilities
Mine is what this build understands, taken from crdt rather than from a list here that would have to be kept in step with it.
func (Capabilities) Accepts ¶ added in v0.38.0
func (c Capabilities) Accepts(name Capability, version byte) bool
Accepts reports whether the peer said it understands this version of this capability.
A capability the peer never named is not accepted. Saying nothing about something is not the same as accepting it, and treating it as acceptance is how a negotiation becomes a formality -- the peer that says nothing at all is exactly the old build this exists to protect.
func (Capabilities) MarshalBinary ¶ added in v0.38.0
func (c Capabilities) MarshalBinary() ([]byte, error)
MarshalBinary writes what a peer says it understands.
Canonical: capabilities by name ascending, versions ascending within each, neither repeated and neither empty. Two peers that understand the same things therefore say so in the same bytes, which is what lets a test compare an advertisement with itself and a reader refuse anything that could have been written two ways.
A capability with no versions is left out rather than written empty. Saying "I understand this, in no version" is saying nothing, and a wire with two ways to say nothing has one too many.
func (*Capabilities) UnmarshalBinary ¶ added in v0.38.0
func (c *Capabilities) UnmarshalBinary(data []byte) error
UnmarshalBinary reads what a peer said, and refuses anything that is not what this would have written.
A peer naming a capability this build has never heard of is not an error: it is kept and answered about with "not accepted", which is the whole reason these are named rather than numbered. Only a malformed encoding is refused.
type Capability ¶ added in v0.38.0
type Capability string
A Capability names something a peer may understand. Named rather than numbered so that a peer meeting one it has never heard of can say nothing about it and carry on, which is what makes this extensible without a registry everyone has to agree on first.
const ( CapText Capability = "snapshot.text" CapList Capability = "snapshot.list" CapMap Capability = "snapshot.map" CapComposite Capability = "snapshot.composite" // CapDigest says this build reads a digest block on a welcome, and is what // lets one be sent without a flag day: a server appends one only to a peer // that asked for it, so a peer built before this never sees a message it // would refuse. // // It is not a snapshot format and carries no format versions. Its value is // the version of the COMPARISON -- what is hashed and in what order -- which // is the thing that has to match for two digests to mean anything, and which // crdt owns rather than this package. CapDigest Capability = "digest.welcome" )
The capabilities this build knows how to talk about. Each is a snapshot encoding, and each moves at its own pace: a change to how a text is written does not disturb a map.
This list will grow. The next entry is already known -- a text operation kind a peer refuses if it does not understand it (crdt#80) -- and it is the reason this says capabilities with versions rather than one number for the snapshot format. A mechanism built for the one and stretched to the other afterwards is how a protocol ends up with two of everything.
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
A Client is one participant's view of a document: a replica that edits locally and is kept in step with everyone else.
It is safe for concurrent use. It builds for js/wasm, so a browser tab runs this code and the server's merge logic unchanged.
func Join ¶
Join opens a session over transport and returns once the document has arrived, so the client is usable the moment it is returned.
Use WebSocket unless there is a reason not to; it is what the same code compiled for a browser can afford. GRPC is there for a native peer that wants what gRPC brings with it.
The session lives until ctx is cancelled or Client.Close is called.
func JoinWithRetry ¶ added in v0.28.0
func JoinWithRetry(ctx context.Context, dial Dialer, cfg ClientConfig, policy RetryPolicy) (*Client, error)
JoinWithRetry joins a document and keeps a participant in it, opening a new session each time the one it has ends.
Why this exists ¶
Config.Backlog says a participant that falls too far behind is disconnected and "rejoins and is caught up from its version vector". That is what the protocol supports and, until now, nothing did it: the material was here — a replica that survives the session and a join that says what it already holds — and the loop was left to whoever held the session, the same way it was for Server.Follow before Server.FollowWithRetry.
Leaving it there turned out to matter. A document busier than the backlog disconnects everyone in it in one instant, however well they are keeping up: one edit by each of P participants is P-1 messages into every other queue. Measured on one server with 800 participants editing at once, a backlog of 256 disconnected 99% of them — and an application without this loop leaves 99% of a room staring at a document that has stopped moving.
What it does that Join does not ¶
The client it returns is the same Client, and the handles taken from it stay valid across every reconnection: a handle holds a name and looks the part up under the client's lock, which is what makes replacing the session underneath it invisible.
Each attempt rejoins with what this replica already holds, so the server sends the difference rather than a snapshot, and nothing edited while disconnected is lost — it is pushed as soon as there is somewhere to push it.
An edit made while there is no session does not fail. It is in the replica, and telling an editor that their keystroke failed when it did not is how an application starts undoing work that was never lost. The outage is reported through RetryPolicy.Notify instead, which is where an operator looks.
What it does not do ¶
It does not resume presence: a cursor is ephemeral, and one from before an outage is a guess about where somebody was. It publishes nothing of its own on reconnection beyond the operations the server is missing.
Client.Done closes when this gives up for good — because the context ended, because Client.Close was called, or because RetryPolicy.Permanent said so — and not when a session ends.
func (*Client) Changes ¶
func (c *Client) Changes() <-chan struct{}
Changes receives a value whenever the document or the participants changed. It coalesces: a reader that is slow sees one wake-up, not a queue of them.
func (*Client) Close ¶
Close ends the session. The local document is left intact, so its Client.Snapshot can resume later.
func (*Client) Done ¶
func (c *Client) Done() <-chan struct{}
Done is closed when the session has ended, whatever the reason.
func (*Client) Err ¶
Err returns why the session ended, or nil while it is still running. Once Client.Done is closed it is never nil: a session that was closed deliberately reports ErrClosed rather than the transport's cancellation.
func (*Client) Parts ¶ added in v0.10.0
Parts returns the parts this replica holds, in the canonical order. A part that has never been written to is not among them.
func (*Client) Peers ¶
Peers returns the other participants and where their cursors are, ordered by site.
func (*Client) SetCursor ¶
SetCursor publishes where this participant is. meta carries whatever the editor wants shown — a display name, a colour — and is not interpreted here.
Cursor positions are ephemeral and are never persisted.
func (*Client) Snapshot ¶
Snapshot returns the document in a form ClientConfig.Resume accepts, which is how a participant keeps its place across a disconnection.
func (*Client) TakeChanges ¶ added in v0.7.0
func (c *Client) TakeChanges() []crdt.PartChange
TakeChanges returns the edits made by everyone else since it was last called, in the order a view of the text has to make them, and forgets them.
It pairs with Client.Changes: that says something happened, this says what. A view that only ever applies these holds what the document holds — see crdt.Change.
Local edits are not reported. A caller that made them already knows.
func (*Client) Text ¶
Text returns a handle on the text part with this name, which is created the first time anybody writes to it. The name is arbitrary UTF-8 and is expected to carry structure — "file:src/main.tex". An empty or invalid name is refused; see crdt.Part.
func (*Client) Version ¶
func (c *Client) Version() crdt.CompositeVersion
Version returns what this participant holds, for ClientConfig.Resume or for diagnostics.
type ClientConfig ¶
type ClientConfig struct {
// Document names the document to join. It is created if it does not exist.
Document string
// Site is this participant's replica identity, and must differ from every
// other participant's in the document. See [crdt.DeriveSiteID].
Site crdt.SiteID
// Resume is a snapshot from an earlier session, obtained from
// [Client.Snapshot]. When set, the participant keeps the work it did while
// disconnected and is sent only what it missed, rather than the whole
// document.
//
// Come back as the site that wrote it. These operations carry the Site the
// participant had when it made them, so a server running [OwnSiteOnly] --
// or any policy about who may speak for whom -- refuses them under a fresh
// identity, and refuses them quietly enough that the tab goes on showing
// work nowhere else has.
//
// A server that can no longer say what this snapshot missed refuses the
// join rather than answering with a history that has a hole in it: that is
// a document purged past this version, and the error carries
// [crdt.ErrPurged]. Nothing is lost by it — these bytes are still the
// caller's — but the participant cannot be caught up from that server and
// has to be reseeded from it instead.
Resume []byte
// MaxOperations bounds how many operations ONE message from the server may
// ask this participant to reserve. Zero, the default, is no bound.
//
// It is the mirror of [Config.MaxOperations], and it exists because the two
// ends of a session are not always the same kind of thing. A participant
// joining a server its own operator runs is trusting that server with far
// more than a message size -- it reads every document, holds the store and
// decides who may join -- so a bound there protects against nothing.
//
// The peer-to-peer carriers are the case this is for. Over [DataChannel] the
// far end is ANOTHER PERSON's browser, and over [JoinBroadcastChannel] it is
// whichever tab won the election. The page that hosts builds a [Server] and
// can bound what it is sent; without this the page that joins could not, so
// one session had a bound in one direction only.
//
// The unit is operations rather than bytes, and that is the whole of why it
// is worth a field. crdt allocates a flat 80 bytes an operation -- 6.2 times
// a realistic single-character insert and 20 times the four-byte floor one
// can encode in -- so what a message costs is decided by how many operations
// it claims, and a bound written in bytes would be a different bound for
// every peer. A message over the bound is refused with [crdt.ErrTooManyOps]
// and the session ends.
//
// Set it far above an honest catch-up: a document holds roughly one
// operation per character ever typed into it, so a bound a real document
// could reach refuses the session it exists to protect.
MaxOperations int
}
ClientConfig describes a participant joining a document.
type ConditionalStore ¶ added in v0.67.0
type ConditionalStore interface {
Store
// LoadToken is [Store.Load] and also the token naming what it returned, so a
// caller can later ask for a save that only lands while that is still what is
// there. nil and nil is a document nobody has written.
LoadToken(ctx context.Context, document string) ([]byte, Token, error)
// SaveIf records the snapshot only while the store still holds the version
// expect names, and returns the token naming what it wrote.
//
// It reports [ErrChanged] if the store holds something else, which is the same
// answer [Archivable.Release] gives to the same question. A nil expect asks for
// a document that is not there yet, so a second server opening the same new
// document is refused rather than racing.
SaveIf(ctx context.Context, document string, snapshot []byte, expect Token) (Token, error)
}
A ConditionalStore is a Store that can refuse a save over a document somebody else has written since it was read.
Store.Save replaces, which is what makes two servers over one store lose the earlier save — measured in TestTwoServersOverOneStoreLoseTheEarlierSave, where the store ends holding one server's work and no error is returned to anybody. This is the interface that lets a store say no instead, and it is optional because most stores have no way to: saying no means comparing before writing, under whatever the store uses to keep its own writes apart.
It does not make a shared store a supported way to run two servers. The document each server holds in memory is still not shared, so the second one's replica is still missing what the first one applied; what changes is that the loss is reported rather than silent, which is the difference between an operator finding out and not. See the package documentation on the two shapes that are supported.
What it costs, measured ¶
The condition has to be cheap or it is not worth having, and the shape of it decides that. Compared on PostgreSQL 18.6 with pgstore's own table, medians of nine rounds, each updating a row that is already there:
snapshot blind save conditioned on a token conditioned on md5(snapshot) 1 MiB 2.878ms 3.143ms 5.291ms 8 MiB 13.287ms 12.791ms 23.022ms 32 MiB 51.349ms 50.261ms 88.467ms
A token the store already has costs nothing measurable — the arm carrying it sits inside the blind save's noise, and it was the pessimistic arm, paying an extra read per round that a real caller would not. Hashing the stored snapshot instead nearly doubles every save, because it is not the hash that costs but detoasting a blob to feed it. That is why this is a token and not a digest.
type Config ¶
type Config struct {
// Store keeps documents between sessions. Defaults to a [MemoryStore].
Store Store
// Backlog is how many messages may be queued for one participant.
// A participant that falls further behind than this is disconnected with
// ResourceExhausted rather than served stale state or allowed to stall
// everyone else; it rejoins and is caught up from its version vector.
// Defaults to [DefaultBacklog].
//
// Size it against the number of participants who may edit at once, not
// against how fast one of them types: one edit by each of P participants is
// P-1 messages into every queue, so a burst on a busy document reaches the
// limit in an instant however well everyone is keeping up. See
// [DefaultBacklog] for what that costs, measured.
//
// Note what "it rejoins" asks of a caller, because it is not free.
// [JoinWithRetry] does it: it opens a session again whenever the one it has
// ends, rejoining with what its replica already holds so that nothing
// edited in between is lost. A caller that uses [Join] instead gets one
// session, and a burst on a busy document then becomes a disconnection
// nobody recovers from.
Backlog int
// PersistEvery, when set, saves every document that has changed at this
// interval, whoever is connected. Without it a document is saved when its
// last participant leaves and when [Server.Flush] is called, so a server
// restarted while anybody was still editing loses everything since the
// document was opened.
//
// It bounds what a crash costs to this interval, which is a number an
// operator can choose. A server that sets it must be closed with
// [Server.Close], which stops the housekeeping and saves what is left.
PersistEvery time.Duration
// EvictAfter, when set, persists a document nobody has been in for this long
// and lets go of it. Without it a long-lived server holds every document it
// has ever served.
//
// A document is reloaded from the store the next time somebody joins it, so
// evicting costs a read rather than anything anybody wrote.
EvictAfter time.Duration
// CollectEvery, when set, gives back what every participant of a document
// has certainly seen: the map tombstones nobody can be confused by any
// more. See [crdt.Map.Collect].
//
// The map parts, and only those. A text and a list had a collection too and
// it was withdrawn in crdt v0.35.0, because it re-pointed survivors at
// something relative to what each replica happened to hold and left two of
// them holding different documents. So a document of text gives nothing
// back here yet.
//
// It is off by default, and being off is not a failure of nerve. Collecting
// asks for a version every replica has delivered, and a server can only
// know one because participants tell it what they hold. If any of them has
// gone quiet, the answer is nothing and nothing is collected — which is the
// behaviour the literature calls for, and the reason it is acceptable is
// that collecting affects size and never correctness:
//
// When these requirements are not met, GC may block. We consider this to
// be acceptable, as GC does not impact correctness (only performance), and
// the normal operations in the object's interface remain live.
// — Shapiro, Preguiça, Baquero and Zawirski, "A comprehensive study of
// Convergent and Commutative Replicated Data Types", §4.1
//
// A map costs nothing to collect beyond the asking: it re-points nothing,
// it gives back a sequence number without its operation as a matter of
// course, and a peer catches up either way. There is no participant to
// refuse and none to re-seed.
//
// In a federation it works, and what a link promises is the whole of why.
// A link is a participant of both documents — it joins the server it
// follows, and [Server.Follow] joins its own — so for a while it was a
// participant that never said anything, and neither end could collect at
// all. It says something now, and what it says is not what its replica
// holds but the meet its own server could collect against: a peer may give
// back only what everybody behind that link has certainly seen.
//
// The difference is somebody's work. Promising the replica's version let a
// peer collect past a participant that was away on the follower — and
// federation exists so that such a participant can come back on the peer,
// where it was then answered with a superseded run, went on showing a value
// everybody else had removed, and had a version equal to the peer's, so no
// rejoin would ever repair it.
//
// So a quiet participant anywhere holds the whole federation's collection
// back, which is this same paragraph's rule one hop further out. See
// TestAFederatedServerCollectsOnceTheLinkSaysWhatItHolds and
// TestAParticipantBehindALinkIsNotCollectedPast, which are the two halves.
//
// All of which holds only for the version it is asked with. "Gone quiet"
// above means gone quiet, not gone off the air: a participant whose carrier
// dropped is out of the open sessions within milliseconds and is exactly the
// one that has not delivered what the others have. So the meet is taken over
// every site this document has seen, present or not — see
// [document.collectable] — and a participant that never comes back holds
// collection back for good. That is the shape of the trade the passage above
// describes, and it is the one this server takes: giving back space is worth
// nothing beside two participants who will never agree again.
CollectEvery time.Duration
// OnEvictError, when set, is told about a document that could not be saved
// as it was evicted. There is nobody left to return an error to, and the
// document cannot be kept — a session may already have opened a fresh
// replica of it — so this is the only place that failure can be seen.
OnEvictError func(document string, err error)
// OnPersistError, when set, is called for every document whose periodic
// save failed, with the error the store gave.
//
// Without it a server that cannot write is silent about it. The saves go on
// being attempted and go on failing, participants go on editing and are
// told nothing, and the work is there until the process stops and then is
// not. A disk that filled up, a name a filesystem will not take, a
// credential that expired: all of them look exactly like a server that is
// working, which is the worst way for durability to fail.
//
// It is called from the housekeeping goroutine, once per document per pass,
// so an implementation that blocks delays the next pass. Counting or
// logging is what it is for; recovery is the operator's.
//
// [Server.Flush] does not call it — a caller that asks for a flush is given
// the error to handle.
OnPersistError func(document string, err error)
// MaxOperations, if set, is the most operations one message may claim, and a
// message claiming more is refused before anything is reserved for it. Zero
// is no limit, which is what this package did before the field existed.
//
// It is a second bound beside Backlog and the transport's own message size,
// and it bounds a different thing: how much a message may ask this server to
// RESERVE. The counted headers in [github.com/go-crdt/crdt] already refuse a
// claim larger than the bytes that follow it, so what is left is a message
// that is honest and enormous -- the worst claim they permit still reserves
// sixteen to twenty-four times the input, that being the ratio of a record in
// memory to the smallest encoded one. A gibibyte of operations is therefore
// sixteen to twenty-four gibibytes, for an honest sender as much as a hostile
// one.
//
// The number is here rather than in crdt for the reason crdt gives: a ceiling
// chosen by a library is a guess about somebody's machine. This is the layering
// HPACK uses, where the decoder offers the knob and the server turns it.
//
// What it costs when it fires: the message is refused whole, the session is
// told, and [Config.OnOperationsRefused] is called. It is not a disconnection
// on its own -- a binding decides that, as it does for any refused batch.
//
// Sizing it is arithmetic rather than taste. Multiply by the in-memory record
// size, which is 80 to 96 bytes, and compare against what this server may
// spend on one message: a million operations is 80 to 96 MB.
//
// See go-crdt/collab#169.
MaxOperations int
// OnOperationsRefused, when set, is called for every batch a document
// refused, with the site of the session that sent it and the reason.
//
// Without it a refusal is told to the offending session and to nobody else,
// which is the wrong room for two of the three reasons a batch is refused.
// A client that sends rubbish deserves the error it gets and no more. The
// other two are the operator's business:
//
// - [Config.AuthorizeOperations] refused it. On a federating server that
// is a peer carrying sites it may not speak for, and the operator is the
// only one who can take that up with the other operator.
// - [github.com/go-crdt/crdt.ErrCollidingID]: an operation wore the name
// of one this replica had already applied and said something else. Two
// replicas chose the same site, by accident as easily as on purpose, and
// an accident here means the identities a deployment hands out are not
// unique -- which nobody can discover from inside a session.
//
// The error wraps the cause, so an implementation asks errors.Is rather than
// reading the wording.
//
// It is called with no lock held, after the batch has been dealt with, so an
// implementation that blocks delays only the session it belongs to and not
// everybody editing the document. Counting or logging is what it is for.
OnOperationsRefused func(document string, from crdt.SiteID, err error)
// Clock is what [Config.EvictAfter] measures with. It defaults to time.Now,
// and exists because a caller that wants a monotonic source, or a test that
// wants to reach an hour of idleness without waiting an hour, has nowhere
// else to say so. It is read from more than one goroutine, so it must be
// safe for concurrent use and must be given here rather than set afterwards.
Clock func() time.Time
// Authorize, when set, decides whether a participant may open a document.
// It is asked once per session, after the join arrives and before the
// document is touched, so a refused session neither reads the store nor
// reveals whether the document exists.
//
// This belongs here rather than in a gRPC interceptor, which is where one
// would first look for it: an interceptor sees the method and the request
// metadata, and the document being joined is in neither — it arrives in the
// stream's first message. Anything deciding per document has to run after
// that message, which means here. Authentication, which is per connection
// rather than per document, still belongs in an interceptor; ctx carries
// whatever it put there.
//
// Returning a gRPC status error passes that status to the participant
// unchanged; any other error is reported as PermissionDenied.
Authorize func(ctx context.Context, document string, site crdt.SiteID) error
// AuthorizeOperations, if set, is asked about every batch a session sends,
// and refusing ends the session.
//
// Authorize runs once, when somebody joins, and decides whether that site
// may be in this document. That is the whole story for an HONEST
// participant: it speaks for itself, and the site it joined as is the site
// its operations carry.
//
// It is not enforced. A session may hand over operations another site made,
// and with no policy here the server applies them: measured, a session
// joined as site 1 wrote a document as site 2 and advanced the server's
// version vector for site 2, which never joined. That is not only wrong
// attribution -- a site identity is half of an operation's name, so two
// writers on one identity produce different characters with the same ID and
// two replicas given both silently hold different text. [OwnSiteOnly] is
// the policy that refuses it, in one line, for a deployment that does not
// federate.
//
// It is not the whole story for a link. [Server.Follow] joins as one site
// and then relays the work of everyone on the server it follows, so what
// arrives on a link names sites this server never authorised — thousands of
// them, belonging to an institution rather than to a person. Inside one
// deployment that is exactly right and there is nothing to decide. Between
// two, it is the decision: whether this link may speak for those sites.
//
// from is the site the session joined as. batches carry the operations it
// is asking to add, each naming the site that made it, so a policy can be
// written about the relationship between the two — "this link may carry
// operations for sites derived within lyon.ac.example" — which is what an
// interfederation has scopes for.
//
// It runs after the operations have been decoded and before any of them is
// applied, so a refused batch changes nothing. Returning a gRPC status
// error passes that status on unchanged, as Authorize does.
AuthorizeOperations func(ctx context.Context, document string, from crdt.SiteID, batches []crdt.PartOps) error
}
Config configures a Server.
type Dialer ¶ added in v0.24.0
A Dialer produces a route to the peer, one per attempt.
It is a function rather than a Transport because a Transport is a way to reach a peer and an attempt needs a fresh one: a gRPC client connection whose server has gone stays broken, a WebSocket that closed cannot be reopened, and a link that held on to the first one would redial nothing. What is stable across attempts is the knowledge of where the peer is and how to authenticate to it, and that is what a closure holds.
It is called on the link's own goroutine, once per attempt, with the context the link was given, so a dialler that blocks blocks the retry loop and a dialler that respects ctx is what makes cancellation prompt during a dial. Returning a transport it has already returned is allowed where the transport itself redials, as WebSocket does.
type DirStore ¶ added in v0.13.0
type DirStore struct {
// contains filtered or unexported fields
}
A DirStore keeps documents as files in one directory. It is what a server wants when the documents belong to whatever else is on that disk — a project whose files are already there, backed up with them and restored with them — and it needs nothing running beside it.
A file per document, named after nothing ¶
A document name is arbitrary UTF-8 and is expected to carry structure, so the names a real consumer uses are "project:default" and "project:ods:chapter one.ods". Those are not file names: a colon is a path separator on one system this package supports, a slash is one everywhere, and "." and ".." name something else entirely. Escaping the awkward characters would leave the question of which ones, on which system, and a name that escapes to the same file as another is two documents sharing a snapshot.
So the file is named after the encoding of the name rather than the name: base64, in the alphabet made for file names, which is total and reversible and has no awkward character in it. It is unreadable at the shell, which is what DirStore.Documents is for.
What a reader may see ¶
A snapshot is written to a temporary file and renamed over the old one, so a reader sees the whole of one version or the whole of the one before. A crash during a save leaves the previous snapshot intact and a temporary file behind; the next NewDirStore on that directory clears those away.
func NewDirStore ¶ added in v0.13.0
NewDirStore returns a store keeping documents in dir, creating it if it is not there, and clears away any temporary file a previous run left behind.
func (*DirStore) Documents ¶ added in v0.13.0
Documents returns the names of the documents held, which is what a caller needs to inspect a store whose file names are an encoding rather than a name.
A file whose name is not one this store wrote is skipped rather than reported: a directory shared with anything else would otherwise turn every stray file into an error nobody can act on.
func (*DirStore) Idle ¶ added in v0.26.1
Idle returns the documents whose file has not been written for longer than d, which is a DirStore's half of Archivable.
func (*DirStore) Load ¶ added in v0.13.0
Load returns the snapshot for a document, or nil if there is none yet.
func (*DirStore) LoadSites ¶ added in v0.36.0
LoadSites returns what SaveSites last wrote, or nil for a document nobody has been recorded in. See SiteStore.
func (*DirStore) Release ¶ added in v0.26.1
Release forgets a document if this store still holds exactly want.
The comparison and the removal are one step under the same lock as a save, so a document written while it was being archived is not deleted by the release that follows.
func (*DirStore) Save ¶ added in v0.13.0
Save records the snapshot, replacing any previous one.
It writes a temporary file, flushes it, and renames it over the old one, so that a reader sees one whole version or the other and never half of either. On every system this package supports, a rename within a directory replaces the destination in one step.
func (*DirStore) SaveSites ¶ added in v0.36.0
SaveSites records the participants, replaced whole and atomically, exactly as a snapshot is: a half-written one would be read back as unreadable and take the document down with it.
The bytes are stored as they arrive, and not for want of anything to compress: measured, brotli takes a two-participant blob from 57 bytes to 36, and a sixty-four-participant one from 15233 to 362. It is a set of version vectors over the same site ids, which is the most repetitive thing in the store. They are stored as they arrive because [encodeSites] promises that two servers holding the same thing write the same bytes, and a compressor inside that promise weakens it into "as long as nobody upgraded" — which is why dircompress.go exists on the storage side rather than in crdt.
Their integrity is not this store's business either. What arrives already carries its own magic and its own CRC32C, put there by the encoder that owns the format, so a participants file that rots is refused wherever it was kept rather than only here.
type GRPCServer ¶ added in v0.18.0
type GRPCServer struct {
collabpb.UnimplementedCollabServer
// contains filtered or unexported fields
}
GRPCServer presents a Server as the generated gRPC service.
It exists because the Server itself no longer does. The session logic speaks the wire format in wire.go — four small types, hand-written — and that is what let it stop depending on the generated protobuf code. The reason is a measurement: compiling the server for the browser with protobuf attached takes the WebAssembly binding from 5.3 MB to 19.3, because gRPC and protobuf come with it. A browser holding a document for a colleague on another continent cannot pay that, and it is the same reason wire.go exists at all.
So gRPC is a binding rather than a foundation: this type converts, and the document logic never sees a protobuf message.
func GRPCService ¶ added in v0.18.0
func GRPCService(s *Server) *GRPCServer
GRPC presents a Server over gRPC. Register the result with collabpb.RegisterCollabServer on a grpc.Server built with GRPCServerOptions -- without them the server keeps gRPC's four-mebibyte receive limit and refuses exactly the messages this protocol is largest in.
func (*GRPCServer) Session ¶ added in v0.18.0
func (g *GRPCServer) Session(stream collabpb.Collab_SessionServer) error
Session is the service method: one bidirectional stream, one participant, one document.
type LinkStatus ¶ added in v0.24.0
type LinkStatus struct {
// Up is true for the one report that says the link is established and the
// local replica has been caught up. The rest of the fields are zero.
Up bool
// Err is why the attempt ended, for a report that is not Up.
Err error
// Attempt counts the consecutive failures since the link was last up, so
// the first report of an outage carries 1. It is what tells a second
// failure from a hundredth.
Attempt int
// RetryIn is how long the link will wait before trying again — the jitter
// already applied, so it is the real interval and not the policy's idea of
// it.
RetryIn time.Duration
// DownFor is how long this outage has lasted: zero on the first report,
// growing with each one. It is the number a health check thresholds on,
// because "down for six seconds" and "down for six hours" want different
// people woken up.
DownFor time.Duration
}
A LinkStatus is what a reconnecting link tells RetryPolicy.Notify: whether it is up, and if it is not, why, since when and for how much longer.
It answers the two questions an operator has about a link that is not working — is it down, and how long has it been down — without their having to keep the state themselves, because the obvious mistake is to report only the failure and leave "still failing" indistinguishable from "failed once".
type List ¶ added in v0.10.0
type List struct {
// contains filtered or unexported fields
}
A List is a handle on one list part — comments, a change log, the messages beside a document.
func (*List) Append ¶ added in v0.10.0
Append adds values after the last one, which is what a chat or a log does.
func (*List) Insert ¶ added in v0.10.0
Insert adds values at index pos, locally and then everywhere.
func (*List) Values ¶ added in v0.10.0
Values returns copies of every value present, in order. It is what a view of a list reads when it is told the list changed; see crdt.PartChange.
type Map ¶ added in v0.10.0
type Map struct {
// contains filtered or unexported fields
}
A Map is a handle on one map part, such as the cells of a sheet.
func (*Map) Delete ¶ added in v0.10.0
Delete removes key, locally and then everywhere. It writes a tombstone whether or not this replica holds the key; see crdt.Map.
func (*Map) Edit ¶ added in v0.29.0
Edit runs fn against the map part itself and sends whatever operations it produced to everyone else.
It is what a structured type built on a map is driven through. A github.com/go-crdt/crdt/structured.Sequence, a Tree or a RecordMap is a binding over a crdt.Map rather than a thing of its own, and each of its methods changes the map here and hands back the operations that change:
err := m.Edit(func(mp *crdt.Map) ([]crdt.MapOp, error) {
_, ops, err := structured.SequenceOf(mp).Insert(after, value)
return ops, err
})
The lock is held for the whole of fn, so what it reads and what it writes are one moment. fn must not call back into the client.
func (*Map) Get ¶ added in v0.10.0
Get returns a copy of the value at key, and whether the key is present. It is what a view reads for each key a crdt.PartChange names.
func (*Map) Len ¶ added in v0.10.0
Len returns how many keys are present, not counting deleted ones.
func (*Map) Read ¶ added in v0.29.0
Read runs fn against the map part for reading. It sends nothing, and anything fn writes to the map would be this replica's alone — so do not: use Map.Edit.
The lock is held for the whole of fn, so a view built from several reads sees one moment rather than a moment per read.
type MemoryStore ¶
type MemoryStore struct {
// contains filtered or unexported fields
}
MemoryStore keeps documents in memory. It is the default, it is what the tests use, and it is enough for a single process that does not need to survive a restart. Anything else — Postgres, object storage — implements Store.
func (*MemoryStore) Documents ¶
func (s *MemoryStore) Documents() []string
Documents returns the names of the documents held, which is what a caller needs to inspect or migrate a store.
func (*MemoryStore) Idle ¶ added in v0.26.1
Idle returns the documents this store has not been asked to write for longer than d, which is a MemoryStore's half of Archivable.
func (*MemoryStore) Load ¶
Load returns a copy of the stored snapshot, or nil if the document is new.
func (*MemoryStore) LoadSites ¶ added in v0.36.0
LoadSites returns a copy of what SaveSites last wrote, or nil. See SiteStore.
func (*MemoryStore) Release ¶ added in v0.26.1
Release forgets a document if this store still holds exactly want.
The comparison and the removal are one step under the same lock, so a save that lands while a document is being archived cannot be deleted by the release that follows it.
type MultiStore ¶ added in v0.24.0
type MultiStore struct {
// contains filtered or unexported fields
}
A MultiStore keeps every document in several stores at once.
What it is for ¶
The stores in this project answer different questions. A database answers "what is the document now", quickly, which is what a server restarting needs. A git repository answers "what did it say last Tuesday, and who wrote this sentence", which is what a person needs. Neither answers the other's question, and an operator who wants both has until now had to choose.
The alternative already in use is worse than choosing: writing to one store from the server and to the other from somewhere else — a browser, a sync job — which is two sources of truth held together by whichever of them happens to run last.
Reading merges rather than picking ¶
[Load] reads every store and merges what they return. The obvious design is to read the first store that has the document and stop, and it is wrong here for a reason particular to this problem: a save that failed halfway leaves the stores holding different documents, and reading only the first would quietly drop whatever only the second had. Merging is the only answer that loses nothing, and a CRDT is what makes it available.
It has a consequence worth having on purpose: adding a store to a running server backfills it. The new store returns nothing, the merge is the other store's document unchanged, and the next save writes it across.
The cost is that opening a document reads every store instead of one, and merges when more than one has content. That is paid once per document, when it is opened, and not per edit.
A store that cannot be read makes the document unavailable ¶
If any store fails to read, [Load] fails. It does not fall back to the stores that answered, because what came back would be a document that is missing whatever the unreadable store alone held — and the next save would then write that shortened document over the store that was merely unreachable. Serving a document that is quietly missing a paragraph is worse than serving none: an error stops at one document and an operator can fix it, while silent loss is discovered by the person who wrote the paragraph.
A merge that refuses fails the same way and for the same reason. Two stores that have each discarded what the other still needs give ErrUnmergeable, and one that holds a write the other can no longer accept gives crdt.ErrStranded; either way the document does not open, rather than opening as whichever of the two [Load] happened to reach first. It takes a store left behind by a purge or a collect to reach that at all — stores written together hold the same bytes, and merging those carries nothing.
Reading stops at the first store that refuses, and that is on purpose ¶
A member whose Load fails -- a file that rotted, a database that is down -- fails the whole read. Redundancy here buys durability, not availability: one unreadable replica takes the document down even though a good copy is beside it, and an operator meeting that has to repair or remove the bad store rather than wait for a failover that is not coming.
Serving the members that did answer would be worse than it looks. A snapshot is a set of operations and the merge is their union, so a member left out is not a smaller document, it is a document missing whatever only that member held -- and the [Save] that follows writes the merge of the others over it, which makes the loss permanent. Refusing keeps the operator's options open; answering closes them silently.
Writing tries every store, and fails if any refused ¶
[Save] writes to all of them even after one has failed, so that a store being down does not stop the others from being written, and then reports every failure together. It returns an error if any store refused, because a caller that gets nil back has to be able to believe the document is durable in all of them.
They are written one after another rather than at the same time because there are two of them, not two hundred.
It does not implement SiteStore. Go has no way to implement an interface only when what is underneath does, so a MultiStore that declared the methods would keep nothing whenever its stores could not — silently, which is worse than not offering it. A server given one falls back to the participants a document names, which is everyone who has written; see SiteStore for what that costs.
func NewMultiStore ¶ added in v0.24.0
func NewMultiStore(stores ...Store) *MultiStore
NewMultiStore returns a store that writes to all of the given stores and reads from all of them.
It panics if given none: a store that silently keeps nothing would look like it was working, and there is no configuration in which that is what somebody meant. One is allowed, and behaves as that store does.
type PipeConn ¶ added in v0.21.0
type PipeConn struct {
// contains filtered or unexported fields
}
PipeConn is the server end of a Pipe, the side Server.ServePipe serves. It is an opaque handle: the session is spoken through it, not by the caller.
type RetryPolicy ¶ added in v0.24.0
type RetryPolicy struct {
// Wait is how long to wait before the first attempt after a link drops.
// Each further failure doubles it, up to Ceiling. Defaults to
// [DefaultRetryWait]; it may not be negative, and may not exceed Ceiling.
Wait time.Duration
// Ceiling is the longest the link will ever wait between attempts.
// Defaults to [DefaultRetryCeiling]; it may not be negative.
Ceiling time.Duration
// Permanent, when set, is asked about the error that ended each attempt and
// stops the link by returning true — the error it was asked about is then
// what [Server.FollowWithRetry] returns.
//
// It exists because this package cannot answer the question honestly for
// everybody. The errors that are genuinely permanent are refused before the
// loop is ever entered, and everything that can then end an attempt came
// off a network or off a peer, where "permanent" is a judgement about a
// deployment and not about an error value. A peer refusing the link is the
// case that decides the shape: it is a policy decision, policy is edited
// and credentials are rotated, so a link that gave up on it would need
// somebody to notice and restart a process — while a link that keeps asking
// costs one attempt per Ceiling, which is nothing. So the default is to
// retry it, and an operator who disagrees says so here, with
// [errors.Is] over whatever their peer returns.
Permanent func(error) bool
// Notify, when set, is told every time the link changes state: down, with
// why and for how long, and up again. See [LinkStatus].
//
// A library has no business choosing where that goes. Writing to stdout
// would put a federation link's troubles into the middle of whatever the
// process's own output is, in a format nobody asked for, and a link that
// says nothing at all is a link nobody can operate — an outage would be
// visible only as a replica that had quietly stopped converging. So it is
// handed over instead, and the operator's logger, metric or health check
// decides.
//
// It is called on the link's own goroutine, in order, and blocking in it
// blocks the link — including the moment it is trying to come back. It must
// not call back into the link.
Notify func(LinkStatus)
}
A RetryPolicy is how a link comes back, which Server.Follow says is the operator's to decide. This is where they decide it.
The zero value is a working policy: DefaultRetryWait doubling to DefaultRetryCeiling, jittered, retrying everything, telling nobody.
type Role ¶ added in v0.26.0
type Role int
A Role is which side of the protocol a tab took, decided by [electRole].
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
A Server hosts documents. Register it with collabpb.RegisterCollabServer on any grpc.Server — over github.com/grpc-transports/websocket for browsers, over plain TCP for anything else.
Documents stay in memory once opened, so a long-lived server holds every document it has served. Call Server.Flush to persist them.
func (*Server) Close ¶ added in v0.13.0
Close stops the housekeeping Config.PersistEvery and Config.EvictAfter ask for, and saves everything that has changed. It does not end the sessions in progress: those belong to whatever is serving them, and stopping that is the caller's to do first.
Calling it twice is harmless. A server that asked for neither still has one, so a caller need not know which kind it configured.
func (*Server) Flush ¶
Flush persists every document that has changed since it was last written. A server that wants durability without waiting for participants to leave calls this on a timer, or before shutting down.
func (*Server) Follow ¶ added in v0.17.0
Follow makes this server a participant in another server's copy of a document, so that the two converge.
What it is for ¶
Not capacity. One server holds a document for a thousand participants at about three kilobytes and two and a half microseconds each, flat, which is twelve percent of a core for a document five people are typing in — see BenchmarkFanOut. A second server earns its place for two other reasons: a participant far from the first pays the round trip on every keystroke echo, and a site that goes down takes its documents with it until somebody brings them back.
Both are answered by a replica near each participant rather than by splitting a document across servers, which is what this is. It is also what the CRDT is for: two replicas that have seen the same operations hold the same document, in any order, with no agreement about the order and nothing to coordinate on the write path. There is no leader here and no consensus, which is why it works between datacentres without paying a round trip per edit.
What a link is ¶
A participant. The server being followed cannot tell the difference and does not need to: a link joins its document, is sent what it is missing, and is broadcast to like anybody else. Everything the local document learns is sent out, and everything that arrives is applied and broadcast onwards — to everyone except the link it arrived on, which is the loop prevention the subscriber machinery already had.
That prevention is not enough on its own, and the missing half is in applyOperations: two servers that follow each other would otherwise pass an operation back and forth forever, each applying it harmlessly and telling the other again. Operations that do not advance the version are not passed on.
Per document ¶
A link follows one document. The alternative — a link that mirrors a whole store — is simpler to operate and replicates documents nobody is looking at, which between continents is bandwidth spent on nothing. Idle documents are evicted here already, and a link is what keeps one alive, so the set of documents a server replicates is the set somebody is using.
What this does not do ¶
It does not reconnect. A link that drops stays dropped, and the error is returned to whoever called Follow, because the policy for coming back — immediately, with a backoff, never — belongs to the operator and not to a library. Server.FollowWithRetry does not overturn that: it is one such policy, written down and opted into by an operator whose answer is "with a backoff", and Follow behaves exactly as it did for everyone else. It does not discover peers. It does not replicate presence: cursors are ephemeral and a link that carried them would have to decide what a cursor in another datacentre means when the link is a second behind.
A peer that has purged ¶
A peer that has discarded the characters this replica would need can send it the whole document instead of the difference, and a link with nobody behind it yet takes that and converges — which is how a fresh datacentre follows a document that has been purged. A link into a server that already holds participants cannot: a session already welcomed has no way to be re-seeded. Then this returns, and the error carries crdt.ErrPurged. So does the one from the other direction, a replica that has purged past the peer it was asked to follow.
Neither is worth retrying until somebody reseeds a store, which is why RetryPolicy.Permanent is where an operator answers it: errors.Is(err, crdt.ErrPurged) is the whole test.
func (*Server) FollowWithRetry ¶ added in v0.24.0
func (s *Server) FollowWithRetry(ctx context.Context, dial Dialer, document string, as crdt.SiteID, policy RetryPolicy) error
FollowWithRetry follows document on the peer dial reaches, exactly as Server.Follow does, and re-establishes the link when it drops — waiting longer after each failure, never longer than the policy's ceiling, and never the same interval as anybody else.
It returns when ctx is cancelled, returning ctx's error, or when RetryPolicy.Permanent says an attempt's failure is not worth another, returning that failure. It returns immediately, without dialling, if the call itself is wrong: no dialler, no document name, a link claiming the server's own replica, or a policy that waits for a negative time or longer than its own ceiling.
Why this is here rather than in every caller ¶
Server.Follow's reasoning stands: the policy belongs to the operator. What does not follow from it is leaving everyone to write the loop, because it is the same loop every time and it is usually written twice wrong. Without jitter, every link in a datacentre that lost the same peer waits the same interval and returns together, so the peer coming back up meets the whole fleet at once and goes down again — the retry itself becomes the outage. Without a ceiling, doubling either arrives somewhere absurd or, more often, is bounded by an attempt counter that gives up, and a federation link that gives up is a replica that stops converging while the process it lives in carries on looking healthy.
So the loop is written here, once, and nothing acquires it by accident. Server.Follow is unchanged and remains what a caller with a different answer builds on.
What is retried, and what is not ¶
Everything the loop can see, and that is a deliberate line rather than a shrug. The two failures that are genuinely permanent — a document with no name, and a link claiming [serverSite] — are decided once, before the loop is entered, so they are returned to the caller rather than re-asked forever. After that, an attempt can only end because a dialler could not reach the peer, because a carrier broke, because the peer refused the link or spoke something unexpected, or because ctx ended. The first two are what a link between datacentres does on a normal day. The third is a deployment's business, not this package's, which is what RetryPolicy.Permanent is for. The last ends the loop rather than being retried.
A peer that has purged past this replica is the third kind and the one worth naming, because retrying it can never work: it is answered by reseeding a store and not by waiting. It is not decided here — a purge is a fact about a deployment and can be fixed while the loop runs — so the error reaches RetryPolicy.Notify through LinkStatus.Err carrying crdt.ErrPurged, and what to do about it is RetryPolicy.Permanent's to say. See Server.Follow.
Jitter ¶
The delay is drawn uniformly from the half-open band between half the current interval and the whole of it, rather than from zero to it: full jitter decorrelates best but can draw a delay near zero many times running, which is the hammering this exists to avoid, while a floor of half the interval bounds the attempt rate and still leaves no two links in step.
It is not a knob. An operator has something to say about how long to wait and how stale they will tolerate, and nothing to say about a jitter fraction — but offered the field they would be able to set it to zero, which is the single mistake this whole function exists to prevent.
func (*Server) ServePipe ¶ added in v0.21.0
ServePipe runs one session over the server end of a Pipe, with this server holding the document. It returns when the session ends — because the client end closed, because ctx was cancelled, or because the session itself failed.
It is the in-process counterpart of Server.ServeWebSocket and, in a browser, [Server.ServeDataChannel]: there is no request to upgrade and no origin to check, because there is no boundary to cross. Give it the PipeConn that Pipe returned beside the Transport the local editor joined over.
func (*Server) ServeWebSocket ¶ added in v0.11.0
ServeWebSocket returns an http.Handler that runs sessions over WebSockets — the carrier a browser can afford, and the one WebSocket dials.
Mount it where the browser will reach it. Everything else is the same server: the same documents, the same store, the same Config.Authorize, and a participant here edits the same document as one arriving over gRPC.
origins, when not empty, are the Origin header values allowed to open a session, which is the check that stops another site's page from opening one with the visitor's cookies. An empty list allows only same-origin requests.
func (*Server) Stable ¶ added in v0.32.0
func (s *Server) Stable(name string) (crdt.CompositeVersion, bool)
What every participant has certainly seen.
A replica may drop a tombstone once every replica has delivered the deletion that made it — see crdt.Doc.Collect, which asks for such a version and cannot compute one. A server can: it is the thing every operation passes through, and participants tell it what they have applied.
This is the telling. A participant sends its version after it applies what the server sent; the server keeps the last one from each, and the meet of them — the element-wise minimum — is what everybody here has. Nothing depends on an acknowledgement arriving: one that is late or lost holds the answer back, and holding it back is the safe direction.
What it is not ¶
It is the meet over the participants **connected now**. A replica that is offline holding work of its own is not in it and cannot be: the server has never heard of what it did. So this is not yet a version anything may be collected against — deciding that a replica is gone is a policy, and it is not one a version vector can make. What this gives is the measurement that policy would have to be worth making: whether, in a room that is actually being used, the meet advances at all.
Stable returns the version every participant of the named document has acknowledged, and false if the document is not open or nobody has said anything yet.
type SiteStore ¶ added in v0.36.0
type SiteStore interface {
// LoadSites returns what SaveSites last wrote for a document, or nil if
// there is none. Nil is not an error: it is how a store says it has never
// been told about this one.
LoadSites(ctx context.Context, document string) ([]byte, error)
// SaveSites records the participants, replacing whatever was there.
SaveSites(ctx context.Context, document string, sites []byte) error
}
A SiteStore keeps, beside a document, the participants that have been in it.
Store is enough to hold a document and not enough to collect one. Collecting asks for a version every participant has delivered, and the participant set a server can derive from a document alone is the sites its own version vector names — which is everyone who has *written*. Somebody who has only ever read is in no version vector at all, so a document that is evicted and loaded again comes back not knowing they were here, and the floor moves past them.
A store that can keep a little more says so by implementing this. What it is given is opaque: the server owns the encoding, a store keeps the bytes and gives them back. They carry their own magic and their own checksum, so a store that gives back bytes that rotted is caught rather than believed, and a store that cannot checksum what it keeps has nothing to add here.
A store that does not implement it loses nothing it had. The server falls back to the sites the document names, which is what every store did before this existed.
type Store ¶
type Store interface {
// Load returns the snapshot for a document, or nil if there is none yet.
//
// nil, and only nil, means "none yet". A zero-length snapshot is not an
// empty document — a document has a header — it is a torn write, and a
// store must refuse it rather than answer nil: the server would open an
// empty replica and the next save would make the loss permanent.
// Returning nil is how a store says "new document", and is not an error.
Load(ctx context.Context, document string) ([]byte, error)
// Save records the current snapshot, replacing any previous one.
//
// Replacing, and not merging: what a server hands over is its whole
// replica, and the store is not asked to combine it with what is there.
// That is why two servers holding one document over one store lose the
// earlier save — see the package documentation on why a shared store is not
// how to run two of them, and [MergeSnapshots] for the combining a caller
// does when it has two snapshots and means to keep both.
Save(ctx context.Context, document string, snapshot []byte) error
}
A Store keeps documents between sessions. It holds snapshots, which are self-contained: a document restored from one can still serve a participant that has been away, because the snapshot carries the whole history.
Implementations must be safe for concurrent use.
type Text ¶ added in v0.10.0
type Text struct {
// contains filtered or unexported fields
}
A Text is a handle on one text part: the buffer of a file, and what an editor binds to.
func (*Text) Anchor ¶ added in v0.10.0
Anchor returns the identity of the character at rune offset pos, which keeps naming that character however the text moves around it. It is what a comment or a stored selection should hold; see crdt.Doc.Anchor.
func (*Text) AnchorUTF16 ¶ added in v0.12.0
AnchorUTF16 is Text.Anchor with pos counted in UTF-16 code units, the units a page counts in. An offset falling between the two units of one character is refused rather than rounded; see crdt.ErrSurrogateBoundary.
func (*Text) AuthorRuns ¶ added in v0.10.0
AuthorRuns splits the visible text into stretches by who wrote them, which is what colouring a document by author needs.
func (*Text) AuthorRunsUTF16 ¶ added in v0.12.0
AuthorRunsUTF16 is Text.AuthorRuns with every offset and length counted in UTF-16 code units, so that a page can colour the string it holds without converting anything by hand.
func (*Text) Delete ¶ added in v0.10.0
Delete removes length runes at rune offset pos, locally and then everywhere.
func (*Text) DeleteUTF16 ¶ added in v0.10.0
DeleteUTF16 removes length code units at an offset counted in the same units.
func (*Text) Insert ¶ added in v0.10.0
Insert adds text at rune offset pos, locally and then everywhere.
func (*Text) InsertUTF16 ¶ added in v0.10.0
InsertUTF16 adds text at an offset counted in UTF-16 code units.
func (*Text) LenUTF16 ¶ added in v0.10.0
LenUTF16 returns the length a browser would report, counting UTF-16 code units. Its companions InsertUTF16 and DeleteUTF16 take offsets in the same units, so a caller in the browser never converts by hand; see crdt.Doc.
func (*Text) Part ¶ added in v0.10.0
Part names this handle's part, which is what a crdt.PartChange from Client.TakeChanges carries.
func (*Text) Position ¶ added in v0.10.0
Position returns where the character an anchor names sits now — or where it was, if it has been deleted. See crdt.Doc.Position.
func (*Text) PositionUTF16 ¶ added in v0.12.0
PositionUTF16 is Text.Position with the offset reported in UTF-16 code units. ok is false for an anchor this replica has never seen, exactly as it is there — which is not the same question as whether the character is still in the text; that one is Text.Visible.
type Tiered ¶ added in v0.26.1
type Tiered struct {
// contains filtered or unexported fields
}
A Tiered store keeps documents somebody is using in one store and documents nobody has opened for a long time in another.
What it is for ¶
A server holds every document it has ever served until it is told to let go — see Config.EvictAfter — and a store holds every document it has ever been given, with nothing to let go of it at all. A year of a busy service is a directory of documents that were opened once, and there has been no way to move them anywhere cheaper without deleting them.
Nothing is deleted that is not already somewhere else ¶
Archiving is three steps in one order: read the hot store, write the cold one, then ask the hot one to release exactly what was read. Every way it can fail leaves the document somewhere. A cold store that will not take it releases nothing. A document that was written between the read and the release is not released, because Archivable.Release compares before it removes; it is archived on a later pass, when it has gone quiet again.
An archived document is not a missing one ¶
Store.Load returning nil means "there is no such document, start a new one", and a server acts on it: it opens an empty document and, at its next save, writes it over whatever was there. So a Tiered store never answers nil for a document the cold store has, and never answers nil because it could not reach the cold store — it fails instead. A document that is unreachable is not a document that does not exist, and confusing the two is how a store loses what it was given.
Reading an archived document brings it back: it is written to the hot store on the way past, so the next read does not go looking again and the next save has somewhere to land.
It does not implement SiteStore. Go has no way to implement an interface only when what is underneath does, so a Tiered that declared the methods would keep nothing whenever its tiers could not — silently, which is worse than not offering it. A server given one falls back to the participants a document names, which is everyone who has written; see SiteStore for what that costs.
func NewTiered ¶ added in v0.26.1
func NewTiered(hot Archivable, cold Store) *Tiered
NewTiered returns a store that writes to hot and falls back to cold.
It panics if either is nil, which is a mistake in the call rather than a state anything could recover from.
func (*Tiered) Archive ¶ added in v0.26.1
Archive moves every document the hot store has not been asked to write for idleFor into the cold store, and returns how many it moved.
It is a method rather than a timer of its own because a server already has housekeeping running on a schedule an operator chose, and a store that starts goroutines is a store that has to be closed.
A document that could not be archived does not stop the ones after it: the errors are returned together, and what was moved is reported whatever else happened.
type Token ¶ added in v0.67.0
type Token []byte
A Token names a version of a stored document, opaquely: a caller keeps one and hands it back, and only the store that issued it knows what is inside.
nil means "nothing was there". A caller that loaded a document nobody had written holds nil, and passing nil to ConditionalStore.SaveIf asks for a save that succeeds only while that is still true.
type Transport ¶ added in v0.11.0
type Transport interface {
// contains filtered or unexported methods
}
A Transport is how a participant reaches a server. WebSocket works anywhere, a browser included; GRPC works outside one.
There are two because of what they cost where they run. Outside a browser a carrier costs nothing anybody notices, and gRPC brings deadlines, interceptors and everything already built around them. Inside one it is paid for on every load: protobuf alone is six times the size of the whole CRDT compiled to wasm — see wire.go for the measurements — so the browser gets a framing of four message kinds over a plain WebSocket instead.
Both carry the same session, byte for byte in the fields that matter, because every field in these messages is something github.com/go-crdt/crdt encoded and will check on arrival. A participant on one and a participant on the other can edit the same document.
func GRPC ¶ added in v0.11.0
func GRPC(conn grpc.ClientConnInterface) Transport
GRPC returns a transport that opens sessions on a gRPC connection.
It is deliberately not the default. Everything it carries is bytes some encoder in github.com/go-crdt/crdt produced, so protobuf is describing fields nobody reads through it — and compiled for a browser it costs six times the CRDT itself. Outside a browser that does not matter, and gRPC brings deadlines, interceptors and the tooling built around them, which is reason enough to keep it. See Transport and WebSocket.
func WebSocket ¶ added in v0.11.0
func WebSocket(url string, opts ...WebSocketOption) Transport
WebSocket returns a transport that opens sessions at url, which is "ws://" or "wss://" and the path the server's handler is mounted at.
This is the transport a browser uses, and the one to reach for by default: it is what the same code compiled to wasm can afford. See Transport.
type WebSocketOption ¶ added in v0.11.0
type WebSocketOption func(*wsTransport)
A WebSocketOption configures WebSocket.
func WithHTTPHeader ¶ added in v0.11.0
func WithHTTPHeader(h http.Header) WebSocketOption
WithHTTPHeader sends these headers with the opening handshake, which is where a cookie or a bearer token goes when the participant is not a browser.
It does not exist for a browser, because a page cannot put a header on a WebSocket handshake — and does not need to, since the browser sends the cookies for that origin itself. Code meant to run in both places should let the cookie do the work; see Config.Authorize.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Command browsertest is the real-browser half of the WebRTC proof.
|
Command browsertest is the real-browser half of the WebRTC proof. |
|
Command electtest is the real-browser proof that a host election decided by a Web Lock leaves exactly one host, where the window it replaces would leave two.
|
Command electtest is the real-browser proof that a host election decided by a Web Lock leaves exactly one host, where the window it replaces would leave two. |
|
gitstore
module
|
|
|
migrate
module
|
|
|
pgstore
module
|
|
|
Package storetest is the contract every collab.Store keeps, written once and run by each implementation against itself.
|
Package storetest is the contract every collab.Store keeps, written once and run by each implementation against itself. |
|
Command wasmtest is the browser half of the end-to-end proof.
|
Command wasmtest is the browser half of the end-to-end proof. |
