diff --git a/README.md b/README.md
index 6de1841..9e681a1 100644
--- a/README.md
+++ b/README.md
@@ -742,6 +742,8 @@ attempt, including durable command receipt lookup. Restore is the exception. It
across Begin, source transfer, Finish, cleanup, and stale session replacement. Status, invalidation, and backup export
streams are intentionally unbounded.
+The page to owner `ReplicaRpc` protocol version is 5. A version mismatch is rejected before the session is admitted.
+
Browser restore protocol version 4 begins with a short unary request. The owner returns a restore nonce and dedicated
`MessagePort`. The page immediately starts a pending `FinishRestoreBackupV4` result RPC and sends `Start`. After the
owner has both signals, it drives archive input with one outstanding `Pull` credit. The page sends one bounded chunk
@@ -1622,25 +1624,25 @@ composition recipes.
### `@lucas-barake/effect-local`
-| Namespace | Public API |
-| ------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
-| `Backup` | `FormatVersion`, `Header`, `ExportOptions`, `RestoreOptions`, `ExportedDocument` |
-| `Canonical` | `stringify`, `hash`, `digest` |
-| `CommandOutcome` | `Rejected`, `DurablyCommittedLocal`, `OutcomeUnknown`, `CommandOutcomeUnknown`, `CommandOutcome`, `schema`, `rejected`, `durablyCommitted`, `unknown`, `match`, `committedOrFail` |
-| `Commit` | `Heads`, `Commit` and their inferred types |
-| `Document` | `WireSchema`, `AutomergeEncoded`, `DocumentSchema`, `Document`, `Any`, `make`, `isAutomergeValue`, `decode`, `encode` |
-| `DocumentSet` | `DocumentSet`, `make`, `get` |
-| `Identity` | Schemas and types for `ReplicaId`, `ReplicaIncarnation`, `SessionId`, `DocumentId`, `CommandId`, `WriterGeneration`, `CommitSequence`, `PeerId`, `ProjectionVersion`, and `DocumentLineage`. `genesisLineage`, `makeReplicaId`, `makeSessionId`, `makeDocumentId`, `makeCommandId`, `makePeerId`, `makeDocumentLineage`, `documentIdFromCommandId` |
-| `Mutation` | `DraftValue`, `Draft`, `SuccessResult`, `HandlerResult`, `HandlerOptions`, `Handler`, `HandlerService`, `Mutation`, `Any`, `make`; definitions expose `payloadSchema`, `successSchema`, `errorSchema`, `of`, and `toLayer`; `toLayer` accepts a handler or an Effect that builds one |
-| `PeerTransport` | `Capabilities`, `Connection`, `ConnectOptions`, `PeerTransport` service |
-| `Projection` | `Projection`, `Any`, `make`, `assertUniqueKeys`, `evaluate` |
-| `Query` | `Handler`, `HandlerService`, `Query`, `Any`, `make`; definitions expose `payloadSchema`, `successSchema`, `errorSchema`, `of`, and `toLayer`; `toLayer` accepts a handler or an Effect that builds one |
-| `Replica` | `Replica` service |
-| `ReplicaDefinition` | `ReplicaDefinition`, `Any`, `invalidationKeys`, `make` |
-| `ReplicaError` | Reason schemas `DocumentNotFound`, `DocumentDecodeError`, `DocumentEncodeError`, `UnsupportedDocumentVersion`, `ProjectionBlocked`, `CommandIdConflict`, `ReceiptOperationMismatch`, `StorageUnavailable`, `CanonicalEncodeError`, `StorageCorrupt`, `QuotaExceeded`, `MigrationFailed`, `BackupInvalid`, `BackupTooLarge`, `RestoreBusy`, `RestoreFailed`, `ProtocolMismatch`, `ReplicaFenced`, `OperationTimeout`, `UnsupportedStorageFormatVersion`, `CheckpointSuperseded`, `DocumentLineageChanged`. Causal reasons use `Schema.Defect()` for transportable arbitrary failures. `Reason`, tagged `ReplicaError` |
-| `ReplicaLimits` | `Values`, `minimumRestoreErrorBytes`, `ReplicaLimits` service, `make`, `layer` |
-| `ReplicaStatus` | `Starting`, `Ready`, `ReadOnly`, `Degraded`, `ProjectionBlocked`, `Restoring`, `Failed`, `ReplicaStatus` schemas and types |
-| `Snapshot` | `ProjectionState`, `Snapshot`, `FromDocument` |
+| Namespace | Public API |
+| ------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
+| `Backup` | `FormatVersion`, `Header`, `ExportOptions`, `RestoreOptions`, `ExportedDocument` |
+| `Canonical` | `stringify`, `hash`, `digest` |
+| `CommandOutcome` | `Rejected`, `DurablyCommittedLocal`, `OutcomeUnknown`, `CommandOutcomeUnknown`, `CommandOutcome`, `schema`, `rejected`, `durablyCommitted`, `unknown`, `match`, `committedOrFail` |
+| `Commit` | `Heads`, `Commit` and their inferred types |
+| `Document` | `WireSchema`, `AutomergeEncoded`, `DocumentSchema`, `Document`, `Any`, `make`, `isAutomergeValue`, `decode`, `encode` |
+| `DocumentSet` | `DocumentSet`, `make`, `get` |
+| `Identity` | Schemas and types for `ReplicaId`, `ReplicaIncarnation`, `SessionId`, `DocumentId`, `CommandId`, `WriterGeneration`, `CommitSequence`, `PeerId`, `ProjectionVersion`, and `DocumentLineage`. `genesisLineage`, `makeReplicaId`, `makeSessionId`, `makeDocumentId`, `makeCommandId`, `makePeerId`, `makeDocumentLineage`, `documentIdFromCommandId` |
+| `Mutation` | `DraftValue`, `Draft`, `SuccessResult`, `HandlerResult`, `HandlerOptions`, `Handler`, `HandlerService`, `Mutation`, `Any`, `make`; definitions expose `payloadSchema`, `successSchema`, `errorSchema`, `of`, and `toLayer`; `toLayer` accepts a handler or an Effect that builds one |
+| `PeerTransport` | `Capabilities`, `Connection`, `ConnectOptions`, `PeerTransport` service |
+| `Projection` | `Projection`, `Any`, `make`, `assertUniqueKeys`, `evaluate` |
+| `Query` | `Handler`, `HandlerService`, `Query`, `Any`, `make`; definitions expose `payloadSchema`, `successSchema`, `errorSchema`, `of`, and `toLayer`; `toLayer` accepts a handler or an Effect that builds one |
+| `Replica` | `Replica` service |
+| `ReplicaDefinition` | `ReplicaDefinition`, `Any`, `invalidationKeys`, `make` |
+| `ReplicaError` | Reason schemas `DocumentNotFound`, `DocumentDecodeError`, `DocumentEncodeError`, `UnsupportedDocumentVersion`, `ProjectionBlocked`, `CommandIdConflict`, `ReceiptOperationMismatch`, `StorageUnavailable`, `CanonicalEncodeError`, `StorageCorrupt`, `ReplicaMetadataMissing`, `QuotaExceeded`, `MigrationFailed`, `BackupInvalid`, `BackupTooLarge`, `RestoreBusy`, `RestoreFailed`, `ProtocolMismatch`, `ReplicaFenced`, `OperationTimeout`, `UnsupportedStorageFormatVersion`, `CheckpointSuperseded`, `DocumentLineageChanged`. Causal reasons use `Schema.Defect()` for transportable arbitrary failures. `Reason`, tagged `ReplicaError` |
+| `ReplicaLimits` | `Values`, `minimumRestoreErrorBytes`, `ReplicaLimits` service, `make`, `layer` |
+| `ReplicaStatus` | `Starting`, `Ready`, `ReadOnly`, `Degraded`, `ProjectionBlocked`, `Restoring`, `Failed`, `ReplicaStatus` schemas and types |
+| `Snapshot` | `ProjectionState`, `Snapshot`, `FromDocument` |
`Replica.Replica` is the application capability. Its methods are `create`, `get`, `mutate`, `delete`, `query`,
`lookupMutation`, `lookupCreate`, `lookupDelete`, `flush`, `status`, `exportBackup`, `restoreBackup`,
@@ -1677,6 +1679,7 @@ Core data shapes:
| `CommandIdConflict` | `commandId` |
| `ReceiptOperationMismatch` | `commandId`, `expected`, `observed` |
| `StorageUnavailable`, `CanonicalEncodeError`, `StorageCorrupt`, `BackupInvalid`, `RestoreFailed` | `cause` |
+| `ReplicaMetadataMissing` | None |
| `QuotaExceeded` | `resource`, `limit` |
| `MigrationFailed` | `migration`, `cause` |
| `BackupTooLarge` | `limit`, `observed` |
@@ -1751,7 +1754,7 @@ Public SQL service methods:
| `ProjectionStore` | `clear`, `replace(binding, snapshot, destinationTable)`, `replaceDocument(document, snapshot, commitSequence)` |
| `QueryExecutor` | `execute(query, payload)` and reactive `reactive(query, payload)` |
| `Recovery` | `recover(document, documentId)`, `recoverWithPermit(document, documentId, permit)`, `exportRaw(documentId)` |
-| `ReplicaGate` | `current`, scoped `shared`, bounded scoped `admit`, exclusive `claim(use)`, `refresh`, `validate(expectedPermit)` |
+| `ReplicaGate` | `current`, scoped `shared`, bounded scoped `admit`, exclusive `claim(use)`, `refresh`, `preflight(expectedPermit)`, `validate(expectedPermit)` |
| `CompactionWorkflow` | `execute(operationId)`, `poll(execution)`, `interrupt(execution)`, `resume(execution)` |
| `HistoryRewriteWorkflow` | `execute(documentId, operationId)`, `poll(execution)`, `interrupt(execution)`, `resume(execution)`. Handles are `DocumentExecution`, which adds `documentId` to `Execution` |
@@ -1759,33 +1762,40 @@ Public SQL service methods:
persisted, and run with the configured SQL transaction annotation. `Create`, `Mutate`, and `Delete` are client
uninterruptible. `ApplySync` is uninterruptible.
-| Procedure | Payload additions | Success |
-| ----------- | -------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------- |
-| `Create` | Common command fields plus schema coded JSON bytes in `payload` | Schema coded JSON bytes |
-| `Mutate` | Common command fields plus `mutationTag` and schema coded JSON bytes in `payload` | Schema coded JSON bytes |
-| `Delete` | Common command fields only | Schema coded JSON bytes |
-| `ApplySync` | `replicaIncarnation`, peer and connection epochs, receive sequence, document type, message hash, message bytes | `{ reply, heads, acceptedHeads, commitSequence, observedByPeer, durableConfirmation: false, duplicate }` |
+| Procedure | Payload additions | Success |
+| ----------- | --------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------- |
+| `Create` | Common command fields plus schema coded JSON bytes in `payload` | Schema coded JSON bytes |
+| `Mutate` | Common command fields plus `mutationTag` and schema coded JSON bytes in `payload` | Schema coded JSON bytes |
+| `Delete` | Common command fields only | Schema coded JSON bytes |
+| `ApplySync` | `replicaIncarnation`, peer and connection epochs, receive sequence, document type, message hash, message bytes, document lineage, writer provenance | `{ reply, heads, acceptedHeads, commitSequence, observedByPeer, durableConfirmation: false, duplicate }` |
The common command fields are `replicaIncarnation`, `writerGeneration`, `commandId`, `documentType`, and
`requestHash`. Command persistence keys include the replica incarnation, command ID, and request hash. Sync
-persistence keys include the incarnation, peer, connection epoch, receive sequence, and message hash.
+persistence keys include the incarnation, peer, connection epoch, receive sequence, message hash, canonical writer
+provenance, and normalized document lineage.
+
+`ReplicaMetadataMissing` reports loss of the replica wide metadata singleton. Public command and inbound sync
+admission returns this typed failure before durable dispatch. If metadata disappears after admission, waiter visibility
+depends on whether the handler was registered before the typed failure. Domain changes and the attempted reply roll
+back without a terminal cached reply. After repair, the same persistence key is safe for storage redelivery or an
+application retry.
SQL composition contracts:
-| Module | Required services or ownership |
-| ------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
-| `SqlReplica` | `layerWithBindings` is the normal constructor. The application supplies platform `SqlClient`, `Crypto`, `ReplicaLimits`, generated mutation and query handler Layers, and all declared `SqlProjection` bindings. `layer` accepts bindings directly. `layerFromServices` is advanced assembly |
-| `ReplicaBootstrap` and `ReplicaGate` | Bootstrap runs migrations and establishes replica identity, incarnation, definition hash, and writer generation. Gate requires `ReplicaLimits` and provides shared operation permits and exclusive restore fencing. `shared` and `claim` always wait. `admit` also waits, except that while an exclusive holder such as restore holds or is waiting for the gate it sheds acquirers past `maxQueuedPermits` with `QuotaExceeded`, so a long restore cannot accumulate an unbounded backlog |
-| `DocumentStore` | Requires `Crypto`, `SqlClient`, and `ReplicaGate`. Callers of loaded native Automerge documents own `free` |
-| `ProjectionStore` | Requires generated binding services and `SqlClient`. It replaces and rebuilds deterministic derived tables |
-| `CommandExecutor` | Requires stores, gate, SQL, Crypto, and every mutation handler. Canonical state, projections, receipt, and commit outbox are one transaction |
-| `QueryExecutor` | Requires gate, projection readiness, SQL, and every query handler. Declared query errors remain typed |
-| `CommitPublisher` | Requires `Reactivity` and SQL. `publishPending` durably retries until publication into a bounded sliding stream. Scoped subscribers receive best effort refresh notifications and must recover from later sequence gaps by rereading canonical state |
-| `DurableRuntime` | Builds Cluster and Workflow over the same SQL storage. `layerWith` accepts additional Workflow registrations |
-| `EntityReplica` | Adapts durable Cluster commands, stores, queries, backup, status, and commit publication to core `Replica` |
-| `BackupStore` | Streams and restores definition bound archives. Export is a snapshot. Restore validates envelope, checksum, bounds, definition, foreign keys, recovery, and projections before installation |
-| `Compaction` and `Recovery` | Compaction publishes only verified checkpoints and prunes changes only with retained safety evidence. It also reclaims command receipts from superseded incarnations, which no current permit can read. Recovery validates checkpoints, changes, heads, and tombstones. `Compaction.rewriteHistory` is the one destructive member. It is reached through `ReplicaWorkflow.HistoryRewriteWorkflow`, which is the only composition that owns the required permit for it |
-| `ReplicaWorkflow` | Registers and executes the scoped compaction and history rewrite Workflows. Workflow durability does not make peer or backup transport durable |
+| Module | Required services or ownership |
+| ------------------------------------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
+| `SqlReplica` | `layerWithBindings` is the normal constructor. The application supplies platform `SqlClient`, `Crypto`, `ReplicaLimits`, generated mutation and query handler Layers, and all declared `SqlProjection` bindings. `layer` accepts bindings directly. `layerFromServices` is advanced assembly |
+| `ReplicaBootstrap` and `ReplicaGate` | Bootstrap runs migrations and establishes replica identity, incarnation, definition hash, and writer generation. Gate requires `ReplicaLimits` and provides shared operation permits and exclusive restore fencing. `shared` and `claim` always wait. `admit` also waits, except that while an exclusive holder such as restore holds or is waiting for the gate it sheds acquirers past `maxQueuedPermits` with `QuotaExceeded`, so a long restore cannot accumulate an unbounded backlog. `preflight` compares durable ownership without adopting it. `validate` is the transaction serialized write authority |
+| `DocumentStore` | Requires `Crypto`, `SqlClient`, and `ReplicaGate`. Callers of loaded native Automerge documents own `free` |
+| `ProjectionStore` | Requires generated binding services and `SqlClient`. It replaces and rebuilds deterministic derived tables |
+| `CommandExecutor` | Requires stores, gate, SQL, Crypto, and every mutation handler. Canonical state, projections, receipt, and commit outbox are one transaction |
+| `QueryExecutor` | Requires gate, projection readiness, SQL, and every query handler. Declared query errors remain typed |
+| `CommitPublisher` | Requires `Reactivity` and SQL. `publishPending` durably retries until publication into a bounded sliding stream. Scoped subscribers receive best effort refresh notifications and must recover from later sequence gaps by rereading canonical state |
+| `DurableRuntime` | Builds Cluster and Workflow over the same SQL storage. `layerWith` accepts additional Workflow registrations |
+| `EntityReplica` | Adapts durable Cluster commands, stores, queries, backup, status, and commit publication to core `Replica` |
+| `BackupStore` | Streams and restores definition bound archives. Export is a snapshot. Restore validates envelope, checksum, bounds, definition, foreign keys, recovery, and projections before installation |
+| `Compaction` and `Recovery` | Compaction publishes only verified checkpoints and prunes changes only with retained safety evidence. It also reclaims command receipts from superseded incarnations, which no current permit can read. Recovery validates checkpoints, changes, heads, and tombstones. `Compaction.rewriteHistory` is the one destructive member. It is reached through `ReplicaWorkflow.HistoryRewriteWorkflow`, which is the only composition that owns the required permit for it |
+| `ReplicaWorkflow` | Registers and executes the scoped compaction and history rewrite Workflows. Workflow durability does not make peer or backup transport durable |
`PeerSession.PeerSession` exposes `peerId`, `connectionEpoch`, `markDirty`, `flush`, `observedByPeer`, and
`durableConfirmation`. `SupervisedPeerSession` adds `awaitDisconnect`. `make` owns transport and receive lifetime.
diff --git a/packages/local-browser/src/ReplicaRpc.ts b/packages/local-browser/src/ReplicaRpc.ts
index f938105..fa8164c 100644
--- a/packages/local-browser/src/ReplicaRpc.ts
+++ b/packages/local-browser/src/ReplicaRpc.ts
@@ -24,7 +24,7 @@ const ExportedDocument = Schema.Struct({
const JsonOutcome = CommandOutcome.schema(Schema.Json, Schema.Json)
const DocumentIdOutcome = CommandOutcome.schema(Identity.DocumentId, Schema.Never)
-export const protocolVersion = 4
+export const protocolVersion = 5
const SessionLease = Schema.Struct({ leaseMillis: Schema.Int })
export const SessionHandshake = Schema.Struct({
leaseMillis: Schema.Int,
diff --git a/packages/local-browser/src/internal/restoreProtocol.ts b/packages/local-browser/src/internal/restoreProtocol.ts
index 72767a3..3e0954f 100644
--- a/packages/local-browser/src/internal/restoreProtocol.ts
+++ b/packages/local-browser/src/internal/restoreProtocol.ts
@@ -58,6 +58,7 @@ const CanonicalEncodeError = Schema.TaggedStruct("CanonicalEncodeError", {
const StorageCorrupt = Schema.TaggedStruct("StorageCorrupt", {
cause: BoundedErrorDescription
})
+const ReplicaMetadataMissing = Schema.TaggedStruct("ReplicaMetadataMissing", {})
const QuotaExceeded = Schema.TaggedStruct("QuotaExceeded", {
resource: Schema.String,
limit: Schema.Int
@@ -116,6 +117,7 @@ export const RestoreWireError = Schema.Union([
StorageUnavailable,
CanonicalEncodeError,
StorageCorrupt,
+ ReplicaMetadataMissing,
QuotaExceeded,
MigrationFailed,
BackupInvalid,
@@ -216,6 +218,7 @@ export const restoreWireErrorFields = fieldMetadata({
StorageUnavailable,
CanonicalEncodeError,
StorageCorrupt,
+ ReplicaMetadataMissing,
QuotaExceeded,
MigrationFailed,
BackupInvalid,
@@ -427,6 +430,12 @@ export const encodeReplicaError = (
{ _tag: reason._tag, cause: emptyErrorDescription },
({ defect }) => ({ _tag: reason._tag, cause: defect(reason.cause) })
)
+ case "ReplicaMetadataMissing":
+ return encodeWithinBudget(
+ maxBytes,
+ { _tag: reason._tag },
+ () => ({ _tag: reason._tag })
+ )
case "QuotaExceeded":
return encodeWithinBudget(
maxBytes,
@@ -599,6 +608,9 @@ export const replicaErrorFromWire = (wire: RestoreWireError): ReplicaError.Repli
case "StorageCorrupt":
reason = new ReplicaError.StorageCorrupt({ cause: decodeDefect(wire.cause) })
break
+ case "ReplicaMetadataMissing":
+ reason = new ReplicaError.ReplicaMetadataMissing()
+ break
case "QuotaExceeded":
reason = new ReplicaError.QuotaExceeded({ resource: wire.resource, limit: wire.limit })
break
diff --git a/packages/local-browser/test/ReplicaClient.test.ts b/packages/local-browser/test/ReplicaClient.test.ts
index 4bcfee1..ee0943b 100644
--- a/packages/local-browser/test/ReplicaClient.test.ts
+++ b/packages/local-browser/test/ReplicaClient.test.ts
@@ -100,6 +100,10 @@ it.layer(NodeCrypto.layer)("ReplicaClient", (it) => {
typeof ReplicaRpc.group,
RpcClientError.RpcClientError
>
+ const resolveRpcResponse = (
+ effect: Effect.Effect, E, R>
+ ) =>
+ Effect.flatMap(effect, (value) => Deferred.isDeferred(value) ? Deferred.await(value) : Effect.succeed(value))
type FinishRestore = TestReplicaRpcClient["FinishRestoreBackupV4"]
const dropTerminalReady = (
rpc: TestReplicaRpcClient,
@@ -236,6 +240,8 @@ it.layer(NodeCrypto.layer)("ReplicaClient", (it) => {
it.effect("decodes and rejects owners using an older protocol", () =>
Effect.scoped(Effect.gen(function*() {
+ const sessions = yield* SessionManager.SessionManager
+ assert.strictEqual(ReplicaRpc.protocolVersion, 5)
const open = ReplicaRpc.group.requests.get("OpenSession")
if (open?._tag !== "OpenSession") return yield* Effect.die(new Error("OpenSession RPC not found"))
yield* Schema.decodeUnknownEffect(open.successSchema)({
@@ -263,8 +269,64 @@ it.layer(NodeCrypto.layer)("ReplicaClient", (it) => {
})
const error = yield* Effect.flip(ReplicaClient.fromRpcClient(definition, older))
assert.strictEqual(error.reason._tag, "ProtocolMismatch")
+ if (error.reason._tag === "ProtocolMismatch") {
+ assert.strictEqual(error.reason.expected, "protocol version 5")
+ assert.strictEqual(error.reason.observed, "protocol version 4")
+ }
+ assert.strictEqual(yield* sessions.activeCount, 0)
+ })).pipe(Effect.provide(Owner)))
+
+ it.effect("decodes and rejects owners using a newer protocol", () =>
+ Effect.scoped(Effect.gen(function*() {
+ const sessions = yield* SessionManager.SessionManager
+ const rpc = yield* RpcTest.makeClient(ReplicaRpc.group)
+ const newer = new Proxy(rpc, {
+ get(target, property, receiver) {
+ const value = Reflect.get(target, property, receiver)
+ if (property !== "OpenSession") return value
+ return (payload: never) =>
+ value(payload).pipe(Effect.map((lease) => ({
+ ...(lease as {
+ readonly leaseMillis: number
+ readonly protocolVersion: number
+ readonly definitionHash: string
+ readonly ownerEpoch: string
+ }),
+ protocolVersion: ReplicaRpc.protocolVersion + 1
+ })))
+ }
+ })
+ const error = yield* Effect.flip(ReplicaClient.fromRpcClient(definition, newer))
+ assert.strictEqual(error.reason._tag, "ProtocolMismatch")
+ if (error.reason._tag === "ProtocolMismatch") {
+ assert.strictEqual(error.reason.expected, "protocol version 5")
+ assert.strictEqual(error.reason.observed, "protocol version 6")
+ }
+ assert.strictEqual(yield* sessions.activeCount, 0)
})).pipe(Effect.provide(Owner)))
+ it.effect("owner rejects an explicit older client protocol before opening a session", () =>
+ Effect.gen(function*() {
+ const open = yield* ReplicaRpc.group.accessHandler("OpenSession")
+ const sessions = yield* SessionManager.SessionManager
+ const response = open({
+ sessionId: yield* Identity.makeSessionId,
+ protocolVersion: 4,
+ definitionHash: definition.hash
+ }, {
+ client: new Rpc.ServerClient(1),
+ requestId: RequestId("old-client"),
+ headers: Headers.empty
+ })
+ const error = yield* Effect.flip(resolveRpcResponse(response))
+ assert.strictEqual(error.reason._tag, "ProtocolMismatch")
+ if (error.reason._tag === "ProtocolMismatch") {
+ assert.strictEqual(error.reason.expected, `5:${definition.hash}`)
+ assert.strictEqual(error.reason.observed, `4:${definition.hash}`)
+ }
+ assert.strictEqual(yield* sessions.activeCount, 0)
+ }).pipe(Effect.provide(Owner)))
+
it.effect("recovers ambiguous commands through typed receipt lookup", () =>
Effect.scoped(Effect.gen(function*() {
const rpc = yield* RpcTest.makeClient(ReplicaRpc.group)
diff --git a/packages/local-browser/test/ReplicaOwnerRestoreV4.test.ts b/packages/local-browser/test/ReplicaOwnerRestoreV4.test.ts
index 237c500..fe84139 100644
--- a/packages/local-browser/test/ReplicaOwnerRestoreV4.test.ts
+++ b/packages/local-browser/test/ReplicaOwnerRestoreV4.test.ts
@@ -99,7 +99,7 @@ it.layer(NodeCrypto.layer)("ReplicaOwner restore version 4", (it) => {
protocolVersion: ReplicaRpc.protocolVersion,
definitionHash: definition.hash
})
- assert.strictEqual(handshake.protocolVersion, 4)
+ assert.strictEqual(handshake.protocolVersion, 5)
assert.strictEqual(handshake.maxChunkBytes, limits.maxChunkBytes)
assert.strictEqual(handshake.maxRestoreCoalesceMillis, limits.maxRestoreCoalesceMillis)
assert.strictEqual(handshake.maxRestoreErrorBytes, limits.maxRestoreErrorBytes)
diff --git a/packages/local-browser/test/RestoreClientProtocol.test.ts b/packages/local-browser/test/RestoreClientProtocol.test.ts
index 9fd504c..dbc076e 100644
--- a/packages/local-browser/test/RestoreClientProtocol.test.ts
+++ b/packages/local-browser/test/RestoreClientProtocol.test.ts
@@ -1665,13 +1665,13 @@ it.layer(NodeCrypto.layer)("RestoreClientProtocol", (it) => {
}
}))
- it.effect("rejects a version 3 handshake before opening the source transport", () =>
+ it.effect("rejects a version 4 handshake before opening the source transport", () =>
Effect.scoped(Effect.gen(function*() {
const oldRpc = {
OpenSession: () =>
Effect.succeed({
leaseMillis: 10_000,
- protocolVersion: 3,
+ protocolVersion: 4,
definitionHash: definition.hash,
ownerEpoch: "owner"
}),
@@ -1680,8 +1680,8 @@ it.layer(NodeCrypto.layer)("RestoreClientProtocol", (it) => {
const error = yield* ReplicaClient.fromRpcClient(definition, oldRpc).pipe(Effect.flip)
assert.strictEqual(error.reason._tag, "ProtocolMismatch")
if (error.reason._tag === "ProtocolMismatch") {
- assert.strictEqual(error.reason.expected, "protocol version 4")
- assert.strictEqual(error.reason.observed, "protocol version 3")
+ assert.strictEqual(error.reason.expected, "protocol version 5")
+ assert.strictEqual(error.reason.observed, "protocol version 4")
}
})))
})
diff --git a/packages/local-browser/test/RestoreProtocol.test.ts b/packages/local-browser/test/RestoreProtocol.test.ts
index 3efe624..7c11e31 100644
--- a/packages/local-browser/test/RestoreProtocol.test.ts
+++ b/packages/local-browser/test/RestoreProtocol.test.ts
@@ -23,6 +23,16 @@ it.effect("round trips a typed restore error through its wire schema", () =>
assert.deepStrictEqual(RestoreProtocol.replicaErrorFromWire(decoded), original)
}))
+it.effect("reconstructs ReplicaMetadataMissing from the wire", () =>
+ Effect.gen(function*() {
+ const original = new ReplicaError.ReplicaError({
+ reason: new ReplicaError.ReplicaMetadataMissing()
+ })
+ const encoded = RestoreProtocol.encodeReplicaError(original, 4_096)
+ const decoded = yield* Schema.decodeUnknownEffect(RestoreProtocol.RestoreWireError)(encoded)
+ assert.deepStrictEqual(RestoreProtocol.replicaErrorFromWire(decoded), original)
+ }))
+
it.effect("round trips a superseded checkpoint reason through the restore wire", () =>
Effect.gen(function*() {
const original = new ReplicaError.ReplicaError({
@@ -237,6 +247,7 @@ it.effect("encodes every restore error and defect at the minimum configured budg
new ReplicaError.StorageUnavailable({ cause }),
new ReplicaError.CanonicalEncodeError({ cause }),
new ReplicaError.StorageCorrupt({ cause }),
+ new ReplicaError.ReplicaMetadataMissing(),
new ReplicaError.QuotaExceeded({ resource: "resource".repeat(32), limit: 1 }),
new ReplicaError.MigrationFailed({ migration: "migration".repeat(32), cause }),
new ReplicaError.BackupInvalid({ cause }),
diff --git a/packages/local-rpc/test/PeerRpcServer.test.ts b/packages/local-rpc/test/PeerRpcServer.test.ts
index bf4990f..f4b3a30 100644
--- a/packages/local-rpc/test/PeerRpcServer.test.ts
+++ b/packages/local-rpc/test/PeerRpcServer.test.ts
@@ -263,6 +263,7 @@ const makeFixture = (options: {
admit: Effect.acquireRelease(Effect.succeed(permit), () => Effect.void),
claim: (use) => use(permit),
refresh: Effect.succeed(permit),
+ preflight: () => Effect.void,
validate: () => Effect.void
})
const sync = PeerSync.PeerSync.of({
diff --git a/packages/local-sql/src/Compaction.ts b/packages/local-sql/src/Compaction.ts
index 459c7e6..159e4fe 100644
--- a/packages/local-sql/src/Compaction.ts
+++ b/packages/local-sql/src/Compaction.ts
@@ -3,6 +3,7 @@ import * as Canonical from "@lucas-barake/effect-local/Canonical"
import type * as Document from "@lucas-barake/effect-local/Document"
import * as Identity from "@lucas-barake/effect-local/Identity"
import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
+import * as Cause from "effect/Cause"
import * as Context from "effect/Context"
import * as Crypto from "effect/Crypto"
import * as DateTime from "effect/DateTime"
@@ -14,6 +15,7 @@ import * as Schema from "effect/Schema"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as SqlSchema from "effect/unstable/sql/SqlSchema"
import * as InternalAutomerge from "./internal/automerge.js"
+import * as InternalReplicaError from "./internal/replicaError.js"
import * as WriterProvenance from "./internal/writerProvenance.js"
import * as Recovery from "./Recovery.js"
import * as ReplicaGate from "./ReplicaGate.js"
@@ -1058,7 +1060,7 @@ export const layer: Layer.Layer<
}
const definitionHash = yield* findDefinitionHash(undefined).pipe(
Effect.map((row) => row.definition_hash),
- Effect.catchTag("NoSuchElementError", () => Effect.die(new Error("Replica metadata was not initialized")))
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing)
)
// The stored schema version, not `document.version`: `recover` decodes at the version the
// row records and does not migrate, so the value the re-rooted change carries is encoded at
@@ -1166,7 +1168,7 @@ export const layer: Layer.Layer<
}
const commitSequence = yield* nextCommitSequence(undefined).pipe(
Effect.map((row) => row.commit_sequence),
- Effect.catchTag("NoSuchElementError", () => Effect.die(new Error("Replica metadata was not initialized")))
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing)
)
yield* sql`DELETE FROM effect_local_changes WHERE document_id = ${documentId}`
yield* sql`DELETE FROM effect_local_checkpoints WHERE document_id = ${documentId}`
diff --git a/packages/local-sql/src/DocumentStore.ts b/packages/local-sql/src/DocumentStore.ts
index 5ce89a6..319750b 100644
--- a/packages/local-sql/src/DocumentStore.ts
+++ b/packages/local-sql/src/DocumentStore.ts
@@ -4,6 +4,7 @@ import * as Identity from "@lucas-barake/effect-local/Identity"
import type * as Mutation from "@lucas-barake/effect-local/Mutation"
import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
import type * as Snapshot from "@lucas-barake/effect-local/Snapshot"
+import * as Cause from "effect/Cause"
import * as Context from "effect/Context"
import type * as Crypto from "effect/Crypto"
import * as DateTime from "effect/DateTime"
@@ -14,6 +15,7 @@ import * as Schema from "effect/Schema"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as SqlSchema from "effect/unstable/sql/SqlSchema"
import * as InternalAutomerge from "./internal/automerge.js"
+import * as InternalReplicaError from "./internal/replicaError.js"
import * as Rows from "./internal/rows.js"
import * as Recovery from "./Recovery.js"
import * as ReplicaGate from "./ReplicaGate.js"
@@ -93,8 +95,8 @@ export const layer: Layer.Layer<
WHERE singleton = 1 RETURNING commit_sequence`
})(undefined).pipe(
Effect.map((row) => row.commit_sequence),
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing),
Effect.catchTags({
- NoSuchElementError: () => Effect.die(new Error("Replica metadata was not initialized")),
SchemaError: (cause) =>
Effect.fail(
new ReplicaError.ReplicaError({
@@ -107,8 +109,8 @@ export const layer: Layer.Layer<
)
const currentDefinitionHash = findDefinitionHash(undefined).pipe(
Effect.map((row) => row.definition_hash),
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing),
Effect.catchTags({
- NoSuchElementError: () => Effect.die(new Error("Replica metadata was not initialized")),
SchemaError: (cause) =>
Effect.fail(
new ReplicaError.ReplicaError({
diff --git a/packages/local-sql/src/EntityReplica.ts b/packages/local-sql/src/EntityReplica.ts
index bd9898b..ed01201 100644
--- a/packages/local-sql/src/EntityReplica.ts
+++ b/packages/local-sql/src/EntityReplica.ts
@@ -97,7 +97,7 @@ export const layer = (definition: ReplicaDefinition.Any): Layer.Layer<
})
})
),
- Effect.flatMap((lock) => lock.withPermit(f(permit)))
+ Effect.flatMap((lock) => lock.withPermit(gate.preflight(permit).pipe(Effect.andThen(f(permit)))))
)
)
diff --git a/packages/local-sql/src/PeerSession.ts b/packages/local-sql/src/PeerSession.ts
index 2911a7b..925b56a 100644
--- a/packages/local-sql/src/PeerSession.ts
+++ b/packages/local-sql/src/PeerSession.ts
@@ -453,6 +453,7 @@ const makeWithTerminal = (
Effect.gen(function*() {
const incarnation = yield* Effect.scoped(Effect.gen(function*() {
const permit = yield* gate.shared
+ yield* gate.preflight(permit)
if (permit.incarnation !== session.replicaIncarnation) {
return yield* new ReplicaError.ReplicaError({
reason: new ReplicaError.ProtocolMismatch({
diff --git a/packages/local-sql/src/PeerSync.ts b/packages/local-sql/src/PeerSync.ts
index a73bbfd..dfb0702 100644
--- a/packages/local-sql/src/PeerSync.ts
+++ b/packages/local-sql/src/PeerSync.ts
@@ -4,6 +4,7 @@ import type * as Document from "@lucas-barake/effect-local/Document"
import * as Identity from "@lucas-barake/effect-local/Identity"
import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
import * as ReplicaLimits from "@lucas-barake/effect-local/ReplicaLimits"
+import * as Cause from "effect/Cause"
import * as Clock from "effect/Clock"
import * as Context from "effect/Context"
import * as Crypto from "effect/Crypto"
@@ -18,6 +19,7 @@ import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as SqlSchema from "effect/unstable/sql/SqlSchema"
import * as DocumentStore from "./DocumentStore.js"
import * as InternalAutomerge from "./internal/automerge.js"
+import * as InternalReplicaError from "./internal/replicaError.js"
import * as WriterProvenance from "./internal/writerProvenance.js"
import * as ReplicaBootstrap from "./ReplicaBootstrap.js"
import * as ReplicaGate from "./ReplicaGate.js"
@@ -140,10 +142,6 @@ const CheckpointHashRow = Schema.Struct({
checkpoint_hash: Schema.String
})
-const CommitSequenceRow = Schema.Struct({
- commit_sequence: Schema.Number
-})
-
const CountRow = Schema.Struct({
count: Schema.Number
})
@@ -486,14 +484,14 @@ export const layer: Layer.Layer<
Schema.encodeSync(WriterProvenance.StoredChangeProvenances)(request.writerProvenance)
}`
})
- const findCommitSequence = SqlSchema.findAll({
+ const findCommitSequence = SqlSchema.findOne({
Request: Schema.Void,
- Result: CommitSequenceRow,
+ Result: Schema.Struct({ commit_sequence: Identity.CommitSequence }),
execute: () => sql`SELECT commit_sequence FROM effect_local_metadata WHERE singleton = 1`
})
- const incrementCommitSequence = SqlSchema.findAll({
+ const incrementCommitSequence = SqlSchema.findOne({
Request: Schema.Void,
- Result: CommitSequenceRow,
+ Result: Schema.Struct({ commit_sequence: Identity.CommitSequence }),
execute: () =>
sql`UPDATE effect_local_metadata SET commit_sequence = commit_sequence + 1
WHERE singleton = 1 RETURNING commit_sequence`
@@ -808,17 +806,15 @@ export const layer: Layer.Layer<
AND accepted_at < ${cutoff}`
}))
- const nextSequence = incrementCommitSequence(undefined).pipe(Effect.flatMap((rows) =>
- rows[0] === undefined
- ? Effect.die(new Error("Replica metadata was not initialized"))
- : Effect.succeed(Identity.CommitSequence.make(rows[0].commit_sequence))
- ))
+ const nextSequence = incrementCommitSequence(undefined).pipe(
+ Effect.map((row) => row.commit_sequence),
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing)
+ )
- const currentSequence = findCommitSequence(undefined).pipe(Effect.flatMap((rows) =>
- rows[0] === undefined
- ? Effect.die(new Error("Replica metadata was not initialized"))
- : Effect.succeed(Identity.CommitSequence.make(rows[0].commit_sequence))
- ))
+ const currentSequence = findCommitSequence(undefined).pipe(
+ Effect.map((row) => row.commit_sequence),
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing)
+ )
const loadWriterProvenance = (documentId: Identity.DocumentId, message: Uint8Array) =>
Effect.gen(function*() {
@@ -1177,15 +1173,6 @@ export const layer: Layer.Layer<
yield* validateSession(permit, session)
const nowMillis = yield* Clock.currentTimeMillis
const acceptedAt = new Date(nowMillis).toISOString()
- yield* quotaLock.withPermit(Effect.gen(function*() {
- yield* validateSessionGeneration(generation, sessionGeneration)
- yield* expirePending(
- receiptSession,
- documentId,
- acceptedAt,
- new Date(nowMillis - limits.maxPendingAgeMillis).toISOString()
- ).pipe(Effect.catchTag("SqlError", failStorageUnavailable))
- }))
const messageHash = yield* digest(message)
const validateReceipt = (receipt: typeof ReceiptRow.Type) =>
Effect.gen(function*() {
@@ -1214,21 +1201,35 @@ export const layer: Layer.Layer<
})
}
})
- const receiptRows = yield* findReceipts({
- replicaIncarnation: receiptSession.replicaIncarnation,
- peerId: receiptSession.peerId,
- connectionEpoch: receiptSession.connectionEpoch,
- receiveSequence
- }).pipe(
- Effect.catchTags({
- SqlError: failStorageUnavailable,
- SchemaError: failStorageCorrupt
- })
+ const receipt = yield* quotaLock.withPermit(
+ sql.withTransaction(
+ Effect.gen(function*() {
+ yield* validateSessionGeneration(generation, sessionGeneration)
+ yield* gate.validate(permit)
+ yield* expirePending(
+ receiptSession,
+ documentId,
+ acceptedAt,
+ new Date(nowMillis - limits.maxPendingAgeMillis).toISOString()
+ ).pipe(Effect.catchTag("SqlError", failStorageUnavailable))
+ const receiptRows = yield* findReceipts({
+ replicaIncarnation: receiptSession.replicaIncarnation,
+ peerId: receiptSession.peerId,
+ connectionEpoch: receiptSession.connectionEpoch,
+ receiveSequence
+ }).pipe(
+ Effect.catchTags({
+ SqlError: failStorageUnavailable,
+ SchemaError: failStorageCorrupt
+ })
+ )
+ const receipt = receiptRows[0]
+ if (receipt !== undefined) yield* validateReceipt(receipt)
+ return receipt
+ })
+ ).pipe(Effect.catchTag("SqlError", failStorageUnavailable))
)
- const receipt = receiptRows[0]
if (receipt !== undefined) {
- yield* validateReceipt(receipt)
- yield* quotaLock.withPermit(validateSessionGeneration(generation, sessionGeneration))
return receivedFromReceipt(documentId, receipt)
}
const validateReceiptQuota = Effect.gen(function*() {
@@ -1785,6 +1786,7 @@ export const layer: Layer.Layer<
const result = yield* quotaLock.withPermit(Effect.gen(function*() {
const result = yield* sql.withTransaction(Effect.gen(function*() {
yield* validateSessionGeneration(generation, sessionGeneration)
+ yield* gate.validate(permit)
const receiptRows = yield* findReceipts({
replicaIncarnation: receiptSession.replicaIncarnation,
peerId: receiptSession.peerId,
@@ -1805,7 +1807,6 @@ export const layer: Layer.Layer<
}))
})
const committedChangeMap = yield* validateExistingChanges(committedChanges)
- yield* gate.validate(permit)
const commitSequence = transition ? yield* nextSequence : yield* currentSequence
for (let index = 0; index < changes.length; index++) {
const change = changes[index]!
diff --git a/packages/local-sql/src/Recovery.ts b/packages/local-sql/src/Recovery.ts
index e4c32dd..bac78a1 100644
--- a/packages/local-sql/src/Recovery.ts
+++ b/packages/local-sql/src/Recovery.ts
@@ -198,10 +198,10 @@ export const make = Effect.gen(function*() {
) =>
Effect.gen(function*() {
const { changes, checkpoints, option } = yield* sql.withTransaction(Effect.gen(function*() {
+ yield* gate.validate(permit)
const option = yield* findDocument(documentId)
const checkpoints = yield* findVerifiedCheckpoints(documentId)
const changes = yield* findChanges(documentId)
- yield* gate.validate(permit)
return { changes, checkpoints, option }
})).pipe(
Effect.catchTags({
diff --git a/packages/local-sql/src/ReplicaBootstrap.ts b/packages/local-sql/src/ReplicaBootstrap.ts
index 350b101..e02b8c7 100644
--- a/packages/local-sql/src/ReplicaBootstrap.ts
+++ b/packages/local-sql/src/ReplicaBootstrap.ts
@@ -3,16 +3,19 @@ import * as DocumentSet from "@lucas-barake/effect-local/DocumentSet"
import * as Identity from "@lucas-barake/effect-local/Identity"
import type * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition"
import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
+import * as Cause from "effect/Cause"
import * as Context from "effect/Context"
import type * as Crypto from "effect/Crypto"
import * as DateTime from "effect/DateTime"
import * as Effect from "effect/Effect"
import * as Layer from "effect/Layer"
+import * as Option from "effect/Option"
import * as Schema from "effect/Schema"
import type * as Migrator from "effect/unstable/sql/Migrator"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import type * as SqlError from "effect/unstable/sql/SqlError"
import * as SqlSchema from "effect/unstable/sql/SqlSchema"
+import * as InternalReplicaError from "./internal/replicaError.js"
import { storageFormatVersion } from "./internal/schema.js"
import * as Migrations from "./Migrations.js"
@@ -27,14 +30,6 @@ export class ReplicaBootstrap extends Context.Service()
"@lucas-barake/effect-local-sql/ReplicaBootstrap"
) {}
-// A function, not a shared instance, so each failure captures its own stack at the site that observed it.
-const metadataMissing = () =>
- new ReplicaError.ReplicaError({
- reason: new ReplicaError.StorageCorrupt({
- cause: new Error("Replica metadata is missing")
- })
- })
-
const failStorageCorrupt = (cause: unknown) =>
Effect.fail(
new ReplicaError.ReplicaError({
@@ -63,12 +58,7 @@ export const make = (definition: ReplicaDefinition.Any) =>
Result: Schema.Struct({ storage_format_version: Schema.Int }),
execute: () => sql`SELECT storage_format_version FROM effect_local_metadata WHERE singleton = 1`
})
- // Migration 1 creates the metadata table together with every canonical store table, so whenever the
- // metadata table exists these are readable too. A populated replica whose metadata singleton is gone is
- // corrupt, and must be rejected before the migrator touches it rather than after. The peer tables are
- // deliberately absent here because migration 2 creates them and they may not exist yet; findPopulated
- // below covers them and stays the authoritative check.
- const findPopulatedBeforeMigrating = SqlSchema.findOneOption({
+ const findCanonicalPopulation = SqlSchema.findOne({
Request: Schema.Void,
Result: Schema.Struct({ populated: Schema.Int }),
execute: () =>
@@ -85,31 +75,63 @@ export const make = (definition: ReplicaDefinition.Any) =>
UNION ALL SELECT 1 FROM effect_local_backup_installations
) AS populated`
})
- // The storage format version decides whether this build may touch the database at all, so it has to be
- // checked before Migrations.run, which commits its own transaction. Checking afterwards would mean a build
- // that refuses to open a replica has already migrated it. A database with no tables yet is a fresh one and
- // is left to the migrator. This runs outside the bootstrap transaction below, so it maps its own SchemaError.
- yield* Effect.gen(function*() {
- const table = yield* findMetadataTable(undefined)
- if (table.length === 0) return
- const stored = (yield* findStorageFormat(undefined))[0]
- if (stored === undefined) {
- const populated = yield* findPopulatedBeforeMigrating(undefined)
- if (populated._tag === "Some" && populated.value.populated === 1) {
- return yield* metadataMissing()
- }
- return
- }
- if (stored.storage_format_version === storageFormatVersion) return
- return yield* unsupportedStorageFormat(stored.storage_format_version)
- }).pipe(Effect.catchTag("SchemaError", failStorageCorrupt))
- yield* Migrations.run
+ const findOptionalTables = SqlSchema.findAll({
+ Request: Schema.Void,
+ Result: Schema.Struct({
+ name: Schema.Literals([
+ "effect_local_peer_receipts",
+ "effect_local_peer_outbox",
+ "effect_local_history_rewrites"
+ ])
+ }),
+ execute: () =>
+ sql`SELECT name FROM sqlite_master
+ WHERE type = 'table'
+ AND name IN (
+ 'effect_local_peer_receipts',
+ 'effect_local_peer_outbox',
+ 'effect_local_history_rewrites'
+ )`
+ })
+ const findPeerReceiptPopulation = SqlSchema.findOne({
+ Request: Schema.Void,
+ Result: Schema.Struct({ populated: Schema.Int }),
+ execute: () => sql`SELECT EXISTS (SELECT 1 FROM effect_local_peer_receipts) AS populated`
+ })
+ const findPeerOutboxPopulation = SqlSchema.findOne({
+ Request: Schema.Void,
+ Result: Schema.Struct({ populated: Schema.Int }),
+ execute: () => sql`SELECT EXISTS (SELECT 1 FROM effect_local_peer_outbox) AS populated`
+ })
+ const findHistoryRewritePopulation = SqlSchema.findOne({
+ Request: Schema.Void,
+ Result: Schema.Struct({ populated: Schema.Int }),
+ execute: () => sql`SELECT EXISTS (SELECT 1 FROM effect_local_history_rewrites) AS populated`
+ })
+ const isPopulated = (
+ effect: Effect.Effect<{ readonly populated: number }, E, R>
+ ) =>
+ effect.pipe(
+ Effect.catchNoSuchElement,
+ Effect.flatMap((option) =>
+ Option.isNone(option)
+ ? failStorageCorrupt(new TypeError("Population query returned no row"))
+ : Effect.succeed(option.value)
+ ),
+ Effect.flatMap((row) =>
+ row.populated === 0
+ ? Effect.succeed(false)
+ : row.populated === 1
+ ? Effect.succeed(true)
+ : failStorageCorrupt(new TypeError(`Invalid population result: ${row.populated}`))
+ )
+ )
const findMetadata = SqlSchema.findAll({
Request: Schema.Void,
Result: Schema.Struct({ singleton: Schema.Int }),
execute: () => sql`SELECT singleton FROM effect_local_metadata WHERE singleton = 1`
})
- const findPopulated = SqlSchema.findOneOption({
+ const findPopulated = SqlSchema.findOne({
Request: Schema.Void,
Result: Schema.Struct({ populated: Schema.Int }),
execute: () =>
@@ -124,6 +146,7 @@ export const make = (definition: ReplicaDefinition.Any) =>
UNION ALL SELECT 1 FROM effect_local_commit_outbox
UNION ALL SELECT 1 FROM effect_local_quarantine
UNION ALL SELECT 1 FROM effect_local_backup_installations
+ UNION ALL SELECT 1 FROM effect_local_history_rewrites
UNION ALL SELECT 1 FROM effect_local_peer_receipts
UNION ALL SELECT 1 FROM effect_local_peer_outbox
) AS populated`
@@ -155,12 +178,45 @@ export const make = (definition: ReplicaDefinition.Any) =>
sql`SELECT replica_id, replica_incarnation, writer_generation
FROM effect_local_metadata WHERE singleton = 1`
})
+ yield* sql.withTransaction(Effect.gen(function*() {
+ const table = yield* findMetadataTable(undefined)
+ if (table.length !== 0) {
+ const stored = (yield* findStorageFormat(undefined))[0]
+ if (stored === undefined) {
+ if (yield* isPopulated(findCanonicalPopulation(undefined))) {
+ return yield* InternalReplicaError.metadataMissing()
+ }
+ const optionalTables = new Set((yield* findOptionalTables(undefined)).map((row) => row.name))
+ if (
+ optionalTables.has("effect_local_peer_receipts") &&
+ (yield* isPopulated(findPeerReceiptPopulation(undefined)))
+ ) {
+ return yield* InternalReplicaError.metadataMissing()
+ }
+ if (
+ optionalTables.has("effect_local_peer_outbox") &&
+ (yield* isPopulated(findPeerOutboxPopulation(undefined)))
+ ) {
+ return yield* InternalReplicaError.metadataMissing()
+ }
+ if (
+ optionalTables.has("effect_local_history_rewrites") &&
+ (yield* isPopulated(findHistoryRewritePopulation(undefined)))
+ ) {
+ return yield* InternalReplicaError.metadataMissing()
+ }
+ } else if (stored.storage_format_version !== storageFormatVersion) {
+ return yield* unsupportedStorageFormat(stored.storage_format_version)
+ }
+ }
+ yield* Migrations.run
+ })).pipe(Effect.catchTag("SchemaError", failStorageCorrupt))
+
return yield* sql.withTransaction(Effect.gen(function*() {
const metadata = yield* findMetadata(undefined)
if (metadata.length === 0) {
- const populated = yield* findPopulated(undefined)
- if (populated._tag === "Some" && populated.value.populated === 1) {
- return yield* metadataMissing()
+ if (yield* isPopulated(findPopulated(undefined))) {
+ return yield* InternalReplicaError.metadataMissing()
}
yield* sql`INSERT INTO effect_local_metadata (
singleton,
@@ -182,7 +238,7 @@ export const make = (definition: ReplicaDefinition.Any) =>
}
const format = (yield* findFormat(undefined))[0]
if (format === undefined) {
- return yield* metadataMissing()
+ return yield* InternalReplicaError.metadataMissing()
}
if (format.storage_format_version !== storageFormatVersion) {
return yield* unsupportedStorageFormat(format.storage_format_version)
@@ -207,7 +263,7 @@ export const make = (definition: ReplicaDefinition.Any) =>
}
yield* sql`UPDATE effect_local_metadata SET writer_generation = writer_generation + 1 WHERE singleton = 1`
const row = yield* findPermit(undefined).pipe(
- Effect.catchTag("NoSuchElementError", () => metadataMissing())
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing)
)
yield* sql`INSERT INTO effect_local_writer_generations (generation, claimed_at)
VALUES (${row.writer_generation}, ${DateTime.formatIso(yield* DateTime.now)})`
diff --git a/packages/local-sql/src/ReplicaGate.ts b/packages/local-sql/src/ReplicaGate.ts
index b42ba43..beacc07 100644
--- a/packages/local-sql/src/ReplicaGate.ts
+++ b/packages/local-sql/src/ReplicaGate.ts
@@ -1,6 +1,7 @@
import * as Identity from "@lucas-barake/effect-local/Identity"
import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
import * as ReplicaLimits from "@lucas-barake/effect-local/ReplicaLimits"
+import * as Cause from "effect/Cause"
import * as Context from "effect/Context"
import * as DateTime from "effect/DateTime"
import * as Deferred from "effect/Deferred"
@@ -15,6 +16,7 @@ import * as TxReentrantLock from "effect/TxReentrantLock"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import type * as SqlError from "effect/unstable/sql/SqlError"
import * as SqlSchema from "effect/unstable/sql/SqlSchema"
+import * as InternalReplicaError from "./internal/replicaError.js"
import * as ReplicaBootstrap from "./ReplicaBootstrap.js"
export interface Permit {
@@ -51,6 +53,7 @@ export class ReplicaGate extends Context.Service Effect.Effect
) => Effect.Effect
readonly refresh: Effect.Effect
+ readonly preflight: (expected: Permit) => Effect.Effect
readonly validate: (expected: Permit) => Effect.Effect
}>()("@lucas-barake/effect-local-sql/ReplicaGate") {}
@@ -274,8 +277,8 @@ export const layer: Layer.Layer<
incarnation: row.replica_incarnation,
writerGeneration: row.writer_generation
})),
+ Effect.catchIf(Cause.isNoSuchElementError, InternalReplicaError.metadataMissing),
Effect.catchTags({
- NoSuchElementError: () => Effect.die(new Error("Replica metadata was not initialized")),
SchemaError: (cause) =>
Effect.fail(
new ReplicaError.ReplicaError({
@@ -310,6 +313,15 @@ export const layer: Layer.Layer<
AND writer_generation = ${expected.writerGeneration}
RETURNING replica_incarnation, writer_generation`
})
+ const fenced = (expected: Permit, observed: Permit) =>
+ Effect.fail(
+ new ReplicaError.ReplicaError({
+ reason: new ReplicaError.ReplicaFenced({
+ expectedGeneration: expected.writerGeneration,
+ observedGeneration: observed.writerGeneration
+ })
+ })
+ )
// A shed acquirer never reaches `readLock` or `release`, because `Effect.acquireUseRelease` installs
// its release only after `acquire` succeeds. So the gate's own occupancy is untouched by a shed.
const sharedWith = (admission: Effect.Effect) => {
@@ -331,6 +343,15 @@ export const layer: Layer.Layer<
// here proves `state` already carries the generation of the last claim that touched the database.
claiming: Ref.get(writer).pipe(Effect.map((owner) => owner !== null)),
refresh: readState.pipe(Effect.tap((next) => Ref.set(state, next))),
+ preflight: (expected) =>
+ readState.pipe(
+ Effect.flatMap((observed) =>
+ observed.incarnation === expected.incarnation &&
+ observed.writerGeneration === expected.writerGeneration
+ ? Effect.void
+ : fenced(expected, observed)
+ )
+ ),
shared: sharedWith(acquire),
admit: sharedWith(acquireBounded),
claim: (use) =>
@@ -381,16 +402,7 @@ export const layer: Layer.Layer<
validateState(expected).pipe(
Effect.flatMap((rows) =>
rows.length === 1 ? Effect.void : readState.pipe(
- Effect.flatMap((observed) =>
- Effect.fail(
- new ReplicaError.ReplicaError({
- reason: new ReplicaError.ReplicaFenced({
- expectedGeneration: expected.writerGeneration,
- observedGeneration: observed.writerGeneration
- })
- })
- )
- )
+ Effect.flatMap((observed) => fenced(expected, observed))
)
),
Effect.catchTags({
diff --git a/packages/local-sql/src/internal/replicaError.ts b/packages/local-sql/src/internal/replicaError.ts
new file mode 100644
index 0000000..bf6a368
--- /dev/null
+++ b/packages/local-sql/src/internal/replicaError.ts
@@ -0,0 +1,9 @@
+import * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
+import * as Effect from "effect/Effect"
+
+export const metadataMissing = () =>
+ Effect.fail(
+ new ReplicaError.ReplicaError({
+ reason: new ReplicaError.ReplicaMetadataMissing()
+ })
+ )
diff --git a/packages/local-sql/test/Compaction.test.ts b/packages/local-sql/test/Compaction.test.ts
index 4f55b51..c236a78 100644
--- a/packages/local-sql/test/Compaction.test.ts
+++ b/packages/local-sql/test/Compaction.test.ts
@@ -6,8 +6,10 @@ import * as Document from "@lucas-barake/effect-local/Document"
import * as DocumentSet from "@lucas-barake/effect-local/DocumentSet"
import * as Identity from "@lucas-barake/effect-local/Identity"
import * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition"
+import * as Cause from "effect/Cause"
import * as Effect from "effect/Effect"
import * as Layer from "effect/Layer"
+import * as Option from "effect/Option"
import * as Schema from "effect/Schema"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as CommandExecutor from "../src/CommandExecutor.js"
@@ -49,6 +51,50 @@ describe("Compaction", () => {
Layer.provide(Layer.mergeAll(Base, Gate, StoreService, Projections))
)
const Services = Layer.mergeAll(Base, Gate, StoreService, RecoveryService, CompactionService, Executor)
+ const servicesWithDefinitionHashRace = Effect.sync(() => {
+ let armed = false
+ const sqlite = SqliteClient.layer({ filename: ":memory:", disableWAL: true })
+ const instrumented = Layer.effect(
+ SqlClient.SqlClient,
+ Effect.map(SqlClient.SqlClient, (sql) =>
+ new Proxy(sql, {
+ apply(target, thisArg, args: Array) {
+ const statement = Reflect.apply(target as never, thisArg, args) as Effect.Effect<
+ unknown,
+ unknown,
+ never
+ >
+ const strings = args[0]
+ if (
+ !armed ||
+ !Array.isArray(strings) ||
+ !(strings as ReadonlyArray).join("").includes(
+ "SELECT definition_hash FROM effect_local_metadata"
+ )
+ ) {
+ return statement
+ }
+ armed = false
+ return sql`DELETE FROM effect_local_metadata WHERE singleton = 1`.pipe(
+ Effect.andThen(statement)
+ )
+ }
+ }) as typeof sql)
+ ).pipe(Layer.provide(sqlite))
+ const database = Layer.merge(instrumented, NodeCrypto.layer)
+ const bootstrap = ReplicaBootstrap.layer(definition).pipe(Layer.provide(database))
+ const base = Layer.merge(database, bootstrap)
+ const gate = ReplicaGate.layer.pipe(withGateLimits, Layer.provide(base))
+ const store = DocumentStore.layer.pipe(Layer.provide(Layer.merge(base, gate)))
+ const recovery = Recovery.layer.pipe(Layer.provide(Layer.mergeAll(base, gate)))
+ const compaction = Compaction.layer.pipe(Layer.provide(Layer.mergeAll(base, gate, recovery)))
+ return {
+ arm: Effect.sync(() => {
+ armed = true
+ }),
+ services: Layer.mergeAll(base, gate, store, recovery, compaction)
+ }
+ })
/** Commits a real create command under the current incarnation and reports the receipt it wrote. */
const commitReceipt = Effect.gen(function*() {
@@ -537,6 +583,84 @@ describe("Compaction", () => {
assert.deepStrictEqual(yield* listReceipts, byKey([superseded, live]))
}).pipe(Effect.provide(Services)))
+ it.effect("reports ReplicaMetadataMissing when metadata disappears after rewrite recovery", () =>
+ Effect.gen(function*() {
+ const harness = yield* servicesWithDefinitionHashRace
+ yield* Effect.gen(function*() {
+ const compaction = yield* Compaction.Compaction
+ const sql = yield* SqlClient.SqlClient
+ const store = yield* DocumentStore.DocumentStore
+ const documentId = yield* Identity.makeDocumentId
+ const created = yield* store.create(Task, documentId, { title: "one", labels: [] })
+ InternalAutomerge.free(created.automerge)
+ yield* compaction.compact(Task, documentId)
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ yield* Effect.addFinalizer(() =>
+ sql`INSERT OR IGNORE INTO effect_local_metadata ${sql.insert(metadata[0]!)}`.pipe(
+ Effect.asVoid,
+ Effect.orDie
+ )
+ )
+ const documentBefore = yield* documentRowOf(documentId)
+ const checkpointsBefore = yield* checkpointsOf(documentId)
+ const changesBefore = yield* changesOf(documentId)
+ const outboxBefore = yield* commitOutboxOf(documentId)
+ const markersBefore = yield* rewriteMarkers
+ yield* harness.arm
+
+ const exit = yield* Effect.exit(compaction.rewriteHistory(Task, documentId, operationId))
+ assert.strictEqual(exit._tag, "Failure")
+ if (exit._tag !== "Failure") return
+ assert.isFalse(Cause.hasDies(exit.cause))
+ const error = Option.getOrThrow(Cause.findErrorOption(exit.cause))
+ assert.strictEqual(error.reason._tag, "ReplicaMetadataMissing")
+ assert.deepStrictEqual(yield* documentRowOf(documentId), documentBefore)
+ assert.deepStrictEqual(yield* checkpointsOf(documentId), checkpointsBefore)
+ assert.deepStrictEqual(yield* changesOf(documentId), changesBefore)
+ assert.deepStrictEqual(yield* commitOutboxOf(documentId), outboxBefore)
+ assert.deepStrictEqual(yield* rewriteMarkers, markersBefore)
+ }).pipe(Effect.scoped, Effect.provide(harness.services))
+ }))
+
+ it.effect("reports ReplicaMetadataMissing when metadata disappears before the rewrite commit sequence", () =>
+ Effect.gen(function*() {
+ const compaction = yield* Compaction.Compaction
+ const sql = yield* SqlClient.SqlClient
+ const store = yield* DocumentStore.DocumentStore
+ const documentId = yield* Identity.makeDocumentId
+ const created = yield* store.create(Task, documentId, { title: "one", labels: [] })
+ InternalAutomerge.free(created.automerge)
+ yield* compaction.compact(Task, documentId)
+ const metadataBefore = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ const documentBefore = yield* documentRowOf(documentId)
+ const checkpointsBefore = yield* checkpointsOf(documentId)
+ const changesBefore = yield* changesOf(documentId)
+ const outboxBefore = yield* commitOutboxOf(documentId)
+ const markersBefore = yield* rewriteMarkers
+ yield* sql`CREATE TRIGGER drop_metadata_before_rewrite_sequence
+ AFTER UPDATE ON effect_local_documents
+ BEGIN
+ DELETE FROM effect_local_metadata WHERE singleton = 1;
+ END`
+
+ const exit = yield* Effect.exit(compaction.rewriteHistory(Task, documentId, operationId))
+ assert.strictEqual(exit._tag, "Failure")
+ if (exit._tag !== "Failure") return
+ assert.isFalse(Cause.hasDies(exit.cause))
+ const error = Option.getOrThrow(Cause.findErrorOption(exit.cause))
+ assert.strictEqual(error.reason._tag, "ReplicaMetadataMissing")
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`,
+ metadataBefore
+ )
+ assert.deepStrictEqual(yield* documentRowOf(documentId), documentBefore)
+ assert.deepStrictEqual(yield* checkpointsOf(documentId), checkpointsBefore)
+ assert.deepStrictEqual(yield* changesOf(documentId), changesBefore)
+ assert.deepStrictEqual(yield* commitOutboxOf(documentId), outboxBefore)
+ assert.deepStrictEqual(yield* rewriteMarkers, markersBefore)
+ }).pipe(Effect.provide(Services)))
+
it.effect("rewrites a high-churn document down to a single change and checkpoint", () =>
Effect.gen(function*() {
const compaction = yield* Compaction.Compaction
diff --git a/packages/local-sql/test/DocumentEntity.test.ts b/packages/local-sql/test/DocumentEntity.test.ts
index d738216..5210faa 100644
--- a/packages/local-sql/test/DocumentEntity.test.ts
+++ b/packages/local-sql/test/DocumentEntity.test.ts
@@ -118,6 +118,7 @@ describe("DocumentEntity", () => {
admit: Effect.die("unused"),
claim: (use) => use(permit),
refresh: Effect.succeed(permit),
+ preflight: () => Effect.void,
validate: () => Effect.void
})
diff --git a/packages/local-sql/test/DocumentStore.test.ts b/packages/local-sql/test/DocumentStore.test.ts
index d7cacb1..aa2df59 100644
--- a/packages/local-sql/test/DocumentStore.test.ts
+++ b/packages/local-sql/test/DocumentStore.test.ts
@@ -9,6 +9,7 @@ import * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition
import type * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
import * as Effect from "effect/Effect"
import * as Layer from "effect/Layer"
+import * as Result from "effect/Result"
import * as Schema from "effect/Schema"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import { vi } from "vitest"
@@ -80,6 +81,46 @@ describe("DocumentStore", () => {
assert.throws(() => InternalAutomerge.heads(leaked!))
}).pipe(Effect.provide(Store)))
+ it.effect("fails create with ReplicaMetadataMissing when metadata is absent", () =>
+ Effect.gen(function*() {
+ const store = yield* DocumentStore.DocumentStore
+ const sql = yield* SqlClient.SqlClient
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+
+ const result = yield* Effect.result(
+ store.create(Task, yield* Identity.makeDocumentId, { title: "missing", labels: [] })
+ )
+
+ assert.isTrue(Result.isFailure(result))
+ if (!Result.isFailure(result)) return
+ assert.strictEqual(result.failure.reason._tag as string, "ReplicaMetadataMissing")
+ const rows = yield* sql<{
+ readonly changes: number
+ readonly documents: number
+ readonly outbox: number
+ }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_changes) AS changes,
+ (SELECT COUNT(*) FROM effect_local_documents) AS documents,
+ (SELECT COUNT(*) FROM effect_local_commit_outbox) AS outbox`
+ assert.deepStrictEqual(rows[0], { changes: 0, documents: 0, outbox: 0 })
+ }).pipe(Effect.provide(Store)))
+
+ it.effect("fails load with ReplicaMetadataMissing when metadata is absent", () =>
+ Effect.gen(function*() {
+ const store = yield* DocumentStore.DocumentStore
+ const sql = yield* SqlClient.SqlClient
+ const documentId = yield* Identity.makeDocumentId
+ const created = yield* store.create(Task, documentId, { title: "one", labels: [] })
+ InternalAutomerge.free(created.automerge)
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+
+ const result = yield* Effect.result(store.load(Task, documentId))
+
+ assert.isTrue(Result.isFailure(result))
+ if (!Result.isFailure(result)) return
+ assert.strictEqual(result.failure.reason._tag as string, "ReplicaMetadataMissing")
+ }).pipe(Effect.provide(Store)))
+
it.effect("rolls application rows back with an outer transaction", () =>
Effect.gen(function*() {
const store = yield* DocumentStore.DocumentStore
diff --git a/packages/local-sql/test/DurableRuntime.test.ts b/packages/local-sql/test/DurableRuntime.test.ts
index 52adc4e..adf2085 100644
--- a/packages/local-sql/test/DurableRuntime.test.ts
+++ b/packages/local-sql/test/DurableRuntime.test.ts
@@ -20,12 +20,14 @@ import * as Fiber from "effect/Fiber"
import * as Latch from "effect/Latch"
import * as Layer from "effect/Layer"
import * as Option from "effect/Option"
+import * as PrimaryKey from "effect/PrimaryKey"
import * as Ref from "effect/Ref"
import * as Schema from "effect/Schema"
import * as TestClock from "effect/testing/TestClock"
import * as MessageStorage from "effect/unstable/cluster/MessageStorage"
import * as Runners from "effect/unstable/cluster/Runners"
import * as Sharding from "effect/unstable/cluster/Sharding"
+import * as Snowflake from "effect/unstable/cluster/Snowflake"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as DurableClock from "effect/unstable/workflow/DurableClock"
import * as Workflow from "effect/unstable/workflow/Workflow"
@@ -39,6 +41,7 @@ import * as DocumentEntity from "../src/DocumentEntity.js"
import * as DocumentStore from "../src/DocumentStore.js"
import * as DurableRuntime from "../src/DurableRuntime.js"
import * as InternalAutomerge from "../src/internal/automerge.js"
+import * as WriterProvenance from "../src/internal/writerProvenance.js"
import * as Recovery from "../src/Recovery.js"
import * as ReplicaBootstrap from "../src/ReplicaBootstrap.js"
import * as ReplicaGate from "../src/ReplicaGate.js"
@@ -61,11 +64,29 @@ describe("DurableRuntime", () => {
projections: [],
queries: []
})
+ const provenanceFor = (message: Uint8Array) =>
+ WriterProvenance.syncMessageChangeHashes(message).map((changeHash) => ({
+ changeHash,
+ writerSchemaVersion: Task.version,
+ writerDefinitionHash: definition.hash
+ }))
+ const keyOf = (payload: unknown) => {
+ if (!PrimaryKey.isPrimaryKey(payload)) throw new TypeError("Expected a primary key payload")
+ return PrimaryKey.value(payload)
+ }
+ type MetadataRow = {
+ readonly commit_sequence: Identity.CommitSequence
+ readonly definition_hash: string
+ readonly replica_id: Identity.ReplicaId
+ readonly replica_incarnation: Identity.ReplicaIncarnation
+ readonly singleton: number
+ readonly storage_format_version: number
+ readonly writer_generation: Identity.WriterGeneration
+ }
const Database = Layer.merge(
SqliteClient.layer({ filename: ":memory:", disableWAL: true }),
NodeCrypto.layer
)
- const Bootstrap = ReplicaBootstrap.layer(definition).pipe(Layer.provide(Database))
const Executor = Layer.succeed(
CommandExecutor.CommandExecutor,
CommandExecutor.CommandExecutor.of({
@@ -78,7 +99,7 @@ describe("DurableRuntime", () => {
lookupDelete: (id) => Effect.succeed(CommandOutcome.unknown(id))
})
)
- const Limits = ReplicaLimits.layer({
+ const limits: ReplicaLimits.Values = {
maxBackupBytes: 1_000_000,
maxChunkBytes: 64_000,
maxArchiveRecords: 1_000,
@@ -109,14 +130,110 @@ describe("DurableRuntime", () => {
maxRestorePullMillis: 10_000,
maxRestoreCoalesceMillis: 25,
maxRestoreErrorBytes: 4_096
+ }
+ const Limits = ReplicaLimits.layer(limits)
+ const servicesWithDatabase = (
+ database: Layer.Layer
+ ) => {
+ const bootstrap = ReplicaBootstrap.layer(definition).pipe(Layer.provide(database))
+ const gate = ReplicaGate.layer.pipe(Layer.provide(Limits), Layer.provide(Layer.merge(database, bootstrap)))
+ const store = DocumentStore.layer.pipe(Layer.provide(Layer.merge(database, gate)))
+ const recovery = Recovery.layer.pipe(Layer.provide(Layer.mergeAll(database, gate)))
+ const compaction = Compaction.layer.pipe(Layer.provide(Layer.mergeAll(database, gate, recovery)))
+ const inputs = Layer.mergeAll(database, bootstrap, Executor, Limits, gate, store, recovery, compaction)
+ const live = DurableRuntime.layer(definition).pipe(Layer.provide(inputs))
+ return { bootstrap, compaction, gate, inputs, layer: Layer.merge(inputs, live), live, recovery, store }
+ }
+ const Default = servicesWithDatabase(Database)
+ const Bootstrap = Default.bootstrap
+ const Gate = Default.gate
+ const Store = Default.store
+ const RecoveryService = Default.recovery
+ const CompactionService = Default.compaction
+ const Live = Default.live
+ const Services = Default.layer
+
+ const rollbackServices = Effect.gen(function*() {
+ const completed = yield* Deferred.make>()
+ const rolledBack = yield* Deferred.make()
+ let armed = false
+ let targetConnection: unknown | undefined
+ const baseDatabase = SqliteClient.layer({ filename: ":memory:", disableWAL: true })
+ const instrumentedDatabase = Layer.effect(
+ SqlClient.SqlClient,
+ Effect.gen(function*() {
+ const sql = yield* SqlClient.SqlClient
+ const call = ((...args: ReadonlyArray) => {
+ const statement = (sql as any)(...args)
+ const text = Array.isArray(args[0]) ? args[0].join("") : ""
+ if (!text.includes("UPDATE effect_local_metadata SET writer_generation = writer_generation")) {
+ return statement
+ }
+ return Effect.serviceOption(sql.transactionService).pipe(
+ Effect.tap((transaction) =>
+ Effect.sync(() => {
+ if (armed && Option.isSome(transaction)) targetConnection = transaction.value[0]
+ })
+ ),
+ Effect.andThen(statement)
+ )
+ }) as SqlClient.SqlClient
+ return Object.assign(
+ call,
+ sql,
+ {
+ withTransaction: (effect: Effect.Effect) =>
+ Effect.suspend(() => {
+ let connection: unknown | undefined
+ return Effect.serviceOption(sql.transactionService).pipe(
+ Effect.flatMap((parent) =>
+ sql.withTransaction(
+ Effect.serviceOption(sql.transactionService).pipe(
+ Effect.tap((transaction) =>
+ Effect.sync(() => {
+ if (Option.isSome(transaction)) connection = transaction.value[0]
+ })
+ ),
+ Effect.andThen(effect)
+ )
+ ).pipe(
+ Effect.onExit((exit) => {
+ if (
+ !armed ||
+ Option.isSome(parent) ||
+ connection === undefined ||
+ connection !== targetConnection
+ ) {
+ return Effect.void
+ }
+ return Deferred.succeed(completed, exit).pipe(
+ Effect.andThen(
+ Exit.isFailure(exit) ? Deferred.succeed(rolledBack, undefined) : Effect.void
+ )
+ )
+ })
+ )
+ )
+ )
+ })
+ }
+ )
+ })
+ ).pipe(Layer.provideMerge(baseDatabase))
+ const database = Layer.merge(instrumentedDatabase, NodeCrypto.layer)
+ return {
+ arm: Effect.sync(() => {
+ armed = true
+ }),
+ completed,
+ database,
+ disarm: Effect.sync(() => {
+ armed = false
+ }),
+ rolledBack,
+ services: servicesWithDatabase(database).layer
+ }
})
- const Gate = ReplicaGate.layer.pipe(Layer.provide(Limits), Layer.provide(Layer.merge(Database, Bootstrap)))
- const Store = DocumentStore.layer.pipe(Layer.provide(Layer.merge(Database, Gate)))
- const RecoveryService = Recovery.layer.pipe(Layer.provide(Layer.mergeAll(Database, Gate)))
- const CompactionService = Compaction.layer.pipe(Layer.provide(Layer.mergeAll(Database, Gate, RecoveryService)))
- const Inputs = Layer.mergeAll(Database, Bootstrap, Executor, Limits, Gate, Store, RecoveryService, CompactionService)
- const Live = DurableRuntime.layer(definition).pipe(Layer.provide(Inputs))
- const Services = Layer.merge(Inputs, Live)
const servicesAtWith = (filename: string, workflowRegistrations: Layer.Layer) => {
const database = Layer.merge(SqliteClient.layer({ filename, disableWAL: true }), NodeCrypto.layer)
@@ -412,6 +529,358 @@ describe("DurableRuntime", () => {
assert.strictEqual((yield* Effect.exit(runtime.interrupt(forged)))._tag, "Failure")
}).pipe(Effect.provide(Services)))
+ it.effect("keeps missing metadata ApplySync messages retryable", () =>
+ Effect.gen(function*() {
+ const harness = yield* rollbackServices
+ yield* Effect.gen(function*() {
+ const entity = yield* DocumentEntity.DocumentEntity.client
+ const gate = yield* ReplicaGate.ReplicaGate
+ const sharding = yield* Sharding.Sharding
+ const sql = yield* SqlClient.SqlClient
+ const store = yield* DocumentStore.DocumentStore
+ const documentId = yield* Identity.makeDocumentId
+ const peerId = yield* Identity.makePeerId
+ const created = yield* store.create(Task, documentId, { title: "local" })
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(created.automerge)))
+ let remote = Automerge.change(
+ Automerge.clone(created.automerge, { actor: "1".repeat(32) }),
+ (draft) => {
+ ;(draft.value as { title: string }).title = "remote"
+ }
+ )
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(remote)))
+ const emptyPeer = Automerge.init>()
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(emptyPeer)))
+ const localHandshake = Automerge.generateSyncMessage(emptyPeer, Automerge.initSyncState())[1]!
+ const receivedHandshake = Automerge.receiveSyncMessage(
+ remote,
+ Automerge.initSyncState(),
+ localHandshake
+ )
+ remote = receivedHandshake[0]
+ const message = Automerge.generateSyncMessage(remote, receivedHandshake[1])[1]!
+ assert.isAbove(Automerge.decodeSyncMessage(message).changes.length, 0)
+ const messageHash = yield* Canonical.digest(message)
+ const writerProvenance = provenanceFor(message)
+ const permit = yield* gate.current
+ const payload = {
+ replicaIncarnation: permit.incarnation,
+ peerId,
+ connectionEpoch: "remote-epoch",
+ localConnectionEpoch: "local-epoch",
+ receiveSequence: 0,
+ documentType: Task.name,
+ messageHash,
+ message,
+ writerProvenance
+ } as const
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ const before = yield* sql<{
+ readonly changes: number
+ readonly commitSequence: number
+ readonly outbox: number
+ }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_changes) AS changes,
+ (SELECT commit_sequence FROM effect_local_metadata WHERE singleton = 1) AS commitSequence,
+ (SELECT COUNT(*) FROM effect_local_commit_outbox) AS outbox`
+ const primaryKey = keyOf(DocumentEntity.ApplySync.payloadSchema.make(payload))
+ const messageId = `EffectLocal/Document/${documentId}/ApplySync/${primaryKey}`
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+ yield* harness.arm
+
+ const first = yield* entity(documentId).ApplySync(payload).pipe(
+ Effect.result,
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(first).pipe(
+ Effect.andThen(Fiber.await(first)),
+ Effect.asVoid
+ )
+ )
+ const transactionExit = yield* Deferred.await(harness.completed)
+ assert.isTrue(Exit.isFailure(transactionExit))
+ if (!Exit.isFailure(transactionExit)) return
+ assert.isTrue(Cause.hasFails(transactionExit.cause))
+ assert.isFalse(Cause.hasDies(transactionExit.cause))
+ const failure = Option.getOrThrow(Cause.findErrorOption(transactionExit.cause))
+ assert.strictEqual((failure as { readonly _tag?: string })._tag, "ReplicaError")
+ if ((failure as { readonly _tag?: string })._tag !== "ReplicaError") return
+ assert.strictEqual(
+ (failure as { readonly reason: { readonly _tag: string } }).reason._tag,
+ "ReplicaMetadataMissing"
+ )
+ yield* Deferred.await(harness.rolledBack)
+
+ const pending = yield* sql<{
+ readonly missing_reply: number
+ readonly processed: number
+ readonly replies: number
+ }>`SELECT processed, last_reply_id IS NULL AS missing_reply, (
+ SELECT COUNT(*) FROM effect_local_cluster_replies
+ WHERE request_id = effect_local_cluster_messages.id
+ ) AS replies
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ assert.lengthOf(pending, 1)
+ assert.strictEqual(pending[0]!.processed, 0)
+ assert.strictEqual(pending[0]!.missing_reply, 1)
+ assert.strictEqual(pending[0]!.replies, 0)
+ const requestIds = yield* sql<{ readonly id: bigint }>`SELECT id
+ FROM effect_local_cluster_messages
+ WHERE message_id = ${messageId}`.pipe(
+ Effect.provideService(SqlClient.SafeIntegers, true)
+ )
+ assert.lengthOf(requestIds, 1)
+
+ yield* sql`INSERT INTO effect_local_metadata ${sql.insert(metadata[0]!)}`
+ yield* harness.disarm
+ assert.isTrue(yield* sharding.reset(Snowflake.Snowflake(requestIds[0]!.id)))
+ yield* sharding.pollStorage
+ yield* Effect.gen(function*() {
+ while (true) {
+ const rows = yield* sql<{
+ readonly processed: number
+ readonly replies: number
+ }>`SELECT processed, (
+ SELECT COUNT(*) FROM effect_local_cluster_replies
+ WHERE request_id = effect_local_cluster_messages.id
+ ) AS replies
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ if (rows[0]?.processed === 1 && rows[0]?.replies === 1) return
+ }
+ }).pipe(Effect.timeout("5 seconds"), TestClock.withLive)
+ const retriedResult = yield* entity(documentId).ApplySync(payload).pipe(
+ Effect.timeout("5 seconds"),
+ TestClock.withLive
+ )
+ assert.isFalse(retriedResult.duplicate)
+
+ const receipts = yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_peer_receipts
+ WHERE replica_incarnation = ${permit.incarnation}
+ AND peer_id = ${peerId}
+ AND connection_epoch = ${payload.connectionEpoch}
+ AND receive_sequence = ${payload.receiveSequence}`
+ assert.strictEqual(receipts[0]?.count, 1)
+ const completed = yield* sql<{
+ readonly has_reply: number
+ readonly processed: number
+ }>`SELECT processed, EXISTS (
+ SELECT 1 FROM effect_local_cluster_replies
+ WHERE id = effect_local_cluster_messages.last_reply_id
+ ) AS has_reply
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ assert.lengthOf(completed, 1)
+ assert.strictEqual(completed[0]!.processed, 1)
+ assert.strictEqual(completed[0]!.has_reply, 1)
+ const replies = yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_replies
+ WHERE request_id = (
+ SELECT id FROM effect_local_cluster_messages WHERE message_id = ${messageId}
+ )`
+ assert.strictEqual(replies[0]?.count, 1)
+ const after = yield* sql<{
+ readonly changes: number
+ readonly commitSequence: number
+ readonly outbox: number
+ }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_changes) AS changes,
+ (SELECT commit_sequence FROM effect_local_metadata WHERE singleton = 1) AS commitSequence,
+ (SELECT COUNT(*) FROM effect_local_commit_outbox) AS outbox`
+ assert.deepStrictEqual(after[0], {
+ changes: before[0]!.changes + 1,
+ commitSequence: before[0]!.commitSequence + 1,
+ outbox: before[0]!.outbox + 1
+ })
+ }).pipe(Effect.scoped, Effect.provide(harness.services))
+ }))
+
+ it.effect("rejects duplicate ApplySync when metadata is absent", () =>
+ Effect.gen(function*() {
+ const harness = yield* rollbackServices
+ yield* Effect.gen(function*() {
+ const entity = yield* DocumentEntity.DocumentEntity.client
+ const gate = yield* ReplicaGate.ReplicaGate
+ const sharding = yield* Sharding.Sharding
+ const sql = yield* SqlClient.SqlClient
+ const store = yield* DocumentStore.DocumentStore
+ const documentId = yield* Identity.makeDocumentId
+ const peerId = yield* Identity.makePeerId
+ const created = yield* store.create(Task, documentId, { title: "local" })
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(created.automerge)))
+ const remote = Automerge.change(
+ Automerge.clone(created.automerge, { actor: "2".repeat(32) }),
+ (draft) => {
+ ;(draft.value as { title: string }).title = "remote"
+ }
+ )
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(remote)))
+ const generated = Automerge.generateSyncMessage(remote, Automerge.initSyncState())
+ assert.isNotNull(generated[1])
+ const message = generated[1]!
+ const messageHash = yield* Canonical.digest(message)
+ const writerProvenance = provenanceFor(message)
+ const permit = yield* gate.current
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ const documentBefore = yield* sql`SELECT * FROM effect_local_documents WHERE document_id = ${documentId}`
+ assert.lengthOf(documentBefore, 1)
+ const document = documentBefore[0]!
+ const expiredAt = new Date(0).toISOString()
+ const activeAt = new Date(limits.maxPendingAgeMillis + 1).toISOString()
+
+ yield* sql`INSERT INTO effect_local_changes (
+ change_hash, document_id, document_type, writer_schema_version, writer_definition_hash,
+ actor, sequence, dependencies, bytes, applied, peer_id, accepted_at, commit_sequence
+ ) VALUES (
+ ${"e".repeat(64)}, ${documentId}, ${Task.name}, ${Task.version}, ${definition.hash},
+ ${"f".repeat(32)}, 100, '[]', ${new Uint8Array([9])}, 0, ${peerId}, ${expiredAt},
+ ${metadata[0]!.commit_sequence}
+ )`
+ yield* sql`INSERT INTO effect_local_peer_receipts (
+ replica_incarnation, peer_id, connection_epoch, receive_sequence, document_id,
+ message_hash, reply, reply_hash, pending_message, heads,
+ accepted_heads, commit_sequence, accepted_at, writer_provenance
+ ) VALUES (
+ ${permit.incarnation}, ${peerId}, 'remote-epoch', 1, ${documentId},
+ 'expired-pending', NULL, NULL, ${new Uint8Array([8])}, ${document.materialized_heads},
+ ${document.accepted_heads}, ${metadata[0]!.commit_sequence}, ${expiredAt}, '[]'
+ )`
+ yield* sql`INSERT INTO effect_local_peer_receipts (
+ replica_incarnation, peer_id, connection_epoch, receive_sequence, document_id,
+ message_hash, reply, reply_hash, pending_message, heads,
+ accepted_heads, commit_sequence, accepted_at, writer_provenance
+ ) VALUES (
+ ${permit.incarnation}, ${peerId}, 'remote-epoch', 0, ${documentId},
+ ${messageHash}, NULL, NULL, NULL, ${document.materialized_heads},
+ ${document.accepted_heads}, ${metadata[0]!.commit_sequence}, ${activeAt},
+ ${Schema.encodeSync(WriterProvenance.StoredChangeProvenances)(writerProvenance)}
+ )`
+ const expiredChangeBefore = yield* sql`SELECT * FROM effect_local_changes
+ WHERE change_hash = ${"e".repeat(64)}`
+ const receiptsBefore = yield* sql`SELECT * FROM effect_local_peer_receipts
+ WHERE replica_incarnation = ${permit.incarnation}
+ AND peer_id = ${peerId}
+ AND connection_epoch = 'remote-epoch'
+ ORDER BY receive_sequence`
+ const outboxBefore = yield* sql`SELECT * FROM effect_local_commit_outbox ORDER BY commit_sequence`
+ const quarantineBefore = yield* sql`SELECT * FROM effect_local_quarantine ORDER BY id`
+ yield* TestClock.setTime(limits.maxPendingAgeMillis + 1)
+ const payload = {
+ replicaIncarnation: permit.incarnation,
+ peerId,
+ connectionEpoch: "remote-epoch",
+ localConnectionEpoch: "local-epoch",
+ receiveSequence: 0,
+ documentType: Task.name,
+ messageHash,
+ message,
+ writerProvenance
+ } as const
+ const primaryKey = keyOf(DocumentEntity.ApplySync.payloadSchema.make(payload))
+ const messageId = `EffectLocal/Document/${documentId}/ApplySync/${primaryKey}`
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+ yield* harness.arm
+
+ const requested = yield* entity(documentId).ApplySync(payload).pipe(
+ Effect.result,
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(requested).pipe(
+ Effect.andThen(Fiber.await(requested)),
+ Effect.asVoid
+ )
+ )
+ const transactionExit = yield* Deferred.await(harness.completed)
+ assert.isTrue(Exit.isFailure(transactionExit))
+ if (!Exit.isFailure(transactionExit)) return
+ assert.isTrue(Cause.hasFails(transactionExit.cause))
+ assert.isFalse(Cause.hasDies(transactionExit.cause))
+ const failure = Option.getOrThrow(Cause.findErrorOption(transactionExit.cause))
+ assert.strictEqual((failure as { readonly _tag?: string })._tag, "ReplicaError")
+ if ((failure as { readonly _tag?: string })._tag !== "ReplicaError") return
+ assert.strictEqual(
+ (failure as { readonly reason: { readonly _tag: string } }).reason._tag,
+ "ReplicaMetadataMissing"
+ )
+ yield* Deferred.await(harness.rolledBack)
+
+ const pending = yield* sql<{
+ readonly missing_reply: number
+ readonly processed: number
+ readonly replies: number
+ }>`SELECT processed, last_reply_id IS NULL AS missing_reply, (
+ SELECT COUNT(*) FROM effect_local_cluster_replies
+ WHERE request_id = effect_local_cluster_messages.id
+ ) AS replies
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ assert.lengthOf(pending, 1)
+ assert.strictEqual(pending[0]!.processed, 0)
+ assert.strictEqual(pending[0]!.missing_reply, 1)
+ assert.strictEqual(pending[0]!.replies, 0)
+
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`,
+ []
+ )
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_changes WHERE change_hash = ${"e".repeat(64)}`,
+ expiredChangeBefore
+ )
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_peer_receipts
+ WHERE replica_incarnation = ${permit.incarnation}
+ AND peer_id = ${peerId}
+ AND connection_epoch = 'remote-epoch'
+ ORDER BY receive_sequence`,
+ receiptsBefore
+ )
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_documents WHERE document_id = ${documentId}`,
+ documentBefore
+ )
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_commit_outbox ORDER BY commit_sequence`,
+ outboxBefore
+ )
+ assert.deepStrictEqual(
+ yield* sql`SELECT * FROM effect_local_quarantine ORDER BY id`,
+ quarantineBefore
+ )
+ const requestIds = yield* sql<{ readonly id: bigint }>`SELECT id
+ FROM effect_local_cluster_messages
+ WHERE message_id = ${messageId}`.pipe(
+ Effect.provideService(SqlClient.SafeIntegers, true)
+ )
+ assert.lengthOf(requestIds, 1)
+
+ yield* sql`INSERT INTO effect_local_metadata ${sql.insert(metadata[0]!)}`
+ yield* harness.disarm
+ assert.isTrue(yield* sharding.reset(Snowflake.Snowflake(requestIds[0]!.id)))
+ yield* sharding.pollStorage
+ yield* Effect.gen(function*() {
+ while (true) {
+ const rows = yield* sql<{
+ readonly processed: number
+ readonly replies: number
+ }>`SELECT processed, (
+ SELECT COUNT(*) FROM effect_local_cluster_replies
+ WHERE request_id = effect_local_cluster_messages.id
+ ) AS replies
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ if (rows[0]?.processed === 1 && rows[0]?.replies === 1) return
+ }
+ }).pipe(Effect.timeout("5 seconds"), TestClock.withLive)
+ const retried = yield* entity(documentId).ApplySync(payload).pipe(
+ Effect.timeout("5 seconds"),
+ TestClock.withLive
+ )
+ assert.isTrue(retried.duplicate)
+ }).pipe(Effect.scoped, Effect.provide(harness.services))
+ }))
+
it.effect("serves ApplySync without holding the connection across the gate", () =>
Effect.gen(function*() {
const atGate = yield* Deferred.make()
diff --git a/packages/local-sql/test/PeerSession.test.ts b/packages/local-sql/test/PeerSession.test.ts
index 98ffde5..7100f47 100644
--- a/packages/local-sql/test/PeerSession.test.ts
+++ b/packages/local-sql/test/PeerSession.test.ts
@@ -117,6 +117,7 @@ it.layer(NodeCrypto.layer)("PeerSession", (it) => {
admit: Effect.acquireRelease(Effect.succeed(permit), () => Effect.void),
claim: (use) => use(permit),
refresh: Effect.succeed(permit),
+ preflight: () => Effect.void,
validate: () => Effect.void
})
@@ -1566,6 +1567,7 @@ it.layer(NodeCrypto.layer)("PeerSession", (it) => {
),
claim: (use) => Ref.get(current).pipe(Effect.flatMap(use)),
refresh: Ref.get(current),
+ preflight: () => Effect.void,
validate: () => Effect.void
})
const sync = PeerSync.PeerSync.of({
diff --git a/packages/local-sql/test/PeerSessionCoverage.test.ts b/packages/local-sql/test/PeerSessionCoverage.test.ts
index 7592b64..82ae213 100644
--- a/packages/local-sql/test/PeerSessionCoverage.test.ts
+++ b/packages/local-sql/test/PeerSessionCoverage.test.ts
@@ -62,6 +62,7 @@ it.layer(NodeCrypto.layer)("PeerSession coverage", (it) => {
admit: Effect.acquireRelease(Effect.succeed(permit), () => Effect.void),
claim: (use) => use(permit),
refresh: Effect.succeed(permit),
+ preflight: () => Effect.void,
validate: () => Effect.void
})
const result = {
diff --git a/packages/local-sql/test/PeerSync.test.ts b/packages/local-sql/test/PeerSync.test.ts
index 7a7475a..2004f12 100644
--- a/packages/local-sql/test/PeerSync.test.ts
+++ b/packages/local-sql/test/PeerSync.test.ts
@@ -10,6 +10,7 @@ import * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition
import type * as ReplicaError from "@lucas-barake/effect-local/ReplicaError"
import * as ReplicaLimits from "@lucas-barake/effect-local/ReplicaLimits"
import * as Cause from "effect/Cause"
+import * as Context from "effect/Context"
import * as Deferred from "effect/Deferred"
import * as Effect from "effect/Effect"
import type * as Exit from "effect/Exit"
@@ -1705,6 +1706,99 @@ describe("PeerSync", () => {
InternalAutomerge.free(created.automerge)
}).pipe(Effect.provide(TestLayer)))
+ it.effect("fences a raced duplicate when metadata disappears before commit", () =>
+ Effect.scoped(
+ Effect.gen(function*() {
+ const sql = yield* SqlClient.SqlClient
+ const store = yield* DocumentStore.DocumentStore
+ const documentId = yield* Identity.makeDocumentId
+ const peerId = yield* Identity.makePeerId
+ const created = yield* store.create(Task, documentId, { title: "one", labels: [] })
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(created.automerge)))
+ const remote = Automerge.change(
+ Automerge.clone(created.automerge, { actor: "7".repeat(32) }),
+ (draft) => {
+ ;(draft.value as { title: string }).title = "two"
+ }
+ )
+ yield* Effect.addFinalizer(() => Effect.sync(() => InternalAutomerge.free(remote)))
+ const message = Automerge.generateSyncMessage(remote, Automerge.initSyncState())[1]!
+ const input = {
+ remoteConnectionEpoch: "raced-remote",
+ receiveSequence: 0,
+ lineage: Identity.genesisLineage,
+ message,
+ writerProvenance: provenanceFor(message, Task.version, definition.hash)
+ }
+ const loaded = yield* Deferred.make()
+ const release = yield* Deferred.make()
+ yield* Effect.addFinalizer(() => Deferred.succeed(release, undefined).pipe(Effect.asVoid))
+ const blockingStore = new Proxy(store, {
+ get(target, property, receiver) {
+ if (property !== "load") return Reflect.get(target, property, receiver)
+ const load: typeof store.load = (document, documentId) =>
+ store.load(document, documentId).pipe(
+ Effect.tap(() =>
+ Deferred.succeed(loaded, undefined).pipe(
+ Effect.andThen(Deferred.await(release))
+ )
+ )
+ )
+ return load
+ }
+ })
+ const loserContext = yield* Layer.build(
+ Layer.fresh(PeerSync.layer).pipe(
+ Layer.provide(Layer.succeed(DocumentStore.DocumentStore, blockingStore)),
+ Layer.provide(Services)
+ )
+ )
+ const winnerContext = yield* Layer.build(
+ Layer.fresh(PeerSync.layer).pipe(Layer.provide(Services))
+ )
+ const loser = Context.get(loserContext, PeerSync.PeerSync)
+ const winner = Context.get(winnerContext, PeerSync.PeerSync)
+ const loserSession = yield* loser.open(peerId)
+ const winnerSession = yield* winner.open(peerId)
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ yield* Effect.addFinalizer(() =>
+ sql`INSERT OR IGNORE INTO effect_local_metadata ${sql.insert(metadata[0]!)}`.pipe(Effect.orDie)
+ )
+
+ const losing = yield* loser.receive(Task, documentId, loserSession, input).pipe(
+ Effect.result,
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(losing).pipe(
+ Effect.andThen(Fiber.await(losing)),
+ Effect.asVoid
+ )
+ )
+ yield* Deferred.await(loaded)
+ const won = yield* winner.receive(Task, documentId, winnerSession, input)
+ assert.isFalse(won.duplicate)
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+ yield* Deferred.succeed(release, undefined)
+
+ const lost = yield* Fiber.join(losing)
+ assert.isTrue(Result.isFailure(lost))
+ if (!Result.isFailure(lost)) return
+ assert.strictEqual(lost.failure._tag, "ReplicaError")
+ if (lost.failure._tag === "ReplicaError") {
+ assert.strictEqual(lost.failure.reason._tag as string, "ReplicaMetadataMissing")
+ }
+ const receipts = yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_peer_receipts
+ WHERE replica_incarnation = ${loserSession.replicaIncarnation}
+ AND peer_id = ${peerId}
+ AND connection_epoch = ${input.remoteConnectionEpoch}
+ AND receive_sequence = ${input.receiveSequence}`
+ assert.strictEqual(receipts[0]?.count, 1)
+ }).pipe(Effect.provide(Services))
+ ))
+
it.effect("serializes one document across sessions without blocking independent documents", () =>
Effect.gen(function*() {
const store = yield* DocumentStore.DocumentStore
diff --git a/packages/local-sql/test/ReplicaBootstrap.test.ts b/packages/local-sql/test/ReplicaBootstrap.test.ts
index 1d32ba3..cb5aa22 100644
--- a/packages/local-sql/test/ReplicaBootstrap.test.ts
+++ b/packages/local-sql/test/ReplicaBootstrap.test.ts
@@ -4,12 +4,19 @@ import { assert, describe, it } from "@effect/vitest"
import * as Document from "@lucas-barake/effect-local/Document"
import * as DocumentSet from "@lucas-barake/effect-local/DocumentSet"
import * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition"
+import * as Context from "effect/Context"
+import * as Deferred from "effect/Deferred"
import * as Effect from "effect/Effect"
+import * as Fiber from "effect/Fiber"
+import * as Latch from "effect/Latch"
import * as Layer from "effect/Layer"
import * as Result from "effect/Result"
import * as Schema from "effect/Schema"
import * as Migrator from "effect/unstable/sql/Migrator"
import * as SqlClient from "effect/unstable/sql/SqlClient"
+import { mkdtempSync, rmSync } from "node:fs"
+import { tmpdir } from "node:os"
+import { join } from "node:path"
import * as Migrations from "../src/Migrations.js"
import * as ReplicaBootstrap from "../src/ReplicaBootstrap.js"
@@ -150,7 +157,7 @@ describe("ReplicaBootstrap", () => {
if (!Result.isFailure(result)) return
assert.strictEqual(result.failure._tag, "ReplicaError")
if (result.failure._tag !== "ReplicaError") return
- assert.strictEqual(result.failure.reason._tag, "StorageCorrupt")
+ assert.strictEqual(result.failure.reason._tag as string, "ReplicaMetadataMissing")
const generations = yield* sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM effect_local_writer_generations
`
@@ -159,6 +166,30 @@ describe("ReplicaBootstrap", () => {
Effect.provide(Layer.merge(SqliteClient.layer({ filename: ":memory:", disableWAL: true }), NodeCrypto.layer))
))
+ it.effect("rejects a surviving history rewrite marker when metadata is absent", () =>
+ Effect.gen(function*() {
+ const first = yield* ReplicaBootstrap.make(definition)
+ const sql = yield* SqlClient.SqlClient
+ yield* sql`INSERT INTO effect_local_history_rewrites (
+ replica_incarnation, operation_id, document_id, lineage, rewritten_at
+ ) VALUES (
+ ${first.incarnation}, 'rewrite-operation', 'document-1', 'lineage-1', '2026-01-01T00:00:00.000Z'
+ )`
+ const marker = yield* sql`SELECT * FROM effect_local_history_rewrites`
+ yield* sql`DELETE FROM effect_local_writer_generations`
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+
+ const result = yield* Effect.result(ReplicaBootstrap.make(definition))
+ assert.isTrue(Result.isFailure(result))
+ if (!Result.isFailure(result) || result.failure._tag !== "ReplicaError") return
+ assert.strictEqual(result.failure.reason._tag, "ReplicaMetadataMissing")
+ assert.deepStrictEqual(yield* sql`SELECT * FROM effect_local_history_rewrites`, marker)
+ assert.deepStrictEqual(yield* sql`SELECT * FROM effect_local_metadata`, [])
+ assert.deepStrictEqual(yield* sql`SELECT * FROM effect_local_writer_generations`, [])
+ }).pipe(
+ Effect.provide(Layer.merge(SqliteClient.layer({ filename: ":memory:", disableWAL: true }), NodeCrypto.layer))
+ ))
+
it.effect("does not migrate a populated replica whose metadata row is missing", () =>
Effect.gen(function*() {
const sql = yield* SqlClient.SqlClient
@@ -185,7 +216,7 @@ describe("ReplicaBootstrap", () => {
if (!Result.isFailure(result)) return
assert.strictEqual(result.failure._tag, "ReplicaError")
if (result.failure._tag !== "ReplicaError") return
- assert.strictEqual(result.failure.reason._tag, "StorageCorrupt")
+ assert.strictEqual(result.failure.reason._tag as string, "ReplicaMetadataMissing")
// a build that refuses to open the replica must not have migrated it on the way to refusing
const applied = yield* sql<{ readonly migration_id: number }>`
@@ -200,6 +231,189 @@ describe("ReplicaBootstrap", () => {
Effect.provide(Layer.merge(SqliteClient.layer({ filename: ":memory:", disableWAL: true }), NodeCrypto.layer))
))
+ it.effect("opens empty migration 1 storage without peer tables", () =>
+ Effect.gen(function*() {
+ yield* Migrator.make({})({
+ loader: Effect.map(Migrations.loader, (migrations) => migrations.slice(0, 1)),
+ table: "effect_local_migrations"
+ })
+ const state = yield* ReplicaBootstrap.make(definition)
+ assert.strictEqual(state.writerGeneration, 1)
+ }).pipe(
+ Effect.provide(Layer.merge(SqliteClient.layer({ filename: ":memory:", disableWAL: true }), NodeCrypto.layer))
+ ))
+
+ it.effect("rejects peer only migration 2 storage before migrating", () =>
+ Effect.forEach(["receipt", "outbox"] as const, (fixture) =>
+ Effect.gen(function*() {
+ const sql = yield* SqlClient.SqlClient
+ yield* sql`PRAGMA foreign_keys = OFF`
+ assert.strictEqual(
+ (yield* sql<{ readonly foreign_keys: number }>`PRAGMA foreign_keys`)[0]?.foreign_keys,
+ 0
+ )
+ yield* Migrator.make({})({
+ loader: Effect.map(Migrations.loader, (migrations) => migrations.slice(0, 2)),
+ table: "effect_local_migrations"
+ })
+ const canonical = (yield* sql<{ readonly count: number }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_writer_generations) +
+ (SELECT COUNT(*) FROM effect_local_documents) +
+ (SELECT COUNT(*) FROM effect_local_changes) +
+ (SELECT COUNT(*) FROM effect_local_checkpoints) +
+ (SELECT COUNT(*) FROM effect_local_command_receipts) +
+ (SELECT COUNT(*) FROM effect_local_projection_registry) +
+ (SELECT COUNT(*) FROM effect_local_document_projections) +
+ (SELECT COUNT(*) FROM effect_local_commit_outbox) +
+ (SELECT COUNT(*) FROM effect_local_quarantine) +
+ (SELECT COUNT(*) FROM effect_local_backup_installations) AS count`)[0]!.count
+ assert.strictEqual(canonical, 0)
+ if (fixture === "receipt") {
+ yield* sql`INSERT INTO effect_local_peer_receipts (
+ replica_incarnation, peer_id, connection_epoch, receive_sequence, document_id,
+ message_hash, reply, reply_hash, pending_message, heads, accepted_heads,
+ commit_sequence, accepted_at
+ ) VALUES (
+ 0, 'peer_1', 'connection_1', 1, 'doc_missing', 'message_1',
+ NULL, NULL, NULL, '[]', '[]', 1, '2026-01-01T00:00:00.000Z'
+ )`
+ } else {
+ yield* sql`INSERT INTO effect_local_peer_outbox (
+ replica_incarnation, peer_id, connection_epoch, document_id, send_sequence,
+ message, message_hash, heads, status
+ ) VALUES (
+ 0, 'peer_1', 'connection_1', 'doc_missing', 1,
+ x'01', 'message_1', '[]', 'Pending'
+ )`
+ }
+ const table = fixture === "receipt" ? "effect_local_peer_receipts" : "effect_local_peer_outbox"
+ const before = fixture === "receipt"
+ ? yield* sql>>`SELECT * FROM effect_local_peer_receipts`
+ : yield* sql>>`SELECT * FROM effect_local_peer_outbox`
+
+ const result = yield* Effect.result(ReplicaBootstrap.make(definition))
+ assert.isTrue(Result.isFailure(result))
+ if (!Result.isFailure(result) || result.failure._tag !== "ReplicaError") return
+ assert.strictEqual(result.failure.reason._tag as string, "ReplicaMetadataMissing")
+ const applied = yield* sql<{ readonly migration_id: number }>`
+ SELECT migration_id FROM effect_local_migrations ORDER BY migration_id
+ `
+ assert.deepStrictEqual(applied.map((row) => row.migration_id), [1, 2])
+ const after = table === "effect_local_peer_receipts"
+ ? yield* sql>>`SELECT * FROM effect_local_peer_receipts`
+ : yield* sql>>`SELECT * FROM effect_local_peer_outbox`
+ assert.deepStrictEqual(after, before)
+ const outboxColumns = yield* sql<{ readonly name: string }>`PRAGMA table_info(effect_local_peer_outbox)`
+ assert.notInclude(outboxColumns.map((column) => column.name), "created_at")
+ assert.deepStrictEqual(
+ (yield* sql<{ readonly migration_id: number }>`
+ SELECT migration_id FROM effect_local_migration_catalog ORDER BY migration_id
+ `).map((row) => row.migration_id),
+ [1, 2]
+ )
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_metadata) +
+ (SELECT COUNT(*) FROM effect_local_writer_generations) AS count`)[0]!.count,
+ 0
+ )
+ }).pipe(
+ Effect.provide(Layer.merge(SqliteClient.layer({ filename: ":memory:", disableWAL: true }), NodeCrypto.layer))
+ ), { discard: true }))
+
+ it.effect("serializes peer validation with migrations", () =>
+ Effect.scoped(Effect.gen(function*() {
+ const directory = yield* Effect.acquireRelease(
+ Effect.sync(() => mkdtempSync(join(tmpdir(), "effect-local-bootstrap-"))),
+ (path) => Effect.sync(() => rmSync(path, { force: true, recursive: true }))
+ )
+ const filename = join(directory, "replica.sqlite")
+ const atMigration = yield* Deferred.make()
+ const releaseMigration = yield* Latch.make()
+ let armed = false
+ const pausingClient = Layer.effect(
+ SqlClient.SqlClient,
+ Effect.map(SqlClient.SqlClient, (sql) =>
+ new Proxy(sql, {
+ apply(target, thisArg, args: Array) {
+ const statement = Reflect.apply(target as never, thisArg, args) as Effect.Effect<
+ unknown,
+ unknown,
+ never
+ >
+ const strings = args[0]
+ if (!armed || !Array.isArray(strings)) return statement
+ const text = (strings as ReadonlyArray).join("?").replace(/\s+/g, " ").trim()
+ if (
+ !text.includes("CREATE TABLE IF NOT EXISTS ?") ||
+ !text.includes("migration_id integer PRIMARY KEY") ||
+ !text.includes("created_at datetime")
+ ) {
+ return statement
+ }
+ armed = false
+ return Deferred.succeed(atMigration, undefined).pipe(
+ Effect.andThen(releaseMigration.await),
+ Effect.andThen(statement)
+ )
+ }
+ }) as typeof sql)
+ )
+ const databaseA = Layer.merge(
+ pausingClient.pipe(Layer.provide(SqliteClient.layer({ filename }))),
+ NodeCrypto.layer
+ )
+ const databaseB = Layer.merge(SqliteClient.layer({ filename }), NodeCrypto.layer)
+ const contextA = yield* Layer.build(databaseA)
+ const contextB = yield* Layer.build(databaseB)
+ const sqlA = Context.get(contextA, SqlClient.SqlClient)
+ const sqlB = Context.get(contextB, SqlClient.SqlClient)
+ yield* Migrator.make({})({
+ loader: Effect.map(Migrations.loader, (migrations) => migrations.slice(0, 2)),
+ table: "effect_local_migrations"
+ }).pipe(Effect.provide(contextA))
+ yield* sqlB`PRAGMA foreign_keys = OFF`
+ armed = true
+ const opening = yield* ReplicaBootstrap.make(definition).pipe(
+ Effect.provide(contextA),
+ Effect.result,
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Deferred.await(atMigration)
+ yield* sqlB`INSERT INTO effect_local_peer_receipts (
+ replica_incarnation, peer_id, connection_epoch, receive_sequence, document_id,
+ message_hash, reply, reply_hash, pending_message, heads, accepted_heads,
+ commit_sequence, accepted_at
+ ) VALUES (
+ 0, 'peer_1', 'connection_1', 1, 'doc_missing', 'message_1',
+ NULL, NULL, NULL, '[]', '[]', 1, '2026-01-01T00:00:00.000Z'
+ )`
+ yield* releaseMigration.open
+ const raced = yield* Fiber.join(opening).pipe(Effect.ensuring(Fiber.interrupt(opening)))
+ assert.isTrue(Result.isFailure(raced))
+ assert.deepStrictEqual(
+ (yield* sqlB<{ readonly migration_id: number }>`
+ SELECT migration_id FROM effect_local_migrations ORDER BY migration_id
+ `).map((row) => row.migration_id),
+ [1, 2]
+ )
+ assert.strictEqual(
+ (yield* sqlB<{ readonly count: number }>`
+ SELECT COUNT(*) AS count FROM effect_local_peer_receipts
+ `)[0]!.count,
+ 1
+ )
+ const retried = yield* Effect.result(ReplicaBootstrap.make(definition).pipe(Effect.provide(contextB)))
+ assert.isTrue(Result.isFailure(retried))
+ if (Result.isFailure(retried) && retried.failure._tag === "ReplicaError") {
+ assert.strictEqual(retried.failure.reason._tag as string, "ReplicaMetadataMissing")
+ }
+ assert.strictEqual(
+ (yield* sqlA<{ readonly count: number }>`SELECT COUNT(*) AS count FROM effect_local_metadata`)[0]!.count,
+ 0
+ )
+ })))
+
it.effect("rejects an incompatible replica definition without modifying metadata", () =>
Effect.gen(function*() {
const first = yield* ReplicaBootstrap.make(definition)
diff --git a/packages/local-sql/test/ReplicaGate.test.ts b/packages/local-sql/test/ReplicaGate.test.ts
index 49a8bed..317bc6d 100644
--- a/packages/local-sql/test/ReplicaGate.test.ts
+++ b/packages/local-sql/test/ReplicaGate.test.ts
@@ -14,6 +14,7 @@ import * as Option from "effect/Option"
import * as Ref from "effect/Ref"
import * as Result from "effect/Result"
import * as Schema from "effect/Schema"
+import { TestClock } from "effect/testing"
import * as SqlClient from "effect/unstable/sql/SqlClient"
import * as ReplicaBootstrap from "../src/ReplicaBootstrap.js"
import * as ReplicaGate from "../src/ReplicaGate.js"
@@ -669,4 +670,74 @@ describe("ReplicaGate", () => {
assert.isTrue(Result.isFailure(result))
if (Result.isFailure(result)) assert.strictEqual(result.failure.reason._tag, "StorageCorrupt")
}).pipe(Effect.provide(Gate)))
+
+ it.effect("releases claim admission after ReplicaMetadataMissing", () =>
+ Effect.gen(function*() {
+ const gate = yield* ReplicaGate.ReplicaGate
+ const sql = yield* SqlClient.SqlClient
+ const initial = yield* gate.current
+ const metadata = (yield* sql<{
+ readonly commit_sequence: number
+ readonly definition_hash: string
+ readonly replica_id: string
+ readonly replica_incarnation: number
+ readonly storage_format_version: number
+ readonly writer_generation: number
+ }>`SELECT
+ commit_sequence, definition_hash, replica_id, replica_incarnation,
+ storage_format_version, writer_generation
+ FROM effect_local_metadata WHERE singleton = 1`)[0]!
+ const generationsBefore = (yield* sql<{ readonly count: number }>`
+ SELECT COUNT(*) AS count FROM effect_local_writer_generations
+ `)[0]!.count
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+
+ const validation = yield* Effect.result(gate.validate(initial))
+ assert.isTrue(Result.isFailure(validation))
+ if (Result.isFailure(validation)) {
+ assert.strictEqual(validation.failure.reason._tag as string, "ReplicaMetadataMissing")
+ }
+ const failedClaim = yield* Effect.result(gate.claim(() => Effect.void))
+ assert.isTrue(Result.isFailure(failedClaim))
+ if (Result.isFailure(failedClaim) && failedClaim.failure._tag === "ReplicaError") {
+ assert.strictEqual(failedClaim.failure.reason._tag as string, "ReplicaMetadataMissing")
+ }
+ assert.isFalse(yield* gate.claiming)
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`
+ SELECT COUNT(*) AS count FROM effect_local_writer_generations
+ `)[0]!.count,
+ generationsBefore
+ )
+
+ yield* sql`INSERT INTO effect_local_metadata (
+ singleton, storage_format_version, replica_id, replica_incarnation,
+ writer_generation, definition_hash, commit_sequence
+ ) VALUES (
+ 1, ${metadata.storage_format_version}, ${metadata.replica_id}, ${metadata.replica_incarnation},
+ ${metadata.writer_generation}, ${metadata.definition_hash}, ${metadata.commit_sequence}
+ )`
+ const callback = yield* Deferred.make()
+ const claim = yield* gate.claim((permit) => Deferred.succeed(callback, permit).pipe(Effect.as(permit))).pipe(
+ Effect.result,
+ Effect.timeoutOption("1 second"),
+ Effect.forkChild({ startImmediately: true })
+ )
+ const claimed = yield* Effect.gen(function*() {
+ yield* TestClock.adjust("1 second")
+ return yield* Fiber.join(claim)
+ }).pipe(Effect.ensuring(Fiber.interrupt(claim)))
+ assert.isTrue(Option.isSome(claimed))
+ if (Option.isSome(claimed)) {
+ assert.isTrue(Result.isSuccess(claimed.value))
+ }
+ assert.isTrue(Option.isSome(yield* Deferred.poll(callback)))
+ assert.isFalse(yield* gate.claiming)
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`
+ SELECT COUNT(*) AS count FROM effect_local_writer_generations
+ `)[0]!.count,
+ generationsBefore + 1
+ )
+ }).pipe(Effect.provide(Gate)))
})
diff --git a/packages/local-sql/test/SqlReplica.test.ts b/packages/local-sql/test/SqlReplica.test.ts
index 761bff4..8ae7443 100644
--- a/packages/local-sql/test/SqlReplica.test.ts
+++ b/packages/local-sql/test/SqlReplica.test.ts
@@ -1,11 +1,13 @@
import { NodeCrypto } from "@effect/platform-node"
import { SqliteClient } from "@effect/sql-sqlite-node"
import { assert, describe, it } from "@effect/vitest"
+import * as Canonical from "@lucas-barake/effect-local/Canonical"
import * as CommandOutcome from "@lucas-barake/effect-local/CommandOutcome"
import * as Document from "@lucas-barake/effect-local/Document"
import * as DocumentSet from "@lucas-barake/effect-local/DocumentSet"
import * as Identity from "@lucas-barake/effect-local/Identity"
import * as Mutation from "@lucas-barake/effect-local/Mutation"
+import * as PeerTransport from "@lucas-barake/effect-local/PeerTransport"
import * as Projection from "@lucas-barake/effect-local/Projection"
import * as Replica from "@lucas-barake/effect-local/Replica"
import * as ReplicaDefinition from "@lucas-barake/effect-local/ReplicaDefinition"
@@ -18,6 +20,7 @@ import * as Latch from "effect/Latch"
import * as Layer from "effect/Layer"
import * as Option from "effect/Option"
import * as Queue from "effect/Queue"
+import * as Result from "effect/Result"
import * as Schema from "effect/Schema"
import * as Stream from "effect/Stream"
import { TestClock } from "effect/testing"
@@ -28,6 +31,7 @@ import * as CommandExecutor from "../src/CommandExecutor.js"
import * as CommitPublisher from "../src/CommitPublisher.js"
import * as DocumentStore from "../src/DocumentStore.js"
import * as ClusterStorage from "../src/internal/clusterStorage.js"
+import * as PeerSession from "../src/PeerSession.js"
import * as ProjectionStore from "../src/ProjectionStore.js"
import * as QueryExecutor from "../src/QueryExecutor.js"
import * as Recovery from "../src/Recovery.js"
@@ -76,6 +80,15 @@ describe("SqlReplica", () => {
projections: [],
queries: []
})
+ type MetadataRow = {
+ readonly commit_sequence: Identity.CommitSequence
+ readonly definition_hash: string
+ readonly replica_id: Identity.ReplicaId
+ readonly replica_incarnation: Identity.ReplicaIncarnation
+ readonly singleton: number
+ readonly storage_format_version: number
+ readonly writer_generation: Identity.WriterGeneration
+ }
const limits: ReplicaLimits.Values = {
maxBackupBytes: 1024 * 1024,
maxChunkBytes: 64 * 1024,
@@ -256,6 +269,286 @@ describe("SqlReplica", () => {
assert.deepStrictEqual(rows[0], { changes: 6, clusterMessages: 8, documents: 4, receipts: 7 })
}).pipe(Effect.provide(Live), Effect.provide(Database), TestClock.withLive))
+ it.effect("keeps missing metadata Mutate messages retryable", () =>
+ Effect.gen(function*() {
+ const replica = yield* Replica.Replica
+ const sql = yield* SqlClient.SqlClient
+ const created = yield* replica.create(Task, {
+ commandId: yield* Identity.makeCommandId,
+ value: { title: "before" }
+ })
+ assert.strictEqual(created._tag, "DurablyCommittedLocal")
+ if (created._tag !== "DurablyCommittedLocal") return
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ const before = yield* sql<{
+ readonly changes: number
+ readonly commitSequence: number
+ }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_changes) AS changes,
+ (SELECT commit_sequence FROM effect_local_metadata WHERE singleton = 1) AS commitSequence`
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+ yield* Effect.addFinalizer(() =>
+ sql`INSERT OR IGNORE INTO effect_local_metadata ${sql.insert(metadata[0]!)}`.pipe(Effect.orDie)
+ )
+
+ const commandId = yield* Identity.makeCommandId
+ const requestHash = yield* CommandExecutor.mutationRequestHash({
+ incarnation: metadata[0]!.replica_incarnation,
+ commandId,
+ documentId: created.value,
+ mutation: Rename,
+ payload: "after"
+ })
+ const messageId = `EffectLocal/Document/${created.value}/Mutate/${
+ metadata[0]!.replica_incarnation
+ }:${commandId}:${requestHash}`
+ const first = yield* replica.mutate(Rename, {
+ commandId,
+ documentId: created.value,
+ payload: "after"
+ }).pipe(
+ Effect.result,
+ Effect.timeoutOption("1 second"),
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(first).pipe(
+ Effect.andThen(Fiber.await(first)),
+ Effect.asVoid
+ )
+ )
+ let dispatched = false
+ yield* Effect.gen(function*() {
+ while (!first.pollUnsafe()) {
+ const rows = yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ if (rows[0]?.count === 1) {
+ dispatched = true
+ return
+ }
+ yield* sql`SELECT 1`
+ }
+ }).pipe(Effect.timeout("2 seconds"), TestClock.withLive)
+ if (dispatched) yield* TestClock.adjust("1 second")
+ const timed = yield* Fiber.join(first)
+ assert.isTrue(Option.isSome(timed))
+ if (!Option.isSome(timed)) return
+ assert.isTrue(Result.isFailure(timed.value))
+ if (!Result.isFailure(timed.value)) return
+ assert.strictEqual(timed.value.failure._tag, "ReplicaError")
+ if (timed.value.failure._tag !== "ReplicaError") return
+ assert.strictEqual(timed.value.failure.reason._tag, "ReplicaMetadataMissing")
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`)[0]?.count,
+ 0
+ )
+
+ yield* sql`INSERT INTO effect_local_metadata ${sql.insert(metadata[0]!)}`
+ const retried = yield* replica.mutate(Rename, {
+ commandId,
+ documentId: created.value,
+ payload: "after"
+ }).pipe(Effect.timeout("5 seconds"), TestClock.withLive)
+ assert.deepStrictEqual(retried, CommandOutcome.durablyCommitted(commandId, undefined))
+ assert.strictEqual((yield* replica.get(Task, created.value)).value.title, "after")
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_command_receipts
+ WHERE replica_incarnation = ${metadata[0]!.replica_incarnation}
+ AND command_id = ${commandId}`)[0]?.count,
+ 1
+ )
+ const after = yield* sql<{
+ readonly changes: number
+ readonly commitSequence: number
+ }>`SELECT
+ (SELECT COUNT(*) FROM effect_local_changes) AS changes,
+ (SELECT commit_sequence FROM effect_local_metadata WHERE singleton = 1) AS commitSequence`
+ assert.deepStrictEqual(after[0], {
+ changes: before[0]!.changes + 1,
+ commitSequence: before[0]!.commitSequence + 1
+ })
+ const newest = yield* sql<{
+ readonly commitSequence: number
+ readonly documentId: Identity.DocumentId
+ }>`SELECT commit_sequence AS commitSequence, document_id AS documentId
+ FROM effect_local_commit_outbox ORDER BY commit_sequence DESC LIMIT 1`
+ assert.deepStrictEqual(newest[0], {
+ commitSequence: after[0]!.commitSequence,
+ documentId: created.value
+ })
+ }).pipe(Effect.scoped, Effect.provide(Live), Effect.provide(Database)))
+
+ it.effect("does not adopt a foreign writer during command preflight", () =>
+ Effect.gen(function*() {
+ const gate = yield* ReplicaGate.ReplicaGate
+ const replica = yield* Replica.Replica
+ const sql = yield* SqlClient.SqlClient
+ const created = yield* replica.create(Task, {
+ commandId: yield* Identity.makeCommandId,
+ value: { title: "before" }
+ })
+ assert.strictEqual(created._tag, "DurablyCommittedLocal")
+ if (created._tag !== "DurablyCommittedLocal") return
+ const local = yield* gate.current
+ yield* sql`UPDATE effect_local_metadata SET writer_generation = writer_generation + 1 WHERE singleton = 1`
+ yield* Effect.addFinalizer(() =>
+ sql`UPDATE effect_local_metadata SET writer_generation = ${local.writerGeneration}
+ WHERE singleton = 1`.pipe(Effect.orDie)
+ )
+
+ const attempt = (commandId: Identity.CommandId, payload: string) =>
+ Effect.gen(function*() {
+ const requestHash = yield* CommandExecutor.mutationRequestHash({
+ incarnation: local.incarnation,
+ commandId,
+ documentId: created.value,
+ mutation: Rename,
+ payload
+ })
+ const messageId =
+ `EffectLocal/Document/${created.value}/Mutate/${local.incarnation}:${commandId}:${requestHash}`
+ const fiber = yield* replica.mutate(Rename, {
+ commandId,
+ documentId: created.value,
+ payload
+ }).pipe(
+ Effect.result,
+ Effect.timeoutOption("1 second"),
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(fiber).pipe(
+ Effect.andThen(Fiber.await(fiber)),
+ Effect.asVoid
+ )
+ )
+ let dispatched = false
+ yield* Effect.gen(function*() {
+ while (!fiber.pollUnsafe()) {
+ const rows = yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`
+ if (rows[0]?.count === 1) {
+ dispatched = true
+ return
+ }
+ yield* sql`SELECT 1`
+ }
+ }).pipe(Effect.timeout("2 seconds"), TestClock.withLive)
+ if (dispatched) yield* TestClock.adjust("1 second")
+ return [yield* Fiber.join(fiber), messageId] as const
+ })
+
+ const [first, firstMessageId] = yield* attempt(yield* Identity.makeCommandId, "first")
+ const [second, secondMessageId] = yield* attempt(yield* Identity.makeCommandId, "second")
+ for (const result of [first, second]) {
+ assert.isTrue(Option.isSome(result))
+ if (!Option.isSome(result)) return
+ assert.isTrue(Result.isFailure(result.value))
+ if (!Result.isFailure(result.value)) return
+ assert.strictEqual(result.value.failure._tag, "ReplicaError")
+ if (result.value.failure._tag !== "ReplicaError") return
+ assert.strictEqual(result.value.failure.reason._tag, "ReplicaFenced")
+ }
+ assert.deepStrictEqual(yield* gate.current, local)
+ for (const messageId of [firstMessageId, secondMessageId]) {
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_messages WHERE message_id = ${messageId}`)[0]?.count,
+ 0
+ )
+ }
+ }).pipe(Effect.scoped, Effect.provide(Live), Effect.provide(Database)))
+
+ it.effect("rejects missing replica metadata before dispatching inbound ApplySync", () =>
+ Effect.gen(function*() {
+ const replica = yield* Replica.Replica
+ const sql = yield* SqlClient.SqlClient
+ const created = yield* replica.create(Task, {
+ commandId: yield* Identity.makeCommandId,
+ value: { title: "before" }
+ })
+ assert.strictEqual(created._tag, "DurablyCommittedLocal")
+ if (created._tag !== "DurablyCommittedLocal") return
+ const peerId = yield* Identity.makePeerId
+ const inbound = yield* Queue.unbounded()
+ yield* Effect.addFinalizer(() => Queue.shutdown(inbound))
+ const transport = PeerTransport.PeerTransport.of({
+ capabilities: { storeAndForward: false },
+ connect: () =>
+ Effect.succeed({
+ peerId,
+ capabilities: { storeAndForward: false },
+ receive: Stream.fromQueue(inbound),
+ send: () => Effect.void,
+ close: Effect.void
+ })
+ })
+ const session = yield* PeerSession.makeSupervised({
+ peerId,
+ documents: [{ document: Task, documentId: created.value }]
+ }).pipe(
+ Effect.provideService(PeerTransport.PeerTransport, transport),
+ Effect.provideService(ReplicaLimits.ReplicaLimits, limits)
+ )
+ const disconnected = yield* Effect.result(session.awaitDisconnect).pipe(
+ Effect.forkChild({ startImmediately: true })
+ )
+ yield* Effect.addFinalizer(() =>
+ Fiber.interrupt(disconnected).pipe(
+ Effect.andThen(Fiber.await(disconnected)),
+ Effect.asVoid
+ )
+ )
+ const metadata = yield* sql`SELECT * FROM effect_local_metadata WHERE singleton = 1`
+ assert.lengthOf(metadata, 1)
+ const before = yield* sql<{ readonly sequence: number }>`SELECT
+ COALESCE(MAX(commit_sequence), 0) AS sequence FROM effect_local_commit_outbox`
+ yield* sql`DELETE FROM effect_local_metadata WHERE singleton = 1`
+ yield* Effect.addFinalizer(() =>
+ sql`INSERT OR IGNORE INTO effect_local_metadata ${sql.insert(metadata[0]!)}`.pipe(Effect.orDie)
+ )
+ const message = Uint8Array.of(1)
+ const messageHash = yield* Canonical.digest(message)
+ const encoded = yield* Schema.encodeEffect(
+ Schema.fromJsonString(Schema.toCodecJson(PeerSession.SyncEnvelope))
+ )({
+ connectionEpoch: "remote-epoch",
+ sequence: 0,
+ documentId: created.value,
+ documentType: Task.name,
+ messageHash,
+ message,
+ lineage: Identity.genesisLineage,
+ writerProvenance: []
+ }).pipe(Effect.map((value) => new TextEncoder().encode(value)))
+ yield* Queue.offer(inbound, encoded)
+
+ const timed = yield* Fiber.join(disconnected).pipe(
+ Effect.timeoutOption("2 seconds"),
+ TestClock.withLive
+ )
+ assert.isTrue(Option.isSome(timed))
+ if (!Option.isSome(timed)) return
+ assert.isTrue(Result.isFailure(timed.value))
+ if (!Result.isFailure(timed.value)) return
+ assert.strictEqual(timed.value.failure._tag, "ReplicaError")
+ assert.strictEqual(timed.value.failure.reason._tag, "ReplicaMetadataMissing")
+ assert.strictEqual(
+ (yield* sql<{ readonly count: number }>`SELECT COUNT(*) AS count
+ FROM effect_local_cluster_messages
+ WHERE message_id LIKE ${`EffectLocal/Document/${created.value}/ApplySync/%`}`)[0]?.count,
+ 0
+ )
+ assert.deepStrictEqual(
+ yield* sql<{ readonly sequence: number }>`SELECT
+ COALESCE(MAX(commit_sequence), 0) AS sequence FROM effect_local_commit_outbox`,
+ before
+ )
+ }).pipe(Effect.scoped, Effect.provide(Live), Effect.provide(Database)))
+
it.effect("provides nonempty projection bindings", () =>
Effect.gen(function*() {
const replica = yield* Replica.Replica
diff --git a/packages/local/src/ReplicaError.ts b/packages/local/src/ReplicaError.ts
index 79be5be..c36709d 100644
--- a/packages/local/src/ReplicaError.ts
+++ b/packages/local/src/ReplicaError.ts
@@ -62,6 +62,10 @@ export class StorageCorrupt extends Schema.TaggedErrorClass(
"@lucas-barake/effect-local/ReplicaError/StorageCorrupt"
)("StorageCorrupt", { cause: Schema.Defect() }) {}
+export class ReplicaMetadataMissing extends Schema.TaggedErrorClass(
+ "@lucas-barake/effect-local/ReplicaError/ReplicaMetadataMissing"
+)("ReplicaMetadataMissing", {}) {}
+
export class QuotaExceeded extends Schema.TaggedErrorClass(
"@lucas-barake/effect-local/ReplicaError/QuotaExceeded"
)("QuotaExceeded", {
@@ -149,6 +153,7 @@ export const Reason = Schema.Union([
StorageUnavailable,
CanonicalEncodeError,
StorageCorrupt,
+ ReplicaMetadataMissing,
QuotaExceeded,
MigrationFailed,
BackupInvalid,
diff --git a/packages/local/test/ReplicaError.test.ts b/packages/local/test/ReplicaError.test.ts
index 139af79..ac2c89b 100644
--- a/packages/local/test/ReplicaError.test.ts
+++ b/packages/local/test/ReplicaError.test.ts
@@ -14,6 +14,15 @@ describe("ReplicaError", () => {
assert.deepStrictEqual(Schema.decodeUnknownSync(ReplicaError.ReplicaError)(encoded), error)
})
+ it("round trips ReplicaMetadataMissing", () => {
+ const error = new ReplicaError.ReplicaError({
+ reason: new ReplicaError.ReplicaMetadataMissing()
+ })
+ const encoded = Schema.encodeSync(ReplicaError.ReplicaError)(error)
+ assert.deepStrictEqual(encoded.reason, { _tag: "ReplicaMetadataMissing" })
+ assert.deepStrictEqual(Schema.decodeUnknownSync(ReplicaError.ReplicaError)(encoded), error)
+ })
+
it("round trips arbitrary defect causes", () => {
const error = new ReplicaError.ReplicaError({
reason: new ReplicaError.StorageUnavailable({