From 4bdc5a9e20bcfcddeed3f96aa485580f6d006722 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Thu, 17 Sep 2026 12:17:30 +0200 Subject: [PATCH 1/4] fix(datatables): drop a DuckDB data table secret once its ATTACH has used it Co-Authored-By: Claude Opus 5 (1M context) --- backend/windmill-worker/src/duckdb_executor.rs | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/backend/windmill-worker/src/duckdb_executor.rs b/backend/windmill-worker/src/duckdb_executor.rs index bc5660b5fb..aecdb738b3 100644 --- a/backend/windmill-worker/src/duckdb_executor.rs +++ b/backend/windmill-worker/src/duckdb_executor.rs @@ -2718,6 +2718,11 @@ fn pg_secret_attach_statements(db_resource: Value, alias_name: &str) -> Result' (TYPE postgres, SECRET …)` would reach a database nobody + // authorized this job for, as this role. + format!("DROP TEMPORARY SECRET {secret_name};"), ]) } @@ -3945,6 +3950,8 @@ mod tests { stmts[3], format!("ATTACH 'sslmode=require' AS dt (TYPE postgres, SECRET {secret_name});") ); + assert_eq!(stmts[4], format!("DROP TEMPORARY SECRET {secret_name};")); + assert_eq!(stmts.len(), 5); } #[test] From ee176e24d5bb97bc6e23bc562f2b11c39ebb0c05 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Thu, 17 Sep 2026 13:04:26 +0200 Subject: [PATCH 2/4] perf(datatables): resolve a workspace's data tables per pointer hop, not per entry Co-Authored-By: Claude Opus 5 (1M context) --- .../windmill-api-workspaces/src/workspaces.rs | 18 ++-- backend/windmill-common/src/workspaces.rs | 90 +++++++++++++++++++ 2 files changed, 98 insertions(+), 10 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 65fdc88212..6bd7d7ea7e 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -2204,17 +2204,15 @@ async fn list_datatables( Extension(db): Extension, Path(w_id): Path, ) -> JsonResult> { - let names = list_datatable_names(&db, &w_id).await?; + // A pointer entry owns no database, so what it resolves to is the only truthful answer here. + // One that resolves to nothing — a pointer whose workspace was deleted — is dropped rather than + // listed with a database it does not have; what happened is named where it is actionable + // instead: by the delete that stranded it, and by any attempt to use it. + let resolved = + windmill_common::workspaces::resolve_workspace_governing_datatables(&db, &w_id).await?; - let mut items = Vec::with_capacity(names.len()); - for name in names { - // A pointer entry owns no database, so what it resolves to is the only truthful answer - // here. One that resolves to nothing — a pointer whose workspace was deleted — is dropped - // rather than listed with a database it does not have; what happened is named where it is - // actionable instead: by the delete that stranded it, and by any attempt to use it. - let Ok(governing) = resolve_governing_datatable(&db, &w_id, &name).await else { - continue; - }; + let mut items = Vec::with_capacity(resolved.len()); + for (name, governing) in resolved { let database = governing .datatable .database diff --git a/backend/windmill-common/src/workspaces.rs b/backend/windmill-common/src/workspaces.rs index 3e8dd0aae5..1bf6314e5e 100644 --- a/backend/windmill-common/src/workspaces.rs +++ b/backend/windmill-common/src/workspaces.rs @@ -1522,6 +1522,96 @@ pub async fn resolve_governing_datatable( ))) } +/// Every entry of a workspace resolved as [`resolve_governing_datatable`] resolves one, in stored +/// order, reading the settings rows one pointer hop at a time rather than once per entry. An entry +/// that does not resolve — malformed, a dangling pointer, a loop — is left out. Same authorization +/// contract as the single resolution: it checks nothing. +pub async fn resolve_workspace_governing_datatables( + db: &DB, + w_id: &str, +) -> Result> { + type Entries = + std::collections::HashMap>; + async fn load(db: &DB, workspaces: &[String], entries: &mut Entries) -> Result> { + let rows: Vec<(String, String, serde_json::Value)> = sqlx::query_as( + "SELECT ws.workspace_id, dt.key, dt.value FROM workspace_settings ws + CROSS JOIN LATERAL jsonb_each(COALESCE(ws.datatable->'datatables', '{}'::jsonb)) dt + WHERE ws.workspace_id = ANY($1)", + ) + .bind(workspaces) + .fetch_all(db) + .await?; + for ws in workspaces { + entries.entry(ws.clone()).or_default(); + } + let mut keys = Vec::with_capacity(rows.len()); + for (ws, key, value) in rows { + keys.push(key.clone()); + entries.entry(ws).or_default().insert(key, value); + } + Ok(keys) + } + + let mut entries = Entries::new(); + let listed = load(db, &[w_id.to_string()], &mut entries).await?; + // (index into `listed`, workspace, entry name) still to be followed. + let mut cursors: Vec<(usize, String, String)> = listed + .iter() + .enumerate() + .map(|(i, name)| (i, w_id.to_string(), name.clone())) + .collect(); + let mut resolved: Vec<(usize, GoverningDatatable)> = vec![]; + + for _ in 0..DATATABLE_REFERENCE_MAX_DEPTH { + let mut next = vec![]; + for (i, ws, name) in cursors.drain(..) { + let Some(value) = entries + .get(&ws) + .and_then(|m| m.get(&name)) + .filter(|v| !v.is_null()) + else { + continue; + }; + let Ok(datatable) = serde_json::from_value::(value.clone()) else { + continue; + }; + if validate_datatable_shape(&name, &datatable).is_err() { + continue; + } + match &datatable.reference { + None => { + resolved.push((i, GoverningDatatable { workspace_id: ws, name, datatable })) + } + Some(reference) => next.push(( + i, + reference.workspace_id.clone(), + reference.datatable.clone(), + )), + } + } + if next.is_empty() { + break; + } + let to_load: Vec = next + .iter() + .map(|(_, ws, _)| ws.clone()) + .filter(|ws| !entries.contains_key(ws)) + .collect::>() + .into_iter() + .collect(); + if !to_load.is_empty() { + load(db, &to_load, &mut entries).await?; + } + cursors = next; + } + + resolved.sort_by_key(|(i, _)| *i); + Ok(resolved + .into_iter() + .map(|(i, governing)| (listed[i].clone(), governing)) + .collect()) +} + /// Build the `admin` connection for a governing entry: `custom_instance_user` for an instance /// database, the user's own resource for a BYO-postgres one. async fn resolve_datatable_connection_unchecked( From 8a7f364cfe0c3309c08131497655e800d1b9abca Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Thu, 17 Sep 2026 13:29:35 +0200 Subject: [PATCH 3/4] fix(datatables): hold the parent's settings while a fork points at its data tables Co-Authored-By: Claude Opus 5 (1M context) --- .../tests/datatable_roles.rs | 54 +++++++++++++++++++ .../windmill-api-workspaces/src/workspaces.rs | 8 +++ .../windmill-common/src/datatable_roles.rs | 3 ++ backend/windmill-common/src/workspaces.rs | 4 ++ backend/windmill-trigger-postgres/src/lib.rs | 5 +- 5 files changed, 73 insertions(+), 1 deletion(-) diff --git a/backend/windmill-api-integration-tests/tests/datatable_roles.rs b/backend/windmill-api-integration-tests/tests/datatable_roles.rs index f12daf1953..57b45088bd 100644 --- a/backend/windmill-api-integration-tests/tests/datatable_roles.rs +++ b/backend/windmill-api-integration-tests/tests/datatable_roles.rs @@ -1071,6 +1071,60 @@ async fn an_alias_saved_elsewhere_waits_for_roles_going_on_for_its_database( Ok(()) } +#[sqlx::test(migrations = "../migrations", fixtures("base", "datatable_roles"))] +async fn a_fork_waits_for_a_rename_of_the_data_table_it_keeps( + db: Pool, +) -> anyhow::Result<()> { + initialize_tracing().await; + // A rename in flight holds the parent's settings row and moves only the pointers it can see; a + // fork being created is invisible to it, so the fork has to read the name the rename commits. + let mut renaming = db.begin().await?; + sqlx::query( + "SELECT 1 FROM workspace_settings WHERE workspace_id = 'test-workspace' FOR UPDATE", + ) + .execute(&mut *renaming) + .await?; + sqlx::query( + "UPDATE workspace_settings SET datatable = jsonb_set(datatable #- '{datatables,main}', + '{datatables,renamed}', datatable->'datatables'->'main') + WHERE workspace_id = 'test-workspace'", + ) + .execute(&mut *renaming) + .await?; + + let server = ApiServer::start(db.clone()).await?; + let url = format!( + "http://localhost:{}/api/w/test-workspace/workspaces/create_fork", + server.addr.port() + ); + let fork = tokio::spawn( + authed(client().post(&url), "SECRET_TOKEN") + .json(&json!({ "id": "wm-fork-race", "name": "race", "color": "#0000ff" })) + .send(), + ); + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + assert!( + !fork.is_finished(), + "the fork copied the parent's data tables while a rename held them" + ); + renaming.commit().await?; + + let resp = fork.await??; + assert!(resp.status().is_success(), "{}", resp.text().await?); + let datatables: Option = sqlx::query_scalar( + "SELECT datatable->'datatables' FROM workspace_settings WHERE workspace_id = 'wm-fork-race'", + ) + .fetch_one(&db) + .await?; + let datatables = datatables.unwrap(); + assert_eq!( + datatables["renamed"]["reference"], + json!({ "workspace_id": "test-workspace", "datatable": "renamed" }), + "{datatables}" + ); + Ok(()) +} + #[cfg(not(all(feature = "private", feature = "enterprise")))] const ENTERPRISE_REFUSAL: &str = "Data table roles are a Windmill Enterprise Edition feature"; diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 6bd7d7ea7e..97ab07c0a5 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -8549,6 +8549,14 @@ async fn create_workspace_fork( .execute(&mut *tx) .await?; + // The pointers this fork writes to the parent's data tables stay invisible until it commits, so + // a rename of one of them cannot carry them. Holding the parent's settings row makes such a + // rename wait for this commit, and makes the copy below read one that committed first. + sqlx::query("SELECT 1 FROM workspace_settings WHERE workspace_id = $1 FOR SHARE") + .bind(&parent_workspace_id) + .execute(&mut *tx) + .await?; + // Clone all data from the parent workspace using Rust implementation if let Err(e) = clone_workspace_data(&mut tx, &db, &parent_workspace_id, &forked_id, &authed).await diff --git a/backend/windmill-common/src/datatable_roles.rs b/backend/windmill-common/src/datatable_roles.rs index 4dd9b06fde..a9195055ac 100644 --- a/backend/windmill-common/src/datatable_roles.rs +++ b/backend/windmill-common/src/datatable_roles.rs @@ -233,6 +233,9 @@ pub fn role_id_by_name<'a>(catalog: &'a DatatableRoleCatalog, name: &str) -> Res /// Every instance database the registry knows about. Role provisioning has to reach all of them: /// a role that cannot `CONNECT` to a database is refused by Postgres before any grant matters. +/// +/// Authorization: checks nothing, and names every instance database across all workspaces. Callers +/// MUST be superadmin-gated or keep the names server-side; never return them to a workspace caller. pub async fn registered_instance_databases(db: &DB) -> Result> { crate::datatable_roles_oss::registered_instance_databases(db).await } diff --git a/backend/windmill-common/src/workspaces.rs b/backend/windmill-common/src/workspaces.rs index 1bf6314e5e..711b4dac29 100644 --- a/backend/windmill-common/src/workspaces.rs +++ b/backend/windmill-common/src/workspaces.rs @@ -1967,6 +1967,10 @@ pub fn strip_datatable_permissions( /// As [`parse_datatable_ref`], except that an entry whose stored name itself contains `?` — which /// names could before they were restricted — resolves by that exact name, without a role. It is /// looked up first, so `sales?role=x` never reaches a different entry than the one stored so. +/// +/// Authorization: checks nothing, and its answer reveals whether `w_id` stores that exact name. +/// Callers MUST already act for `w_id` — a job of it, or a caller authenticated into it — and +/// MUST still pass the name to [`get_datatable_resource_from_db`] or an admin-access check. pub async fn parse_datatable_ref_for( db: &DB, w_id: &str, diff --git a/backend/windmill-trigger-postgres/src/lib.rs b/backend/windmill-trigger-postgres/src/lib.rs index 25b19e6c1b..b92860662d 100644 --- a/backend/windmill-trigger-postgres/src/lib.rs +++ b/backend/windmill-trigger-postgres/src/lib.rs @@ -377,7 +377,10 @@ pub async fn get_raw_postgres_connection( /// A replication stream reads every row of every table whatever the data table's roles grant, so /// the two don't mix: a data table under roles takes no triggers or captures, and roles cannot be /// turned on while one is enabled on it. -pub async fn ensure_not_under_roles( +/// +/// Authorization: checks nothing, and its refusal says whether `w_id`'s data table is under roles. +/// Callers MUST have established that the caller may manage triggers in `w_id` first. +pub(crate) async fn ensure_not_under_roles( db: &DB, w_id: &str, postgres_resource_path: &str, From b68768084e5cebbeeb0019b530eee739b3ce2809 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Thu, 17 Sep 2026 18:16:48 +0200 Subject: [PATCH 4/4] fix(datatables): bind fork database copies to their workspace, and count every use before dropping one Co-Authored-By: Claude Opus 5 (1M context) --- .../windmill-api-workspaces/src/workspaces.rs | 20 +++++- .../src/workspaces_extra.rs | 26 ++++++- .../windmill-common/src/instance_config.rs | 3 + backend/windmill-common/src/lib.rs | 45 ++++++++++++- backend/windmill-common/src/workspaces.rs | 67 +++++++++++++++++++ 5 files changed, 156 insertions(+), 5 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 9f8425cc69..2d840d9e51 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -3437,8 +3437,13 @@ async fn create_pg_database( } if is_instance_datatable_source(&db, &w_id, &req.source).await? { - windmill_common::create_custom_instance_database(&db, &req.target_dbname, "datatable") - .await?; + windmill_common::create_custom_instance_database( + &db, + &req.target_dbname, + "datatable", + Some(&w_id), + ) + .await?; } else { let source_pg = resolve_pg_source_checked(&db, &user_db, &authed, &w_id, &req.source).await?; @@ -3592,6 +3597,10 @@ async fn import_pg_database( .to_string(), )); } + if is_instance_datatable_source(&db, &w_id, &req.target).await? { + windmill_common::ensure_fork_database_available_to(&db, override_dbname, &w_id) + .await?; + } } target_pg.dbname = override_dbname.clone(); } @@ -8030,6 +8039,7 @@ async fn point_kept_datatables_at_parent( forked_w_id: &str, cloned: &[ForkedDatatableInfo], ) -> Result<()> { + windmill_common::workspaces::lock_fork_datatables(tx, parent_w_id).await?; let settings: Option = sqlx::query_scalar!( "SELECT datatable FROM workspace_settings WHERE workspace_id = $1", forked_w_id @@ -8198,6 +8208,12 @@ async fn apply_forked_datatable( })?, }; + if database.resource_type == DataTableCatalogResourceType::Instance + && !windmill_api_auth::is_super_admin_authed(db, authed).await? + { + windmill_common::ensure_fork_database_available_to(db, &fdt.new_dbname, parent_w_id) + .await?; + } if database.resource_type == DataTableCatalogResourceType::Instance { // The whole `database` object, not just its `resource_path`: a pointer entry has none to // patch. `reference` goes with it — exactly one of the two may be set. diff --git a/backend/windmill-api-workspaces/src/workspaces_extra.rs b/backend/windmill-api-workspaces/src/workspaces_extra.rs index 1cb5188e5b..480f815ad8 100644 --- a/backend/windmill-api-workspaces/src/workspaces_extra.rs +++ b/backend/windmill-api-workspaces/src/workspaces_extra.rs @@ -1420,7 +1420,31 @@ pub async fn drop_forked_datatable_databases( )); continue; } - if let Err(e) = windmill_common::drop_custom_instance_database(&db, db_to_drop).await { + // The fork's own entry is what is going away; anything else still reaching the copy, + // a child fork's pointer at this entry included, keeps it. The lock keeps a child fork + // from gaining such a pointer before the drop. + let dropped = async { + let mut tx = db.begin().await?; + windmill_common::workspaces::lock_fork_datatables(&mut tx, &w_id).await?; + let uses = windmill_common::workspaces::managed_database_uses( + &mut tx, + windmill_common::workspaces::DataTableCatalogResourceType::Instance, + db_to_drop, + Some((&w_id, dt_name)), + ) + .await?; + if !uses.is_empty() { + return Err(Error::BadRequest(format!( + "it is still used by {}", + uses.join(", ") + ))); + } + windmill_common::drop_custom_instance_database(&db, db_to_drop).await?; + tx.commit().await?; + Ok::<_, Error>(()) + } + .await; + if let Err(e) = dropped { errors.push(format!( "Could not drop instance database '{}' for datatable://{}: {}", db_to_drop, dt_name, e diff --git a/backend/windmill-common/src/instance_config.rs b/backend/windmill-common/src/instance_config.rs index de3fd7684a..afb2eab436 100644 --- a/backend/windmill-common/src/instance_config.rs +++ b/backend/windmill-common/src/instance_config.rs @@ -809,6 +809,9 @@ pub struct CustomInstanceDb { pub error: Option, #[serde(skip_serializing_if = "Option::is_none")] pub tag: Option, + /// The workspace a member created this fork copy for. Absent when a superadmin created it. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace_id: Option, } /// Setup log entries for a custom instance database. diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index a67bb6bfa1..ab3a1dc46b 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -1568,11 +1568,13 @@ pub async fn ensure_instance_db_grant_options_unchecked( } /// Create a custom instance database: CREATE DATABASE, grant permissions, register in global_settings. -/// The `tag` is stored in global_settings metadata (e.g. "datatable" or "ducklake"). +/// The `tag` is stored in global_settings metadata (e.g. "datatable" or "ducklake"). `for_workspace` +/// is the workspace a member creates a fork copy for; see [`ensure_fork_database_available_to`]. pub async fn create_custom_instance_database( db: &DB, dbname: &str, tag: &str, + for_workspace: Option<&str>, ) -> error::Result<()> { let dbname = dbname.trim(); validate_dbname(dbname)?; @@ -1626,7 +1628,8 @@ pub async fn create_custom_instance_database( }, "success": true, "error": null, - "tag": tag + "tag": tag, + "workspace_id": for_workspace, }); sqlx::query!( r#"UPDATE global_settings SET value = jsonb_set(value, '{databases}', (COALESCE(value->'databases', '{}'::jsonb) || to_jsonb($1::json))) WHERE name = 'custom_instance_pg_databases'"#, @@ -1646,6 +1649,44 @@ pub async fn create_custom_instance_database( Ok(()) } +/// Refuse a workspace member writing a fork copy into, or pointing a fork at, the instance database +/// `dbname`, unless `w_id` created it for that ([`create_custom_instance_database`]) and nothing uses +/// it yet. The `wm_fork_` prefix is no authorization: every instance database answers to the same +/// `custom_instance_user`, so a name is all it takes to reach another workspace's copy. +pub async fn ensure_fork_database_available_to( + db: &DB, + dbname: &str, + w_id: &str, +) -> error::Result<()> { + let created_for = sqlx::query_scalar::<_, Option>( + "SELECT value->'databases'->$1->>'workspace_id' FROM global_settings + WHERE name = 'custom_instance_pg_databases'", + ) + .bind(dbname) + .fetch_optional(db) + .await? + .flatten(); + if created_for.as_deref() != Some(w_id) { + return Err(Error::BadRequest(format!( + "Database '{dbname}' was not created for a fork of workspace '{w_id}'" + ))); + } + let uses = workspaces::managed_database_uses( + &mut *db.acquire().await?, + workspaces::DataTableCatalogResourceType::Instance, + dbname, + None, + ) + .await?; + if !uses.is_empty() { + return Err(Error::BadRequest(format!( + "Database '{dbname}' is already in use: {}", + uses.join(", ") + ))); + } + Ok(()) +} + /// Connection options parsed from a database URL. /// /// The only place a database URL becomes `PgConnectOptions`. Providers that mint the password diff --git a/backend/windmill-common/src/workspaces.rs b/backend/windmill-common/src/workspaces.rs index 352eaf943b..ed54b6f481 100644 --- a/backend/windmill-common/src/workspaces.rs +++ b/backend/windmill-common/src/workspaces.rs @@ -1471,6 +1471,73 @@ pub struct GoverningDatatable { pub datatable: DataTable, } +/// Everything still using the Windmill-managed database `dbname`, one description per use: data +/// table entries naming it, fork entries pointing at those, Ducklake catalogs on it, and fork +/// Ducklake metadata schemas there that cleanup has not dropped yet. `exempt` is the one data table +/// entry, `(workspace_id, name)`, the caller is about to stop using it through; pointers at that +/// entry still count, since dropping the database would leave them resolving to nothing. +/// +/// Authorization: reads every workspace's settings and checks nothing. Callers MUST only turn the +/// answer into a refusal for someone allowed to administer `dbname`. +pub async fn managed_database_uses( + conn: &mut sqlx::PgConnection, + kind: DataTableCatalogResourceType, + dbname: &str, + exempt: Option<(&str, &str)>, +) -> Result> { + let (exempt_workspace, exempt_name) = exempt.unzip(); + Ok(sqlx::query_scalar::<_, String>( + "WITH entries AS ( + SELECT ws.workspace_id::text AS workspace_id, dt.key AS name, dt.value + FROM workspace_settings ws + CROSS JOIN LATERAL jsonb_each( + CASE WHEN jsonb_typeof(ws.datatable->'datatables') = 'object' + THEN ws.datatable->'datatables' ELSE '{}'::jsonb END) dt + ), naming AS ( + SELECT workspace_id, name FROM entries + WHERE value->'database'->>'resource_type' = $1 + AND value->'database'->>'resource_path' = $2 + ) + SELECT format('data table ''%s'' in workspace ''%s''', name, workspace_id) FROM naming + WHERE $3::text IS NULL OR NOT (workspace_id = $3 AND name = $4) + UNION ALL + SELECT format('data table ''%s'' in workspace ''%s'', which points at the one in ''%s''', + e.name, e.workspace_id, n.workspace_id) + FROM entries e JOIN naming n + ON e.value->'reference'->>'workspace_id' = n.workspace_id + AND e.value->'reference'->>'datatable' = n.name + UNION ALL + SELECT format('Ducklake ''%s'' in workspace ''%s''', dl.key, ws.workspace_id) + FROM workspace_settings ws + CROSS JOIN LATERAL jsonb_each( + CASE WHEN jsonb_typeof(ws.ducklake->'ducklakes') = 'object' + THEN ws.ducklake->'ducklakes' ELSE '{}'::jsonb END) dl + WHERE dl.value->'catalog'->>'resource_type' = $1 + AND dl.value->'catalog'->>'resource_path' = $2 + UNION ALL + SELECT format('the Ducklake namespace of fork ''%s'', not cleaned up yet', workspace_id) + FROM fork_ducklake_namespace + WHERE catalog = $1 || ':' || $2 AND NOT schema_dropped + ORDER BY 1", + ) + .bind(kind.as_ref()) + .bind(dbname) + .bind(exempt_workspace) + .bind(exempt_name) + .fetch_all(&mut *conn) + .await?) +} + +/// Held by fork cleanup of `w_id`'s data tables and by forking `w_id`, which can hand the new fork +/// pointers at them, so a pointer cannot appear between cleanup's check and its drop. +pub async fn lock_fork_datatables(conn: &mut sqlx::PgConnection, w_id: &str) -> Result<()> { + sqlx::query("SELECT pg_advisory_xact_lock(hashtext('fork_datatables:' || $1))") + .bind(w_id) + .execute(&mut *conn) + .await?; + Ok(()) +} + impl GoverningDatatable { /// Backed by the Windmill instance's own Postgres, which is the only substrate data table /// roles apply to.