Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 20 additions & 12 deletions apps/cli/src/server/checkpoints.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,9 @@ import {
} from "./durable-files"
import {
eventingControlSnapshotPath,
LocalEventingControlStore,
openControlStore,
restoreControlSnapshot,
validateControlSnapshot,
type EventingControlSnapshotValidation,
} from "./eventing/control-store"
import { CURRENT_LOCAL_SCHEMA, SCHEMA_FINGERPRINT } from "./schema-identity"
Expand Down Expand Up @@ -885,7 +887,7 @@ const resolveCheckpointById = async (
const controlSha256 = await sha256File(controlPath)
if (controlSha256 !== manifest.controlSha256)
throw new Error("checkpoint control-store digest mismatch")
const controlValidation = LocalEventingControlStore.validateSnapshot(controlPath)
const controlValidation = await Effect.runPromise(validateControlSnapshot(controlPath))
if (!controlValidationMatches(manifest.controlValidation, controlValidation))
throw new Error("checkpoint control-store validation does not match its manifest")
} else if (existsSync(controlPath)) {
Expand Down Expand Up @@ -933,13 +935,12 @@ const restoreResolvedInto = async (
"SETTINGS allow_different_database_def=1",
)
if (resolvedCheckpoint.manifest.formatVersion === MANIFEST_FORMAT_VERSION) {
await LocalEventingControlStore.restoreSnapshot(
join(resolvedCheckpoint.snapshotDir, "control.sqlite"),
targetDataDir,
await Effect.runPromise(
restoreControlSnapshot(join(resolvedCheckpoint.snapshotDir, "control.sqlite"), targetDataDir),
)
} else {
const controlStore = await LocalEventingControlStore.open(targetDataDir)
controlStore.close()
// A legacy checkpoint has no control snapshot: open and close to create an empty store.
await Effect.runPromise(Effect.scoped(Effect.asVoid(openControlStore(targetDataDir))))
}
return { db, validation: validateRestoredDatabase(db) }
} catch (error) {
Expand Down Expand Up @@ -1687,15 +1688,22 @@ const createCheckpointTraced = Effect.fn("CheckpointService.create")(function* (
: createError(error),
),
)
return yield* Effect.tryPromise({
const controlPath = eventingControlSnapshotPath(options.dataDir, checkpointId)
yield* Effect.tryPromise({
try: async () => {
const { oldState, snapshot, startedAt } = prepared
let { operation } = prepared
await syncTree(snapshotBackupDir(options.dataDir, checkpointId))
const controlPath = eventingControlSnapshotPath(options.dataDir, checkpointId)
await assertNoSymlink(checkpointSnapshotsRoot(options.dataDir), controlPath)
await assertRealFile(controlPath, "checkpoint eventing control snapshot")
const controlValidation = LocalEventingControlStore.validateSnapshot(controlPath)
},
catch: createError,
})
const controlValidation = yield* validateControlSnapshot(controlPath).pipe(
Effect.mapError(createError),
)
return yield* Effect.tryPromise({
try: async () => {
const { oldState, snapshot, startedAt } = prepared
let { operation } = prepared
operation = { ...operation, phase: "backup-complete" }
await writeOperation(options.dataDir, operation, options.faults)
const provisionalManifest: CheckpointManifest = {
Expand Down
17 changes: 15 additions & 2 deletions apps/cli/src/server/eventing/consumer-auth.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,20 @@
import { resolve } from "node:path"
import { Effect, Schema } from "effect"
import { ensureLocalToken, localTokenMatches } from "../local-token"

export class EventConsumerTokenError extends Schema.TaggedError<EventConsumerTokenError>()(
"@maple/cli/eventing/EventConsumerTokenFailed",
{ message: Schema.String, cause: Schema.Defect() },
) {}

export const eventConsumerTokenPath = (dataDir: string): string => `${resolve(dataDir)}.event-consumer-token`
export const ensureEventConsumerToken = (dataDir: string): Promise<string> =>
ensureLocalToken(eventConsumerTokenPath(dataDir), "event consumer token")
export const ensureEventConsumerToken = (dataDir: string): Effect.Effect<string, EventConsumerTokenError> =>
Effect.tryPromise({
try: () => ensureLocalToken(eventConsumerTokenPath(dataDir), "event consumer token"),
catch: (cause) =>
new EventConsumerTokenError({
message: cause instanceof Error ? cause.message : String(cause),
cause,
}),
})
export const eventConsumerTokenMatches = localTokenMatches
Loading
Loading