diff --git a/backend/.sqlx/query-6ac00d65b3b7707cba9456171d4332c696e12b52066a5ae2ab10fe3b5bd0f3d0.json b/backend/.sqlx/query-6ac00d65b3b7707cba9456171d4332c696e12b52066a5ae2ab10fe3b5bd0f3d0.json new file mode 100644 index 0000000000..3a1dc9da28 --- /dev/null +++ b/backend/.sqlx/query-6ac00d65b3b7707cba9456171d4332c696e12b52066a5ae2ab10fe3b5bd0f3d0.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT name, code_down FROM datatable_migrations WHERE workspace_id = $1 AND datatable = $2 AND timestamp = $3", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "name", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "code_down", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Int8" + ] + }, + "nullable": [ + false, + true + ] + }, + "hash": "6ac00d65b3b7707cba9456171d4332c696e12b52066a5ae2ab10fe3b5bd0f3d0" +} diff --git a/backend/.sqlx/query-98f03bde861caf75a4a98457fe9dc1c25362de35ef086731095bac0e7e679a48.json b/backend/.sqlx/query-98f03bde861caf75a4a98457fe9dc1c25362de35ef086731095bac0e7e679a48.json new file mode 100644 index 0000000000..8cf1c0b2ee --- /dev/null +++ b/backend/.sqlx/query-98f03bde861caf75a4a98457fe9dc1c25362de35ef086731095bac0e7e679a48.json @@ -0,0 +1,19 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO datatable_migrations (workspace_id, datatable, timestamp, name, code_up, code_down) SELECT $1, * FROM UNNEST($2::varchar[], $3::bigint[], $4::varchar[], $5::text[], $6::text[])", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "VarcharArray", + "Int8Array", + "VarcharArray", + "TextArray", + "TextArray" + ] + }, + "nullable": [] + }, + "hash": "98f03bde861caf75a4a98457fe9dc1c25362de35ef086731095bac0e7e679a48" +} diff --git a/backend/.sqlx/query-a618f0dbb1a4913fca09bdb2b01544c760d4c302f1513ce8c6b904f43868c755.json b/backend/.sqlx/query-a618f0dbb1a4913fca09bdb2b01544c760d4c302f1513ce8c6b904f43868c755.json deleted file mode 100644 index 412e9032b6..0000000000 --- a/backend/.sqlx/query-a618f0dbb1a4913fca09bdb2b01544c760d4c302f1513ce8c6b904f43868c755.json +++ /dev/null @@ -1,18 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "INSERT INTO datatable_migrations (workspace_id, timestamp, name, code_up, code_down) SELECT $1, * FROM UNNEST($2::bigint[], $3::varchar[], $4::text[], $5::text[])", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Varchar", - "Int8Array", - "VarcharArray", - "TextArray", - "TextArray" - ] - }, - "nullable": [] - }, - "hash": "a618f0dbb1a4913fca09bdb2b01544c760d4c302f1513ce8c6b904f43868c755" -} diff --git a/backend/.sqlx/query-d9fbb752a49f80bfaf890f18295ba83238b59b306b741221e4c19bb40ca93712.json b/backend/.sqlx/query-a8cec061756c7be61af829f4a5207ba6333ad5feb4557a9341427e9e68608025.json similarity index 63% rename from backend/.sqlx/query-d9fbb752a49f80bfaf890f18295ba83238b59b306b741221e4c19bb40ca93712.json rename to backend/.sqlx/query-a8cec061756c7be61af829f4a5207ba6333ad5feb4557a9341427e9e68608025.json index 5f8e722f0b..c262f3d1f6 100644 --- a/backend/.sqlx/query-d9fbb752a49f80bfaf890f18295ba83238b59b306b741221e4c19bb40ca93712.json +++ b/backend/.sqlx/query-a8cec061756c7be61af829f4a5207ba6333ad5feb4557a9341427e9e68608025.json @@ -1,25 +1,30 @@ { "db_name": "PostgreSQL", - "query": "SELECT timestamp, name, code_up, code_down FROM datatable_migrations WHERE workspace_id = $1 ORDER BY timestamp ASC", + "query": "SELECT datatable, timestamp, name, code_up, code_down FROM datatable_migrations WHERE workspace_id = $1 ORDER BY datatable, timestamp ASC", "describe": { "columns": [ { "ordinal": 0, + "name": "datatable", + "type_info": "Varchar" + }, + { + "ordinal": 1, "name": "timestamp", "type_info": "Int8" }, { - "ordinal": 1, + "ordinal": 2, "name": "name", "type_info": "Varchar" }, { - "ordinal": 2, + "ordinal": 3, "name": "code_up", "type_info": "Text" }, { - "ordinal": 3, + "ordinal": 4, "name": "code_down", "type_info": "Text" } @@ -33,8 +38,9 @@ false, false, false, + false, true ] }, - "hash": "d9fbb752a49f80bfaf890f18295ba83238b59b306b741221e4c19bb40ca93712" + "hash": "a8cec061756c7be61af829f4a5207ba6333ad5feb4557a9341427e9e68608025" } diff --git a/backend/.sqlx/query-5103cb8f00056b1d8656d11672378aa93afc3d1a88ec8c0ed5093b5a57075dd3.json b/backend/.sqlx/query-fcb0cf55f5c9047066a6fdd62bf6f85238cbf189ccf18191493d6357e269d12d.json similarity index 76% rename from backend/.sqlx/query-5103cb8f00056b1d8656d11672378aa93afc3d1a88ec8c0ed5093b5a57075dd3.json rename to backend/.sqlx/query-fcb0cf55f5c9047066a6fdd62bf6f85238cbf189ccf18191493d6357e269d12d.json index 47a00b9032..c630f5dd05 100644 --- a/backend/.sqlx/query-5103cb8f00056b1d8656d11672378aa93afc3d1a88ec8c0ed5093b5a57075dd3.json +++ b/backend/.sqlx/query-fcb0cf55f5c9047066a6fdd62bf6f85238cbf189ccf18191493d6357e269d12d.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT timestamp, name, code_up FROM datatable_migrations WHERE workspace_id = $1 ORDER BY timestamp ASC", + "query": "SELECT timestamp, name, code_up FROM datatable_migrations WHERE workspace_id = $1 AND datatable = $2 ORDER BY timestamp ASC", "describe": { "columns": [ { @@ -21,6 +21,7 @@ ], "parameters": { "Left": [ + "Text", "Text" ] }, @@ -30,5 +31,5 @@ false ] }, - "hash": "5103cb8f00056b1d8656d11672378aa93afc3d1a88ec8c0ed5093b5a57075dd3" + "hash": "fcb0cf55f5c9047066a6fdd62bf6f85238cbf189ccf18191493d6357e269d12d" } diff --git a/backend/migrations/20260617081932_datatable_migrations.up.sql b/backend/migrations/20260617081932_datatable_migrations.up.sql index c1d9c98254..72c7dc5ec6 100644 --- a/backend/migrations/20260617081932_datatable_migrations.up.sql +++ b/backend/migrations/20260617081932_datatable_migrations.up.sql @@ -1,17 +1,21 @@ --- SQL migrations defined per workspace for data tables. --- `timestamp` is the migration version (YYYYMMDDHHMMSS), recorded as `version` --- in the data table's `_wm_migrations` table once applied. +-- SQL migrations defined per data table within a workspace. +-- `datatable` is the target data table name, `name` is the migration name +-- (e.g. add_index_to_customers), and `timestamp` is the migration version +-- (YYYYMMDDHHMMSS), recorded as `version` in the data table's `_wm_migrations` +-- table once applied. CREATE TABLE datatable_migrations ( workspace_id VARCHAR(50) NOT NULL REFERENCES workspace(id) ON DELETE CASCADE, + datatable VARCHAR(255) NOT NULL, timestamp BIGINT NOT NULL, name VARCHAR(255) NOT NULL, code_up TEXT NOT NULL, code_down TEXT, - PRIMARY KEY (workspace_id, timestamp) + PRIMARY KEY (workspace_id, datatable, timestamp) ); --- No standalone index on workspace_id: the (workspace_id, timestamp) primary-key --- btree already serves `WHERE workspace_id = $1` lookups via its leading column. +-- No standalone index: the (workspace_id, datatable, timestamp) primary-key btree +-- already serves both `WHERE workspace_id = $1` and `WHERE workspace_id = $1 AND +-- datatable = $2` lookups via its leading columns. GRANT ALL ON datatable_migrations TO windmill_user; GRANT ALL ON datatable_migrations TO windmill_admin; diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index f553038cf0..1125c9ad5f 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -130,6 +130,10 @@ pub fn workspaced_service() -> Router { "/run_datatable_migrations/{datatable_name}", post(run_datatable_migrations), ) + .route( + "/rollback_datatable_migrations/{datatable_name}", + post(rollback_datatable_migrations), + ) .route("/list_datatable_migrations", get(list_datatable_migrations)) .route( "/update_datatable_migrations", @@ -2416,8 +2420,9 @@ async fn run_datatable_migrations( let migrations = sqlx::query!( "SELECT timestamp, name, code_up FROM datatable_migrations \ - WHERE workspace_id = $1 ORDER BY timestamp ASC", - &w_id + WHERE workspace_id = $1 AND datatable = $2 ORDER BY timestamp ASC", + &w_id, + &datatable_name, ) .fetch_all(&db) .await?; @@ -2484,8 +2489,124 @@ async fn run_datatable_migrations( Ok(Json(RunDatatableMigrationsResult { applied })) } +#[derive(Serialize)] +struct RolledBackMigration { + version: i64, + name: String, +} + +#[derive(Serialize)] +struct RollbackDatatableMigrationsResult { + rolled_back: Vec, +} + +/// Roll back the most recently applied migration on a given data table (one +/// step): run its `code_down` and drop its `_wm_migrations` row, atomically. +async fn rollback_datatable_migrations( + authed: ApiAuthed, + Extension(db): Extension, + Path((w_id, datatable_name)): Path<(String, String)>, +) -> JsonResult { + require_admin(authed.is_admin, &authed.username)?; + + audit_log( + &db, + &authed, + "workspaces.rollback_datatable_migrations", + ActionKind::Update, + &w_id, + Some(datatable_name.as_str()), + None, + ) + .await?; + + // Connect to the data table's database + let db_resource = get_datatable_resource_from_db_unchecked(&db, &w_id, &datatable_name).await?; + let pg_db: PgDatabase = serde_json::from_value(db_resource) + .map_err(|e| Error::internal_err(format!("Failed to parse database credentials: {}", e)))?; + let (mut client, connection) = pg_db.connect(Some(&db)).await?; + tokio::spawn(async move { + if let Err(e) = connection.await { + tracing::error!("Datatable connection error: {}", e); + } + }); + + client + .batch_execute( + "CREATE TABLE IF NOT EXISTS _wm_migrations (\ + version BIGINT PRIMARY KEY, \ + installed_at TIMESTAMPTZ NOT NULL DEFAULT now())", + ) + .await + .map_err(|e| { + Error::internal_err(format!("Failed to ensure _wm_migrations table: {}", e)) + })?; + + let latest = client + .query_opt( + "SELECT version FROM _wm_migrations ORDER BY version DESC LIMIT 1", + &[], + ) + .await + .map_err(|e| Error::internal_err(format!("Failed to read _wm_migrations: {}", e)))?; + + let version: i64 = match latest { + Some(row) => row.get::<_, i64>(0), + None => { + return Ok(Json(RollbackDatatableMigrationsResult { + rolled_back: vec![], + })) + } + }; + + let definition = sqlx::query!( + "SELECT name, code_down FROM datatable_migrations \ + WHERE workspace_id = $1 AND datatable = $2 AND timestamp = $3", + &w_id, + &datatable_name, + version + ) + .fetch_optional(&db) + .await? + .ok_or_else(|| { + Error::BadRequest(format!( + "Cannot roll back migration {version}: its definition no longer exists" + )) + })?; + + let code_down = definition.code_down.ok_or_else(|| { + Error::BadRequest(format!( + "Cannot roll back migration {} ({}): it has no down migration", + version, definition.name + )) + })?; + + // Run the down migration and forget its version atomically. + let tx = client + .transaction() + .await + .map_err(|e| Error::internal_err(format!("Failed to start transaction: {}", e)))?; + tx.batch_execute(&code_down).await.map_err(|e| { + Error::internal_err(format!( + "Failed to roll back migration {} ({}): {}", + version, definition.name, e + )) + })?; + tx.execute("DELETE FROM _wm_migrations WHERE version = $1", &[&version]) + .await + .map_err(|e| Error::internal_err(format!("Failed to drop migration record: {}", e)))?; + tx.commit().await.map_err(|e| { + Error::internal_err(format!("Failed to commit rollback of {}: {}", version, e)) + })?; + + Ok(Json(RollbackDatatableMigrationsResult { + rolled_back: vec![RolledBackMigration { version, name: definition.name }], + })) +} + #[derive(Serialize, Deserialize)] pub struct DatatableMigration { + pub datatable: String, pub timestamp: i64, pub name: String, pub code_up: String, @@ -2500,8 +2621,8 @@ async fn list_datatable_migrations( ) -> JsonResult> { let migrations = sqlx::query_as!( DatatableMigration, - "SELECT timestamp, name, code_up, code_down FROM datatable_migrations \ - WHERE workspace_id = $1 ORDER BY timestamp ASC", + "SELECT datatable, timestamp, name, code_up, code_down FROM datatable_migrations \ + WHERE workspace_id = $1 ORDER BY datatable, timestamp ASC", &w_id ) .fetch_all(&db) @@ -2525,6 +2646,11 @@ async fn update_datatable_migrations( ) -> Result { require_admin(authed.is_admin, &authed.username)?; + let datatables: Vec = payload + .migrations + .iter() + .map(|m| m.datatable.clone()) + .collect(); let timestamps: Vec = payload.migrations.iter().map(|m| m.timestamp).collect(); let names: Vec = payload.migrations.iter().map(|m| m.name.clone()).collect(); let code_ups: Vec = payload @@ -2548,9 +2674,10 @@ async fn update_datatable_migrations( .await?; sqlx::query!( - "INSERT INTO datatable_migrations (workspace_id, timestamp, name, code_up, code_down) \ - SELECT $1, * FROM UNNEST($2::bigint[], $3::varchar[], $4::text[], $5::text[])", + "INSERT INTO datatable_migrations (workspace_id, datatable, timestamp, name, code_up, code_down) \ + SELECT $1, * FROM UNNEST($2::varchar[], $3::bigint[], $4::varchar[], $5::text[], $6::text[])", &w_id, + &datatables, ×tamps, &names, &code_ups, diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 9966715f22..fc23022483 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -4385,6 +4385,40 @@ paths: name: type: string + /w/{workspace}/workspaces/rollback_datatable_migrations/{datatable_name}: + post: + summary: roll back the most recently applied migration on a datatable + operationId: rollbackDatatableMigrations + tags: + - workspace + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - name: datatable_name + in: path + required: true + schema: + type: string + responses: + "200": + description: rolled back migrations + content: + application/json: + schema: + type: object + required: [rolled_back] + properties: + rolled_back: + type: array + items: + type: object + required: [version, name] + properties: + version: + type: integer + format: int64 + name: + type: string + /w/{workspace}/workspaces/list_datatable_migrations: get: summary: list datatable migrations for a workspace @@ -27761,8 +27795,10 @@ components: additionalProperties: true DatatableMigration: type: object - required: [timestamp, name, code_up] + required: [datatable, timestamp, name, code_up] properties: + datatable: + type: string timestamp: type: integer format: int64 diff --git a/cli/src/commands/datatable/datatable.ts b/cli/src/commands/datatable/datatable.ts index c86c971a90..a5cfe08bf3 100644 --- a/cli/src/commands/datatable/datatable.ts +++ b/cli/src/commands/datatable/datatable.ts @@ -9,6 +9,10 @@ import { GlobalOptions } from "../../types.ts"; import { runCatalogQuery } from "../../utils/catalog.ts"; import { psql as psqlDatatable } from "./psql.ts"; import { serve as serveDatatable } from "./serve.ts"; +import { + rollbackMigrations, + runMigrations, +} from "../datatable_migrations.ts"; const DEFAULT_DATATABLE_NAME = "main"; @@ -41,6 +45,27 @@ async function run( await runCatalogQuery(opts, "datatable", name, sql); } +async function migrateUp(opts: GlobalOptions, name?: string) { + const workspace = await resolveWorkspace(opts); + await requireLogin(opts); + await runMigrations(workspace.workspaceId, name ?? DEFAULT_DATATABLE_NAME); +} + +async function migrateDown(opts: GlobalOptions, name?: string) { + const workspace = await resolveWorkspace(opts); + await requireLogin(opts); + await rollbackMigrations(workspace.workspaceId, name ?? DEFAULT_DATATABLE_NAME); +} + +const migrateCommand = new Command() + .description("apply or roll back datatable migrations") + .command("up", "apply all pending migrations to a datatable") + .arguments("[name:string]") + .action(migrateUp as any) + .command("down", "roll back the most recent migration on a datatable") + .arguments("[name:string]") + .action(migrateDown as any); + async function serve( opts: GlobalOptions & { port?: number; host?: string; password?: string }, ) { @@ -69,6 +94,7 @@ const command = new Command() "Output only the final result as JSON. Useful for scripting.", ) .action(run as any) + .command("migrate", migrateCommand) .command( "serve", "Serve all datatables as a Postgres-wire endpoint (psql, DBeaver, pgAdmin); the client picks the datatable via the database name in its connection string", diff --git a/cli/src/commands/datatable_migrations.ts b/cli/src/commands/datatable_migrations.ts index eb963eeec2..966f63fa11 100644 --- a/cli/src/commands/datatable_migrations.ts +++ b/cli/src/commands/datatable_migrations.ts @@ -5,9 +5,63 @@ import { colors } from "@cliffy/ansi/colors"; import * as wmill from "../../gen/services.gen.ts"; import type { DatatableMigration } from "../../gen/types.gen.ts"; import { readTextFile, readTextFileSync } from "../utils/utils.ts"; +import { Confirm } from "@cliffy/prompt/confirm"; +// Migrations live under /datatable_migrations//, one folder per +// target data table, as `_.up.sql` (and optional `.down.sql`). const MIGRATIONS_DIR = "datatable_migrations"; +/** + * Apply the workspace's pending migrations to a data table (forwards migrations + * recorded in `_wm_migrations`). Mirrors `wmill datatable migrate up`. + */ +export async function runMigrations( + workspace: string, + datatableName: string, +): Promise { + const result = await wmill.runDatatableMigrations({ + workspace, + datatableName, + }); + const applied = result.applied ?? []; + if (applied.length === 0) { + log.info(colors.gray(`No pending migrations to run on '${datatableName}'`)); + return; + } + log.info( + colors.green(`Applied ${applied.length} migration(s) to '${datatableName}':`), + ); + for (const m of applied) { + log.info(colors.gray(` ${m.version} ${m.name}`)); + } +} + +/** + * Roll back the most recently applied migration on a data table (one step). + * Mirrors `wmill datatable migrate down`. + */ +export async function rollbackMigrations( + workspace: string, + datatableName: string, +): Promise { + const result = await wmill.rollbackDatatableMigrations({ + workspace, + datatableName, + }); + const rolledBack = result.rolled_back ?? []; + if (rolledBack.length === 0) { + log.info( + colors.gray(`No applied migrations to roll back on '${datatableName}'`), + ); + return; + } + for (const m of rolledBack) { + log.info( + colors.green(`Rolled back migration ${m.version} ${m.name} on '${datatableName}'`), + ); + } +} + function upFileName(m: { timestamp: number; name: string }): string { return `${m.timestamp}_${m.name}.up.sql`; } @@ -28,6 +82,13 @@ function parseMigrationFileName( return { timestamp: Number(m[1]), name: m[2], kind: m[3] as "up" | "down" }; } +function sortMigrations(migrations: DatatableMigration[]): DatatableMigration[] { + return migrations.sort( + (a, b) => + a.datatable.localeCompare(b.datatable) || a.timestamp - b.timestamp, + ); +} + /** * Push the local /datatable_migrations/ folder to the workspace, replacing * the remote set. Returns true if a push was performed, false if the folder is @@ -35,47 +96,55 @@ function parseMigrationFileName( */ export async function pushDatatableMigrations( workspace: string, + opts?: { yes?: boolean; jsonOutput?: boolean }, ): Promise { const localDir = path.join(process.cwd(), MIGRATIONS_DIR); if (!fs.existsSync(localDir)) { return false; } - // Group .up.sql / .down.sql files by their (timestamp, name) key. - const byKey = new Map< - string, - { timestamp: number; name: string; code_up?: string; code_down?: string } - >(); - for (const file of fs.readdirSync(localDir)) { - const parsed = parseMigrationFileName(file); - if (!parsed) continue; - const key = `${parsed.timestamp}_${parsed.name}`; - const entry = - byKey.get(key) ?? { timestamp: parsed.timestamp, name: parsed.name }; - const content = await readTextFile(path.join(localDir, file)); - if (parsed.kind === "up") entry.code_up = content; - else entry.code_down = content; - byKey.set(key, entry); - } - + // Each immediate subfolder is a target data table; read its migration files. const migrations: DatatableMigration[] = []; - for (const entry of byKey.values()) { - if (entry.code_up === undefined) { - log.warn( - colors.yellow( - `Skipping migration ${entry.timestamp}_${entry.name}: missing .up.sql file`, - ), - ); - continue; + for (const datatable of fs.readdirSync(localDir)) { + const dtDir = path.join(localDir, datatable); + if (!fs.statSync(dtDir).isDirectory()) continue; + + // Group .up.sql / .down.sql files by their (timestamp, name) key. + const byKey = new Map< + string, + { timestamp: number; name: string; code_up?: string; code_down?: string } + >(); + for (const file of fs.readdirSync(dtDir)) { + const parsed = parseMigrationFileName(file); + if (!parsed) continue; + const key = `${parsed.timestamp}_${parsed.name}`; + const entry = + byKey.get(key) ?? { timestamp: parsed.timestamp, name: parsed.name }; + const content = await readTextFile(path.join(dtDir, file)); + if (parsed.kind === "up") entry.code_up = content; + else entry.code_down = content; + byKey.set(key, entry); + } + + for (const entry of byKey.values()) { + if (entry.code_up === undefined) { + log.warn( + colors.yellow( + `Skipping ${datatable}/${entry.timestamp}_${entry.name}: missing .up.sql file`, + ), + ); + continue; + } + migrations.push({ + datatable, + timestamp: entry.timestamp, + name: entry.name, + code_up: entry.code_up, + ...(entry.code_down !== undefined ? { code_down: entry.code_down } : {}), + }); } - migrations.push({ - timestamp: entry.timestamp, - name: entry.name, - code_up: entry.code_up, - ...(entry.code_down !== undefined ? { code_down: entry.code_down } : {}), - }); } - migrations.sort((a, b) => a.timestamp - b.timestamp); + sortMigrations(migrations); // Skip if the remote set is already identical. let remote: DatatableMigration[] = []; @@ -94,17 +163,62 @@ export async function pushDatatableMigrations( requestBody: { migrations }, }); log.info( - colors.green( - `Pushed ${migrations.length} datatable migration(s)`, - ), + colors.green(`Pushed ${migrations.length} datatable migration(s)`), ); + + // Migrations newly introduced by this push (absent from the remote set before), + // keyed by (datatable, timestamp) since timestamps are unique only per table. + const remoteKeys = new Set(remote.map((m) => `${m.datatable}\0${m.timestamp}`)); + const newMigrations = migrations.filter( + (m) => !remoteKeys.has(`${m.datatable}\0${m.timestamp}`), + ); + if (newMigrations.length > 0) { + await offerToRunNewMigrations(workspace, newMigrations, opts); + } + return true; } /** - * Pull the workspace's datatable migrations into /datatable_migrations/ as - * `_.up.sql` (and `.down.sql` when a down migration exists). - * Files no longer present remotely are removed locally. + * After a push that introduced new migrations, list them and (interactively) + * offer to run them, equivalent to `wmill datatable migrate up` on each affected + * data table. + */ +async function offerToRunNewMigrations( + workspace: string, + newMigrations: DatatableMigration[], + opts?: { yes?: boolean; jsonOutput?: boolean }, +): Promise { + log.info(colors.green("New migrations were pushed:")); + for (const m of newMigrations) { + log.info(colors.gray(` ${m.datatable}: ${m.timestamp} ${m.name}`)); + } + + // Running migrations mutates the data tables, so skip the prompt in + // non-interactive contexts (--yes, --json, no TTY). + const interactive = !opts?.jsonOutput && !opts?.yes && !!process.stdin.isTTY; + if (!interactive) { + return; + } + + const shouldRun = await Confirm.prompt({ + message: "New migrations were pushed, run them?", + default: false, + }); + if (!shouldRun) { + return; + } + + for (const datatable of new Set(newMigrations.map((m) => m.datatable))) { + await runMigrations(workspace, datatable); + } +} + +/** + * Pull the workspace's datatable migrations into + * /datatable_migrations// as `_.up.sql` (and + * `.down.sql` when a down migration exists). Files no longer present remotely + * are removed locally. */ export async function pullDatatableMigrations( workspace: string, @@ -121,35 +235,41 @@ export async function pullDatatableMigrations( if (migrations.length === 0 && !fs.existsSync(localDir)) { return false; } - fs.mkdirSync(localDir, { recursive: true }); + // Relative paths (`/`) that should exist after the pull. const known = new Set(); for (const m of migrations) { + const dtDir = path.join(localDir, m.datatable); + fs.mkdirSync(dtDir, { recursive: true }); const up = upFileName(m); - known.add(up); - writeIfChanged(path.join(localDir, up), m.code_up); + known.add(`${m.datatable}/${up}`); + writeIfChanged(path.join(dtDir, up), m.code_up); if (m.code_down !== undefined && m.code_down !== null) { const down = downFileName(m); - known.add(down); - writeIfChanged(path.join(localDir, down), m.code_down); + known.add(`${m.datatable}/${down}`); + writeIfChanged(path.join(dtDir, down), m.code_down); } } - // Delete locally-orphaned migration files. - for (const file of fs.readdirSync(localDir)) { - if (!parseMigrationFileName(file)) continue; - if (!known.has(file)) { - try { - fs.unlinkSync(path.join(localDir, file)); - } catch { - // ignore + // Delete locally-orphaned migration files across all datatable subfolders. + if (fs.existsSync(localDir)) { + for (const datatable of fs.readdirSync(localDir)) { + const dtDir = path.join(localDir, datatable); + if (!fs.statSync(dtDir).isDirectory()) continue; + for (const file of fs.readdirSync(dtDir)) { + if (!parseMigrationFileName(file)) continue; + if (!known.has(`${datatable}/${file}`)) { + try { + fs.unlinkSync(path.join(dtDir, file)); + } catch { + // ignore + } + } } } } - log.info( - colors.green(`Pulled ${migrations.length} datatable migration(s)`), - ); + log.info(colors.green(`Pulled ${migrations.length} datatable migration(s)`)); return true; } diff --git a/cli/src/commands/sync/sync.ts b/cli/src/commands/sync/sync.ts index 59c717c326..505726ee51 100644 --- a/cli/src/commands/sync/sync.ts +++ b/cli/src/commands/sync/sync.ts @@ -4928,7 +4928,10 @@ export async function push( log.warn(`Failed to push shared UI folder: ${e}`); } try { - await pushDatatableMigrations(workspace.workspaceId); + await pushDatatableMigrations(workspace.workspaceId, { + yes: opts.yes, + jsonOutput: opts.jsonOutput, + }); } catch (e) { log.warn(`Failed to push datatable migrations folder: ${e}`); } @@ -4977,7 +4980,10 @@ export async function push( log.warn(`Failed to push shared UI folder: ${e}`); } try { - await pushDatatableMigrations(workspace.workspaceId); + await pushDatatableMigrations(workspace.workspaceId, { + yes: opts.yes, + jsonOutput: opts.jsonOutput, + }); } catch (e) { log.warn(`Failed to push datatable migrations folder: ${e}`); }