feat: track datatable table DDL changes in workspace_diff
Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -5044,6 +5044,14 @@ async fn compare_workspaces(
|
|||||||
compare_two_folders(&db, &source_workspace_id, &fork_workspace_id, &item.path)
|
compare_two_folders(&db, &source_workspace_id, &fork_workspace_id, &item.path)
|
||||||
.await?,
|
.await?,
|
||||||
),
|
),
|
||||||
|
"datatable_table" => {
|
||||||
|
// No stateless comparison yet — assume changes are real
|
||||||
|
Some(ItemComparison {
|
||||||
|
has_changes: true,
|
||||||
|
exists_in_source: true,
|
||||||
|
exists_in_fork: true,
|
||||||
|
})
|
||||||
|
}
|
||||||
k => {
|
k => {
|
||||||
tracing::error!("Received unrecognized item kind `{k}` with path: `{}` while computing diff of {fork_workspace_id} and {source_workspace_id} workspaces. Skipping this item", item.path);
|
tracing::error!("Received unrecognized item kind `{k}` with path: `{}` while computing diff of {fork_workspace_id} and {source_workspace_id} workspaces. Skipping this item", item.path);
|
||||||
None
|
None
|
||||||
@@ -5284,6 +5292,8 @@ async fn query_visible_items<'c>(
|
|||||||
.fetch_all(&mut **tx)
|
.fetch_all(&mut **tx)
|
||||||
.await?
|
.await?
|
||||||
}
|
}
|
||||||
|
// Datatable tables are always visible (no permission gating)
|
||||||
|
"datatable_table" => paths_vec,
|
||||||
_ => vec![], // Unknown kind
|
_ => vec![], // Unknown kind
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -24351,7 +24351,7 @@ components:
|
|||||||
kind:
|
kind:
|
||||||
type: string
|
type: string
|
||||||
enum:
|
enum:
|
||||||
["script", "flow", "app", "raw_app", "resource", "variable", "resource_type"]
|
["script", "flow", "app", "raw_app", "resource", "variable", "resource_type", "datatable_table"]
|
||||||
description: Type of the item
|
description: Type of the item
|
||||||
path:
|
path:
|
||||||
type: string
|
type: string
|
||||||
|
|||||||
@@ -230,6 +230,12 @@ fn wrap_ducklake_query(query: &str, ducklake: &str) -> String {
|
|||||||
format!("{}{}{}", &query[..insert_pos], attach, &query[insert_pos..])
|
format!("{}{}{}", &query[..insert_pos], attach, &query[insert_pos..])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Describes a DDL operation that modifies a table structure.
|
||||||
|
#[derive(Debug, Clone, PartialEq)]
|
||||||
|
pub struct DdlOperation {
|
||||||
|
pub table_name: String,
|
||||||
|
}
|
||||||
|
|
||||||
/// Result of expanding a WM_INTERNAL_DB marker.
|
/// Result of expanding a WM_INTERNAL_DB marker.
|
||||||
#[derive(Debug, PartialEq)]
|
#[derive(Debug, PartialEq)]
|
||||||
pub struct ExpandedQuery {
|
pub struct ExpandedQuery {
|
||||||
@@ -237,6 +243,8 @@ pub struct ExpandedQuery {
|
|||||||
/// If set, the worker should use this language instead of the original.
|
/// If set, the worker should use this language instead of the original.
|
||||||
/// Used when the expanded code is in a different language (e.g. BigQuery all-tables uses Bun).
|
/// Used when the expanded code is in a different language (e.g. BigQuery all-tables uses Bun).
|
||||||
pub language_override: Option<ScriptLang>,
|
pub language_override: Option<ScriptLang>,
|
||||||
|
/// If set, this expansion was a DDL operation (CREATE/ALTER/DROP TABLE).
|
||||||
|
pub ddl_operation: Option<DdlOperation>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Checks if a SQL script is a WM_INTERNAL_DB marker and expands it into real SQL.
|
/// Checks if a SQL script is a WM_INTERNAL_DB marker and expands it into real SQL.
|
||||||
@@ -279,9 +287,18 @@ pub fn try_expand_internal_db_query(
|
|||||||
"INSERT" => expand_insert(json_str, db_type).map(ExpandedQuery::sql),
|
"INSERT" => expand_insert(json_str, db_type).map(ExpandedQuery::sql),
|
||||||
"UPDATE" => expand_update(json_str, db_type).map(ExpandedQuery::sql),
|
"UPDATE" => expand_update(json_str, db_type).map(ExpandedQuery::sql),
|
||||||
// Schema DDL operations
|
// Schema DDL operations
|
||||||
"DROP_TABLE" => expand_drop_table(json_str, db_type).map(ExpandedQuery::sql),
|
"DROP_TABLE" => {
|
||||||
"CREATE_TABLE" => expand_create_table(json_str, db_type).map(ExpandedQuery::sql),
|
let table_name = extract_ddl_table_name(json_str, "table");
|
||||||
"ALTER_TABLE" => expand_alter_table(json_str, db_type).map(ExpandedQuery::sql),
|
expand_drop_table(json_str, db_type).map(|code| ExpandedQuery::ddl(code, table_name))
|
||||||
|
}
|
||||||
|
"CREATE_TABLE" => {
|
||||||
|
let table_name = extract_ddl_table_name(json_str, "name");
|
||||||
|
expand_create_table(json_str, db_type).map(|code| ExpandedQuery::ddl(code, table_name))
|
||||||
|
}
|
||||||
|
"ALTER_TABLE" => {
|
||||||
|
let table_name = extract_ddl_table_name(json_str, "name");
|
||||||
|
expand_alter_table(json_str, db_type).map(|code| ExpandedQuery::ddl(code, table_name))
|
||||||
|
}
|
||||||
"CREATE_SCHEMA" => expand_create_schema(json_str, db_type).map(ExpandedQuery::sql),
|
"CREATE_SCHEMA" => expand_create_schema(json_str, db_type).map(ExpandedQuery::sql),
|
||||||
"DROP_SCHEMA" => expand_drop_schema(json_str, db_type).map(ExpandedQuery::sql),
|
"DROP_SCHEMA" => expand_drop_schema(json_str, db_type).map(ExpandedQuery::sql),
|
||||||
// Metadata queries
|
// Metadata queries
|
||||||
@@ -297,13 +314,26 @@ pub fn try_expand_internal_db_query(
|
|||||||
Some(result)
|
Some(result)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Extract the table name from a DDL JSON payload. The key varies by operation:
|
||||||
|
/// DROP_TABLE uses "table", CREATE_TABLE and ALTER_TABLE use "name".
|
||||||
|
fn extract_ddl_table_name(json_str: &str, key: &str) -> String {
|
||||||
|
serde_json::from_str::<serde_json::Value>(json_str)
|
||||||
|
.ok()
|
||||||
|
.and_then(|v| v.get(key).and_then(|v| v.as_str()).map(|s| s.to_string()))
|
||||||
|
.unwrap_or_default()
|
||||||
|
}
|
||||||
|
|
||||||
impl ExpandedQuery {
|
impl ExpandedQuery {
|
||||||
fn sql(code: String) -> Self {
|
fn sql(code: String) -> Self {
|
||||||
Self { code, language_override: None }
|
Self { code, language_override: None, ddl_operation: None }
|
||||||
|
}
|
||||||
|
|
||||||
|
fn ddl(code: String, table_name: String) -> Self {
|
||||||
|
Self { code, language_override: None, ddl_operation: Some(DdlOperation { table_name }) }
|
||||||
}
|
}
|
||||||
|
|
||||||
fn with_language(code: String, lang: ScriptLang) -> Self {
|
fn with_language(code: String, lang: ScriptLang) -> Self {
|
||||||
Self { code, language_override: Some(lang) }
|
Self { code, language_override: Some(lang), ddl_operation: None }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ pub enum DeployedObject {
|
|||||||
Settings { setting_type: String },
|
Settings { setting_type: String },
|
||||||
Key { key_type: String },
|
Key { key_type: String },
|
||||||
WorkspaceDependencies { path: String },
|
WorkspaceDependencies { path: String },
|
||||||
|
DatatableTable { datatable_name: String, table_name: String },
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DeployedObject {
|
impl DeployedObject {
|
||||||
@@ -71,6 +72,9 @@ impl DeployedObject {
|
|||||||
DeployedObject::Settings { .. } => "settings.yaml".to_string(),
|
DeployedObject::Settings { .. } => "settings.yaml".to_string(),
|
||||||
DeployedObject::Key { .. } => "encryption_key.yaml".to_string(),
|
DeployedObject::Key { .. } => "encryption_key.yaml".to_string(),
|
||||||
DeployedObject::WorkspaceDependencies { path, .. } => path.to_owned(),
|
DeployedObject::WorkspaceDependencies { path, .. } => path.to_owned(),
|
||||||
|
DeployedObject::DatatableTable { datatable_name, table_name } => {
|
||||||
|
format!("{datatable_name}/{table_name}")
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -81,7 +85,8 @@ impl DeployedObject {
|
|||||||
| Self::ResourceType { .. }
|
| Self::ResourceType { .. }
|
||||||
| Self::Settings { .. }
|
| Self::Settings { .. }
|
||||||
| Self::Key { .. }
|
| Self::Key { .. }
|
||||||
| Self::WorkspaceDependencies { .. } => true,
|
| Self::WorkspaceDependencies { .. }
|
||||||
|
| Self::DatatableTable { .. } => true,
|
||||||
_ => false,
|
_ => false,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -111,6 +116,7 @@ impl DeployedObject {
|
|||||||
DeployedObject::Settings { .. } => None,
|
DeployedObject::Settings { .. } => None,
|
||||||
DeployedObject::Key { .. } => None,
|
DeployedObject::Key { .. } => None,
|
||||||
DeployedObject::WorkspaceDependencies { .. } => None,
|
DeployedObject::WorkspaceDependencies { .. } => None,
|
||||||
|
DeployedObject::DatatableTable { .. } => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -139,6 +145,7 @@ impl DeployedObject {
|
|||||||
DeployedObject::Settings { .. } => "settings",
|
DeployedObject::Settings { .. } => "settings",
|
||||||
DeployedObject::Key { .. } => "key",
|
DeployedObject::Key { .. } => "key",
|
||||||
DeployedObject::WorkspaceDependencies { .. } => "workspace_dependencies",
|
DeployedObject::WorkspaceDependencies { .. } => "workspace_dependencies",
|
||||||
|
DeployedObject::DatatableTable { .. } => "datatable_table",
|
||||||
}
|
}
|
||||||
.to_string()
|
.to_string()
|
||||||
}
|
}
|
||||||
@@ -272,7 +279,10 @@ mod tests {
|
|||||||
path: "f/folder/script".to_string(),
|
path: "f/folder/script".to_string(),
|
||||||
parent_path: Some("f/folder/old_script".to_string()),
|
parent_path: Some("f/folder/old_script".to_string()),
|
||||||
};
|
};
|
||||||
assert_eq!(obj.get_parent_path(), Some("f/folder/old_script".to_string()));
|
assert_eq!(
|
||||||
|
obj.get_parent_path(),
|
||||||
|
Some("f/folder/old_script".to_string())
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -313,21 +323,13 @@ mod tests {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_get_kind_flow() {
|
fn test_get_kind_flow() {
|
||||||
let obj = DeployedObject::Flow {
|
let obj = DeployedObject::Flow { path: "test".to_string(), parent_path: None, version: 1 };
|
||||||
path: "test".to_string(),
|
|
||||||
parent_path: None,
|
|
||||||
version: 1,
|
|
||||||
};
|
|
||||||
assert_eq!(obj.get_kind(), "flow");
|
assert_eq!(obj.get_kind(), "flow");
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_get_kind_app() {
|
fn test_get_kind_app() {
|
||||||
let obj = DeployedObject::App {
|
let obj = DeployedObject::App { path: "test".to_string(), version: 1, parent_path: None };
|
||||||
path: "test".to_string(),
|
|
||||||
version: 1,
|
|
||||||
parent_path: None,
|
|
||||||
};
|
|
||||||
assert_eq!(obj.get_kind(), "app");
|
assert_eq!(obj.get_kind(), "app");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -346,7 +348,8 @@ mod tests {
|
|||||||
"http_trigger"
|
"http_trigger"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
DeployedObject::WebsocketTrigger { path: "t".to_string(), parent_path: None }.get_kind(),
|
DeployedObject::WebsocketTrigger { path: "t".to_string(), parent_path: None }
|
||||||
|
.get_kind(),
|
||||||
"websocket_trigger"
|
"websocket_trigger"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|||||||
@@ -4212,12 +4212,14 @@ pub async fn run_language_executor(
|
|||||||
// Expand WM_INTERNAL_DB markers into real SQL before dispatching
|
// Expand WM_INTERNAL_DB markers into real SQL before dispatching
|
||||||
let expanded_code: String;
|
let expanded_code: String;
|
||||||
let mut language = language;
|
let mut language = language;
|
||||||
|
let mut ddl_operation: Option<windmill_common::query_builders::DdlOperation> = None;
|
||||||
let code = if let Some(ref lang) = language {
|
let code = if let Some(ref lang) = language {
|
||||||
match windmill_common::query_builders::try_expand_internal_db_query(code, lang) {
|
match windmill_common::query_builders::try_expand_internal_db_query(code, lang) {
|
||||||
Some(Ok(expanded)) => {
|
Some(Ok(expanded)) => {
|
||||||
if let Some(lang_override) = expanded.language_override {
|
if let Some(lang_override) = expanded.language_override {
|
||||||
language = Some(lang_override);
|
language = Some(lang_override);
|
||||||
}
|
}
|
||||||
|
ddl_operation = expanded.ddl_operation;
|
||||||
expanded_code = expanded.code;
|
expanded_code = expanded.code;
|
||||||
&expanded_code
|
&expanded_code
|
||||||
}
|
}
|
||||||
@@ -4232,6 +4234,40 @@ pub async fn run_language_executor(
|
|||||||
} else {
|
} else {
|
||||||
code
|
code
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Tally datatable table DDL change for workspace diff tracking
|
||||||
|
if let Some(ref ddl_op) = ddl_operation {
|
||||||
|
if let Connection::Sql(ref db) = conn {
|
||||||
|
// Extract datatable name from the "database" arg (e.g. "datatable://main")
|
||||||
|
let datatable_name = job
|
||||||
|
.args
|
||||||
|
.as_ref()
|
||||||
|
.and_then(|args| args.get("database"))
|
||||||
|
.and_then(|v| serde_json::from_str::<String>(v.get()).ok())
|
||||||
|
.and_then(|s| s.strip_prefix("datatable://").map(|n| n.to_string()));
|
||||||
|
|
||||||
|
if let Some(dt_name) = datatable_name {
|
||||||
|
if let Err(e) = windmill_git_sync::handle_deployment_metadata(
|
||||||
|
&job.permissioned_as_email,
|
||||||
|
&job.created_by,
|
||||||
|
db,
|
||||||
|
&job.workspace_id,
|
||||||
|
windmill_git_sync::DeployedObject::DatatableTable {
|
||||||
|
datatable_name: dt_name,
|
||||||
|
table_name: ddl_op.table_name.clone(),
|
||||||
|
},
|
||||||
|
None,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
tracing::error!(%e, "error handling datatable DDL deployment metadata");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if let Some(modules) = modules {
|
if let Some(modules) = modules {
|
||||||
#[cfg(feature = "python")]
|
#[cfg(feature = "python")]
|
||||||
let base_dir = if language == Some(ScriptLang::Python3) {
|
let base_dir = if language == Some(ScriptLang::Python3) {
|
||||||
|
|||||||
Reference in New Issue
Block a user