From 27443bc474485b3fe548f59ae1e2df1365cf2fb3 Mon Sep 17 00:00:00 2001 From: Alexei Gorobet Date: Tue, 1 Sep 2026 10:01:05 -0400 Subject: [PATCH 1/2] fix(cli): support remote server commands --- .changeset/remote-cli-commands.md | 8 + README.md | 12 + ...r-boundary-for-local-and-remote-clients.md | 23 +- packages/migrate-sdk/src/cli/command.ts | 360 +++++++++++++++++- packages/migrate-sdk/src/cli/render.ts | 35 +- .../migrate-sdk/src/cli/runs-command.test.ts | 188 ++++++++- packages/migrate-sdk/src/cli/runtime.ts | 22 ++ packages/migrate-sdk/src/client/index.ts | 21 +- .../src/client/internal/http-connection.ts | 55 ++- .../src/client/internal/rpc-client-error.ts | 18 + .../src/client/node/remote-connection.test.ts | 16 +- 11 files changed, 715 insertions(+), 43 deletions(-) create mode 100644 .changeset/remote-cli-commands.md create mode 100644 packages/migrate-sdk/src/client/internal/rpc-client-error.ts diff --git a/.changeset/remote-cli-commands.md b/.changeset/remote-cli-commands.md new file mode 100644 index 0000000..053184e --- /dev/null +++ b/.changeset/remote-cli-commands.md @@ -0,0 +1,8 @@ +--- +"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. 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..b9aaa8c 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 @@ -265,20 +265,27 @@ 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. +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 diff --git a/packages/migrate-sdk/src/cli/command.ts b/packages/migrate-sdk/src/cli/command.ts index 0f06d67..be085ae 100644 --- a/packages/migrate-sdk/src/cli/command.ts +++ b/packages/migrate-sdk/src/cli/command.ts @@ -23,16 +23,20 @@ import { MigrationMessage } from "../domain/message.ts"; import type { MigrationDefinitionRegistry, MigrationDefinitionRegistryMessagesError, + MigrationDefinitionRegistryMessagesReport, MigrationDefinitionRegistryPlanningError, MigrationDefinitionRegistrySelectionInput, MigrationDefinitionRegistryStatusError, MigrationDefinitionRegistryStatusInput, + MigrationDefinitionRegistryStatusReport, } from "../domain/registry.ts"; import type { MigrateAction, + MigrateDashboardRow, MigrateOperationRequest, MigratePreparedOperation, MigrateSelection, + MigrateTarget, } from "../protocol/index.ts"; import { MigrationStore } from "../services/migration-store.ts"; import { SqlMigrationStore } from "../stores/sql/sql-migration-store.ts"; @@ -57,6 +61,8 @@ import { renderPreparedOperationDependencyFailure, renderPreparedOperationPlan, renderPreparedOperationWarnings, + renderRegistryEntriesGraph, + renderRegistryEntriesList, renderRegistryGraph, renderRegistryList, renderRunStopResult, @@ -169,7 +175,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" ); } @@ -327,6 +333,17 @@ const withCliMigrateConnection = ( releaseCliMigrateConnection ); +const reportCliConnectionErrors = ( + effect: Effect.Effect +): Effect.Effect => + effect.pipe( + Effect.catch((error) => + CliError.isCliError(error) + ? Effect.fail(error) + : failReportedCliMessage(renderStoredFailure(error)) + ) + ); + const requireConfiguredSqlStore = (config: MigrationCliConfig) => config.sqlStore === undefined ? failReportedCliMessage( @@ -499,6 +516,24 @@ const hasRegisteredDefinition = ( const listCommand = Command.make("list", {}, () => Effect.gen(function* () { + const root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* reportCliConnectionErrors( + withCliMigrateConnection(({ connection }) => + Effect.gen(function* () { + const snapshot = yield* connection.getDashboard; + yield* Console.log( + renderRegistryEntriesList( + snapshot.dashboard.rows.map((row) => row.entry), + { colors: yield* useColor } + ) + ); + }) + ) + ); + } + const registry = yield* loadConfiguredRegistry; yield* Console.log( @@ -514,9 +549,41 @@ const graphCommand = Command.make( { definition: graphDefinition }, ({ definition }) => Effect.gen(function* () { + const focusedDefinitionId = Option.getOrUndefined(definition); + const root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* reportCliConnectionErrors( + withCliMigrateConnection(({ connection }) => + Effect.gen(function* () { + const snapshot = yield* connection.getDashboard; + const entries = snapshot.dashboard.rows.map((row) => row.entry); + + if ( + focusedDefinitionId !== undefined && + !entries.some((entry) => entry.id === focusedDefinitionId) + ) { + return yield* failReportedCliMessage( + `Migration Definition was not found in the registry: ${focusedDefinitionId}` + ); + } + + yield* Console.log( + renderRegistryEntriesGraph( + entries, + focusedDefinitionId === undefined + ? undefined + : toMigrationDefinitionId(focusedDefinitionId), + { colors: yield* useColor } + ) + ); + }) + ) + ); + } + const loadedConfig = yield* loadConfiguredConfig; const registry = loadedConfig.registry; - const focusedDefinitionId = Option.getOrUndefined(definition); if (focusedDefinitionId !== undefined) { const definitionId = toMigrationDefinitionId(focusedDefinitionId); @@ -830,6 +897,139 @@ const renderMessagesCommandError = ( }) : renderRuntimeError(error); +interface RemoteReadSelection { + readonly operation: MigratePreparedOperation; + readonly selection: MigrateSelection; +} + +function prepareRemoteReadSelection( + connection: MigrationCliServerConnection, + input: CliRegistrySelectionInput, + renderError: (error: unknown) => string +): Effect.Effect { + return Effect.gen(function* () { + const selection = yield* makeMigrateSelection(input); + const operation = yield* connection + .prepareOperation({ + action: "run", + options: { withDependencies: input.withDependencies }, + selection, + }) + .pipe( + Effect.catch((error) => failReportedCliMessage(renderError(error))) + ); + + return { operation, selection }; + }); +} + +function remoteStatusScanTargets( + selection: MigrateSelection, + includedDefinitionIds: readonly MigrationDefinitionId[] +): readonly MigrateTarget[] { + switch (selection.kind) { + case "group": + return [{ groupId: selection.groupId, kind: "group" }]; + case "definitions": + return selection.definitionIds.map((definitionId) => ({ + definitionId, + kind: "migration", + })); + case "all": + return includedDefinitionIds.map((definitionId) => ({ + definitionId, + kind: "migration", + })); + default: { + const unhandled: never = selection; + return unhandled; + } + } +} + +function collectRemoteStatusRows( + connection: MigrationCliServerConnection, + selection: MigrateSelection, + operation: MigratePreparedOperation, + options: { readonly concurrency?: number; readonly scanSource: boolean } +): Effect.Effect< + ReadonlyMap, + unknown +> { + const rowsById = new Map( + operation.planRows.map((row) => [row.entry.id, row]) + ); + + if (!options.scanSource) { + return Effect.succeed(rowsById); + } + + return Effect.gen(function* () { + const includedIds = new Set(operation.plan.includedDefinitionIds); + const targets = remoteStatusScanTargets( + selection, + operation.plan.includedDefinitionIds + ); + + for (const target of targets) { + const dashboard = yield* connection.scanSource({ + ...(options.concurrency === undefined + ? {} + : { concurrency: options.concurrency }), + target, + }); + + for (const row of dashboard.rows) { + if (includedIds.has(row.entry.id) && row.status?.source !== undefined) { + rowsById.set(row.entry.id, row); + } + } + } + + return rowsById; + }); +} + +function makeRemoteStatusReport( + operation: MigratePreparedOperation, + rowsById: ReadonlyMap, + scanSource: boolean +): MigrationDefinitionRegistryStatusReport { + const definitions = operation.plan.includedDefinitionIds.flatMap( + (definitionId) => { + const status = rowsById.get(definitionId)?.status; + return status === undefined ? [] : [status]; + } + ); + + return { + definitions, + includedDefinitionIds: operation.plan.includedDefinitionIds, + notices: operation.plan.notices, + ...(operation.plan.requestedGroup === undefined + ? {} + : { requestedGroup: operation.plan.requestedGroup }), + requestedDefinitionIds: operation.plan.requestedDefinitionIds, + scanSource, + warnings: definitions.flatMap((definition) => definition.warnings), + }; +} + +function makeRemoteMessagesReport( + operation: MigratePreparedOperation, + messages: readonly MigrationMessage[] +): MigrationDefinitionRegistryMessagesReport { + return { + includedDefinitionIds: operation.plan.includedDefinitionIds, + messages, + notices: operation.plan.notices, + ...(operation.plan.requestedGroup === undefined + ? {} + : { requestedGroup: operation.plan.requestedGroup }), + requestedDefinitionIds: operation.plan.requestedDefinitionIds, + }; +} + const statusCommand = Command.make( "status", { @@ -842,10 +1042,60 @@ const statusCommand = Command.make( }, (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 root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* reportCliConnectionErrors( + withCliMigrateConnection(({ connection }) => + Effect.gen(function* () { + const { operation, selection } = + yield* prepareRemoteReadSelection( + connection, + { + all: input.all, + definitionIds: input.definitions, + ...(groupInput === undefined ? {} : { group: groupInput }), + withDependencies: input.withDependencies, + }, + (error) => + isPlanningError(error) + ? renderStatusCommandError(error, { + definitionIds: input.definitions, + ...(groupInput === undefined + ? {} + : { group: groupInput }), + }) + : renderStoredFailure(error) + ); + const rowsById = yield* collectRemoteStatusRows( + connection, + selection, + operation, + { + ...(concurrencyInput === undefined + ? {} + : { concurrency: concurrencyInput }), + scanSource: input.scanSource, + } + ); + const report = makeRemoteStatusReport( + operation, + rowsById, + input.scanSource + ); + + yield* Console.log( + renderStatusReport(report, { colors: yield* useColor }) + ); + }) + ) + ); + } + + const loadedConfig = yield* loadConfiguredConfig; + const registry = toCliExecutableRegistry(loadedConfig.registry); const statusInput = makeStatusInput({ all: input.all, ...(concurrencyInput === undefined @@ -884,9 +1134,59 @@ const messagesCommand = Command.make( }, (input) => Effect.gen(function* () { + const groupInput = Option.getOrUndefined(input.group); + const root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* reportCliConnectionErrors( + withCliMigrateConnection(({ connection }) => + Effect.gen(function* () { + const { operation } = yield* prepareRemoteReadSelection( + connection, + { + all: input.all, + definitionIds: input.definitions, + ...(groupInput === undefined ? {} : { group: groupInput }), + withDependencies: input.withDependencies, + }, + (error) => + isPlanningError(error) + ? renderMessagesCommandError(error, { + definitionIds: input.definitions, + ...(groupInput === undefined + ? {} + : { group: groupInput }), + }) + : renderStoredFailure(error) + ); + const messages = (yield* Effect.forEach( + operation.plan.includedDefinitionIds, + (definitionId) => + connection.getMessages({ + definitionId, + kind: "migration", + }) + )).flat(); + const report = makeRemoteMessagesReport(operation, messages); + + 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 }) + ); + }) + ) + ); + } + const loadedConfig = yield* loadConfiguredConfig; const registry = loadedConfig.registry; - const groupInput = Option.getOrUndefined(input.group); const report = yield* registry .messages( makeRegistrySelectionInput({ @@ -932,9 +1232,57 @@ const unlockCommand = Command.make( { definition: unlockDefinition }, ({ definition }) => Effect.gen(function* () { + const definitionId = toMigrationDefinitionId(definition); + const root = yield* migrateBaseCommand; + + if (Option.isSome(root.server)) { + return yield* reportCliConnectionErrors( + withCliMigrateConnection(({ connection }) => + Effect.gen(function* () { + const snapshot = yield* connection.getDashboard; + const row = snapshot.dashboard.rows.find( + (candidate) => candidate.entry.id === definitionId + ); + + if (row === undefined) { + return yield* failReportedCliMessage( + `Migration Definition was not found in the registry: ${definitionId}` + ); + } + + const lock = row.status?.lock; + + if (lock === undefined || lock === null) { + yield* Console.log( + `Migration Definition lock is already clear: ${definitionId}` + ); + return; + } + + const result = yield* connection.breakLock(lock); + + if (result.kind === "already-clear") { + yield* Console.log( + `Migration Definition lock is already clear: ${definitionId}` + ); + return; + } + + yield* Console.log( + [ + "Migration Definition lock cleared", + `Migration ID ${definitionId}`, + `Owner Run ID ${lock.ownerRunId}`, + `Token ${lock.token}`, + ].join("\n") + ); + }) + ) + ); + } + const loadedConfig = yield* loadConfiguredConfig; const registry = loadedConfig.registry; - const definitionId = toMigrationDefinitionId(definition); const migrationDefinition = Option.getOrUndefined( registry.get(definitionId) ); 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..b558760 100644 --- a/packages/migrate-sdk/src/cli/runs-command.test.ts +++ b/packages/migrate-sdk/src/cli/runs-command.test.ts @@ -16,8 +16,10 @@ import { TestClock, TestConsole } from "effect/testing"; import { CliOutput, Command } from "effect/unstable/cli"; import type { MigrateServerConnectionInput } from "../client/node/index.ts"; import { + toEncodedSourceIdentity, toMigrationDefinitionGroupId, toMigrationDefinitionId, + toMigrationDefinitionLockToken, toMigrationRunId, } from "../domain/ids.ts"; import { @@ -26,12 +28,16 @@ import { } from "../domain/registry.ts"; import type { MigrateActiveRun, + MigrateDashboardRow, MigrateObservationEvent, MigrateOperationRequest, 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 { @@ -52,6 +58,26 @@ const activeRun: MigrateActiveRun = { status: "running", stopSupported: true, }; +const remoteDashboardRow = { + entry: { + dependencies: { optional: [], required: [] }, + hasRollback: true, + id: definitionId, + }, + status: { + definitionId, + discovery: "full", + durable: { failed: 0, migrated: 3, needsUpdate: 0, skipped: 1 }, + lastRun: null, + lock: { + createdAt: new Date("2026-08-29T11:00:00.000Z"), + definitionId, + ownerRunId: runId, + token: toMigrationDefinitionLockToken("lock-remote"), + }, + warnings: [], + }, +} satisfies MigrateDashboardRow; const runTerminalSummary = ( status: "succeeded" | "failed" = "succeeded" @@ -143,14 +169,26 @@ 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"), + }), + getMessages: () => Effect.succeed([]), observeRun: () => Stream.die("Unexpected run observation"), prepareOperation: () => Effect.die("Unexpected operation preparation"), startOperation: () => Effect.die("Unexpected operation start"), + scanSource: () => Effect.die("Unexpected source scan"), stopRun: (requestedRunId) => Effect.succeed({ kind: "requested" as const, @@ -160,6 +198,38 @@ const makeConnection = ( ...overrides, }); +const makeServerBackedCommandConnection = (): MigrationCliServerConnection => + makeConnection({ + breakLock: (lock) => + Effect.succeed({ definitionId: lock.definitionId, kind: "cleared" }), + getDashboard: Effect.succeed({ + dashboard: { + activeRuns: [], + groups: [], + rows: [remoteDashboardRow], + scannedSource: false, + }, + resumeToken: MigrateDashboardResumeToken.make("dashboard-remote"), + }), + getMessages: () => + Effect.succeed([ + { + definitionId, + kind: "skip-reason", + message: "Already migrated remotely", + runId, + severity: "info", + sourceIdentity: toEncodedSourceIdentity("article-1"), + updatedAt: new Date("2026-08-29T11:30:00.000Z"), + }, + ]), + prepareOperation: (request) => + Effect.succeed({ + ...preparedOperation(request), + planRows: [remoteDashboardRow], + }), + }); + const makeLayer = (runtime: MigrationCliRuntimeShape) => Layer.mergeAll( CliOutput.layer(CliOutput.defaultFormatter({ colors: false })), @@ -197,6 +267,122 @@ const interruptRuntime = ( }); describe("migrate runs", () => { + it.effect( + "lists remote Migration Definitions through the shared connection", + () => + Effect.gen(function* () { + const result = yield* runCli( + ["list", "--server", "https://migrate.example/api/migrate"], + { + connectMigrateServer: () => Effect.succeed(makeConnection()), + cwd: "/workspace", + } + ); + + expect(result.exitCode).toBe(0); + expect(result.stderr).toBe(""); + expect(result.stdout).toContain("Migration Definitions"); + }) + ); + + it.effect( + "routes every server-backed read and control command remotely", + () => + Effect.gen(function* () { + const runtime = { + connectMigrateServer: () => + Effect.succeed(makeServerBackedCommandConnection()), + cwd: "/workspace", + } satisfies MigrationCliRuntimeShape; + const server = "https://migrate.example/api/migrate"; + const graph = yield* runCli( + ["graph", "articles", "--server", server], + runtime + ); + const status = yield* runCli( + ["status", "articles", "--server", server], + runtime + ); + const messages = yield* runCli( + ["messages", "articles", "--server", server], + runtime + ); + const unlock = yield* runCli( + ["unlock", "articles", "--server", server], + runtime + ); + + expect(graph.exitCode).toBe(0); + expect(graph.stdout).toContain("Migration Dependency Graph: articles"); + expect(status.exitCode).toBe(0); + expect(status.stdout).toContain("Migration Status"); + expect(status.stdout).toContain("articles"); + expect(status.stdout).toContain("3"); + 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("scans remote status through the Migrate Protocol", () => + Effect.gen(function* () { + let scanInput: + | Parameters[0] + | undefined; + const scannedRow = { + ...remoteDashboardRow, + status: { + ...remoteDashboardRow.status, + source: { + duplicate: 0, + invalid: 0, + orphaned: 0, + total: 5, + unprocessed: 2, + }, + }, + } satisfies MigrateDashboardRow; + const connection = makeServerBackedCommandConnection(); + const result = yield* runCli( + [ + "status", + "articles", + "--scan-source", + "--concurrency", + "2", + "--server", + "https://migrate.example/api/migrate", + ], + { + connectMigrateServer: () => + Effect.succeed({ + ...connection, + scanSource: (input) => { + scanInput = input; + return Effect.succeed({ + activeRuns: [], + groups: [], + rows: [scannedRow], + scannedSource: true, + }); + }, + }), + cwd: "/workspace", + } + ); + + expect(result.exitCode).toBe(0); + expect(result.stdout).toContain("source inventory"); + expect(result.stdout).toContain("Unprocessed"); + expect(scanInput).toEqual({ + concurrency: 2, + target: { definitionId, kind: "migration" }, + }); + }) + ); + it.effect("renders typed remote planning errors with their identifiers", () => Effect.gen(function* () { const cases = [ diff --git a/packages/migrate-sdk/src/cli/runtime.ts b/packages/migrate-sdk/src/cli/runtime.ts index eab6b4b..65572cf 100644 --- a/packages/migrate-sdk/src/cli/runtime.ts +++ b/packages/migrate-sdk/src/cli/runtime.ts @@ -17,13 +17,19 @@ import { type MigrateServerConnectionInput, } from "../client/node/index.ts"; import type { MigrationRunId } from "../domain/ids.ts"; +import type { MigrationDefinitionLock } from "../domain/lock.ts"; +import type { MigrationMessage } from "../domain/message.ts"; import type { MigrateActiveRun, + MigrateBreakLockResult, + MigrateDashboard, + MigrateDashboardSnapshot, MigrateObservationEvent, MigrateOperationRequest, MigratePreparedOperation, MigrateRunStartResult, MigrateRunStopResult, + MigrateTarget, } from "../protocol/index.ts"; import type { SqlMigrationStoreSchemaPlan } from "../stores/sql/sql-migration-store-schema.ts"; import { @@ -47,14 +53,25 @@ export class MigrationCliConnectionError extends Schema.TaggedError Effect.Effect; readonly dispose: () => Promise; readonly getActiveRuns: Effect.Effect; + readonly getDashboard: Effect.Effect; + readonly getMessages: ( + target: MigrateTarget + ) => Effect.Effect; readonly observeRun: ( runId: MigrationRunId ) => Stream.Stream; readonly prepareOperation: ( request: MigrateOperationRequest ) => Effect.Effect; + readonly scanSource: (input: { + readonly concurrency?: number; + readonly target: MigrateTarget; + }) => Effect.Effect; readonly startOperation: (input: { readonly acceptedFingerprint: MigratePreparedOperation["fingerprint"]; readonly request: MigrateOperationRequest; @@ -167,14 +184,19 @@ 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(), + getMessages: (target) => + connection.client.GetMessages({ target }), observeRun: (runId: MigrationRunId) => connection.client.observeRun({ runId }), prepareOperation: (request) => connection.client.PrepareOperation(request), startOperation: (input) => connection.client.StartOperation(input), + scanSource: (input) => connection.client.ScanSource(input), stopRun: (runId: MigrationRunId) => connection.client.StopRun({ runId }), })) 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/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..b36cb8b 100644 --- a/packages/migrate-sdk/src/client/node/remote-connection.test.ts +++ b/packages/migrate-sdk/src/client/node/remote-connection.test.ts @@ -178,12 +178,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, From 747888da4af47beadc0d5dea28b551480f64594c Mon Sep 17 00:00:00 2001 From: Alexei Gorobet Date: Tue, 1 Sep 2026 17:06:59 -0400 Subject: [PATCH 2/2] fix(cli): close remote command gaps --- .changeset/remote-cli-commands.md | 4 +- ...r-boundary-for-local-and-remote-clients.md | 21 +- packages/migrate-sdk/src/cli/command.ts | 737 +++++------------- packages/migrate-sdk/src/cli/inspection.ts | 111 +++ .../src/cli/remote-commands.test.ts | 336 ++++++++ .../migrate-sdk/src/cli/runs-command.test.ts | 196 +---- packages/migrate-sdk/src/cli/runtime.ts | 54 +- .../src/client/internal/client-service.ts | 3 + .../src/client/node/remote-connection.test.ts | 31 + .../migrate-sdk/src/protocol/index.test.ts | 140 ++++ packages/migrate-sdk/src/protocol/index.ts | 85 +- .../src/protocol/registry-selection.ts | 25 + .../migrate-sdk/src/server/handlers.test.ts | 57 ++ packages/migrate-sdk/src/server/handlers.ts | 3 + .../src/server/registry-backend.ts | 7 + .../src/server/registry-runtime.test.ts | 206 +++++ .../src/server/registry-runtime.ts | 27 + .../migrate-sdk/src/server/service.test.ts | 80 ++ packages/migrate-sdk/src/server/service.ts | 24 + .../test/fixtures/remote-server.ts | 46 ++ 20 files changed, 1430 insertions(+), 763 deletions(-) create mode 100644 packages/migrate-sdk/src/cli/inspection.ts create mode 100644 packages/migrate-sdk/src/cli/remote-commands.test.ts create mode 100644 packages/migrate-sdk/src/protocol/registry-selection.ts diff --git a/.changeset/remote-cli-commands.md b/.changeset/remote-cli-commands.md index 053184e..79bfc41 100644 --- a/.changeset/remote-cli-commands.md +++ b/.changeset/remote-cli-commands.md @@ -5,4 +5,6 @@ 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. +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/docs/adr/0007-server-boundary-for-local-and-remote-clients.md b/docs/adr/0007-server-boundary-for-local-and-remote-clients.md index b9aaa8c..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 @@ -280,6 +287,14 @@ lifecycle management: - `migrate runs stop ` explicitly requests durable cooperative cancellation. +`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 @@ -322,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 be085ae..975dd98 100644 --- a/packages/migrate-sdk/src/cli/command.ts +++ b/packages/migrate-sdk/src/cli/command.ts @@ -11,34 +11,26 @@ 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, - MigrationDefinitionRegistryMessagesReport, - MigrationDefinitionRegistryPlanningError, - MigrationDefinitionRegistrySelectionInput, - MigrationDefinitionRegistryStatusError, - MigrationDefinitionRegistryStatusInput, - MigrationDefinitionRegistryStatusReport, +import { + MigrationDefinitionRegistryInvalidSelectionError, + MigrationDefinitionRegistryMissingExplicitRequiredDependenciesError, + type MigrationDefinitionRegistryPlanningError, + MigrationDefinitionRegistryUnknownDefinitionError, + MigrationDefinitionRegistryUnknownGroupError, } from "../domain/registry.ts"; import type { MigrateAction, - MigrateDashboardRow, MigrateOperationRequest, MigratePreparedOperation, MigrateSelection, - MigrateTarget, } 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"; @@ -46,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, @@ -63,8 +60,6 @@ import { renderPreparedOperationWarnings, renderRegistryEntriesGraph, renderRegistryEntriesList, - renderRegistryGraph, - renderRegistryList, renderRunStopResult, renderRuntimeError, renderSqlMigrationStoreSchemaPlan, @@ -191,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; @@ -333,7 +323,7 @@ const withCliMigrateConnection = ( releaseCliMigrateConnection ); -const reportCliConnectionErrors = ( +const reportCliRegistryCommandErrors = ( effect: Effect.Effect ): Effect.Effect => effect.pipe( @@ -344,6 +334,26 @@ const reportCliConnectionErrors = ( ) ); +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( @@ -501,45 +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 root = yield* migrateBaseCommand; - - if (Option.isSome(root.server)) { - return yield* reportCliConnectionErrors( - withCliMigrateConnection(({ connection }) => - Effect.gen(function* () { - const snapshot = yield* connection.getDashboard; - yield* Console.log( - renderRegistryEntriesList( - snapshot.dashboard.rows.map((row) => row.entry), - { colors: yield* useColor } - ) - ); - }) - ) + withCliRegistryOperations((operations) => + Effect.gen(function* () { + const registry = yield* operations.getRegistry; + yield* Console.log( + renderRegistryEntriesList(registry.entries, { + colors: yield* useColor, + }) ); - } - - const registry = yield* loadConfiguredRegistry; - - yield* Console.log( - renderRegistryList(registry, { colors: yield* useColor }) - ); - }) + }) + ) ).pipe(Command.withDescription("List registered Migration Definitions")); const graphDefinition = Argument.string("definition").pipe(Argument.optional); @@ -548,66 +530,31 @@ const graphCommand = Command.make( "graph", { definition: graphDefinition }, ({ definition }) => - Effect.gen(function* () { - const focusedDefinitionId = Option.getOrUndefined(definition); - const root = yield* migrateBaseCommand; - - if (Option.isSome(root.server)) { - return yield* reportCliConnectionErrors( - withCliMigrateConnection(({ connection }) => - Effect.gen(function* () { - const snapshot = yield* connection.getDashboard; - const entries = snapshot.dashboard.rows.map((row) => row.entry); - - if ( - focusedDefinitionId !== undefined && - !entries.some((entry) => entry.id === focusedDefinitionId) - ) { - return yield* failReportedCliMessage( - `Migration Definition was not found in the registry: ${focusedDefinitionId}` - ); - } - - yield* Console.log( - renderRegistryEntriesGraph( - entries, - focusedDefinitionId === undefined - ? undefined - : toMigrationDefinitionId(focusedDefinitionId), - { colors: yield* useColor } - ) - ); - }) - ) - ); - } - - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - - 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( @@ -746,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) + @@ -788,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" ); } @@ -813,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" ); } @@ -826,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; @@ -874,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; @@ -895,140 +798,7 @@ const renderMessagesCommandError = ( ...(input.group === undefined ? {} : { group: input.group }), hasTarget: false, }) - : renderRuntimeError(error); - -interface RemoteReadSelection { - readonly operation: MigratePreparedOperation; - readonly selection: MigrateSelection; -} - -function prepareRemoteReadSelection( - connection: MigrationCliServerConnection, - input: CliRegistrySelectionInput, - renderError: (error: unknown) => string -): Effect.Effect { - return Effect.gen(function* () { - const selection = yield* makeMigrateSelection(input); - const operation = yield* connection - .prepareOperation({ - action: "run", - options: { withDependencies: input.withDependencies }, - selection, - }) - .pipe( - Effect.catch((error) => failReportedCliMessage(renderError(error))) - ); - - return { operation, selection }; - }); -} - -function remoteStatusScanTargets( - selection: MigrateSelection, - includedDefinitionIds: readonly MigrationDefinitionId[] -): readonly MigrateTarget[] { - switch (selection.kind) { - case "group": - return [{ groupId: selection.groupId, kind: "group" }]; - case "definitions": - return selection.definitionIds.map((definitionId) => ({ - definitionId, - kind: "migration", - })); - case "all": - return includedDefinitionIds.map((definitionId) => ({ - definitionId, - kind: "migration", - })); - default: { - const unhandled: never = selection; - return unhandled; - } - } -} - -function collectRemoteStatusRows( - connection: MigrationCliServerConnection, - selection: MigrateSelection, - operation: MigratePreparedOperation, - options: { readonly concurrency?: number; readonly scanSource: boolean } -): Effect.Effect< - ReadonlyMap, - unknown -> { - const rowsById = new Map( - operation.planRows.map((row) => [row.entry.id, row]) - ); - - if (!options.scanSource) { - return Effect.succeed(rowsById); - } - - return Effect.gen(function* () { - const includedIds = new Set(operation.plan.includedDefinitionIds); - const targets = remoteStatusScanTargets( - selection, - operation.plan.includedDefinitionIds - ); - - for (const target of targets) { - const dashboard = yield* connection.scanSource({ - ...(options.concurrency === undefined - ? {} - : { concurrency: options.concurrency }), - target, - }); - - for (const row of dashboard.rows) { - if (includedIds.has(row.entry.id) && row.status?.source !== undefined) { - rowsById.set(row.entry.id, row); - } - } - } - - return rowsById; - }); -} - -function makeRemoteStatusReport( - operation: MigratePreparedOperation, - rowsById: ReadonlyMap, - scanSource: boolean -): MigrationDefinitionRegistryStatusReport { - const definitions = operation.plan.includedDefinitionIds.flatMap( - (definitionId) => { - const status = rowsById.get(definitionId)?.status; - return status === undefined ? [] : [status]; - } - ); - - return { - definitions, - includedDefinitionIds: operation.plan.includedDefinitionIds, - notices: operation.plan.notices, - ...(operation.plan.requestedGroup === undefined - ? {} - : { requestedGroup: operation.plan.requestedGroup }), - requestedDefinitionIds: operation.plan.requestedDefinitionIds, - scanSource, - warnings: definitions.flatMap((definition) => definition.warnings), - }; -} - -function makeRemoteMessagesReport( - operation: MigratePreparedOperation, - messages: readonly MigrationMessage[] -): MigrationDefinitionRegistryMessagesReport { - return { - includedDefinitionIds: operation.plan.includedDefinitionIds, - messages, - notices: operation.plan.notices, - ...(operation.plan.requestedGroup === undefined - ? {} - : { requestedGroup: operation.plan.requestedGroup }), - requestedDefinitionIds: operation.plan.requestedDefinitionIds, - }; -} + : renderStoredFailure(error); const statusCommand = Command.make( "status", @@ -1041,86 +811,46 @@ const statusCommand = Command.make( withDependencies, }, (input) => - Effect.gen(function* () { - const concurrencyInput = Option.getOrUndefined(input.concurrency); - const groupInput = Option.getOrUndefined(input.group); - const root = yield* migrateBaseCommand; - - if (Option.isSome(root.server)) { - return yield* reportCliConnectionErrors( - withCliMigrateConnection(({ connection }) => - Effect.gen(function* () { - const { operation, selection } = - yield* prepareRemoteReadSelection( - connection, - { - all: input.all, - definitionIds: input.definitions, - ...(groupInput === undefined ? {} : { group: groupInput }), - withDependencies: input.withDependencies, - }, - (error) => - isPlanningError(error) - ? renderStatusCommandError(error, { - definitionIds: input.definitions, - ...(groupInput === undefined - ? {} - : { group: groupInput }), - }) - : renderStoredFailure(error) - ); - const rowsById = yield* collectRemoteStatusRows( - connection, - selection, - operation, - { - ...(concurrencyInput === undefined - ? {} - : { concurrency: concurrencyInput }), - scanSource: input.scanSource, - } - ); - const report = makeRemoteStatusReport( - operation, - rowsById, - input.scanSource - ); - - yield* Console.log( - renderStatusReport(report, { colors: yield* useColor }) - ); - }) - ) + 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 loadedConfig = yield* loadConfiguredConfig; - const registry = toCliExecutableRegistry(loadedConfig.registry); - 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 }), - }) - ) - ) - ); + 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( @@ -1133,92 +863,46 @@ const messagesCommand = Command.make( withDependencies, }, (input) => - Effect.gen(function* () { - const groupInput = Option.getOrUndefined(input.group); - const root = yield* migrateBaseCommand; - - if (Option.isSome(root.server)) { - return yield* reportCliConnectionErrors( - withCliMigrateConnection(({ connection }) => - Effect.gen(function* () { - const { operation } = yield* prepareRemoteReadSelection( - connection, - { - all: input.all, - definitionIds: input.definitions, - ...(groupInput === undefined ? {} : { group: groupInput }), - withDependencies: input.withDependencies, - }, - (error) => - isPlanningError(error) - ? renderMessagesCommandError(error, { - definitionIds: input.definitions, - ...(groupInput === undefined - ? {} - : { group: groupInput }), - }) - : renderStoredFailure(error) - ); - const messages = (yield* Effect.forEach( - operation.plan.includedDefinitionIds, - (definitionId) => - connection.getMessages({ - definitionId, - kind: "migration", - }) - )).flat(); - const report = makeRemoteMessagesReport(operation, messages); - - 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 }) - ); - }) - ) - ); - } - - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - 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( @@ -1231,95 +915,34 @@ const unlockCommand = Command.make( "unlock", { definition: unlockDefinition }, ({ definition }) => - Effect.gen(function* () { - const definitionId = toMigrationDefinitionId(definition); - const root = yield* migrateBaseCommand; - - if (Option.isSome(root.server)) { - return yield* reportCliConnectionErrors( - withCliMigrateConnection(({ connection }) => - Effect.gen(function* () { - const snapshot = yield* connection.getDashboard; - const row = snapshot.dashboard.rows.find( - (candidate) => candidate.entry.id === definitionId - ); - - if (row === undefined) { - return yield* failReportedCliMessage( - `Migration Definition was not found in the registry: ${definitionId}` - ); - } - - const lock = row.status?.lock; - - if (lock === undefined || lock === null) { - yield* Console.log( - `Migration Definition lock is already clear: ${definitionId}` - ); - return; - } - - const result = yield* connection.breakLock(lock); - - if (result.kind === "already-clear") { - yield* Console.log( - `Migration Definition lock is already clear: ${definitionId}` - ); - return; - } - - yield* Console.log( - [ - "Migration Definition lock cleared", - `Migration ID ${definitionId}`, - `Owner Run ID ${lock.ownerRunId}`, - `Token ${lock.token}`, - ].join("\n") - ); - }) - ) - ); - } - - const loadedConfig = yield* loadConfiguredConfig; - const registry = loadedConfig.registry; - 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/runs-command.test.ts b/packages/migrate-sdk/src/cli/runs-command.test.ts index b558760..dd54793 100644 --- a/packages/migrate-sdk/src/cli/runs-command.test.ts +++ b/packages/migrate-sdk/src/cli/runs-command.test.ts @@ -14,12 +14,14 @@ 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 { - toEncodedSourceIdentity, toMigrationDefinitionGroupId, toMigrationDefinitionId, - toMigrationDefinitionLockToken, toMigrationRunId, } from "../domain/ids.ts"; import { @@ -28,7 +30,6 @@ import { } from "../domain/registry.ts"; import type { MigrateActiveRun, - MigrateDashboardRow, MigrateObservationEvent, MigrateOperationRequest, MigratePreparedOperation, @@ -49,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" }, @@ -58,26 +63,6 @@ const activeRun: MigrateActiveRun = { status: "running", stopSupported: true, }; -const remoteDashboardRow = { - entry: { - dependencies: { optional: [], required: [] }, - hasRollback: true, - id: definitionId, - }, - status: { - definitionId, - discovery: "full", - durable: { failed: 0, migrated: 3, needsUpdate: 0, skipped: 1 }, - lastRun: null, - lock: { - createdAt: new Date("2026-08-29T11:00:00.000Z"), - definitionId, - ownerRunId: runId, - token: toMigrationDefinitionLockToken("lock-remote"), - }, - warnings: [], - }, -} satisfies MigrateDashboardRow; const runTerminalSummary = ( status: "succeeded" | "failed" = "succeeded" @@ -184,11 +169,12 @@ const makeConnection = ( }, resumeToken: MigrateDashboardResumeToken.make("dashboard-empty"), }), - getMessages: () => Effect.succeed([]), + 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"), - scanSource: () => Effect.die("Unexpected source scan"), stopRun: (requestedRunId) => Effect.succeed({ kind: "requested" as const, @@ -198,38 +184,6 @@ const makeConnection = ( ...overrides, }); -const makeServerBackedCommandConnection = (): MigrationCliServerConnection => - makeConnection({ - breakLock: (lock) => - Effect.succeed({ definitionId: lock.definitionId, kind: "cleared" }), - getDashboard: Effect.succeed({ - dashboard: { - activeRuns: [], - groups: [], - rows: [remoteDashboardRow], - scannedSource: false, - }, - resumeToken: MigrateDashboardResumeToken.make("dashboard-remote"), - }), - getMessages: () => - Effect.succeed([ - { - definitionId, - kind: "skip-reason", - message: "Already migrated remotely", - runId, - severity: "info", - sourceIdentity: toEncodedSourceIdentity("article-1"), - updatedAt: new Date("2026-08-29T11:30:00.000Z"), - }, - ]), - prepareOperation: (request) => - Effect.succeed({ - ...preparedOperation(request), - planRows: [remoteDashboardRow], - }), - }); - const makeLayer = (runtime: MigrationCliRuntimeShape) => Layer.mergeAll( CliOutput.layer(CliOutput.defaultFormatter({ colors: false })), @@ -267,122 +221,6 @@ const interruptRuntime = ( }); describe("migrate runs", () => { - it.effect( - "lists remote Migration Definitions through the shared connection", - () => - Effect.gen(function* () { - const result = yield* runCli( - ["list", "--server", "https://migrate.example/api/migrate"], - { - connectMigrateServer: () => Effect.succeed(makeConnection()), - cwd: "/workspace", - } - ); - - expect(result.exitCode).toBe(0); - expect(result.stderr).toBe(""); - expect(result.stdout).toContain("Migration Definitions"); - }) - ); - - it.effect( - "routes every server-backed read and control command remotely", - () => - Effect.gen(function* () { - const runtime = { - connectMigrateServer: () => - Effect.succeed(makeServerBackedCommandConnection()), - cwd: "/workspace", - } satisfies MigrationCliRuntimeShape; - const server = "https://migrate.example/api/migrate"; - const graph = yield* runCli( - ["graph", "articles", "--server", server], - runtime - ); - const status = yield* runCli( - ["status", "articles", "--server", server], - runtime - ); - const messages = yield* runCli( - ["messages", "articles", "--server", server], - runtime - ); - const unlock = yield* runCli( - ["unlock", "articles", "--server", server], - runtime - ); - - expect(graph.exitCode).toBe(0); - expect(graph.stdout).toContain("Migration Dependency Graph: articles"); - expect(status.exitCode).toBe(0); - expect(status.stdout).toContain("Migration Status"); - expect(status.stdout).toContain("articles"); - expect(status.stdout).toContain("3"); - 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("scans remote status through the Migrate Protocol", () => - Effect.gen(function* () { - let scanInput: - | Parameters[0] - | undefined; - const scannedRow = { - ...remoteDashboardRow, - status: { - ...remoteDashboardRow.status, - source: { - duplicate: 0, - invalid: 0, - orphaned: 0, - total: 5, - unprocessed: 2, - }, - }, - } satisfies MigrateDashboardRow; - const connection = makeServerBackedCommandConnection(); - const result = yield* runCli( - [ - "status", - "articles", - "--scan-source", - "--concurrency", - "2", - "--server", - "https://migrate.example/api/migrate", - ], - { - connectMigrateServer: () => - Effect.succeed({ - ...connection, - scanSource: (input) => { - scanInput = input; - return Effect.succeed({ - activeRuns: [], - groups: [], - rows: [scannedRow], - scannedSource: true, - }); - }, - }), - cwd: "/workspace", - } - ); - - expect(result.exitCode).toBe(0); - expect(result.stdout).toContain("source inventory"); - expect(result.stdout).toContain("Unprocessed"); - expect(scanInput).toEqual({ - concurrency: 2, - target: { definitionId, kind: "migration" }, - }); - }) - ); - it.effect("renders typed remote planning errors with their identifiers", () => Effect.gen(function* () { const cases = [ @@ -592,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], @@ -796,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], @@ -1084,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 65572cf..9baef2d 100644 --- a/packages/migrate-sdk/src/cli/runtime.ts +++ b/packages/migrate-sdk/src/cli/runtime.ts @@ -12,24 +12,28 @@ 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 { MigrationMessage } from "../domain/message.ts"; import type { MigrateActiveRun, MigrateBreakLockResult, - MigrateDashboard, MigrateDashboardSnapshot, MigrateObservationEvent, MigrateOperationRequest, MigratePreparedOperation, + MigrateProtocolError, + MigrateRegistry, + MigrateRegistryMessagesReport, + MigrateRegistryMessagesRequest, + MigrateRegistryStatusReport, + MigrateRegistryStatusRequest, MigrateRunStartResult, MigrateRunStopResult, - MigrateTarget, } from "../protocol/index.ts"; import type { SqlMigrationStoreSchemaPlan } from "../stores/sql/sql-migration-store-schema.ts"; import { @@ -52,33 +56,41 @@ export class MigrationCliConnectionError extends Schema.TaggedError Effect.Effect; + ) => Effect.Effect; readonly dispose: () => Promise; - readonly getActiveRuns: Effect.Effect; - readonly getDashboard: Effect.Effect; - readonly getMessages: ( - target: MigrateTarget - ) => 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; - readonly scanSource: (input: { - readonly concurrency?: number; - readonly target: MigrateTarget; - }) => 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 { @@ -188,15 +200,17 @@ export class MigrationCliRuntime extends Service< dispose: connection.dispose, getActiveRuns: connection.client.GetActiveRuns(), getDashboard: connection.client.GetDashboard(), - getMessages: (target) => - connection.client.GetMessages({ target }), + 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) => connection.client.PrepareOperation(request), startOperation: (input) => connection.client.StartOperation(input), - scanSource: (input) => connection.client.ScanSource(input), stopRun: (runId: MigrationRunId) => connection.client.StopRun({ runId }), })) 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/node/remote-connection.test.ts b/packages/migrate-sdk/src/client/node/remote-connection.test.ts index b36cb8b..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({}) 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) =>