maybeTakeAlignedSnapshot method
Take a snapshot (and prune the confirmed history) when every client
subscribed to documentId has confirmed at least the server's current
state.
Alignment is derived from the version vectors clients report on their pings (and at handshake). The stability frontier is the intersection (per-peer minimum) of those vectors. Snapshotting only when the frontier covers the server's current version guarantees pruning never drops history a client has not yet confirmed — a lagging client still re-syncs from the stored snapshot on its next handshake.
Implementation
Future<void> maybeTakeAlignedSnapshot(String documentId) async {
if (!isRunning) {
return;
}
final subscribed = sessions.values
.where((session) => session.isSubscribedTo(documentId))
.toList();
if (subscribed.isEmpty) {
return;
}
final versionVectors = <VersionVector>[];
for (final session in subscribed) {
final versionVector = session.lastKnownVersionVector;
if (versionVector == null) {
// A subscribed client has not reported its state yet: we cannot know
// how far it has advanced, so pruning would be unsafe.
return;
}
versionVectors.add(versionVector);
}
final frontier = VersionVector.intersection(versionVectors);
final document = await _serverRegistry.getDocument(documentId);
if (document == null) {
return;
}
final serverVersion = document.getVersionVector();
if (serverVersion.isEmpty) {
return;
}
// Every client has confirmed at least the server's current state.
if (!frontier.isStrictlyNewerOrEqualThan(serverVersion)) {
return;
}
// Avoid redundant work: skip if a snapshot already covers this state.
final existing = await _serverRegistry.getLatestSnapshot(documentId);
if (existing != null &&
!serverVersion.isStrictlyNewerThan(existing.versionVector)) {
return;
}
await _serverRegistry.createSnapshot(documentId);
addServerEvent(
ServerEvent(
type: ServerEventType.snapshotCreated,
message: 'All clients aligned on document $documentId: '
'snapshot taken and confirmed history pruned',
data: {
'documentId': documentId,
},
),
);
}