diff --git a/.changeset/remote-cli-commands.md b/.changeset/remote-cli-commands.md new file mode 100644 index 0000000..79bfc41 --- /dev/null +++ b/.changeset/remote-cli-commands.md @@ -0,0 +1,10 @@ +--- +"migrate-sdk": patch +--- + +Support every server-backed CLI command through `--server`, including registry +listing, dependency graphs, status, messages, and lock recovery. Remote HTTP +authentication failures now report 401 and 403 permission errors with a token +configuration hint. Finalize the pre-adoption Migrate Protocol v1 baseline with +static registry metadata and canonical selection-based status and message +reports; servers from the earlier incomplete v1 draft must be redeployed. diff --git a/README.md b/README.md index bdf4822..81081c6 100644 --- a/README.md +++ b/README.md @@ -88,6 +88,18 @@ pnpm exec migrate rollback articles --plan pnpm exec migrate rollback articles ``` +The same server-backed commands can target a remote Migrate Server without a +local migration configuration: + +```sh +MIGRATE_SERVER_TOKEN=secret pnpm exec migrate \ + --server https://migrate.example.com/api/migrate \ + list +``` + +Remote mode supports `list`, `graph`, `status`, `messages`, `unlock`, `run`, +`rollback`, and `runs`. Migration Store schema administration remains local. + ## Full and incremental source discovery Sources default to full discovery. Cursors let interrupted runs resume, but a diff --git a/docs/adr/0007-server-boundary-for-local-and-remote-clients.md b/docs/adr/0007-server-boundary-for-local-and-remote-clients.md index 17d634d..012fd31 100644 --- a/docs/adr/0007-server-boundary-for-local-and-remote-clients.md +++ b/docs/adr/0007-server-boundary-for-local-and-remote-clients.md @@ -79,6 +79,13 @@ the Effect RPC encoding incompatibly requires a new protocol version. An additive wire change may remain in the current protocol only when cross-version contract tests demonstrate compatibility in both client/server directions. +Protocol v1 is finalized before external adoption by including the complete CLI +inspection surface described below. Earlier `0.8.0` development deployments +that advertised the incomplete v1 draft must be redeployed with the corrected +server before using the new remote CLI inspection commands. This is a one-time +baseline correction, not precedent for adding required operations to an adopted +protocol version. Existing TUI operations and payloads remain unchanged. + The currently implemented local connection uses a Node Migrate Server process started by the TUI. The Bun renderer communicates with that process over a reconnectable local Effect RPC socket, while the Node process loads the same @@ -265,20 +272,35 @@ infrastructure-only commands may remain outside that interface. The CLI adopts that seam incrementally through a shared Node Migrate Connection owned by Migrate SDK rather than by the TUI package. Both clients use the same local socket or remote HTTPS connection, version checks, authentication header, -and logical `MigrateClient`; renderer concerns remain in the TUI. The first CLI -surface over this connection is run lifecycle management: - +and logical `MigrateClient`; renderer concerns remain in the TUI. CLI surfaces +over this connection include registry inspection, migration control, and run +lifecycle management: + +- `migrate list` lists registered Migration Definitions. +- `migrate graph` renders Migration Definition dependencies. +- `migrate status` and `migrate messages` inspect durable server state. +- `migrate unlock` requests a lock break through the server. +- `migrate run` and `migrate rollback` plan and start operations. - `migrate runs list` discovers active Migration Runs. - `migrate runs observe ` observes one run until it is terminal or the caller detaches. - `migrate runs stop ` explicitly requests durable cooperative cancellation. -`migrate list` continues to mean registered Migration Definitions. Remote CLI -commands accept a Migrate Server URL and read its bearer token from -`MIGRATE_SERVER_TOKEN`; secrets are not accepted as command-line flags. Local -commands discover the local migration configuration and connect to the same -reconnectable Node Migrate Server used by the TUI. +`GetRegistry` returns static registry entries and groups without reading +Migration Stores. `GetRegistryStatus` and `GetRegistryMessages` accept the same +selection and dependency-expansion inputs as their local registry operations; +the server performs selection validation, source-scan concurrency, dependency +deduplication, and global message ordering authoritatively. Clients render +those canonical reports rather than reconstructing them from dashboard rows or +per-definition requests. + +Remote CLI commands accept a Migrate Server URL and read its bearer token from +`MIGRATE_SERVER_TOKEN`; secrets are not accepted as command-line flags. Store +schema commands remain local infrastructure operations because the Migrate +Protocol does not expose Migration Store administration. Local commands +discover the local migration configuration and may continue using direct SDK +services while the shared server seam is adopted incrementally. While observing in an interactive terminal, the first Ctrl+C offers three distinct choices: detach the observation, stop safely and keep observing while @@ -315,8 +337,10 @@ unknown, active-run discovery is the recovery path. CLI-specific execution progress services, or wait on in-process run handles. - Execution providers remain replaceable behind Migration Executable instead of becoming TUI integrations. -- The TUI and server negotiate protocol and SDK compatibility and exchange - registry and environment identity before operations begin. +- Clients and servers negotiate Migrate Protocol compatibility and exchange + registry, environment, and SDK identity before operations begin. Remote SDK + versions are diagnostic; local socket connections additionally require exact + SDK identity. - Contextual availability, such as rollback support or whether a run can be stopped, is represented by migration and run data rather than server feature negotiation. diff --git a/packages/migrate-sdk/src/cli/command.ts b/packages/migrate-sdk/src/cli/command.ts index 0f06d67..975dd98 100644 --- a/packages/migrate-sdk/src/cli/command.ts +++ b/packages/migrate-sdk/src/cli/command.ts @@ -11,22 +11,19 @@ import { } from "effect"; import { Argument, CliError, Command, Flag } from "effect/unstable/cli"; import type { SqlClient } from "effect/unstable/sql"; -import type { AnySelfContainedMigrationDefinition } from "../domain/definition.ts"; import type { PipelineExecutionConcurrency } from "../domain/execution.ts"; import { - type MigrationDefinitionId, toMigrationDefinitionGroupId, toMigrationDefinitionId, toMigrationRunId, } from "../domain/ids.ts"; import { MigrationMessage } from "../domain/message.ts"; -import type { - MigrationDefinitionRegistry, - MigrationDefinitionRegistryMessagesError, - MigrationDefinitionRegistryPlanningError, - MigrationDefinitionRegistrySelectionInput, - MigrationDefinitionRegistryStatusError, - MigrationDefinitionRegistryStatusInput, +import { + MigrationDefinitionRegistryInvalidSelectionError, + MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError, + type MigrationDefinitionRegistryPlanningError, + MigrationDefinitionRegistryUnknownDefinitionError, + MigrationDefinitionRegistryUnknownGroupError, } from "../domain/registry.ts"; import type { MigrateAction, @@ -34,7 +31,6 @@ import type { MigratePreparedOperation, MigrateSelection, } from "../protocol/index.ts"; -import { MigrationStore } from "../services/migration-store.ts"; import { SqlMigrationStore } from "../stores/sql/sql-migration-store.ts"; import type { SqlMigrationStoreSchemaPlan } from "../stores/sql/sql-migration-store-schema.ts"; import type { MigrationCliConfig } from "./config.ts"; @@ -42,6 +38,11 @@ import { loadMigrationCliConfig, type MigrationCliConfigLoadError, } from "./config-loader.ts"; +import { + type MigrationCliRegistryOperations, + makeLocalMigrationCliRegistryOperations, + makeRemoteMigrationCliRegistryOperations, +} from "./inspection.ts"; import type { ActiveMigrationCliInterrupts } from "./interrupts.ts"; import { type CliObservationProgressMode, @@ -57,8 +58,8 @@ import { renderPreparedOperationDependencyFailure, renderPreparedOperationPlan, renderPreparedOperationWarnings, - renderRegistryGraph, - renderRegistryList, + renderRegistryEntriesGraph, + renderRegistryEntriesList, renderRunStopResult, renderRuntimeError, renderSqlMigrationStoreSchemaPlan, @@ -169,7 +170,7 @@ const loadConfiguredConfig = Effect.gen(function* () { if (serverUrl !== undefined) { return yield* failReportedCliMessage( - "--server is currently supported by migrate run, rollback, and runs commands" + "--server is not supported by local Migration Store schema commands" ); } @@ -185,11 +186,6 @@ const loadConfiguredConfig = Effect.gen(function* () { return loadedConfig; }); -const loadConfiguredRegistry = Effect.map( - loadConfiguredConfig, - (loadedConfig) => loadedConfig.registry -); - interface CliMigrateConnection { readonly connection: MigrationCliServerConnection; readonly listActiveRuns: string; @@ -327,6 +323,37 @@ const withCliMigrateConnection = ( releaseCliMigrateConnection ); +const reportCliRegistryCommandErrors = ( + effect: Effect.Effect +): Effect.Effect => + effect.pipe( + Effect.catch((error) => + CliError.isCliError(error) + ? Effect.fail(error) + : failReportedCliMessage(renderStoredFailure(error)) + ) + ); + +const withCliRegistryOperations = ( + use: ( + operations: MigrationCliRegistryOperations + ) => Effect.Effect +) => + Effect.gen(function* () { + const root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* withCliMigrateConnection(({ connection }) => + use(makeRemoteMigrationCliRegistryOperations(connection)) + ); + } + + const loadedConfig = yield* loadConfiguredConfig; + return yield* use( + makeLocalMigrationCliRegistryOperations(loadedConfig.registry) + ); + }).pipe(reportCliRegistryCommandErrors); + const requireConfiguredSqlStore = (config: MigrationCliConfig) => config.sqlStore === undefined ? failReportedCliMessage( @@ -484,27 +511,17 @@ const storeCommand = Command.make("store").pipe( Command.withSubcommands([schemaCommand]) ); -type CliExecutableRegistry = MigrationDefinitionRegistry< - readonly AnySelfContainedMigrationDefinition[] ->; - -const toCliExecutableRegistry = ( - registry: MigrationDefinitionRegistry -): CliExecutableRegistry => registry as CliExecutableRegistry; - -const hasRegisteredDefinition = ( - registry: MigrationDefinitionRegistry, - definitionId: MigrationDefinitionId -): boolean => registry.list().some((entry) => entry.id === definitionId); - const listCommand = Command.make("list", {}, () => - Effect.gen(function* () { - const registry = yield* loadConfiguredRegistry; - - yield* Console.log( - renderRegistryList(registry, { colors: yield* useColor }) - ); - }) + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const registry = yield* operations.getRegistry; + yield* Console.log( + renderRegistryEntriesList(registry.entries, { + colors: yield* useColor, + }) + ); + }) + ) ).pipe(Command.withDescription("List registered Migration Definitions")); const graphDefinition = Argument.string("definition").pipe(Argument.optional); @@ -513,34 +530,31 @@ const graphCommand = Command.make( "graph", { definition: graphDefinition }, ({ definition }) => - Effect.gen(function* () { - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - const focusedDefinitionId = Option.getOrUndefined(definition); - - if (focusedDefinitionId !== undefined) { - const definitionId = toMigrationDefinitionId(focusedDefinitionId); + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const focusedDefinitionId = Option.getOrUndefined(definition); + const registry = yield* operations.getRegistry; - if (!hasRegisteredDefinition(registry, definitionId)) { + if ( + focusedDefinitionId !== undefined && + !registry.entries.some((entry) => entry.id === focusedDefinitionId) + ) { return yield* failReportedCliMessage( - `Migration Definition was not found in the registry: ${definitionId}` + `Migration Definition was not found in the registry: ${focusedDefinitionId}` ); } yield* Console.log( - renderRegistryGraph(registry, definitionId, { - colors: yield* useColor, - }) + renderRegistryEntriesGraph( + registry.entries, + focusedDefinitionId === undefined + ? undefined + : toMigrationDefinitionId(focusedDefinitionId), + { colors: yield* useColor } + ) ); - return; - } - - yield* Console.log( - renderRegistryGraph(registry, undefined, { - colors: yield* useColor, - }) - ); - }) + }) + ) ).pipe(Command.withDescription("Inspect Migration Definition dependencies")); const plan = Flag.boolean("plan").pipe( @@ -679,40 +693,12 @@ interface CliRegistrySelectionInput { readonly withDependencies: boolean; } -const makeRegistrySelectionInput = ( - input: CliRegistrySelectionInput -): MigrationDefinitionRegistrySelectionInput => { - if (input.group !== undefined) { - return { - group: input.group, - ...(input.all ? { all: true } : {}), - ...(input.definitionIds.length === 0 - ? {} - : { definitionIds: input.definitionIds }), - withDependencies: input.withDependencies, - } as unknown as MigrationDefinitionRegistrySelectionInput; - } - - if (input.all) { - return { - all: true, - ...(input.definitionIds.length === 0 - ? {} - : { definitionIds: input.definitionIds }), - withDependencies: input.withDependencies, - } as MigrationDefinitionRegistrySelectionInput; - } - - return input.definitionIds.length === 0 - ? ({} as MigrationDefinitionRegistrySelectionInput) - : { - definitionIds: input.definitionIds as [string, ...string[]], - withDependencies: input.withDependencies, - }; -}; - const makeMigrateSelection = ( - input: Pick + input: Pick, + messages: { + readonly missing?: string; + readonly mixed?: string; + } = {} ): Effect.Effect => { const selectedKinds = Number(input.all) + @@ -721,13 +707,15 @@ const makeMigrateSelection = ( if (selectedKinds === 0) { return failReportedCliMessage( - "Select migrations with definition IDs, --group, or --all" + messages.missing ?? + "Select migrations with definition IDs, --group, or --all" ); } if (selectedKinds > 1) { return failReportedCliMessage( - "Choose only one migration selection: definition IDs, --group, or --all" + messages.mixed ?? + "Choose only one migration selection: definition IDs, --group, or --all" ); } @@ -746,7 +734,8 @@ const makeMigrateSelection = ( if (firstDefinitionId === undefined) { return failReportedCliMessage( - "Select migrations with definition IDs, --group, or --all" + messages.missing ?? + "Select migrations with definition IDs, --group, or --all" ); } @@ -759,40 +748,25 @@ const makeMigrateSelection = ( }); }; -const makeStatusInput = ( - input: CliRegistrySelectionInput & { - readonly concurrency?: number; - readonly scanSource: boolean; - } -): MigrationDefinitionRegistryStatusInput => - ({ - ...makeRegistrySelectionInput(input), - ...(input.concurrency === undefined - ? {} - : { concurrency: input.concurrency }), - scanSource: input.scanSource, - }) as MigrationDefinitionRegistryStatusInput; +const registryReadSelectionMessages = { + missing: + "Registry planning requires all: true, a Migration Definition group, or at least one Migration Definition id", + mixed: + "Registry planning accepts only one selection form: all: true, Migration Definition ids, or a Migration Definition group", +} as const; const isPlanningError = ( error: unknown -): error is MigrationDefinitionRegistryPlanningError => { - if (typeof error !== "object" || error === null || !("_tag" in error)) { - return false; - } - - switch (error._tag) { - case "MigrationDefinitionRegistryInvalidSelectionError": - case "MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError": - case "MigrationDefinitionRegistryUnknownDefinitionError": - case "MigrationDefinitionRegistryUnknownGroupError": - return true; - default: - return false; - } -}; +): error is MigrationDefinitionRegistryPlanningError => + Schema.is(MigrationDefinitionRegistryInvalidSelectionError)(error) || + Schema.is( + MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError + )(error) || + Schema.is(MigrationDefinitionRegistryUnknownDefinitionError)(error) || + Schema.is(MigrationDefinitionRegistryUnknownGroupError)(error); const renderStatusCommandError = ( - error: MigrationDefinitionRegistryStatusError, + error: unknown, input: { readonly definitionIds: readonly string[]; readonly group?: string; @@ -807,15 +781,11 @@ const renderStatusCommandError = ( }); } - if (error._tag === "MigrationStatusRequestError") { - return error.message; - } - - return renderRuntimeError(error); + return renderStoredFailure(error); }; const renderMessagesCommandError = ( - error: MigrationDefinitionRegistryMessagesError, + error: unknown, input: { readonly definitionIds: readonly string[]; readonly group?: string; @@ -828,7 +798,7 @@ const renderMessagesCommandError = ( ...(input.group === undefined ? {} : { group: input.group }), hasTarget: false, }) - : renderRuntimeError(error); + : renderStoredFailure(error); const statusCommand = Command.make( "status", @@ -841,36 +811,46 @@ const statusCommand = Command.make( withDependencies, }, (input) => - Effect.gen(function* () { - const loadedConfig = yield* loadConfiguredConfig; - const registry = toCliExecutableRegistry(loadedConfig.registry); - const concurrencyInput = Option.getOrUndefined(input.concurrency); - const groupInput = Option.getOrUndefined(input.group); - const statusInput = makeStatusInput({ - all: input.all, - ...(concurrencyInput === undefined - ? {} - : { concurrency: concurrencyInput }), - definitionIds: input.definitions, - ...(groupInput === undefined ? {} : { group: groupInput }), - scanSource: input.scanSource, - withDependencies: input.withDependencies, - }); - const report = yield* registry.status(statusInput).pipe( - Effect.catch((error) => - failReportedCliMessage( - renderStatusCommandError(error, { - definitionIds: input.definitions, - ...(groupInput === undefined ? {} : { group: groupInput }), - }) - ) - ) - ); + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const concurrencyInput = Option.getOrUndefined(input.concurrency); + const groupInput = Option.getOrUndefined(input.group); + const selection = yield* makeMigrateSelection( + { + all: input.all, + definitionIds: input.definitions, + ...(groupInput === undefined ? {} : { group: groupInput }), + }, + registryReadSelectionMessages + ); + const report = yield* operations + .getStatus({ + ...(concurrencyInput === undefined + ? {} + : { concurrency: concurrencyInput }), + scanSource: input.scanSource, + selection, + withDependencies: input.withDependencies, + }) + .pipe( + Effect.catchTag("MigrationStatusRequestError", (error) => + failReportedCliMessage(error.message) + ), + Effect.catch((error) => + failReportedCliMessage( + renderStatusCommandError(error, { + definitionIds: input.definitions, + ...(groupInput === undefined ? {} : { group: groupInput }), + }) + ) + ) + ); - yield* Console.log( - renderStatusReport(report, { colors: yield* useColor }) - ); - }) + yield* Console.log( + renderStatusReport(report, { colors: yield* useColor }) + ); + }) + ) ).pipe(Command.withDescription("Inspect Migration Definition status")); const messagesCommand = Command.make( @@ -883,42 +863,46 @@ const messagesCommand = Command.make( withDependencies, }, (input) => - Effect.gen(function* () { - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - const groupInput = Option.getOrUndefined(input.group); - const report = yield* registry - .messages( - makeRegistrySelectionInput({ + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const groupInput = Option.getOrUndefined(input.group); + const selection = yield* makeMigrateSelection( + { all: input.all, definitionIds: input.definitions, ...(groupInput === undefined ? {} : { group: groupInput }), + }, + registryReadSelectionMessages + ); + const report = yield* operations + .getMessages({ + selection, withDependencies: input.withDependencies, }) - ) - .pipe( - Effect.catch((error) => - failReportedCliMessage( - renderMessagesCommandError(error, { - definitionIds: input.definitions, - ...(groupInput === undefined ? {} : { group: groupInput }), - }) + .pipe( + Effect.catch((error) => + failReportedCliMessage( + renderMessagesCommandError(error, { + definitionIds: input.definitions, + ...(groupInput === undefined ? {} : { group: groupInput }), + }) + ) ) - ) - ); + ); - if (input.json) { - const json = yield* Schema.encodeEffect(MigrationMessagesFromJson)( - report.messages - ).pipe(Effect.orDie); - yield* Console.log(json); - return; - } + if (input.json) { + const json = yield* Schema.encodeEffect(MigrationMessagesFromJson)( + report.messages + ).pipe(Effect.orDie); + yield* Console.log(json); + return; + } - yield* Console.log( - renderMessagesReport(report, { colors: yield* useColor }) - ); - }) + yield* Console.log( + renderMessagesReport(report, { colors: yield* useColor }) + ); + }) + ) ).pipe(Command.withDescription("Inspect durable Migration Messages")); const unlockDefinition = Argument.string("definition").pipe( @@ -931,47 +915,34 @@ const unlockCommand = Command.make( "unlock", { definition: unlockDefinition }, ({ definition }) => - Effect.gen(function* () { - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - const definitionId = toMigrationDefinitionId(definition); - const migrationDefinition = Option.getOrUndefined( - registry.get(definitionId) - ); - - if (migrationDefinition === undefined) { - return yield* failReportedCliMessage( - `Migration Definition was not found in the registry: ${definitionId}` - ); - } + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const definitionId = toMigrationDefinitionId(definition); + const result = yield* operations.breakLock(definitionId); - const lock = yield* Effect.gen(function* () { - const store = yield* MigrationStore; + if (result.kind === "not-found") { + return yield* failReportedCliMessage( + `Migration Definition was not found in the registry: ${definitionId}` + ); + } - return yield* store.breakDefinitionLock(definitionId); - }).pipe( - Effect.provide(migrationDefinition.store), - Effect.catch((error) => - failReportedCliMessage(renderRuntimeError(error)) - ) - ); + if (result.kind === "already-clear") { + yield* Console.log( + `Migration Definition lock is already clear: ${definitionId}` + ); + return; + } - if (lock === null) { yield* Console.log( - `Migration Definition lock is already clear: ${definitionId}` + [ + "Migration Definition lock cleared", + `Migration ID ${definitionId}`, + `Owner Run ID ${result.lock.ownerRunId}`, + `Token ${result.lock.token}`, + ].join("\n") ); - return; - } - - yield* Console.log( - [ - "Migration Definition lock cleared", - `Migration ID ${definitionId}`, - `Owner Run ID ${lock.ownerRunId}`, - `Token ${lock.token}`, - ].join("\n") - ); - }) + }) + ) ).pipe(Command.withDescription("Break a Migration Definition lock")); const runIdArgument = Argument.string("run-id").pipe( diff --git a/packages/migrate-sdk/src/cli/inspection.ts b/packages/migrate-sdk/src/cli/inspection.ts new file mode 100644 index 0000000..a77b666 --- /dev/null +++ b/packages/migrate-sdk/src/cli/inspection.ts @@ -0,0 +1,111 @@ +import { Effect, Option } from "effect"; +import type { AnySelfContainedMigrationDefinition } from "../domain/definition.ts"; +import type { MigrationDefinitionId } from "../domain/ids.ts"; +import type { MigrationDefinitionLock } from "../domain/lock.ts"; +import type { MigrationDefinitionRegistry } from "../domain/registry.ts"; +import type { + MigrateRegistry, + MigrateRegistryMessagesReport, + MigrateRegistryMessagesRequest, + MigrateRegistryStatusReport, + MigrateRegistryStatusRequest, +} from "../protocol/index.ts"; +import { toMigrationDefinitionRegistrySelectionInput } from "../protocol/registry-selection.ts"; +import { MigrationStore } from "../services/migration-store.ts"; +import type { + MigrationCliServerConnection, + MigrationCliServerError, +} from "./runtime.ts"; + +export type MigrationCliUnlockResult = + | { readonly kind: "already-clear" } + | { readonly kind: "cleared"; readonly lock: MigrationDefinitionLock } + | { readonly kind: "not-found" }; + +export interface MigrationCliRegistryOperations { + readonly breakLock: ( + definitionId: MigrationDefinitionId + ) => Effect.Effect; + readonly getMessages: ( + input: MigrateRegistryMessagesRequest + ) => Effect.Effect; + readonly getRegistry: Effect.Effect; + readonly getStatus: ( + input: MigrateRegistryStatusRequest + ) => Effect.Effect; +} + +type CliExecutableRegistry = MigrationDefinitionRegistry< + readonly AnySelfContainedMigrationDefinition[] +>; + +export const makeLocalMigrationCliRegistryOperations = ( + registry: MigrationDefinitionRegistry +): MigrationCliRegistryOperations => { + const executableRegistry = registry as CliExecutableRegistry; + + return { + breakLock: (definitionId) => + Effect.gen(function* () { + const definition = Option.getOrUndefined(registry.get(definitionId)); + + if (definition === undefined) { + return { kind: "not-found" as const }; + } + + const lock = yield* MigrationStore.pipe( + Effect.flatMap((store) => store.breakDefinitionLock(definitionId)), + Effect.provide(definition.store) + ); + + return lock === null + ? { kind: "already-clear" as const } + : { kind: "cleared" as const, lock }; + }), + getMessages: (input) => + registry.messages(toMigrationDefinitionRegistrySelectionInput(input)), + getRegistry: Effect.succeed({ + entries: registry.list(), + groups: registry.groups(), + }), + getStatus: (input) => + executableRegistry.status({ + ...toMigrationDefinitionRegistrySelectionInput(input), + ...(input.concurrency === undefined + ? {} + : { concurrency: input.concurrency }), + scanSource: input.scanSource, + }), + }; +}; + +export const makeRemoteMigrationCliRegistryOperations = ( + connection: MigrationCliServerConnection +): MigrationCliRegistryOperations => ({ + breakLock: (definitionId) => + Effect.gen(function* () { + const snapshot = yield* connection.getDashboard; + const row = snapshot.dashboard.rows.find( + (candidate) => candidate.entry.id === definitionId + ); + + if (row === undefined) { + return { kind: "not-found" as const }; + } + + const lock = row.status?.lock; + + if (lock === undefined || lock === null) { + return { kind: "already-clear" as const }; + } + + const result = yield* connection.breakLock(lock); + + return result.kind === "already-clear" + ? { kind: "already-clear" as const } + : { kind: "cleared" as const, lock }; + }), + getMessages: connection.getRegistryMessages, + getRegistry: connection.getRegistry, + getStatus: connection.getRegistryStatus, +}); diff --git a/packages/migrate-sdk/src/cli/remote-commands.test.ts b/packages/migrate-sdk/src/cli/remote-commands.test.ts new file mode 100644 index 0000000..5ad0488 --- /dev/null +++ b/packages/migrate-sdk/src/cli/remote-commands.test.ts @@ -0,0 +1,336 @@ +import { layer as nodeServicesLayer } from "@effect/platform-node/NodeServices"; +import { describe, expect, it } from "@effect/vitest"; +import { Effect, Layer, Ref, Stdio, Stream } from "effect"; +import { pretty as prettyCause } from "effect/Cause"; +import { isFailure, isSuccess } from "effect/Exit"; +import { TestConsole } from "effect/testing"; +import { CliOutput, Command } from "effect/unstable/cli"; +import { + toEncodedSourceIdentity, + toMigrationDefinitionId, + toMigrationDefinitionLockToken, + toMigrationRunId, +} from "../domain/ids.ts"; +import { MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError } from "../domain/registry.ts"; +import { MigrationStatusRequestError } from "../domain/status.ts"; +import { + MigrateDashboardResumeToken, + type MigrateRegistryStatusRequest, +} from "../protocol/index.ts"; +import { migrateCommand } from "./command.ts"; +import { + MigrationCliRuntime, + type MigrationCliRuntimeShape, + type MigrationCliServerConnection, +} from "./runtime.ts"; + +const authorsId = toMigrationDefinitionId("authors"); +const articlesId = toMigrationDefinitionId("articles"); +const runId = toMigrationRunId("run-remote"); +const lock = { + createdAt: new Date("2026-08-29T11:00:00.000Z"), + definitionId: articlesId, + ownerRunId: runId, + token: toMigrationDefinitionLockToken("lock-remote"), +}; +const authorsEntry = { + dependencies: { optional: [], required: [] }, + hasRollback: true, + id: authorsId, +}; +const articlesEntry = { + dependencies: { optional: [], required: [authorsId] }, + hasRollback: true, + id: articlesId, +}; +const registry = { + entries: [authorsEntry, articlesEntry], + groups: [], +}; +const articlesStatus = { + definitionId: articlesId, + discovery: "full" as const, + durable: { failed: 0, migrated: 3, needsUpdate: 0, skipped: 1 }, + lastRun: null, + lock, + source: { + duplicate: 0, + invalid: 0, + orphaned: 0, + total: 5, + unprocessed: 2, + }, + warnings: [], +}; + +const makeConnection = ( + overrides: Partial = {} +): MigrationCliServerConnection => ({ + breakLock: () => Effect.die("Unexpected lock break"), + dispose: () => Promise.resolve(), + getActiveRuns: Effect.succeed([]), + getDashboard: Effect.succeed({ + dashboard: { + activeRuns: [], + groups: [], + rows: [], + scannedSource: false, + }, + resumeToken: MigrateDashboardResumeToken.make("dashboard-empty"), + }), + getRegistry: Effect.succeed(registry), + getRegistryMessages: () => + Effect.succeed({ + includedDefinitionIds: [], + messages: [], + notices: [], + requestedDefinitionIds: [], + }), + getRegistryStatus: () => + Effect.succeed({ + definitions: [], + includedDefinitionIds: [], + notices: [], + requestedDefinitionIds: [], + scanSource: false, + warnings: [], + }), + observeRun: () => Stream.die("Unexpected run observation"), + prepareOperation: () => Effect.die("Unexpected operation preparation"), + startOperation: () => Effect.die("Unexpected operation start"), + stopRun: () => Effect.die("Unexpected run stop"), + ...overrides, +}); + +const makeLayer = (runtime: MigrationCliRuntimeShape) => + Layer.mergeAll( + CliOutput.layer(CliOutput.defaultFormatter({ colors: false })), + Layer.succeed(MigrationCliRuntime, runtime), + nodeServicesLayer, + Stdio.layerTest({}), + TestConsole.layer + ); + +const runCli = ( + args: readonly string[], + connection: MigrationCliServerConnection +) => + Effect.gen(function* () { + const exit = yield* Effect.exit( + Command.runWith(migrateCommand, { version: "0.0.0" })(args) + ); + + return { + cause: isFailure(exit) ? prettyCause(exit.cause) : "", + exitCode: isSuccess(exit) ? 0 : 1, + stderr: (yield* TestConsole.errorLines).map(String).join("\n"), + stdout: (yield* TestConsole.logLines).map(String).join("\n"), + }; + }).pipe( + Effect.provide( + makeLayer({ + connectMigrateServer: () => Effect.succeed(connection), + cwd: "/workspace", + }) + ) + ); + +const withServer = (args: readonly string[]) => [ + ...args, + "--server", + "https://migrate.example/api/migrate", +]; + +describe("remote CLI inspection commands", () => { + it.effect( + "reads list and graph metadata without loading dashboard state", + () => + Effect.gen(function* () { + const connection = makeConnection({ + getDashboard: Effect.die("Dashboard must not be loaded"), + }); + const list = yield* runCli(withServer(["list"]), connection); + const graph = yield* runCli( + withServer(["graph", "articles"]), + connection + ); + + expect(list.exitCode).toBe(0); + expect(list.stdout).toContain("Migration Definitions"); + expect(list.stdout).toContain("articles"); + expect(graph.exitCode).toBe(0); + expect(graph.stdout).toContain("Migration Dependency Graph: articles"); + expect(graph.stdout).toContain("articles(required) --> authors"); + }) + ); + + it.effect("sends the complete status request to the server", () => + Effect.gen(function* () { + const request = yield* Ref.make( + undefined + ); + const connection = makeConnection({ + getRegistryStatus: (input) => + Ref.set(request, input).pipe( + Effect.as({ + definitions: [articlesStatus], + includedDefinitionIds: [authorsId, articlesId], + notices: [], + requestedDefinitionIds: [articlesId], + scanSource: true, + warnings: [], + }) + ), + }); + const result = yield* runCli( + withServer([ + "status", + "articles", + "--with-dependencies", + "--scan-source", + "--concurrency", + "2", + ]), + connection + ); + + expect(result.exitCode).toBe(0); + expect(result.stdout).toContain("source inventory"); + expect(yield* Ref.get(request)).toEqual({ + concurrency: 2, + scanSource: true, + selection: { definitionIds: [articlesId], kind: "definitions" }, + withDependencies: true, + }); + }) + ); + + it.effect("renders canonical server validation errors", () => + Effect.gen(function* () { + const missingDependency = yield* runCli( + withServer(["status", "articles"]), + makeConnection({ + getRegistryStatus: () => + Effect.fail( + new MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError( + { + definitionId: articlesId, + message: + "Migration Definition selection is missing required dependencies", + missingDependencyIds: [authorsId], + } + ) + ), + }) + ); + const invalidConcurrency = yield* runCli( + withServer(["status", "--all", "--concurrency", "2"]), + makeConnection({ + getRegistryStatus: () => + Effect.fail( + new MigrationStatusRequestError({ + message: + "Status concurrency is only valid when source scanning is enabled", + }) + ), + }) + ); + + expect(missingDependency.exitCode).toBe(1); + expect(missingDependency.stderr).toContain( + "articles is missing required dependencies: authors" + ); + expect(missingDependency.stderr).toContain( + "migrate status --with-dependencies articles" + ); + expect(invalidConcurrency.exitCode).toBe(1); + expect(invalidConcurrency.stderr).toContain( + "Status concurrency is only valid when source scanning is enabled" + ); + expect(invalidConcurrency.stderr).not.toContain( + "MigrationStatusRequestError" + ); + }) + ); + + it.effect("renders server messages and clears a remote lock", () => + Effect.gen(function* () { + const connection = makeConnection({ + breakLock: () => + Effect.succeed({ definitionId: articlesId, kind: "cleared" }), + getDashboard: Effect.succeed({ + dashboard: { + activeRuns: [], + groups: [], + rows: [{ entry: articlesEntry, status: articlesStatus }], + scannedSource: false, + }, + resumeToken: MigrateDashboardResumeToken.make("dashboard-remote"), + }), + getRegistryMessages: () => + Effect.succeed({ + includedDefinitionIds: [articlesId], + messages: [ + { + definitionId: articlesId, + kind: "skip-reason", + message: "Already migrated remotely", + runId, + severity: "info", + sourceIdentity: toEncodedSourceIdentity("article-1"), + updatedAt: new Date("2026-08-29T11:30:00.000Z"), + }, + ], + notices: [], + requestedDefinitionIds: [articlesId], + }), + }); + const messages = yield* runCli( + withServer(["messages", "articles"]), + connection + ); + const unlock = yield* runCli( + withServer(["unlock", "articles"]), + connection + ); + + expect(messages.exitCode).toBe(0); + expect(messages.stdout).toContain("Already migrated remotely"); + expect(unlock.exitCode).toBe(0); + expect(unlock.stdout).toContain("Migration Definition lock cleared"); + expect(unlock.stdout).toContain("lock-remote"); + }) + ); + + it.effect("accepts empty all-selection reports", () => + Effect.gen(function* () { + const connection = makeConnection({ + getRegistry: Effect.succeed({ entries: [], groups: [] }), + getRegistryMessages: () => + Effect.succeed({ + includedDefinitionIds: [], + messages: [], + notices: [], + requestedDefinitionIds: "all" as const, + }), + getRegistryStatus: () => + Effect.succeed({ + definitions: [], + includedDefinitionIds: [], + notices: [], + requestedDefinitionIds: "all" as const, + scanSource: false, + warnings: [], + }), + }); + const status = yield* runCli(withServer(["status", "--all"]), connection); + const messages = yield* runCli( + withServer(["messages", "--all"]), + connection + ); + + expect(status.exitCode).toBe(0); + expect(messages.exitCode).toBe(0); + }) + ); +}); diff --git a/packages/migrate-sdk/src/cli/render.ts b/packages/migrate-sdk/src/cli/render.ts index 7d0aa7f..7f8df72 100644 --- a/packages/migrate-sdk/src/cli/render.ts +++ b/packages/migrate-sdk/src/cli/render.ts @@ -6,7 +6,6 @@ import type { MigrationDefinitionPlanNotice, MigrationDefinitionRegistry, MigrationDefinitionRegistryConstructionIssue, - MigrationDefinitionRegistryEntry, MigrationDefinitionRegistryMessagesReport, MigrationDefinitionRegistryPlanningError, MigrationDefinitionRegistrySelectionReport, @@ -33,6 +32,16 @@ interface MigrationDefinitionGraphEdge { readonly unresolved: boolean; } +interface RegistryEntryProjection { + readonly dependencies: { + readonly optional: readonly MigrationDefinitionId[]; + readonly required: readonly MigrationDefinitionId[]; + }; + readonly group?: MigrationDefinitionGroupId | undefined; + readonly hasRollback: boolean; + readonly id: MigrationDefinitionId; +} + const ansi = { bold: "\x1b[1m", cyan: "\x1b[36m", @@ -272,11 +281,10 @@ const formatOptionalDependencies = ( ) .join(", "); -export const renderRegistryList = ( - registry: MigrationDefinitionRegistry, +export const renderRegistryEntriesList = ( + entries: readonly RegistryEntryProjection[], options: RenderOptions = {} ): string => { - const entries = registry.list(); const registeredIds = new Set(entries.map((entry) => entry.id)); if (entries.length === 0) { @@ -332,8 +340,13 @@ export const renderRegistryList = ( ].join("\n"); }; +export const renderRegistryList = ( + registry: MigrationDefinitionRegistry, + options: RenderOptions = {} +): string => renderRegistryEntriesList(registry.list(), options); + const collectGraphEdges = ( - entries: readonly MigrationDefinitionRegistryEntry[] + entries: readonly RegistryEntryProjection[] ): readonly MigrationDefinitionGraphEdge[] => { const registeredIds = new Set(entries.map((entry) => entry.id)); @@ -374,12 +387,11 @@ const renderGraphEdge = ( return `${edge.fromDefinitionId}(${styledLabel}) --> ${edge.toDefinitionId}`; }; -export const renderRegistryGraph = ( - registry: MigrationDefinitionRegistry, +export const renderRegistryEntriesGraph = ( + entries: readonly RegistryEntryProjection[], focusedDefinitionId?: MigrationDefinitionId, options: RenderOptions = {} ): string => { - const entries = registry.list(); const edges = collectGraphEdges(entries).filter( (edge) => focusedDefinitionId === undefined || @@ -401,6 +413,13 @@ export const renderRegistryGraph = ( ].join("\n"); }; +export const renderRegistryGraph = ( + registry: MigrationDefinitionRegistry, + focusedDefinitionId?: MigrationDefinitionId, + options: RenderOptions = {} +): string => + renderRegistryEntriesGraph(registry.list(), focusedDefinitionId, options); + const renderRequestedDefinitionIdsInline = ( requestedDefinitionIds: MigratePlanProjection["requestedDefinitionIds"] ): string => diff --git a/packages/migrate-sdk/src/cli/runs-command.test.ts b/packages/migrate-sdk/src/cli/runs-command.test.ts index 288ba5a..dd54793 100644 --- a/packages/migrate-sdk/src/cli/runs-command.test.ts +++ b/packages/migrate-sdk/src/cli/runs-command.test.ts @@ -14,6 +14,10 @@ import { pretty as prettyCause } from "effect/Cause"; import { isFailure, isSuccess } from "effect/Exit"; import { TestClock, TestConsole } from "effect/testing"; import { CliOutput, Command } from "effect/unstable/cli"; +import { + RpcClientDefect, + RpcClientError, +} from "effect/unstable/rpc/RpcClientError"; import type { MigrateServerConnectionInput } from "../client/node/index.ts"; import { toMigrationDefinitionGroupId, @@ -31,7 +35,10 @@ import type { MigratePreparedOperation, MigrateTerminalSummary, } from "../protocol/index.ts"; -import { MigratePlanFingerprint } from "../protocol/index.ts"; +import { + MigrateDashboardResumeToken, + MigratePlanFingerprint, +} from "../protocol/index.ts"; import { migrateCommand } from "./command.ts"; import type { ActiveMigrationCliInterrupts } from "./interrupts.ts"; import { @@ -43,6 +50,10 @@ import { const definitionId = toMigrationDefinitionId("articles"); const runId = toMigrationRunId("run-cli-observation"); +const rpcFailure = (message: string) => + new RpcClientError({ + reason: new RpcClientDefect({ cause: message, message }), + }); const activeRun: MigrateActiveRun = { definitionIds: [definitionId], execution: { adapter: "workflow-sdk", executionId: "workflow-1" }, @@ -143,11 +154,24 @@ const makeConnection = ( overrides: Partial = {}, onDispose: () => void = () => undefined ): MigrationCliServerConnection => ({ + breakLock: () => Effect.die("Unexpected lock break"), dispose: () => { onDispose(); return Promise.resolve(); }, getActiveRuns: Effect.succeed([]), + getDashboard: Effect.succeed({ + dashboard: { + activeRuns: [], + groups: [], + rows: [], + scannedSource: false, + }, + resumeToken: MigrateDashboardResumeToken.make("dashboard-empty"), + }), + getRegistry: Effect.succeed({ entries: [], groups: [] }), + getRegistryMessages: () => Effect.die("Unexpected registry messages read"), + getRegistryStatus: () => Effect.die("Unexpected registry status read"), observeRun: () => Stream.die("Unexpected run observation"), prepareOperation: () => Effect.die("Unexpected operation preparation"), startOperation: () => Effect.die("Unexpected operation start"), @@ -406,7 +430,7 @@ describe("migrate runs", () => { const configPath = "/workspace/migrations $HOME/$(touch-nope)/operator's.config.ts"; const connection = makeConnection({ - observeRun: () => Stream.fail(new Error("connection dropped")), + observeRun: () => Stream.fail(rpcFailure("connection dropped")), }); const result = yield* runCli( ["runs", "observe", runId, "--config", configPath], @@ -610,7 +634,7 @@ describe("migrate runs", () => { const unsafeServerUrl = "https://migrate.example/api/$(touch-nope)/$HOME/`id`/operator's"; const connection = makeConnection({ - observeRun: () => Stream.fail(new Error("connection dropped")), + observeRun: () => Stream.fail(rpcFailure("connection dropped")), }); const result = yield* runCli( ["runs", "observe", unsafeRunId, "--server", unsafeServerUrl], @@ -898,13 +922,7 @@ describe("migrate runs", () => { }) ) ).pipe(Stream.concat(Stream.never)), - stopRun: () => - Effect.fail( - new MigrationCliConnectionError({ - cause: "stop request failed", - message: "stop request failed", - }) - ), + stopRun: () => Effect.fail(rpcFailure("stop request failed")), }); const resultFiber = yield* runCli( ["runs", "observe", runId], diff --git a/packages/migrate-sdk/src/cli/runtime.ts b/packages/migrate-sdk/src/cli/runtime.ts index eab6b4b..9baef2d 100644 --- a/packages/migrate-sdk/src/cli/runtime.ts +++ b/packages/migrate-sdk/src/cli/runtime.ts @@ -12,16 +12,26 @@ import { } from "effect"; import { Service } from "effect/Context"; import { Prompt } from "effect/unstable/cli"; +import type { RpcClientError } from "effect/unstable/rpc/RpcClientError"; import { connectMigrateServer, type MigrateServerConnectionInput, } from "../client/node/index.ts"; import type { MigrationRunId } from "../domain/ids.ts"; +import type { MigrationDefinitionLock } from "../domain/lock.ts"; import type { MigrateActiveRun, + MigrateBreakLockResult, + MigrateDashboardSnapshot, MigrateObservationEvent, MigrateOperationRequest, MigratePreparedOperation, + MigrateProtocolError, + MigrateRegistry, + MigrateRegistryMessagesReport, + MigrateRegistryMessagesRequest, + MigrateRegistryStatusReport, + MigrateRegistryStatusRequest, MigrateRunStartResult, MigrateRunStopResult, } from "../protocol/index.ts"; @@ -46,22 +56,41 @@ export class MigrationCliConnectionError extends Schema.TaggedError Effect.Effect; readonly dispose: () => Promise; - readonly getActiveRuns: Effect.Effect; + readonly getActiveRuns: Effect.Effect< + readonly MigrateActiveRun[], + MigrationCliServerError + >; + readonly getDashboard: Effect.Effect< + MigrateDashboardSnapshot, + MigrationCliServerError + >; + readonly getRegistry: Effect.Effect; + readonly getRegistryMessages: ( + input: MigrateRegistryMessagesRequest + ) => Effect.Effect; + readonly getRegistryStatus: ( + input: MigrateRegistryStatusRequest + ) => Effect.Effect; readonly observeRun: ( runId: MigrationRunId - ) => Stream.Stream; + ) => Stream.Stream; readonly prepareOperation: ( request: MigrateOperationRequest - ) => Effect.Effect; + ) => Effect.Effect; readonly startOperation: (input: { readonly acceptedFingerprint: MigratePreparedOperation["fingerprint"]; readonly request: MigrateOperationRequest; - }) => Effect.Effect; + }) => Effect.Effect; readonly stopRun: ( runId: MigrationRunId - ) => Effect.Effect; + ) => Effect.Effect; } export interface MigrationCliRuntimeShape { @@ -167,8 +196,15 @@ export class MigrationCliRuntime extends Service< try: () => connectMigrateServer(input), }).pipe( Effect.map((connection) => ({ + breakLock: (lock) => connection.client.BreakLock({ lock }), dispose: connection.dispose, getActiveRuns: connection.client.GetActiveRuns(), + getDashboard: connection.client.GetDashboard(), + getRegistry: connection.client.GetRegistry(), + getRegistryMessages: (input) => + connection.client.GetRegistryMessages(input), + getRegistryStatus: (input) => + connection.client.GetRegistryStatus(input), observeRun: (runId: MigrationRunId) => connection.client.observeRun({ runId }), prepareOperation: (request) => diff --git a/packages/migrate-sdk/src/client/index.ts b/packages/migrate-sdk/src/client/index.ts index 07604e5..32add67 100644 --- a/packages/migrate-sdk/src/client/index.ts +++ b/packages/migrate-sdk/src/client/index.ts @@ -28,6 +28,7 @@ import { makeMigrateClientService, makeStreamingMigrateClientService, } from "./internal/client-service.ts"; +import { rpcClientHttpStatusCode } from "./internal/rpc-client-error.ts"; type MigrateHttpRpcClient = RpcClient< Rpcs, @@ -38,29 +39,15 @@ export type { MigrateClientService } from "./internal/client-service.ts"; const transientHttpStatuses = new Set([408, 429, 500, 502, 503, 504]); -const statusCodeFromCause = (cause: unknown): number | undefined => { - if ( - typeof cause !== "object" || - cause === null || - !("response" in cause) || - typeof cause.response !== "object" || - cause.response === null || - !("status" in cause.response) || - typeof cause.response.status !== "number" - ) { - return; - } - - return cause.response.status; -}; - const retryableHttpObservationFailure = (cause: unknown): boolean => Schema.is(RpcClientError)(cause) && cause.reason._tag === "HttpError" && (cause.reason.kind === "TransportError" || cause.reason.kind === "EmptyBodyError" || (cause.reason.kind === "StatusCodeError" && - transientHttpStatuses.has(statusCodeFromCause(cause.reason.cause) ?? 0))); + Option.exists(rpcClientHttpStatusCode(cause), (status) => + transientHttpStatuses.has(status) + ))); const leasedObservation = ( client: MigrateHttpRpcClient, diff --git a/packages/migrate-sdk/src/client/internal/client-service.ts b/packages/migrate-sdk/src/client/internal/client-service.ts index be2d0e9..f69870c 100644 --- a/packages/migrate-sdk/src/client/internal/client-service.ts +++ b/packages/migrate-sdk/src/client/internal/client-service.ts @@ -44,6 +44,9 @@ export const makeMigrateClientService = ( GetActiveRuns: client.GetActiveRuns, GetDashboard: client.GetDashboard, GetMessages: client.GetMessages, + GetRegistry: client.GetRegistry, + GetRegistryMessages: client.GetRegistryMessages, + GetRegistryStatus: client.GetRegistryStatus, GetServerInfo: client.GetServerInfo, GetSourceIdentityHistory: client.GetSourceIdentityHistory, GetSourceItemTotals: client.GetSourceItemTotals, diff --git a/packages/migrate-sdk/src/client/internal/http-connection.ts b/packages/migrate-sdk/src/client/internal/http-connection.ts index 642adb9..d9129a6 100644 --- a/packages/migrate-sdk/src/client/internal/http-connection.ts +++ b/packages/migrate-sdk/src/client/internal/http-connection.ts @@ -1,12 +1,14 @@ -import { Effect, Layer, ManagedRuntime } from "effect"; +import { Cause, Effect, Layer, ManagedRuntime, Option, Schema } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { layerProtocolHttp } from "effect/unstable/rpc/RpcClient"; +import { RpcClientError } from "effect/unstable/rpc/RpcClientError"; import { layerNdjson } from "effect/unstable/rpc/RpcSerialization"; import { type MigrateConnection, validateMigrateServerProtocol, } from "../connection.ts"; import { MigrateClient } from "../index.ts"; +import { rpcClientHttpStatusCode } from "./rpc-client-error.ts"; export interface MigrateHttpConnectionOptions { readonly bearerToken?: string | undefined; @@ -15,6 +17,53 @@ export interface MigrateHttpConnectionOptions { readonly url: string; } +class MigrateHttpRpcConnectionError extends Schema.TaggedError()( + "MigrateHttpRpcConnectionError", + { + cause: Schema.Defect(), + message: Schema.String, + } +) {} + +const rpcConnectionFailureMessage = (error: RpcClientError): string => + Option.match(rpcClientHttpStatusCode(error), { + onNone: () => error.message, + onSome: (status) => { + switch (status) { + case 401: + return "Permission denied by Migrate Server (HTTP 401 Unauthorized). Check MIGRATE_SERVER_TOKEN."; + case 403: + return "Permission denied by Migrate Server (HTTP 403 Forbidden)."; + default: + return `Migrate Server returned HTTP ${status}.`; + } + }, + }); + +const mapRpcClientConnectionCause = ( + cause: Cause.Cause +): Effect.Effect => { + const rpcClientError = Cause.findErrorOption(cause).pipe( + Option.orElse(() => + Cause.findDefect(cause).pipe( + Option.getSuccess, + Option.filter(Schema.is(RpcClientError)) + ) + ) + ); + + return Option.match(rpcClientError, { + onNone: () => Effect.failCause(cause), + onSome: (error) => + Effect.fail( + new MigrateHttpRpcConnectionError({ + cause: error, + message: rpcConnectionFailureMessage(error), + }) + ), + }); +}; + const ipv4Part = /^\d{1,3}$/; export const migrateHttpServerUrl = (input: string): string => { @@ -81,7 +130,9 @@ const connectMigrateHttpServerUnsafe = ({ .runPromise( Effect.gen(function* () { const client = yield* MigrateClient; - const serverInfo = yield* client.GetServerInfo(); + const serverInfo = yield* client + .GetServerInfo() + .pipe(Effect.catchCause(mapRpcClientConnectionCause)); validateMigrateServerProtocol(serverInfo); return { diff --git a/packages/migrate-sdk/src/client/internal/rpc-client-error.ts b/packages/migrate-sdk/src/client/internal/rpc-client-error.ts new file mode 100644 index 0000000..b89d14c --- /dev/null +++ b/packages/migrate-sdk/src/client/internal/rpc-client-error.ts @@ -0,0 +1,18 @@ +import { Option } from "effect"; +import { HttpClientError } from "effect/unstable/http"; +import type { RpcClientError } from "effect/unstable/rpc/RpcClientError"; + +export const rpcClientHttpStatusCode = ( + error: RpcClientError +): Option.Option => { + if ( + error.reason._tag !== "HttpError" || + error.reason.kind !== "StatusCodeError" + ) { + return Option.none(); + } + + return error.reason.cause instanceof HttpClientError.StatusCodeError + ? Option.some(error.reason.cause.response.status) + : Option.none(); +}; diff --git a/packages/migrate-sdk/src/client/node/remote-connection.test.ts b/packages/migrate-sdk/src/client/node/remote-connection.test.ts index 42a7522..e16d514 100644 --- a/packages/migrate-sdk/src/client/node/remote-connection.test.ts +++ b/packages/migrate-sdk/src/client/node/remote-connection.test.ts @@ -84,6 +84,37 @@ describe("remote Migrate Server connection", () => { expect( await connection.runPromise(connection.client.GetDashboard()) ).toMatchObject({ dashboard }); + await expect( + connection.runPromise(connection.client.GetRegistry()) + ).resolves.toEqual({ + entries: dashboard.rows.map((row) => row.entry), + groups: dashboard.groups, + }); + await expect( + connection.runPromise( + connection.client.GetRegistryStatus({ + scanSource: false, + selection: { kind: "all" }, + withDependencies: false, + }) + ) + ).resolves.toMatchObject({ + includedDefinitionIds: ["articles"], + requestedDefinitionIds: "all", + scanSource: false, + }); + await expect( + connection.runPromise( + connection.client.GetRegistryMessages({ + selection: { kind: "all" }, + withDependencies: false, + }) + ) + ).resolves.toMatchObject({ + includedDefinitionIds: ["articles"], + messages: [], + requestedDefinitionIds: "all", + }); const dashboardSnapshots = await connection.runPromise( connection.client .observeDashboard({}) @@ -178,12 +209,26 @@ describe("remote Migrate Server connection", () => { fetch: (input, init) => http.handler(new Request(input, init)), url: "https://migrate.example/rpc", }) - ).rejects.toThrow("Unable to connect to Migrate Server"); + ).rejects.toThrow( + "Permission denied by Migrate Server (HTTP 401 Unauthorized). Check MIGRATE_SERVER_TOKEN." + ); } finally { await http.dispose(); } }); + it("reports forbidden remote connections as permission errors", async () => { + await expect( + connectRemoteMigrateServer({ + fetch: () => + Promise.resolve(new Response("Forbidden", { status: 403 })), + url: "https://migrate.example/rpc", + }) + ).rejects.toThrow( + "Permission denied by Migrate Server (HTTP 403 Forbidden)." + ); + }); + it("connects to a remote server with the same protocol and a different SDK version", async () => { const mismatchedServerLayer = Layer.effect( MigrateServer, diff --git a/packages/migrate-sdk/src/protocol/index.test.ts b/packages/migrate-sdk/src/protocol/index.test.ts index 8a6b39a..944df55 100644 --- a/packages/migrate-sdk/src/protocol/index.test.ts +++ b/packages/migrate-sdk/src/protocol/index.test.ts @@ -6,6 +6,9 @@ import { GetActiveRuns, GetDashboard, GetMessages, + GetRegistry, + GetRegistryMessages, + GetRegistryStatus, GetServerInfo, GetSourceIdentityHistory, MIGRATE_PROTOCOL_VERSION, @@ -35,8 +38,13 @@ import { MigratePrepareOptions, MigrateProtocolError, MigrateProtocolVersion, + MigrateRegistry, MigrateRegistryEntry, MigrateRegistryGroup, + MigrateRegistryMessagesReport, + MigrateRegistryMessagesRequest, + MigrateRegistryStatusReport, + MigrateRegistryStatusRequest, MigrateRunStartResult, MigrateSelection, MigrateServerInfo, @@ -87,6 +95,55 @@ const dashboardValue = { scannedSource: false, }; +const registryValue = { + entries: dashboardValue.rows.map((row) => row.entry), + groups: dashboardValue.groups, +}; + +const registryMessagesRequestValue = { + selection: { definitionIds: ["articles"], kind: "definitions" }, + withDependencies: true, +}; + +const registryMessagesReportValue = { + includedDefinitionIds: ["authors", "articles"], + messages: [ + { + ...messageBase, + kind: "skip-reason", + severity: "info", + }, + ], + notices: [], + requestedDefinitionIds: ["articles"], +}; + +const registryStatusRequestValue = { + concurrency: 2, + scanSource: true, + selection: { groupId: "content", kind: "group" }, + withDependencies: false, +}; + +const registryStatusReportValue = { + definitions: [ + { + definitionId: "articles", + discovery: "full", + durable: { failed: 0, migrated: 1, needsUpdate: 0, skipped: 0 }, + lastRun: null, + lock: null, + warnings: [], + }, + ], + includedDefinitionIds: ["articles"], + notices: [], + requestedDefinitionIds: ["articles"], + requestedGroup: "content", + scanSource: true, + warnings: [], +}; + const operationRequestValue = { action: "run", options: { withDependencies: true }, @@ -276,6 +333,31 @@ const contractCases: readonly { schema: MigrateRegistryGroup, value: { definitionIds: ["authors", "articles"], id: "content" }, }, + { + name: "registry", + schema: MigrateRegistry, + value: registryValue, + }, + { + name: "registry messages request", + schema: MigrateRegistryMessagesRequest, + value: registryMessagesRequestValue, + }, + { + name: "registry messages report", + schema: MigrateRegistryMessagesReport, + value: registryMessagesReportValue, + }, + { + name: "registry status request", + schema: MigrateRegistryStatusRequest, + value: registryStatusRequestValue, + }, + { + name: "registry status report", + schema: MigrateRegistryStatusReport, + value: registryStatusReportValue, + }, { name: "dashboard", schema: MigrateDashboard, @@ -524,6 +606,31 @@ const contractCases: readonly { message: "Migration group was not found", }, }, + { + name: "status request protocol error", + schema: MigrateProtocolError, + value: { + _tag: "MigrationStatusRequestError", + message: + "Status concurrency is only valid when source scanning is enabled", + }, + }, + { + name: "migration store protocol error", + schema: MigrateProtocolError, + value: { + _tag: "MigrationStoreError", + message: "Unable to read migration state", + }, + }, + { + name: "source protocol error", + schema: MigrateProtocolError, + value: { + _tag: "SourceError", + message: "Unable to scan source inventory", + }, + }, { name: "item error message", schema: MigrationMessage, @@ -594,6 +701,21 @@ const rpcPayloadCases: readonly { schema: GetDashboard.payloadSchema, value: undefined, }, + { + name: "GetRegistry", + schema: GetRegistry.payloadSchema, + value: undefined, + }, + { + name: "GetRegistryMessages", + schema: GetRegistryMessages.payloadSchema, + value: registryMessagesRequestValue, + }, + { + name: "GetRegistryStatus", + schema: GetRegistryStatus.payloadSchema, + value: registryStatusRequestValue, + }, { name: "GetActiveRuns", schema: GetActiveRuns.payloadSchema, @@ -694,6 +816,21 @@ const rpcUnarySuccessCases: readonly { resumeToken: "sha256:dashboard", }, }, + { + name: "GetRegistry", + schema: GetRegistry.successSchema, + value: registryValue, + }, + { + name: "GetRegistryMessages", + schema: GetRegistryMessages.successSchema, + value: registryMessagesReportValue, + }, + { + name: "GetRegistryStatus", + schema: GetRegistryStatus.successSchema, + value: registryStatusReportValue, + }, { name: "ObserveDashboardLease", schema: ObserveDashboardLease.successSchema, @@ -758,6 +895,9 @@ const rpcUnarySuccessCases: readonly { const rpcErrorCases = [ GetActiveRuns, GetDashboard, + GetRegistry, + GetRegistryMessages, + GetRegistryStatus, ObserveDashboardLease, GetMessages, GetSourceIdentityHistory, diff --git a/packages/migrate-sdk/src/protocol/index.ts b/packages/migrate-sdk/src/protocol/index.ts index b3c6f09..3a1968d 100644 --- a/packages/migrate-sdk/src/protocol/index.ts +++ b/packages/migrate-sdk/src/protocol/index.ts @@ -1,6 +1,7 @@ import { Schema } from "effect"; import { make as makeRpc } from "effect/unstable/rpc/Rpc"; import { make as makeRpcGroup } from "effect/unstable/rpc/RpcGroup"; +import { MigrationStoreError, SourceError } from "../domain/errors.ts"; import { MigrationDefinitionGroupId, MigrationDefinitionId, @@ -17,7 +18,11 @@ import { MigrationDefinitionRegistryUnknownGroupError, } from "../domain/registry.ts"; import { activeMigrationRunHasObservationDefinition } from "../domain/run.ts"; -import { MigrationDefinitionStatus } from "../domain/status.ts"; +import { + MigrationDefinitionStatus, + MigrationStatusRequestError, + MigrationStatusWarning, +} from "../domain/status.ts"; export const MIGRATE_PROTOCOL_VERSION = 1; @@ -154,6 +159,61 @@ export const MigrateRegistryGroup = Schema.Struct({ }); export type MigrateRegistryGroup = typeof MigrateRegistryGroup.Type; +export const MigrateRegistry = Schema.Struct({ + entries: Schema.Array(MigrateRegistryEntry), + groups: Schema.Array(MigrateRegistryGroup), +}); +export type MigrateRegistry = typeof MigrateRegistry.Type; + +const MigrateRegistrySelectionFields = { + selection: MigrateSelection, + withDependencies: Schema.Boolean, +} as const; + +export const MigrateRegistryMessagesRequest = Schema.Struct( + MigrateRegistrySelectionFields +); +export type MigrateRegistryMessagesRequest = + typeof MigrateRegistryMessagesRequest.Type; + +const MigrateRegistryStatusRequestFields = { + ...MigrateRegistrySelectionFields, + concurrency: Schema.optional(Schema.Int), + scanSource: Schema.Boolean, +} as const; + +export const MigrateRegistryStatusRequest = Schema.Struct( + MigrateRegistryStatusRequestFields +); +export type MigrateRegistryStatusRequest = + typeof MigrateRegistryStatusRequest.Type; + +const MigrateRegistrySelectionReportFields = { + includedDefinitionIds: Schema.Array(MigrationDefinitionId), + notices: Schema.Array(MigrationDefinitionPlanNotice), + requestedDefinitionIds: Schema.Union([ + Schema.Literal("all"), + Schema.Array(MigrationDefinitionId), + ]), + requestedGroup: Schema.optionalKey(MigrationDefinitionGroupId), +} as const; + +export const MigrateRegistryMessagesReport = Schema.Struct({ + ...MigrateRegistrySelectionReportFields, + messages: Schema.Array(MigrationMessage), +}); +export type MigrateRegistryMessagesReport = + typeof MigrateRegistryMessagesReport.Type; + +export const MigrateRegistryStatusReport = Schema.Struct({ + ...MigrateRegistrySelectionReportFields, + definitions: Schema.Array(MigrationDefinitionStatus), + scanSource: Schema.Boolean, + warnings: Schema.Array(MigrationStatusWarning), +}); +export type MigrateRegistryStatusReport = + typeof MigrateRegistryStatusReport.Type; + export const MigrateDashboardRow = Schema.Struct({ entry: MigrateRegistryEntry, status: Schema.optional(MigrationDefinitionStatus), @@ -500,6 +560,9 @@ export const MigrateProtocolError = Schema.Union([ MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError, MigrationDefinitionRegistryUnknownDefinitionError, MigrationDefinitionRegistryUnknownGroupError, + MigrationStatusRequestError, + MigrationStoreError, + SourceError, MigrateOperationError, MigratePlanChangedError, ]); @@ -514,6 +577,23 @@ export class GetDashboard extends makeRpc("GetDashboard", { success: MigrateDashboardSnapshot, }) {} +export class GetRegistry extends makeRpc("GetRegistry", { + error: MigrateProtocolError, + success: MigrateRegistry, +}) {} + +export class GetRegistryMessages extends makeRpc("GetRegistryMessages", { + error: MigrateProtocolError, + payload: MigrateRegistrySelectionFields, + success: MigrateRegistryMessagesReport, +}) {} + +export class GetRegistryStatus extends makeRpc("GetRegistryStatus", { + error: MigrateProtocolError, + payload: MigrateRegistryStatusRequestFields, + success: MigrateRegistryStatusReport, +}) {} + export class GetActiveRuns extends makeRpc("GetActiveRuns", { error: MigrateProtocolError, success: Schema.Array(MigrateActiveRun), @@ -630,6 +710,9 @@ export class BreakLock extends makeRpc("BreakLock", { const MigrateControlRpcs = makeRpcGroup( GetServerInfo, GetDashboard, + GetRegistry, + GetRegistryMessages, + GetRegistryStatus, GetActiveRuns, GetMessages, GetSourceIdentityHistory, diff --git a/packages/migrate-sdk/src/protocol/registry-selection.ts b/packages/migrate-sdk/src/protocol/registry-selection.ts new file mode 100644 index 0000000..9e944f4 --- /dev/null +++ b/packages/migrate-sdk/src/protocol/registry-selection.ts @@ -0,0 +1,25 @@ +import type { MigrationDefinitionRegistrySelectionInput } from "../domain/registry.ts"; +import type { MigrateRegistryMessagesRequest } from "./index.ts"; + +export const toMigrationDefinitionRegistrySelectionInput = ( + input: MigrateRegistryMessagesRequest +): MigrationDefinitionRegistrySelectionInput => { + switch (input.selection.kind) { + case "all": + return { all: true, withDependencies: input.withDependencies }; + case "definitions": + return { + definitionIds: input.selection.definitionIds, + withDependencies: input.withDependencies, + }; + case "group": + return { + group: input.selection.groupId, + withDependencies: input.withDependencies, + }; + default: { + const unhandled: never = input.selection; + return unhandled; + } + } +}; diff --git a/packages/migrate-sdk/src/server/handlers.test.ts b/packages/migrate-sdk/src/server/handlers.test.ts index a9c6e14..1161e62 100644 --- a/packages/migrate-sdk/src/server/handlers.test.ts +++ b/packages/migrate-sdk/src/server/handlers.test.ts @@ -64,6 +64,26 @@ const serverLayer = Layer.succeed( getActiveRuns: Effect.succeed(activeRuns), getDashboard: Effect.succeed(dashboardSnapshot), getMessages: () => Effect.succeed([]), + getRegistry: Effect.succeed({ + entries: dashboard.rows.map((row) => row.entry), + groups: dashboard.groups, + }), + getRegistryMessages: () => + Effect.succeed({ + includedDefinitionIds: [MigrationDefinitionId.make("articles")], + messages: [], + notices: [], + requestedDefinitionIds: [MigrationDefinitionId.make("articles")], + }), + getRegistryStatus: () => + Effect.succeed({ + definitions: [], + includedDefinitionIds: [MigrationDefinitionId.make("articles")], + notices: [], + requestedDefinitionIds: [MigrationDefinitionId.make("articles")], + scanSource: false, + warnings: [], + }), getServerInfo: Effect.succeed(info), getSourceIdentityHistory: () => Effect.succeed([]), getSourceItemTotals: () => @@ -124,6 +144,22 @@ const program = Effect.gen(function* () { const client = yield* makeClient(MigrateStreamingRpcs); const serverInfo = yield* client.GetServerInfo(); const currentDashboard = yield* client.GetDashboard(); + const currentRegistry = yield* client.GetRegistry(); + const registryMessages = yield* client.GetRegistryMessages({ + selection: { + definitionIds: [MigrationDefinitionId.make("articles")], + kind: "definitions", + }, + withDependencies: false, + }); + const registryStatus = yield* client.GetRegistryStatus({ + scanSource: false, + selection: { + definitionIds: [MigrationDefinitionId.make("articles")], + kind: "definitions", + }, + withDependencies: false, + }); const currentActiveRuns = yield* client.GetActiveRuns(); const sourceItemTotals = yield* client.GetSourceItemTotals({ definitionIds: [MigrationDefinitionId.make("articles")], @@ -138,7 +174,10 @@ const program = Effect.gen(function* () { return { currentActiveRuns, currentDashboard, + currentRegistry, dashboardSnapshots: [...dashboardSnapshots], + registryMessages, + registryStatus, runEvents: [...runEvents], serverInfo, sourceItemTotals, @@ -170,6 +209,24 @@ describe("Migrate Server RPC handlers", () => { expect(result.serverInfo).toEqual(info); expect(result.currentDashboard).toEqual(dashboardSnapshot); + expect(result.currentRegistry).toEqual({ + entries: dashboard.rows.map((row) => row.entry), + groups: [], + }); + expect(result.registryMessages).toEqual({ + includedDefinitionIds: ["articles"], + messages: [], + notices: [], + requestedDefinitionIds: ["articles"], + }); + expect(result.registryStatus).toEqual({ + definitions: [], + includedDefinitionIds: ["articles"], + notices: [], + requestedDefinitionIds: ["articles"], + scanSource: false, + warnings: [], + }); expect(result.dashboardSnapshots).toEqual([dashboardSnapshot]); expect(result.currentActiveRuns).toEqual(activeRuns); expect(result.sourceItemTotals).toEqual([ diff --git a/packages/migrate-sdk/src/server/handlers.ts b/packages/migrate-sdk/src/server/handlers.ts index 6be0269..86e7cd6 100644 --- a/packages/migrate-sdk/src/server/handlers.ts +++ b/packages/migrate-sdk/src/server/handlers.ts @@ -22,6 +22,9 @@ const controlHandlers = (server: MigrateServerService) => ({ GetActiveRuns: () => server.getActiveRuns, GetDashboard: () => server.getDashboard, GetMessages: server.getMessages, + GetRegistry: () => server.getRegistry, + GetRegistryMessages: server.getRegistryMessages, + GetRegistryStatus: server.getRegistryStatus, GetServerInfo: () => server.getServerInfo, GetSourceIdentityHistory: server.getSourceIdentityHistory, GetSourceItemTotals: server.getSourceItemTotals, diff --git a/packages/migrate-sdk/src/server/registry-backend.ts b/packages/migrate-sdk/src/server/registry-backend.ts index 7d062c2..9ae5ff5 100644 --- a/packages/migrate-sdk/src/server/registry-backend.ts +++ b/packages/migrate-sdk/src/server/registry-backend.ts @@ -7,6 +7,7 @@ import type { MigrateDashboard, MigrateOperationRequest, MigratePreparedOperation, + MigrateRegistry, } from "../protocol/index.ts"; import type { ExecutableMigrationOperation, @@ -126,6 +127,12 @@ export const makeRegistryMigrateServerBackend = ( Effect.map((snapshot) => dashboard(runtime, snapshot)) ), getMessages: (target) => runtime.listMessages(target), + getRegistry: Effect.succeed({ + entries: runtime.entries, + groups: runtime.groups, + } satisfies MigrateRegistry), + getRegistryMessages: runtime.getRegistryMessages, + getRegistryStatus: runtime.getRegistryStatus, getRunProgress: runtime.getRunProgress, getSourceIdentityHistory: runtime.listSourceIdentityHistory, getSourceItemTotals: runtime.getSourceItemTotals, diff --git a/packages/migrate-sdk/src/server/registry-runtime.test.ts b/packages/migrate-sdk/src/server/registry-runtime.test.ts index 7d6d611..babcd93 100644 --- a/packages/migrate-sdk/src/server/registry-runtime.test.ts +++ b/packages/migrate-sdk/src/server/registry-runtime.test.ts @@ -11,6 +11,7 @@ import { SourceIdentity, toMigrationDefinitionId, toMigrationRunId, + toSourceVersion, } from "../index.ts"; import { InMemoryMigrationStore } from "../stores/in-memory/in-memory-migration-store.ts"; import { makeRegistryMigrateServerRuntime } from "./registry-runtime.ts"; @@ -26,6 +27,211 @@ describe("registry migration server runtime", () => { expect(runtime.groups).toEqual([]); }); + it.effect("returns canonical empty registry reports", () => + Effect.gen(function* () { + const runtime = makeRegistryMigrateServerRuntime({ + executable: MigrationExecutable.inlineService, + registry: MigrationDefinitionRegistry.make({ definitions: [] }), + }); + + expect( + yield* runtime.getRegistryStatus({ + scanSource: false, + selection: { kind: "all" }, + withDependencies: false, + }) + ).toEqual({ + definitions: [], + includedDefinitionIds: [], + notices: [], + requestedDefinitionIds: "all", + scanSource: false, + warnings: [], + }); + expect( + yield* runtime.getRegistryMessages({ + selection: { kind: "all" }, + withDependencies: false, + }) + ).toEqual({ + includedDefinitionIds: [], + messages: [], + notices: [], + requestedDefinitionIds: "all", + }); + }) + ); + + it.effect( + "preserves registry status selection and validation semantics", + () => + Effect.gen(function* () { + const identity = SourceIdentity.make({ + id: "registry-server-status@v1", + schema: SourceIdentity.key("id", Schema.NonEmptyString), + }); + const makeDefinition = ( + id: "authors" | "articles", + required: readonly ReturnType[] = [] + ) => + MigrationDefinition.make({ + dependencies: { required }, + id: toMigrationDefinitionId(id), + process: () => Effect.void, + source: Source.make({ + cursorSchema: Schema.Struct({ offset: Schema.Int }), + identity, + lookupStrategy: "direct", + read: () => Effect.succeed({ items: [] }), + readByIdentity: () => Effect.succeed(null), + sourceSchema: Schema.Struct({ title: Schema.String }), + }), + store: InMemoryMigrationStore.layer(), + }); + const authorsId = toMigrationDefinitionId("authors"); + const articlesId = toMigrationDefinitionId("articles"); + const runtime = makeRegistryMigrateServerRuntime({ + executable: MigrationExecutable.inlineService, + registry: MigrationDefinitionRegistry.make({ + definitions: [ + makeDefinition("authors"), + makeDefinition("articles", [authorsId]), + ] as const, + }), + }); + + expect( + yield* runtime + .getRegistryStatus({ + scanSource: false, + selection: { + definitionIds: [articlesId], + kind: "definitions", + }, + withDependencies: false, + }) + .pipe(Effect.flip) + ).toMatchObject({ + _tag: "MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError", + definitionId: articlesId, + missingDependencyIds: [authorsId], + }); + + const expanded = yield* runtime.getRegistryStatus({ + scanSource: false, + selection: { definitionIds: [articlesId], kind: "definitions" }, + withDependencies: true, + }); + expect(expanded.includedDefinitionIds).toEqual([authorsId, articlesId]); + expect( + expanded.definitions.map((status) => status.definitionId) + ).toEqual([authorsId, articlesId]); + + expect( + yield* runtime + .getRegistryStatus({ + concurrency: 2, + scanSource: false, + selection: { kind: "all" }, + withDependencies: false, + }) + .pipe(Effect.flip) + ).toMatchObject({ + _tag: "MigrationStatusRequestError", + message: + "Status concurrency is only valid when source scanning is enabled", + }); + expect( + yield* runtime + .getRegistryStatus({ + concurrency: 0, + scanSource: true, + selection: { kind: "all" }, + withDependencies: false, + }) + .pipe(Effect.flip) + ).toMatchObject({ + _tag: "MigrationStatusRequestError", + message: "Status concurrency must be a positive integer", + }); + }) + ); + + it.effect("sorts registry messages globally across definitions", () => + Effect.gen(function* () { + const identity = SourceIdentity.make({ + id: "registry-server-messages@v1", + schema: SourceIdentity.key("id", Schema.NonEmptyString), + }); + const state = InMemoryMigrationStore.makeState(); + const store = InMemoryMigrationStore.layer(state); + const makeDefinition = (id: "authors" | "articles") => + MigrationDefinition.make({ + id: toMigrationDefinitionId(id), + process: () => Effect.void, + source: Source.make({ + cursorSchema: Schema.Struct({ offset: Schema.Int }), + identity, + lookupStrategy: "direct", + read: () => Effect.succeed({ items: [] }), + readByIdentity: () => Effect.succeed(null), + sourceSchema: Schema.Struct({ title: Schema.String }), + }), + store, + }); + const authorsId = toMigrationDefinitionId("authors"); + const articlesId = toMigrationDefinitionId("articles"); + const runId = toMigrationRunId("run-messages"); + const addSkippedItem = ( + definitionId: typeof authorsId, + sourceIdentity: string, + message: string, + updatedAt: Date + ) => { + state.itemStates.set( + InMemoryMigrationStore.itemStateKey(definitionId, sourceIdentity), + { + definitionId, + lastRunId: runId, + skipReason: message, + sourceIdentity: SourceIdentity.fromKey(identity, sourceIdentity), + sourceVersion: toSourceVersion("v1"), + status: "skipped", + updatedAt, + } + ); + }; + addSkippedItem( + authorsId, + "author-1", + "Older author message", + new Date("2026-08-29T11:00:00.000Z") + ); + addSkippedItem( + articlesId, + "article-1", + "Newer article message", + new Date("2026-08-29T12:00:00.000Z") + ); + const runtime = makeRegistryMigrateServerRuntime({ + executable: MigrationExecutable.inlineService, + registry: MigrationDefinitionRegistry.make({ + definitions: [makeDefinition("authors"), makeDefinition("articles")], + }), + }); + + const report = yield* runtime.getRegistryMessages({ + selection: { kind: "all" }, + withDependencies: false, + }); + + expect(report.messages.map((message) => message.message)).toEqual([ + "Newer article message", + "Older author message", + ]); + }) + ); + it.effect( "reports server-owned execution failures independently of durable failed items", () => { diff --git a/packages/migrate-sdk/src/server/registry-runtime.ts b/packages/migrate-sdk/src/server/registry-runtime.ts index afc007f..d6fe576 100644 --- a/packages/migrate-sdk/src/server/registry-runtime.ts +++ b/packages/migrate-sdk/src/server/registry-runtime.ts @@ -9,8 +9,10 @@ import { type MigrationDefinitionId, type MigrationDefinitionLock, type MigrationDefinitionRegistry, + type MigrationDefinitionRegistryEntry, type MigrationDefinitionRegistryGroup, type MigrationDefinitionRegistryId, + type MigrationDefinitionRegistryMessagesReport, type MigrationDefinitionRegistryStatusReport, type MigrationDefinitionStatus, type MigrationExecutableObservationResult, @@ -40,12 +42,15 @@ import type { MigrateDefinitionSourceItemTotal, MigrateDependencyCheck, MigrateExecutionState, + MigrateRegistryMessagesRequest, + MigrateRegistryStatusRequest, MigrateSelection, MigrateSourceIdentityHistoryEntry, MigrateSourceItemTotal, MigrateTarget, MigrateTerminalSummary, } from "../protocol/index.ts"; +import { toMigrationDefinitionRegistrySelectionInput } from "../protocol/registry-selection.ts"; import { MigrationDefinitionSource } from "../services/migration-definition-source.ts"; import { isTerminalRunState, @@ -253,6 +258,13 @@ export interface RegistryMigrateServerRuntime { readonly breakLock: ( lock: MigrationDefinitionLock ) => Effect.Effect; + readonly entries: readonly MigrationDefinitionRegistryEntry[]; + readonly getRegistryMessages: ( + input: MigrateRegistryMessagesRequest + ) => Effect.Effect; + readonly getRegistryStatus: ( + input: MigrateRegistryStatusRequest + ) => Effect.Effect; readonly getRunProgress: ( runId: MigrationRunId, observationDefinitionId?: MigrationDefinitionId @@ -373,6 +385,18 @@ export const makeRegistryMigrateServerRuntime = ( const providerSettlementGraceMs = input.providerSettlementGraceMs ?? 2000; const terminalPollIntervalMs = input.terminalPollIntervalMs ?? 500; + const getRegistryMessages = (input: MigrateRegistryMessagesRequest) => + registry.messages(toMigrationDefinitionRegistrySelectionInput(input)); + + const getRegistryStatus = (input: MigrateRegistryStatusRequest) => + registry.status({ + ...toMigrationDefinitionRegistrySelectionInput(input), + ...(input.concurrency === undefined + ? {} + : { concurrency: input.concurrency }), + scanSource: input.scanSource, + }); + const readExecutionProgressEffect = ( definitionIds: readonly MigrationDefinitionId[] ): Effect.Effect => { @@ -1546,6 +1570,9 @@ export const makeRegistryMigrateServerRuntime = ( return { breakLock, + entries, + getRegistryMessages, + getRegistryStatus, groups, hasActiveExecutions: () => activeExecutions.size > 0, getRunProgress, diff --git a/packages/migrate-sdk/src/server/service.test.ts b/packages/migrate-sdk/src/server/service.test.ts index 2d1cf5e..bff056d 100644 --- a/packages/migrate-sdk/src/server/service.test.ts +++ b/packages/migrate-sdk/src/server/service.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from "@effect/vitest"; import { Deferred, Effect, Fiber, Option, Queue, Stream } from "effect"; import { TestClock } from "effect/testing"; +import { MigrationStoreError, SourceError } from "../domain/errors.ts"; import { MigrationDefinitionGroupId, MigrationDefinitionId, @@ -12,6 +13,7 @@ import { MigrationDefinitionRegistryUnknownDefinitionError, MigrationDefinitionRegistryUnknownGroupError, } from "../domain/registry.ts"; +import { MigrationStatusRequestError } from "../domain/status.ts"; import { type MigrateActiveRun, type MigrateDashboard, @@ -114,6 +116,9 @@ const makeBackend = (input?: { readonly executeOperation?: MigrateServerBackend["executeOperation"]; readonly getActiveRuns?: MigrateServerBackend["getActiveRuns"]; readonly getDashboard?: MigrateServerBackend["getDashboard"]; + readonly getRegistry?: MigrateServerBackend["getRegistry"]; + readonly getRegistryMessages?: MigrateServerBackend["getRegistryMessages"]; + readonly getRegistryStatus?: MigrateServerBackend["getRegistryStatus"]; readonly getRunProgress?: MigrateServerBackend["getRunProgress"]; readonly getSourceItemTotals?: MigrateServerBackend["getSourceItemTotals"]; readonly observeRun?: MigrateServerBackend["observeRun"]; @@ -148,6 +153,28 @@ const makeBackend = (input?: { scannedSource: false, }), getMessages: () => Effect.succeed([]), + getRegistry: + input?.getRegistry ?? Effect.succeed({ entries: [], groups: [] }), + getRegistryMessages: + input?.getRegistryMessages ?? + (() => + Effect.succeed({ + includedDefinitionIds: [], + messages: [], + notices: [], + requestedDefinitionIds: "all", + })), + getRegistryStatus: + input?.getRegistryStatus ?? + (() => + Effect.succeed({ + definitions: [], + includedDefinitionIds: [], + notices: [], + requestedDefinitionIds: "all", + scanSource: false, + warnings: [], + })), getRunProgress: input?.getRunProgress ?? (() => Effect.sync(() => undefined)), getSourceIdentityHistory: () => Effect.succeed([]), getSourceItemTotals: input?.getSourceItemTotals ?? (() => Effect.succeed([])), @@ -212,6 +239,59 @@ describe("Migrate Server", () => { }) ); + it.effect("preserves typed registry inspection errors", () => + Effect.gen(function* () { + const statusError = new MigrationStatusRequestError({ + message: + "Status concurrency is only valid when source scanning is enabled", + }); + const server = yield* makeServer( + makeBackend({ + getRegistryStatus: () => Effect.fail(statusError), + }) + ); + const error = yield* server + .getRegistryStatus({ + concurrency: 2, + scanSource: false, + selection: { kind: "all" }, + withDependencies: false, + }) + .pipe(Effect.flip); + + expect(error).toEqual(statusError); + + const storeError = new MigrationStoreError({ + message: "Unable to read migration messages", + }); + const sourceError = new SourceError({ + message: "Unable to scan source inventory", + }); + const failingServer = yield* makeServer( + makeBackend({ + getRegistryMessages: () => Effect.fail(storeError), + getRegistryStatus: () => Effect.fail(sourceError), + }) + ); + const messagesFailure = yield* failingServer + .getRegistryMessages({ + selection: { kind: "all" }, + withDependencies: false, + }) + .pipe(Effect.flip); + const statusFailure = yield* failingServer + .getRegistryStatus({ + scanSource: true, + selection: { kind: "all" }, + withDependencies: false, + }) + .pipe(Effect.flip); + + expect(messagesFailure).toEqual(storeError); + expect(statusFailure).toEqual(sourceError); + }) + ); + it.effect( "discovers and observes an active run without the transient execution map", () => diff --git a/packages/migrate-sdk/src/server/service.ts b/packages/migrate-sdk/src/server/service.ts index 07c0263..5085780 100644 --- a/packages/migrate-sdk/src/server/service.ts +++ b/packages/migrate-sdk/src/server/service.ts @@ -45,6 +45,11 @@ import { MigratePlanFingerprint, type MigratePreparedOperation, MigrateProtocolError, + type MigrateRegistry, + type MigrateRegistryMessagesReport, + type MigrateRegistryMessagesRequest, + type MigrateRegistryStatusReport, + type MigrateRegistryStatusRequest, type MigrateRunStartResult, type MigrateRunStopResult, type MigrateServerInfo, @@ -72,6 +77,13 @@ export interface MigrateServerService { readonly getMessages: (input: { readonly target: MigrateTarget; }) => Effect.Effect; + readonly getRegistry: Effect.Effect; + readonly getRegistryMessages: ( + input: MigrateRegistryMessagesRequest + ) => Effect.Effect; + readonly getRegistryStatus: ( + input: MigrateRegistryStatusRequest + ) => Effect.Effect; readonly getServerInfo: Effect.Effect; readonly getSourceIdentityHistory: (input: { readonly definitionId: MigrationDefinitionId; @@ -168,6 +180,13 @@ export interface MigrateServerBackend { readonly getMessages: ( target: MigrateTarget ) => Effect.Effect; + readonly getRegistry: Effect.Effect; + readonly getRegistryMessages: ( + input: MigrateRegistryMessagesRequest + ) => Effect.Effect; + readonly getRegistryStatus: ( + input: MigrateRegistryStatusRequest + ) => Effect.Effect; readonly getRunProgress: ( runId: MigrationRunId, observationDefinitionId?: MigrationDefinitionId @@ -1302,6 +1321,11 @@ const makeMigrationServerServiceWithInvalidationQueue = ( getActiveRuns, getMessages: ({ target }) => backend.getMessages(target).pipe(Effect.mapError(operationError)), + getRegistry: backend.getRegistry.pipe(Effect.mapError(operationError)), + getRegistryMessages: (input) => + backend.getRegistryMessages(input).pipe(Effect.mapError(operationError)), + getRegistryStatus: (input) => + backend.getRegistryStatus(input).pipe(Effect.mapError(operationError)), getServerInfo: Effect.succeed(serverInfo), getSourceIdentityHistory: ({ definitionId }) => backend diff --git a/packages/migrate-sdk/test/fixtures/remote-server.ts b/packages/migrate-sdk/test/fixtures/remote-server.ts index 3ef6e17..6f579ae 100644 --- a/packages/migrate-sdk/test/fixtures/remote-server.ts +++ b/packages/migrate-sdk/test/fixtures/remote-server.ts @@ -13,6 +13,7 @@ import { import type { MigrateActiveRun, MigrateDashboard, + MigrateSelection, MigrateServerInfo, } from "../../src/protocol/index.ts"; import { MIGRATE_PROTOCOL_VERSION } from "../../src/protocol/index.ts"; @@ -111,6 +112,23 @@ export const makeAuthorizedRemoteMigrateServerHttp = ( authorizationMiddleware(authorize) ); +const requestedDefinitionIds = ( + selection: MigrateSelection +): "all" | readonly MigrationDefinitionId[] => { + switch (selection.kind) { + case "all": + return "all"; + case "definitions": + return selection.definitionIds; + case "group": + return [remoteMigrateDefinitionId]; + default: { + const unhandled: never = selection; + return unhandled; + } + } +}; + export const remoteMigrateServerBackend: MigrateServerBackend = { breakLock: (lock: MigrationDefinitionLock) => @@ -138,6 +156,34 @@ export const remoteMigrateServerBackend: MigrateServerBackend Effect.succeed([]), + getRegistry: Effect.succeed({ + entries: remoteMigrateDashboard.rows.map((row) => row.entry), + groups: remoteMigrateDashboard.groups, + }), + getRegistryMessages: (request) => + Effect.succeed({ + includedDefinitionIds: [remoteMigrateDefinitionId], + messages: [], + notices: [], + ...(request.selection.kind === "group" + ? { requestedGroup: request.selection.groupId } + : {}), + requestedDefinitionIds: requestedDefinitionIds(request.selection), + }), + getRegistryStatus: (request) => + Effect.succeed({ + definitions: remoteMigrateDashboard.rows.flatMap((row) => + row.status === undefined ? [] : [row.status] + ), + includedDefinitionIds: [remoteMigrateDefinitionId], + notices: [], + ...(request.selection.kind === "group" + ? { requestedGroup: request.selection.groupId } + : {}), + requestedDefinitionIds: requestedDefinitionIds(request.selection), + scanSource: request.scanSource, + warnings: [], + }), getRunProgress: () => Effect.succeed({ definitions: remoteMigrateDashboard.rows.flatMap((row) =>