Aggregation
Aggregation merges what many writers write under one name into one document that every reader reads the same, tells each writer whether each change was applied, and lists every writer's current value.
Paths
Paths are relative to the participant's root. {name} is any path, the empty one included; {writer} is everything after the first .aggregation segment.
| Written by | Path | Track | Content |
|---|---|---|---|
| A writer | {name}/.aggregation/{writer} |
patch.json |
The writer's changes: each group holds one frame, a JSON object applied as an RFC 7396 merge patch; the group's sequence numbers the change |
| A writer | the same | value.json |
The writer's current value, a JSON snapshot track |
| Aggregation | .aggregation/{name} |
state.json |
The merged document, a JSON snapshot track |
| Aggregation | the same | writers.json |
A JSON snapshot track of {"{writer}": value}, one entry per writer whose value.json is published now |
| Aggregation | the same | applied/{writer} |
One frame per change of that writer: {"applied": n} or {"refused": n, "code": name} |
Only a session whose publish covers {name}/.aggregation/alice writes as alice, so a link with publish board/.aggregation/{id} lets each person write the board as themselves. Reading .aggregation/{name} shows what the writers wrote, so grant it as you grant their paths.
// A shared board through Aggregation: changes go to board/.aggregation/{me}, the merged board and
// every writer's current value come from .aggregation/board.
import { type Connection, json, moq } from 'tablebox.io';
export async function board(connection: Connection, me: string, show: (board: unknown, writers: unknown) => void) {
const mine = connection.publish(`board/.aggregation/${me}`);
const changes = mine.createTrack('patch.json');
const value = new json.Snapshot.Producer({ track: mine.createTrack('value.json') });
// A writer never uses a change number twice, so each broadcast numbers its changes from the time.
let next = Date.now();
const merged = await connection.read('.aggregation/board');
const state = new json.Snapshot.Consumer({ track: merged.track('state.json').subscribe() });
const writers = new json.Snapshot.Consumer({ track: merged.track('writers.json').subscribe() });
let board: unknown;
let everyone: unknown;
void (async () => { for await (board of state) show(board, everyone); })();
void (async () => { for await (everyone of writers) show(board, everyone); })();
return {
// One RFC 7396 merge patch per change: { "note-1": { "text": "Hi" } } sets a key, null deletes it.
change: (patch: Record<string, unknown>) => {
const group = new moq.Group.Producer(next++);
changes.writeGroup(group);
group.writeJson(patch);
group.close();
},
// This writer's own entry in writers.json, such as its name and cursor.
set: (mine: unknown) => value.update(mine),
};
}Merging and acknowledgements
Reading any track of .aggregation/{name} starts Aggregation for that name; it then reads every writer under {name}/.aggregation/. Each writer's changes are applied in that writer's order, and the changes of different writers in the order they reach Aggregation; on every key the last applied change wins. Every reader of state.json reads the same sequence of documents, whichever Relay it is on.
Every change a writer publishes while Aggregation runs gets one frame on applied/{writer}: applied once it is in state.json, or refused with bad-request (the frame is no JSON object), too-large (the document would pass 16 MiB) or missing-variable. n is the change's group sequence. A writer numbers its changes up by one and never uses a number twice: each broadcast starts at the time in milliseconds since the Unix epoch, or above the writer's last number when that is higher, so a frame on applied/{writer} names one change whichever of the writer's broadcasts it came in. A writer reads its applied/{writer} before it writes, and publishes again the changes it holds no acknowledgement for whenever that read starts again (after a reconnect, or after a takeover ended it): applying a merge patch twice gives the same document.
Every writer's value
writers.json holds each writer's newest value.json value, such as a name or a cursor. A writer's entry leaves when its broadcast or its value.json ends, so the list follows who is there. A range too large for everyone to list everyone's broadcasts can keep its people in writers.json instead, read by those who need it.
Saved
state.json starts from the document saved for it and is saved as it changes, through Persistence, so it needs Persistence's storage variables. Aggregation publishes on state.json once Persistence gives the saved document, and starts from {} when none is saved. The document is in your storage at .aggregation/{project path of name}/state.json under STORAGE_PREFIX.
Running and failures
Aggregation runs while its output is read or a writer is there, and stops 20 s after both are gone. Persistence, which saves state.json, reads it until 20 s after the last writer leaves, so a name nobody else reads stops 40 s after its last writer leaves. All the names of one project are served by one instance at a time.
- The storage variables are missing when a name starts:
state.jsonis refused withmissing-variable, and so is every change;writers.jsonworks. Reading again after you set them works. - Persistence unreachable or failing before the document starts:
state.jsonwaits, and changes stay unacknowledged until it answers. - Saving stops after the document started: the document goes on, and changes are applied and acknowledged; it is saved again once Persistence can.
- An instance stops: reads end, and reading again reaches the instance that takes over. That instance starts from the saved document (at most 10 s old), applies again the changes the writers' broadcasts still hold, then the changes it receives; on a key another writer changed since one of those changes, the older value can come back. Changes acknowledged after the last save and no longer held are lost unless their writers publish them again.
- A writer leaves: its applied changes stay in
state.json, and its entry leaveswriters.json. - Instances are added or removed: a project can be served by two instances for a moment, and readers on different Relays can then see different documents until the old instance lets go.