diff --git a/backend/.sqlx/query-2ee6d24b95cdda151585dcff19f8e7c931785fc21f7bbe9c3a82671943ced0ea.json b/backend/.sqlx/query-7947ffe31b8e6f4a38fba9d9caf434fd409ea2806a7c720c2b1d76b2e7db6c32.json similarity index 65% rename from backend/.sqlx/query-2ee6d24b95cdda151585dcff19f8e7c931785fc21f7bbe9c3a82671943ced0ea.json rename to backend/.sqlx/query-7947ffe31b8e6f4a38fba9d9caf434fd409ea2806a7c720c2b1d76b2e7db6c32.json index 6acb31666c..e69a61e65d 100644 --- a/backend/.sqlx/query-2ee6d24b95cdda151585dcff19f8e7c931785fc21f7bbe9c3a82671943ced0ea.json +++ b/backend/.sqlx/query-7947ffe31b8e6f4a38fba9d9caf434fd409ea2806a7c720c2b1d76b2e7db6c32.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3", + "query": "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3 AND enabled = true", "describe": { "columns": [], "parameters": { @@ -12,5 +12,5 @@ }, "nullable": [] }, - "hash": "2ee6d24b95cdda151585dcff19f8e7c931785fc21f7bbe9c3a82671943ced0ea" + "hash": "7947ffe31b8e6f4a38fba9d9caf434fd409ea2806a7c720c2b1d76b2e7db6c32" } diff --git a/backend/.sqlx/query-928767710fb8b7dc0b1edc897d8ce9b6b59ae2f63e684d3db6eecbcffd767711.json b/backend/.sqlx/query-928767710fb8b7dc0b1edc897d8ce9b6b59ae2f63e684d3db6eecbcffd767711.json new file mode 100644 index 0000000000..36e067a79a --- /dev/null +++ b/backend/.sqlx/query-928767710fb8b7dc0b1edc897d8ce9b6b59ae2f63e684d3db6eecbcffd767711.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO trigger_history\n (workspace_id, trigger_kind, path, operation, source, username, changes)\n SELECT $1, $2, p, $3, $4, $5, $6 FROM unnest($7::text[]) p", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Jsonb", + "TextArray" + ] + }, + "nullable": [] + }, + "hash": "928767710fb8b7dc0b1edc897d8ce9b6b59ae2f63e684d3db6eecbcffd767711" +} diff --git a/backend/.sqlx/query-e52a80386b132e53a956458e2eb77d29bd6da4fc8511ed44f21d49f4c0965d36.json b/backend/.sqlx/query-e52a80386b132e53a956458e2eb77d29bd6da4fc8511ed44f21d49f4c0965d36.json new file mode 100644 index 0000000000..26a5b9f932 --- /dev/null +++ b/backend/.sqlx/query-e52a80386b132e53a956458e2eb77d29bd6da4fc8511ed44f21d49f4c0965d36.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO trigger_history\n (workspace_id, trigger_kind, path, operation, source, username, changes)\n VALUES ($1, $2, $3, $4, $5, $6, $7)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "e52a80386b132e53a956458e2eb77d29bd6da4fc8511ed44f21d49f4c0965d36" +} diff --git a/backend/.sqlx/query-fcbbc3b697249c6ee0ca542ea42fadecddead568110dd531fbace439c2263e4f.json b/backend/.sqlx/query-fcbbc3b697249c6ee0ca542ea42fadecddead568110dd531fbace439c2263e4f.json new file mode 100644 index 0000000000..b039aab41c --- /dev/null +++ b/backend/.sqlx/query-fcbbc3b697249c6ee0ca542ea42fadecddead568110dd531fbace439c2263e4f.json @@ -0,0 +1,71 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT id, trigger_kind, path, operation, source, username, created_at, changes\n FROM trigger_history\n WHERE workspace_id = $1\n AND ($2::TEXT IS NULL OR trigger_kind = $2)\n AND ($3::TEXT IS NULL OR path = $3)\n AND ( $6\n OR path = ANY($7)\n OR EXISTS ( SELECT 1 FROM unnest($8::text[]) AS pfx\n WHERE path = pfx\n OR left(path, length(pfx) + 1) = pfx || '/' ) )\n ORDER BY id DESC\n LIMIT $4 OFFSET $5", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "trigger_kind", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "path", + "type_info": "Varchar" + }, + { + "ordinal": 3, + "name": "operation", + "type_info": "Varchar" + }, + { + "ordinal": 4, + "name": "source", + "type_info": "Varchar" + }, + { + "ordinal": 5, + "name": "username", + "type_info": "Varchar" + }, + { + "ordinal": 6, + "name": "created_at", + "type_info": "Timestamptz" + }, + { + "ordinal": 7, + "name": "changes", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Text", + "Int8", + "Int8", + "Bool", + "TextArray", + "TextArray" + ] + }, + "nullable": [ + false, + false, + false, + false, + false, + true, + false, + true + ] + }, + "hash": "fcbbc3b697249c6ee0ca542ea42fadecddead568110dd531fbace439c2263e4f" +} diff --git a/backend/migrations/20260814090221_trigger_history.down.sql b/backend/migrations/20260814090221_trigger_history.down.sql new file mode 100644 index 0000000000..c4aad347b0 --- /dev/null +++ b/backend/migrations/20260814090221_trigger_history.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS trigger_history; diff --git a/backend/migrations/20260814090221_trigger_history.up.sql b/backend/migrations/20260814090221_trigger_history.up.sql new file mode 100644 index 0000000000..1eb4f778a2 --- /dev/null +++ b/backend/migrations/20260814090221_trigger_history.up.sql @@ -0,0 +1,65 @@ +-- Append-only record of every schedule/trigger mutation: who, what changed, and +-- from which kind of client. +CREATE TABLE IF NOT EXISTS trigger_history ( + id BIGSERIAL PRIMARY KEY, + workspace_id VARCHAR(50) NOT NULL REFERENCES workspace(id) ON DELETE CASCADE, + -- 'schedule' or a trigger's TRIGGER_TYPE ('http', 'kafka', ...). Not the + -- TRIGGER_KIND enum: that one is capture-oriented and misses 'schedule'. + trigger_kind VARCHAR(50) NOT NULL, + path VARCHAR(255) NOT NULL, + -- 'create' | 'update' | 'delete' | 'enable' | 'disable' | 'suspend' + operation VARCHAR(20) NOT NULL, + -- 'ui' | 'cli' | 'api' | 'worker' + source VARCHAR(20) NOT NULL, + -- NULL when the server acted on its own (worker auto-disable). + username VARCHAR(255), + created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(), + -- {field: {old, new}} for the fields that actually changed. `old` is + -- absent where it is not known: a create, and the workspace-wide handler + -- override that rewrites every schedule without reading them first. NULL + -- when the operation carries no field-level diff at all (delete). + changes JSONB +); + +CREATE INDEX IF NOT EXISTS idx_trigger_history_workspace_kind_path + ON trigger_history(workspace_id, trigger_kind, path, id DESC); + +CREATE INDEX IF NOT EXISTS idx_trigger_history_workspace_id + ON trigger_history(workspace_id, id DESC); + +GRANT ALL ON TABLE trigger_history TO windmill_user; +GRANT ALL ON TABLE trigger_history TO windmill_admin; +GRANT ALL ON SEQUENCE trigger_history_id_seq TO windmill_user; +GRANT ALL ON SEQUENCE trigger_history_id_seq TO windmill_admin; + +ALTER TABLE trigger_history ENABLE ROW LEVEL SECURITY; + +CREATE POLICY admin_all ON trigger_history FOR ALL TO windmill_admin USING (true) WITH CHECK (true); + +-- Every mutating trigger route writes through the RLS pool, so windmill_user +-- must be able to append. +CREATE POLICY allow_insert ON trigger_history FOR INSERT TO windmill_user WITH CHECK (true); + +-- Reads mirror the path half of the live trigger's own policies: a row can +-- quote a schedule's `args`, so it must not be readable by anyone who could not +-- read the trigger it describes. Deliberately narrower than the live row on one +-- point — the `extra_perms` grants have no counterpart here, since the history +-- does not carry the row's ACL and must survive its deletion. Narrower is the +-- safe direction. +CREATE POLICY see_own ON trigger_history FOR SELECT TO windmill_user +USING ( + SPLIT_PART(path::text, '/', 1) = 'u' + AND SPLIT_PART(path::text, '/', 2) = current_setting('session.user') +); + +CREATE POLICY see_member ON trigger_history FOR SELECT TO windmill_user +USING ( + SPLIT_PART(path::text, '/', 1) = 'g' + AND SPLIT_PART(path::text, '/', 2) = ANY(regexp_split_to_array(current_setting('session.groups'), ',')) +); + +CREATE POLICY see_folder_extra_perms_user ON trigger_history FOR SELECT TO windmill_user +USING ( + SPLIT_PART(path::text, '/', 1) = 'f' + AND SPLIT_PART(path::text, '/', 2) = ANY(regexp_split_to_array(current_setting('session.folders_read'), ',')) +); diff --git a/backend/summarized_schema.txt b/backend/summarized_schema.txt index 44b62ca51e..732f1ece52 100644 --- a/backend/summarized_schema.txt +++ b/backend/summarized_schema.txt @@ -176,6 +176,9 @@ token: token_hash(char), token_prefix(char), token(char), label(char), expiratio FK: (workspace_id) -> workspace(id) token_expiry_notification: token_hash(char), expiration(ts) INDEX: idx_token_expiry_notification_expiration (expiration) +trigger_history: id(bigint), workspace_id(char), trigger_kind(char), path(char), operation(char), source(char), username(char), created_at(ts), changes(jsonb) + FK: (workspace_id) -> workspace(id) + INDEX: idx_trigger_history_workspace_kind_path (workspace_id, trigger_kind, path, id), idx_trigger_history_workspace_id (workspace_id, id) tutorial_progress: email(char), progress(bit64), skipped_all(bool) unique_ext_jwt_token: jwt_hash(bigint), last_used_at(ts), email(text), username(text), is_admin(bool), is_operator(bool), workspace_id(text?), label(text?), scopes(text[]?) usage: id(char), is_workspace(bool), month_(int), usage(int) diff --git a/backend/windmill-api-auth/src/lib.rs b/backend/windmill-api-auth/src/lib.rs index d92c3a5f2d..88eb11093a 100644 --- a/backend/windmill-api-auth/src/lib.rs +++ b/backend/windmill-api-auth/src/lib.rs @@ -67,8 +67,9 @@ pub struct ApiAuthed { /// `label-*` string. Only `username_override_from_label` sets it. pub username_override_is_token_label: bool, /// Whether the request authenticated with the session token minted at browser login. - /// Only `trigger_or_fallback` reads it — see `is_session_label` for why it attributes - /// rather than proves, and must not gate authority. + /// Read by `trigger_or_fallback` and by `TriggerSource::of_request` (which attributes a + /// trigger mutation to the UI) — see `is_session_label` for why it attributes rather than + /// proves, and must not gate authority. pub is_session_token: bool, pub token_prefix: Option, pub read_only: bool, diff --git a/backend/windmill-api-auth/src/scopes.rs b/backend/windmill-api-auth/src/scopes.rs index ed273c7d15..fb3da99ad3 100644 --- a/backend/windmill-api-auth/src/scopes.rs +++ b/backend/windmill-api-auth/src/scopes.rs @@ -274,6 +274,7 @@ pub enum ScopeDomain { // Native trigger domains NativeTriggers, + TriggersHistory, // System domains Audit, @@ -335,6 +336,7 @@ impl ScopeDomain { Self::PostgresTriggers => "postgres_triggers", Self::EmailTriggers => "email_triggers", Self::NativeTriggers => "native_triggers", + Self::TriggersHistory => "triggers_history", Self::Audit => "audit", Self::Settings => "settings", Self::Workers => "workers", @@ -401,6 +403,7 @@ impl ScopeDomain { "indexer" | "srch" => Some(Self::Indexer), "teams" => Some(Self::Teams), "native_triggers" => Some(Self::NativeTriggers), + "triggers_history" => Some(Self::TriggersHistory), "git_sync" | "github_app" => Some(Self::GitSync), "capture" => Some(Self::Capture), "drafts" => Some(Self::Drafts), diff --git a/backend/windmill-api-schedule/src/lib.rs b/backend/windmill-api-schedule/src/lib.rs index b4e9bf269a..ec790ec060 100644 --- a/backend/windmill-api-schedule/src/lib.rs +++ b/backend/windmill-api-schedule/src/lib.rs @@ -27,6 +27,9 @@ use windmill_common::{ db::UserDB, error::{Error, JsonResult, Result}, schedule::Schedule, + trigger_history::{ + self, TriggerHistoryEvent, TriggerOperation, TriggerSource, SCHEDULE_TRIGGER_KIND, + }, user_drafts::{ delete_all_drafts_for_path, fetch_draft_only_list_rows, overlay_or_draft_only, UserDraftItemKind, WithDraftOverlay, WithDraftQuery, @@ -85,6 +88,40 @@ fn resolve_edited_by(authed: &ApiAuthed) -> String { authed.username.clone() } +/// Append this mutation to `trigger_history`, diffing the row against `before`. +/// +/// Call it on the transaction that made the change, after the change: the +/// snapshot it takes is the "after" side of the diff, and the two commit or roll +/// back together. +async fn record_schedule_history( + tx: &mut sqlx::PgConnection, + authed: &ApiAuthed, + w_id: &str, + path: &str, + operation: TriggerOperation, + before: Option, +) -> Result<()> { + let after = trigger_history::snapshot_row(&mut *tx, "schedule", w_id, path).await?; + // Nothing to describe when the row is not there after the write: the same + // guard the trigger side needs, kept here so the two read alike. + if after.is_none() { + return Ok(()); + } + trigger_history::record( + &mut *tx, + TriggerHistoryEvent { + workspace_id: w_id, + trigger_kind: SCHEDULE_TRIGGER_KIND, + path, + operation, + source: TriggerSource::of_request(authed.is_session_token), + username: Some(&authed.username), + changes: trigger_history::summarize_changes(before.as_ref(), after.as_ref()), + }, + ) + .await +} + pub fn workspaced_service() -> Router { Router::new() .route("/list", get(list_schedule)) @@ -417,6 +454,16 @@ async fn create_schedule( .await .map_err(|e| Error::internal_err(format!("inserting schedule in {w_id}: {e:#}")))?; + record_schedule_history( + &mut *tx, + &authed, + &w_id, + &ns.path, + TriggerOperation::Create, + None, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -524,6 +571,8 @@ async fn edit_schedule( authed.email.clone() }; + let before = trigger_history::snapshot_row(&mut *tx, "schedule", &w_id, path).await?; + let schedule = sqlx::query_as!( Schedule, r#" @@ -632,6 +681,16 @@ async fn edit_schedule( // like set_enabled, flow updates, and worker job completions. clear_schedule(&mut tx, path, &w_id).await?; + record_schedule_history( + &mut *tx, + &authed, + &w_id, + path, + TriggerOperation::Update, + before, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -1084,6 +1143,8 @@ pub async fn set_enabled( } } } + let before = trigger_history::snapshot_row(&mut *tx, "schedule", &w_id, path).await?; + // email is still written for backwards compat with old workers that don't know about permissioned_as let schedule_o = sqlx::query_as!( Schedule, @@ -1139,6 +1200,20 @@ pub async fn set_enabled( clear_schedule(&mut tx, path, &w_id).await?; + record_schedule_history( + &mut *tx, + &authed, + &w_id, + path, + if payload.enabled { + TriggerOperation::Enable + } else { + TriggerOperation::Disable + }, + before, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -1285,6 +1360,22 @@ async fn delete_schedule( .await?; } + // No diff: the row is gone, and the trashbin above already keeps its full + // contents for a restore. + trigger_history::record( + &mut *tx, + TriggerHistoryEvent { + workspace_id: &w_id, + trigger_kind: SCHEDULE_TRIGGER_KIND, + path, + operation: TriggerOperation::Delete, + source: TriggerSource::of_request(authed.is_session_token), + username: Some(&authed.username), + changes: None, + }, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -1373,6 +1464,11 @@ async fn set_default_error_handler( } if payload.override_existing { + // The rewrite and its history rows go in one transaction: on separate + // connections a concurrent edit could interleave, leaving the + // id-ordered drawer showing the wrong latest change, and a failed + // insert would leave the schedules rewritten with nothing recording it. + let mut tx = db.begin().await?; let updated_schedules: Vec; match payload.handler_type { HandlerType::Error => { @@ -1386,14 +1482,14 @@ async fn set_default_error_handler( payload.number_of_occurence_exact, w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } else { updated_schedules = sqlx::query_scalar!( "UPDATE schedule SET ws_error_handler_muted = false, on_failure = NULL, on_failure_extra_args = NULL, on_failure_times = NULL, on_failure_exact = NULL WHERE workspace_id = $1 RETURNING path", w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } } @@ -1406,14 +1502,14 @@ async fn set_default_error_handler( payload.number_of_occurence, w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } else { updated_schedules = sqlx::query_scalar!( "UPDATE schedule SET on_recovery = NULL, on_recovery_extra_args = NULL, on_recovery_times = NULL WHERE workspace_id = $1 RETURNING path", w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } } @@ -1425,18 +1521,70 @@ async fn set_default_error_handler( payload.extra_args, w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } else { updated_schedules = sqlx::query_scalar!( "UPDATE schedule SET on_success = NULL, on_success_extra_args = NULL WHERE workspace_id = $1 RETURNING path", w_id, ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; } } } + // One row per schedule the workspace-wide override rewrote, so a handler + // that appeared on a schedule nobody edited is traceable. Every column + // the UPDATE above wrote, not just the handler path: the mute flag and + // the occurrence thresholds are what someone auditing a surprise + // notification change most needs. No `old` side and no + // already-had-this-value filter — the UPDATE rewrites the whole + // workspace unconditionally, so these rows record the write rather than + // a delta. + // Built from the same values the branch that ran actually bound: a reset + // (`payload.path` absent) hardcodes NULL / false in SQL while the request + // still carries the form's other fields, so reading them here would name + // values the write never produced. + let cleared = payload.path.is_none(); + let handler_path = payload.path.clone(); + let extra_args = (!cleared).then(|| payload.extra_args.clone()).flatten(); + let times = (!cleared).then_some(payload.number_of_occurence).flatten(); + let handler_fields = match payload.handler_type { + HandlerType::Error => serde_json::json!({ + "on_failure": { "new": handler_path }, + "on_failure_extra_args": { "new": extra_args }, + "on_failure_times": { "new": times }, + "on_failure_exact": { + "new": (!cleared).then_some(payload.number_of_occurence_exact).flatten() + }, + "ws_error_handler_muted": { + "new": !cleared && payload.workspace_handler_muted.unwrap_or(false) + }, + }), + HandlerType::Recovery => serde_json::json!({ + "on_recovery": { "new": handler_path }, + "on_recovery_extra_args": { "new": extra_args }, + "on_recovery_times": { "new": times }, + }), + HandlerType::Success => serde_json::json!({ + "on_success": { "new": handler_path }, + "on_success_extra_args": { "new": extra_args }, + }), + }; + trigger_history::record_bulk( + &mut tx, + &w_id, + SCHEDULE_TRIGGER_KIND, + &updated_schedules, + TriggerOperation::Update, + TriggerSource::of_request(authed.is_session_token), + Some(&authed.username), + Some(handler_fields), + ) + .await?; + + tx.commit().await?; + for updated_schedule_path in updated_schedules { // managed ducklake-maintenance rows get the handler update (their // failures should reach workspace handlers) but must not be diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index df46ec1ec8..03689f090e 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -20342,6 +20342,36 @@ paths: type: string nullable: true + /w/{workspace}/triggers_history/list: + get: + summary: list the history of schedule and trigger modifications + operationId: listTriggerHistory + tags: + - trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Page" + - $ref: "#/components/parameters/PerPage" + - name: trigger_kind + description: "'schedule' or a trigger type (http, kafka, ...)" + in: query + schema: + type: string + - name: path + description: only return the history of the trigger at this path + in: query + schema: + type: string + responses: + "200": + description: trigger history + content: + application/json: + schema: + type: array + items: + $ref: "#/components/schemas/TriggerHistoryEntry" + /w/{workspace}/folders/list: get: summary: list folders @@ -28245,6 +28275,44 @@ components: is_fileset: type: boolean + TriggerHistoryEntry: + type: object + properties: + id: + type: integer + format: int64 + trigger_kind: + type: string + description: "'schedule' or a trigger type (http, kafka, ...)" + path: + type: string + operation: + type: string + enum: [create, update, delete, enable, disable, suspend] + source: + type: string + description: The kind of client the change came from. `worker` means the server disabled the trigger on its own after a failure. + enum: [ui, cli, api, worker] + username: + type: string + nullable: true + description: Unset when the server acted on its own. + created_at: + type: string + format: date-time + changes: + type: object + nullable: true + additionalProperties: true + description: "{field: {old, new}} for the fields that actually changed. Unset for a delete." + required: + - id + - trigger_kind + - path + - operation + - source + - created_at + Schedule: type: object properties: diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index ce709d3121..44d1d7cb2b 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -178,6 +178,7 @@ mod teams_oss; mod token; mod tracing_init; mod trash; +mod trigger_history; pub mod triggers; mod users; #[cfg(feature = "private")] @@ -280,6 +281,23 @@ async fn set_deploy_origin( windmill_common::deploy_origin::scope(origin, next.run(req)).await } +/// Scope the request in the client kind it declares, so a trigger mutation can +/// be attributed to the CLI rather than to a bare API call. Entered for every +/// request, undeclared ones included: `TriggerSource::of_request` reads the +/// scope's absence as "no request is being served", which is what separates a +/// caller from a worker disabling a trigger on its own. +async fn set_request_client( + req: axum::extract::Request, + next: axum::middleware::Next, +) -> axum::response::Response { + let client = req + .headers() + .get(windmill_common::trigger_history::CLIENT_HEADER) + .and_then(|v| v.to_str().ok()) + .and_then(windmill_common::trigger_history::client_from_header); + windmill_common::trigger_history::scope_client(client, next.run(req)).await +} + #[cfg(not(feature = "tantivy"))] type IndexReader = (); @@ -639,6 +657,7 @@ pub async fn run_server( .nest("/folders_history", folder_history::workspaced_service()) .nest("/groups", groups::workspaced_service()) .nest("/groups_history", group_history::workspaced_service()) + .nest("/triggers_history", trigger_history::workspaced_service()) .nest("/inputs", windmill_api_inputs::workspaced_service()) .nest("/internal_db", internal_db::workspaced_service()) .route("/labels/list", get(list_workspace_labels)) @@ -1136,6 +1155,8 @@ pub async fn run_server( let app = app.layer(axum::middleware::from_fn(set_deploy_origin)); + let app = app.layer(axum::middleware::from_fn(set_request_client)); + let app = app.layer(CatchPanicLayer::custom(|err| { tracing::error!("panic in handler, returning 500: {:?}", err); Response::builder() diff --git a/backend/windmill-api/src/token.rs b/backend/windmill-api/src/token.rs index acb65d6cf0..fc808767ec 100644 --- a/backend/windmill-api/src/token.rs +++ b/backend/windmill-api/src/token.rs @@ -236,6 +236,26 @@ lazy_static! { }], }); + // Read-only: `trigger_history` is append-only and written by the server + // alone, so there is no `triggers_history:write`. Its own domain rather + // than a `schedules`/`*_triggers` alias: one listing spans every kind, + // and a history row quotes the whole trigger row (a schedule's `args` + // included), so reading it is an explicit grant rather than a side + // effect of being able to read the trigger. Path-selectable because the + // route filters rows by the caller's path grants. + groups.push(ScopeDomain { + name: "Trigger History".to_string(), + description: Some( + "Read-only access to the modification history of schedules and triggers" + .to_string(), + ), + scopes: vec![ScopeOption { + value: "triggers_history:read".to_string(), + label: "Read".to_string(), + requires_resource_path: true, + }], + }); + groups.extend(build_standard_scope_domains()); groups.extend(build_trigger_scope_domains()); diff --git a/backend/windmill-api/src/trash.rs b/backend/windmill-api/src/trash.rs index 416546b850..53fc088582 100644 --- a/backend/windmill-api/src/trash.rs +++ b/backend/windmill-api/src/trash.rs @@ -9,6 +9,9 @@ use windmill_common::{ db::UserDB, error::{Error, Result}, trashbin::{self, TrashItem, TrashItemWithData}, + trigger_history::{ + self, TriggerHistoryEvent, TriggerOperation, TriggerSource, SCHEDULE_TRIGGER_KIND, + }, utils::require_admin, }; @@ -88,6 +91,31 @@ async fn restore_trash_item( .execute(&mut *tx) .await?; + // A restore puts the trigger back, so the history has to say so: otherwise + // the last thing it records for a live trigger is its own deletion. The + // trashed row is the snapshot, so this needs no extra read. + let restored_trigger_kind = match item.item_kind.as_str() { + SCHEDULE_TRIGGER_KIND => Some(SCHEDULE_TRIGGER_KIND), + // `_trigger` is what `delete_trigger` trashes it under, and the + // stem is the `TRIGGER_TYPE` the history records against. + kind => kind.strip_suffix("_trigger"), + }; + if let Some(trigger_kind) = restored_trigger_kind { + trigger_history::record( + &mut *tx, + TriggerHistoryEvent { + workspace_id: &w_id, + trigger_kind, + path: &item.item_path, + operation: TriggerOperation::Create, + source: TriggerSource::of_request(authed.is_session_token), + username: Some(&authed.username), + changes: trigger_history::summarize_changes(None, item.item_data.get("row")), + }, + ) + .await?; + } + audit_log( &mut *tx, &authed, diff --git a/backend/windmill-api/src/trigger_history.rs b/backend/windmill-api/src/trigger_history.rs new file mode 100644 index 0000000000..75ca03b90f --- /dev/null +++ b/backend/windmill-api/src/trigger_history.rs @@ -0,0 +1,105 @@ +/* + * Author: Ruben Fiszel + * Copyright: Windmill Labs, Inc 2022 + * This file and its contents are licensed under the AGPLv3 License. + * Please see the included NOTICE for copyright information and + * LICENSE-AGPL for a copy of the license. + */ + +use axum::{ + extract::{Extension, Path, Query}, + routing::get, + Router, +}; +use serde::{Deserialize, Serialize}; +use windmill_api_auth::{build_scope_path_filter, check_scopes, ApiAuthed, ScopePathFilter}; +use windmill_common::{ + db::UserDB, + error::JsonResult, + utils::{paginate, Pagination}, +}; + +pub fn workspaced_service() -> Router { + Router::new().route("/list", get(list_trigger_history)) +} + +#[derive(Serialize)] +pub struct TriggerHistoryEntry { + pub id: i64, + pub trigger_kind: String, + pub path: String, + pub operation: String, + pub source: String, + pub username: Option, + pub created_at: chrono::DateTime, + pub changes: Option, +} + +#[derive(Deserialize)] +pub struct ListTriggerHistoryQuery { + pub page: Option, + pub per_page: Option, + /// `"schedule"` or a trigger type (`"http"`, `"kafka"`, …). + pub trigger_kind: Option, + pub path: Option, +} + +/// Two gates, because they answer different questions: the RLS policies on +/// `trigger_history` bound the rows to what the *user* may read, and +/// `triggers_history:read:` bounds them further to what this *token* may +/// read. Without the second, a token scoped to one path could read the diffs of +/// every trigger its user can see, and a `create` row quotes the whole trigger +/// row, a schedule's `args` included. +async fn list_trigger_history( + authed: ApiAuthed, + Extension(user_db): Extension, + Path(w_id): Path, + Query(query): Query, +) -> JsonResult> { + if let Some(path) = query.path.as_deref() { + check_scopes(&authed, || format!("triggers_history:read:{}", path))?; + } + + // In the WHERE, not a retain after the fetch: the result is paginated, and a + // post-fetch filter would let a page's size report how many rows the token + // may not read — and return short pages that read as "no history". + let (scope_all, scope_exact, scope_prefix) = + match build_scope_path_filter(&authed, "triggers_history", "read") { + ScopePathFilter::AllowAll => (true, Vec::new(), Vec::new()), + ScopePathFilter::Restricted { exact, prefix } => (false, exact, prefix), + }; + + let mut tx = user_db.begin(&authed).await?; + + let (per_page, offset) = paginate(Pagination { page: query.page, per_page: query.per_page }); + + let history = sqlx::query_as!( + TriggerHistoryEntry, + "SELECT id, trigger_kind, path, operation, source, username, created_at, changes + FROM trigger_history + WHERE workspace_id = $1 + AND ($2::TEXT IS NULL OR trigger_kind = $2) + AND ($3::TEXT IS NULL OR path = $3) + AND ( $6 + OR path = ANY($7) + OR EXISTS ( SELECT 1 FROM unnest($8::text[]) AS pfx + WHERE path = pfx + OR left(path, length(pfx) + 1) = pfx || '/' ) ) + ORDER BY id DESC + LIMIT $4 OFFSET $5", + w_id, + query.trigger_kind, + query.path, + per_page as i64, + offset as i64, + scope_all, + &scope_exact[..], + &scope_prefix[..], + ) + .fetch_all(&mut *tx) + .await?; + + tx.commit().await?; + + Ok(axum::Json(history)) +} diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index a8de6d74c5..46af469725 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -121,6 +121,7 @@ pub mod teams_ee; pub mod teams_oss; pub mod tracing_init; pub mod trashbin; +pub mod trigger_history; pub mod triggers; pub mod user_drafts; pub mod usernames; diff --git a/backend/windmill-common/src/trigger_history.rs b/backend/windmill-common/src/trigger_history.rs new file mode 100644 index 0000000000..70c86b9ffd --- /dev/null +++ b/backend/windmill-common/src/trigger_history.rs @@ -0,0 +1,551 @@ +/* + * Author: Ruben Fiszel + * Copyright: Windmill Labs, Inc 2022 + * This file and its contents are licensed under the AGPLv3 License. + * Please see the included NOTICE for copyright information and + * LICENSE-AGPL for a copy of the license. + */ + +//! Append-only history of schedule and trigger mutations (`trigger_history`). +//! +//! Every field of a row is derived by the server at write time: the caller +//! passes what it is doing, never who it claims to be or where it claims to +//! come from. **Who** (the authed username, or nobody for a server-initiated +//! change) and **what** (a field-level diff computed from the row before and +//! after the write) are derived by the server and cannot be forged. **From what +//! kind of client** ([`TriggerSource`]) is weaker on purpose: a first-party +//! client declares itself in a header, so it attributes rather than proves — +//! see [`TriggerSource::of_request`]. +//! +//! # What is recorded +//! +//! Authoring a single trigger through its own surface — create, update, delete, +//! enable/disable/suspend, restore from the trashbin, and the workspace-wide +//! default-handler override — plus the server disabling one after a failure. +//! **Adding a route that authors a trigger means adding a `record` call to it**; +//! nothing enforces that, because the alternative (a database trigger) cannot +//! see who or which client asked, and would fire on every listener ping. +//! +//! Deliberately outside that line, and not a gap to be closed one call site at a +//! time: +//! +//! - **Cascades of renaming or deleting something else** — a script/flow rename +//! rewriting `script_path` (`triggers::update_triggers_script_path`), a user +//! being removed rewriting ownership. The event belongs to the runnable or the +//! user, not to the trigger. +//! - **Workspace-level bulk operations** — archive, fork clone, cross-workspace +//! deploy. They move whole workspaces; a per-trigger row per path would say +//! nothing the workspace event does not. +//! - **Runtime housekeeping** — clearing `paused_until` / `error` after a run, +//! consumer-offset state (`reset_offset`, `server_id`), the managed +//! ducklake-maintenance schedule. The same category as the `server_id` and +//! `last_server_ping` columns the diff already drops. +//! +//! # The server-initiated disables: the disable wins +//! +//! When the server disables a trigger it could not run, two things want to be +//! true and cannot both be guaranteed: the trigger ends up disabled, and the +//! history says who disabled it. The disable wins, every time. +//! +//! A trigger left enabled reads as healthy while never firing again, and for a +//! flow schedule nothing comes back to retry — it arms its next occurrence when +//! the flow *starts*, so once the runnable is gone that code is never reached +//! again. Enabled-and-dead is silent; disabled-without-an-audit-row is not, and +//! the trigger's own `error` column still says why. +//! +//! So each writer puts the disabling `UPDATE` and the record in one +//! transaction, with only the insert inside a savepoint +//! ([`record_in_disable_tx`]). Both land on the same commit, and the trigger's +//! row lock is held across the pair — so the row cannot end up describing a +//! trigger deleted and recreated at that path in between. If the insert alone +//! fails it rolls back to the savepoint, the disable still commits, and the lost +//! row is reported to the workspace error handler and the critical alert +//! channel — loud, never silent. +//! +//! # Authorization contract +//! +//! None of the helpers here authorize anything: they take a connection and +//! write what they are given, exactly like `audit_log`. A caller must already +//! have authorized the mutation *and* performed it, and must derive `username` +//! from the request's `ApiAuthed` and `source` from +//! [`TriggerSource::of_request`] — never from anything the request body +//! carries. Reads are gated separately, by the RLS policies on the table and by +//! the token scopes the listing route checks. + +use sqlx::{Acquire, PgConnection}; + +use crate::error::Result; + +/// Header a first-party client sets to name itself. Only `cli`, `ui` and `api` +/// mean anything; any other value, and the header being absent, falls back to +/// what the credentials say. +pub const CLIENT_HEADER: &str = "x-windmill-client"; + +/// `trigger_kind` a schedule is recorded under. Triggers use their own +/// `TriggerCrud::TRIGGER_TYPE`. +pub const SCHEDULE_TRIGGER_KIND: &str = "schedule"; + +/// The kind of client a trigger mutation came from. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TriggerSource { + /// A browser session in the Windmill app. + Ui, + /// The `wmill` CLI (including the git-sync pull that shells out to it). + Cli, + /// A direct API call with a token: user scripts, CI, third-party clients. + Api, + /// No request at all: a worker or a trigger listener disabling something + /// after a failure. + Worker, +} + +impl TriggerSource { + pub fn as_str(&self) -> &'static str { + match self { + TriggerSource::Ui => "ui", + TriggerSource::Cli => "cli", + TriggerSource::Api => "api", + TriggerSource::Worker => "worker", + } + } + + fn from_client_header(value: &str) -> Option { + match value.trim().to_ascii_lowercase().as_str() { + "cli" => Some(TriggerSource::Cli), + "ui" => Some(TriggerSource::Ui), + "api" => Some(TriggerSource::Api), + _ => None, + } + } + + /// The source of the request currently being served. + /// + /// The declared client wins when it is one we know; otherwise the token + /// decides, and only the session token minted at browser login attributes + /// to the UI. Both inputs are attribution, never authority — nothing reads + /// a history row to make an access decision, so a caller lying about either + /// only mislabels its own row. + pub fn of_request(is_session_token: bool) -> Self { + match REQUEST_CLIENT.try_with(|client| *client) { + Ok(Some(source)) => source, + Ok(None) if is_session_token => TriggerSource::Ui, + Ok(None) => TriggerSource::Api, + // Outside a request there is no caller to attribute to. The + // server-initiated paths pass `Worker` themselves; this is what + // keeps a stray call from inventing one. + Err(_) => TriggerSource::Worker, + } + } +} + +/// What a mutation did to the trigger it is recorded against. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TriggerOperation { + Create, + Update, + Delete, + Enable, + Disable, + Suspend, +} + +impl TriggerOperation { + pub fn as_str(&self) -> &'static str { + match self { + TriggerOperation::Create => "create", + TriggerOperation::Update => "update", + TriggerOperation::Delete => "delete", + TriggerOperation::Enable => "enable", + TriggerOperation::Disable => "disable", + TriggerOperation::Suspend => "suspend", + } + } +} + +tokio::task_local! { + static REQUEST_CLIENT: Option; +} + +/// Run `f` with `client` as the declared client of every trigger mutation it +/// causes. Entered for every request, unmarked ones included, so that having no +/// scope at all means "not serving a request" — which is what +/// [`TriggerSource::Worker`] records. +pub async fn scope_client( + client: Option, + f: F, +) -> F::Output { + REQUEST_CLIENT.scope(client, f).await +} + +/// Parse the declared client of the request being served, if any. +pub fn client_from_header(value: &str) -> Option { + TriggerSource::from_client_header(value) +} + +/// Row fields that say nothing about the change itself: bookkeeping the history +/// row already carries, and listener runtime state that moves on its own. +const IGNORED_FIELDS: &[&str] = &[ + "workspace_id", + "edited_at", + "edited_by", + "extra_perms", + "last_server_ping", + "server_id", + // Listener runtime state like the two above: every trigger update clears it, + // so keeping it here would tag an ordinary edit with the failure it had + // before. The server-initiated disables put the error in `changes` + // themselves, so nothing is lost. + "error", + // Written from the requester on every schedule mutation, purely for workers + // that predate `permissioned_as`; it tracks the editor, not the schedule. + "email", +]; + +/// A `changes` payload bigger than this is replaced by the list of field names +/// it would have held. A schedule's `args` is caller-supplied and bounded only +/// by the API's request-size limit, and a history row is not worth a +/// multi-megabyte write. +const MAX_CHANGES_BYTES: usize = 32 * 1024; + +/// The row at `path` as JSON, or `None` when there is none — which, on an RLS +/// connection, also covers a row the caller cannot see. +/// +/// `FOR UPDATE`, so the preimage and the mutation that follows it see the same +/// row: without the lock another request can commit between the two, and its +/// change then lands in this caller's diff under this caller's name. +/// +/// Two things follow from taking the lock here rather than at the write: +/// +/// - The only row locked is the one the caller is about to write, and the +/// schedule paths reach the job queue only afterwards, so their documented +/// schedule-then-queue order is unchanged. +/// - The lock is held for whatever the caller does before its own `UPDATE`. For +/// `TriggerCrud::update_trigger` that includes the impl's external work — the +/// postgres impl opens a replication slot on a user-supplied host, the gcp and +/// azure impls call their subscription APIs — so a concurrent `setmode`, a +/// listener error write, or a script rename's bulk `script_path` update waits +/// on that call. Bounded by those APIs, not by us; the alternative is a +/// preimage inside each impl next to its own `UPDATE`. +/// +/// `table` is interpolated: pass a compile-time constant, never anything a +/// caller can reach. +pub async fn snapshot_row( + conn: &mut PgConnection, + table: &'static str, + workspace_id: &str, + path: &str, +) -> Result> { + // SAFETY: `table` is a compile-time constant. + let snapshot: Option = sqlx::query_scalar(&format!( + "SELECT to_jsonb(t) FROM {table} t WHERE workspace_id = $1 AND path = $2 FOR UPDATE" + )) + .bind(workspace_id) + .bind(path) + .fetch_optional(&mut *conn) + .await?; + Ok(snapshot) +} + +/// A field-level diff of two row snapshots, as `{field: {"old": …, "new": …}}`, +/// with `"old"` omitted where there is none to report. +/// +/// A create (`before` absent) keeps every non-null column of the new row, which +/// is its initial shape including whatever the column defaults supplied — +/// `to_jsonb` cannot tell a caller-set column from a defaulted one. Returns +/// `None` when nothing meaningful changed. +pub fn summarize_changes( + before: Option<&serde_json::Value>, + after: Option<&serde_json::Value>, +) -> Option { + let empty = serde_json::Map::new(); + let before = before.and_then(|v| v.as_object()).unwrap_or(&empty); + let after = after.and_then(|v| v.as_object())?; + + // Names of the changed fields, and the running size of what has been cloned + // so far. Measured as it goes rather than by serializing the finished map: + // a caller-sized `args` would otherwise be cloned in full and then copied + // again just to learn it was too big. + let mut fields = Vec::new(); + let mut changes = serde_json::Map::new(); + let mut bytes = 0usize; + for (field, new_value) in after { + if IGNORED_FIELDS.contains(&field.as_str()) { + continue; + } + let old_value = before.get(field); + match old_value { + Some(old_value) if old_value == new_value => continue, + None if new_value.is_null() => continue, + _ => {} + } + fields.push(field.clone()); + if bytes <= MAX_CHANGES_BYTES { + bytes += json_len(new_value) + old_value.map_or(0, json_len) + field.len(); + } + if bytes > MAX_CHANGES_BYTES { + continue; + } + let mut entry = serde_json::Map::new(); + if let Some(old_value) = old_value { + entry.insert("old".to_string(), old_value.clone()); + } + entry.insert("new".to_string(), new_value.clone()); + changes.insert(field.clone(), serde_json::Value::Object(entry)); + } + + if fields.is_empty() { + return None; + } + if bytes > MAX_CHANGES_BYTES { + return Some(serde_json::json!({ "truncated_fields": fields })); + } + Some(serde_json::Value::Object(changes)) +} + +/// Serialized size of `value` without building the string for it. +fn json_len(value: &serde_json::Value) -> usize { + struct Counter(usize); + impl std::io::Write for Counter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0 += buf.len(); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let mut counter = Counter(0); + let _ = serde_json::to_writer(&mut counter, value); + counter.0 +} + +/// jsonb rejects `\u0000` inside a string, and `changes` quotes caller-supplied +/// text — a schedule's `args`, a worker's error message. One NUL anywhere in +/// there would fail the insert and cost the row. +fn strip_nuls(value: &mut serde_json::Value) { + match value { + serde_json::Value::String(s) if s.contains('\0') => *s = s.replace('\0', ""), + serde_json::Value::Array(items) => items.iter_mut().for_each(strip_nuls), + serde_json::Value::Object(map) => map.values_mut().for_each(strip_nuls), + _ => {} + } +} + +/// The last word on what reaches the column, applied at the write itself so a +/// hand-built `changes` (the server-initiated disables carry an error string of +/// unknown length and origin) gets it too, not just a computed diff. +fn cap_changes(changes: Option) -> Option { + let mut changes = changes?; + strip_nuls(&mut changes); + if json_len(&changes) <= MAX_CHANGES_BYTES { + return Some(changes); + } + let fields = changes + .as_object() + .map(|o| o.keys().cloned().collect::>()) + .unwrap_or_default(); + Some(serde_json::json!({ "truncated_fields": fields })) +} + +/// One trigger mutation, as it is about to be recorded. +#[derive(Clone)] +pub struct TriggerHistoryEvent<'a> { + pub workspace_id: &'a str, + /// `"schedule"`, or the trigger's `TRIGGER_TYPE` (`"http"`, `"kafka"`, …). + pub trigger_kind: &'a str, + pub path: &'a str, + pub operation: TriggerOperation, + pub source: TriggerSource, + /// `None` when the server acted on its own. + pub username: Option<&'a str>, + pub changes: Option, +} + +impl<'a> TriggerHistoryEvent<'a> { + /// The event for a trigger the server disabled on its own after a failure. + /// + /// `forced_state` is the column the disable wrote, in the same + /// `{field: {old, new}}` shape as a diff — the two disable paths write + /// different columns (`enabled` for a schedule, `mode` for a trigger). + /// + /// Record this only when the disabling `UPDATE` reported an affected row, + /// and only when that `UPDATE` was itself predicated on the trigger still + /// being enabled. The server reads the trigger long before it writes, so + /// without both the row describes a transition a user had already made. + pub fn server_disable( + workspace_id: &'a str, + trigger_kind: &'a str, + path: &'a str, + mut forced_state: serde_json::Value, + error: &str, + ) -> Self { + if let Some(obj) = forced_state.as_object_mut() { + obj.insert("error".to_string(), serde_json::json!({ "new": error })); + } + Self { + workspace_id, + trigger_kind, + path, + operation: TriggerOperation::Disable, + source: TriggerSource::Worker, + username: None, + changes: Some(forced_state), + } + } +} + +/// Record a disable inside the transaction that made it, without letting a +/// failed insert take the disable down with it. +/// +/// The caller's `UPDATE` holds the trigger's row lock until that transaction +/// commits, and this runs inside that window — so the row cannot end up +/// describing a trigger that was deleted and recreated at the same path in +/// between, which is the whole point of doing it here rather than on a second +/// connection afterwards. +/// +/// The insert itself goes in a savepoint. If it fails it rolls back alone, the +/// caller still commits the disable, and the reason comes back here so the +/// caller can alert: a trigger left enabled reads as healthy while never firing +/// again, which is worse than a missing audit row. +pub async fn record_in_disable_tx( + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + event: TriggerHistoryEvent<'_>, +) -> Option { + let mut savepoint = match tx.begin().await { + Ok(savepoint) => savepoint, + Err(e) => return Some(e.to_string()), + }; + match record(&mut savepoint, event).await { + Ok(()) => savepoint.commit().await.err().map(|e| e.to_string()), + Err(e) => { + savepoint.rollback().await.ok(); + Some(e.to_string()) + } + } +} + +/// Append `event` to the history. +/// +/// Pass the same connection as the mutation for the two to commit together. +/// Does not authorize — see the module docs. +pub async fn record(conn: &mut PgConnection, event: TriggerHistoryEvent<'_>) -> Result<()> { + sqlx::query!( + "INSERT INTO trigger_history + (workspace_id, trigger_kind, path, operation, source, username, changes) + VALUES ($1, $2, $3, $4, $5, $6, $7)", + event.workspace_id, + event.trigger_kind, + event.path, + event.operation.as_str(), + event.source.as_str(), + event.username, + cap_changes(event.changes) as _, + ) + .execute(&mut *conn) + .await?; + Ok(()) +} + +/// Append one row per path, all describing the same change. +/// +/// For the workspace-wide operations that rewrite every schedule at once, where +/// a per-path diff would cost a snapshot per row and say the same thing each +/// time. Does not authorize — see the module docs. +pub async fn record_bulk( + conn: &mut PgConnection, + workspace_id: &str, + trigger_kind: &str, + paths: &[String], + operation: TriggerOperation, + source: TriggerSource, + username: Option<&str>, + changes: Option, +) -> Result<()> { + if paths.is_empty() { + return Ok(()); + } + sqlx::query!( + "INSERT INTO trigger_history + (workspace_id, trigger_kind, path, operation, source, username, changes) + SELECT $1, $2, p, $3, $4, $5, $6 FROM unnest($7::text[]) p", + workspace_id, + trigger_kind, + operation.as_str(), + source.as_str(), + username, + cap_changes(changes) as _, + paths, + ) + .execute(&mut *conn) + .await?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + /// The whole worker side of the attribution rests on this: a mutation made + /// outside a request records `worker` without each call site saying so. + #[tokio::test] + async fn client_is_absent_outside_a_request() { + assert_eq!(TriggerSource::of_request(false), TriggerSource::Worker); + assert_eq!( + scope_client(None, async { TriggerSource::of_request(true) }).await, + TriggerSource::Ui + ); + assert_eq!( + scope_client(None, async { TriggerSource::of_request(false) }).await, + TriggerSource::Api + ); + assert_eq!( + scope_client(Some(TriggerSource::Cli), async { + TriggerSource::of_request(true) + }) + .await, + TriggerSource::Cli + ); + } + + /// `error` and `edited_at` stand in for the whole ignore list: every trigger + /// update clears `error`, so without it an ordinary edit would carry the + /// failure the trigger had before it. + #[test] + fn diff_keeps_only_what_changed() { + let before = json!({"schedule": "0 0 * * *", "enabled": true, "edited_at": "a", "error": "boom"}); + let after = json!({"schedule": "0 1 * * *", "enabled": true, "edited_at": "b", "error": null}); + assert_eq!( + summarize_changes(Some(&before), Some(&after)), + Some(json!({"schedule": {"old": "0 0 * * *", "new": "0 1 * * *"}})) + ); + assert_eq!(summarize_changes(Some(&before), Some(&before)), None); + } + + #[test] + fn create_drops_null_columns_and_bookkeeping() { + let after = json!({"schedule": "0 0 * * *", "summary": null, "workspace_id": "w"}); + assert_eq!( + summarize_changes(None, Some(&after)), + Some(json!({"schedule": {"new": "0 0 * * *"}})) + ); + } + + /// A NUL reaching the column fails the insert, and `changes` quotes + /// caller-supplied text — so this is the difference between a recorded + /// disable and a lost one. + #[test] + fn nul_bytes_never_reach_the_column() { + let changes = cap_changes(Some(json!({ "error": { "new": "boom\u{0}tail" } }))); + assert_eq!(changes, Some(json!({ "error": { "new": "boomtail" } }))); + } + + #[test] + fn oversized_changes_keep_the_field_names() { + let after = json!({ "args": "x".repeat(MAX_CHANGES_BYTES + 1) }); + assert_eq!( + summarize_changes(None, Some(&after)), + Some(json!({"truncated_fields": ["args"]})) + ); + } +} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 2b4ebe8bfd..ec84faa2d0 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -2503,6 +2503,68 @@ pub async fn send_success_to_workspace_handler<'a, 'c, T: Serialize + Send + Syn Ok(()) } +/// The event for a schedule the server disabled on its own. +pub fn schedule_auto_disable_event<'a>( + workspace_id: &'a str, + path: &'a str, + error: &str, +) -> windmill_common::trigger_history::TriggerHistoryEvent<'a> { + windmill_common::trigger_history::TriggerHistoryEvent::server_disable( + workspace_id, + windmill_common::trigger_history::SCHEDULE_TRIGGER_KIND, + path, + serde_json::json!({ "enabled": { "old": true, "new": false } }), + error, + ) +} + +/// Disable a schedule the server can no longer arm, and record that it did. +/// +/// Contract on `record_in_disable_tx`. Here `tx` is the job-completion +/// transaction, so the savepoint also keeps a failed insert from poisoning it. +/// +/// Returns `Err` only when the disable itself failed; a lost history row comes +/// back through `history_lost` for the caller to report. +async fn disable_schedule_and_record( + tx: &mut Transaction<'_, Postgres>, + schedule: &Schedule, + err: &Error, + history_lost: &mut Option, +) -> Result { + let disable_result = sqlx::query!( + "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3 AND enabled = true", + err.to_string(), + &schedule.workspace_id, + &schedule.path + ) + .execute(&mut **tx) + .await; + + #[cfg(feature = "failpoints")] + let disable_result = if schedule_failpoints::is_active( + schedule_failpoints::ScheduleFailPoint::ScheduleDisable, + ) { + Err(sqlx::Error::Protocol( + "failpoint: schedule disable".to_string(), + )) + } else { + disable_result + }; + + let rows = disable_result?.rows_affected(); + // Zero rows means a user disabled the schedule first: no transition of ours + // to record. + if rows == 0 { + return Ok(0); + } + + let event = + schedule_auto_disable_event(&schedule.workspace_id, &schedule.path, &err.to_string()); + *history_lost = windmill_common::trigger_history::record_in_disable_tx(tx, event).await; + + Ok(rows) +} + pub async fn try_schedule_next_job<'c>( db: &Pool, mut tx: Transaction<'c, Postgres>, @@ -2657,36 +2719,31 @@ pub async fn try_schedule_next_job<'c>( "Could not push next scheduled job for {}: {err}. Disabling schedule.", schedule.path ); - let disable_result = sqlx::query!( - "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3", - err.to_string(), - &schedule.workspace_id, - &schedule.path - ) - .execute(&mut *tx) - .await; - #[cfg(feature = "failpoints")] - let disable_result = if schedule_failpoints::is_active( - schedule_failpoints::ScheduleFailPoint::ScheduleDisable, - ) { - Err(sqlx::Error::Protocol( - "failpoint: schedule disable".to_string(), - )) - } else { - disable_result - }; - if let Err(disable_err) = disable_result { + let mut history_lost = None; + match disable_schedule_and_record(&mut tx, schedule, err, &mut history_lost).await { + Err(disable_err) => { + report_error_to_workspace_handler_or_critical_side_channel( + job, + db, + format!( + "Could not push next scheduled job for {} and could not disable schedule: {disable_err}", + schedule.path, + ), + ) + .await; + } + Ok(_) => push_err = None, + } + if let Some(history_err) = history_lost { report_error_to_workspace_handler_or_critical_side_channel( job, db, format!( - "Could not push next scheduled job for {} and could not disable schedule: {disable_err}", + "Disabled schedule {} but could not record it in the trigger history: {history_err}", schedule.path, ), ) .await; - } else { - push_err = None; } } } diff --git a/backend/windmill-trigger-http/src/handler.rs b/backend/windmill-trigger-http/src/handler.rs index 17321b7775..14e1b9a4d2 100644 --- a/backend/windmill-trigger-http/src/handler.rs +++ b/backend/windmill-trigger-http/src/handler.rs @@ -307,6 +307,37 @@ pub async fn create_many_http_triggers( .map_err(|err| error_wrapper(&new_http_trigger.config.route_path, err.into()))?; } + // Bulk create is still authoring, so it records like the single-create + // route rather than being the one way to make a trigger appear with no + // history behind it. + let created = windmill_common::trigger_history::snapshot_row( + &mut *tx, + "http_trigger", + &w_id, + &new_http_trigger.base.path, + ) + .await + .map_err(|err| error_wrapper(&new_http_trigger.config.route_path, err))?; + windmill_common::trigger_history::record( + &mut *tx, + windmill_common::trigger_history::TriggerHistoryEvent { + workspace_id: &w_id, + trigger_kind: HttpTrigger::TRIGGER_TYPE, + path: &new_http_trigger.base.path, + operation: windmill_common::trigger_history::TriggerOperation::Create, + source: windmill_common::trigger_history::TriggerSource::of_request( + authed.is_session_token, + ), + username: Some(&authed.username), + changes: windmill_common::trigger_history::summarize_changes( + None, + created.as_ref(), + ), + }, + ) + .await + .map_err(|err| error_wrapper(&new_http_trigger.config.route_path, err))?; + audit_log( &mut *tx, &authed, diff --git a/backend/windmill-trigger/src/handler.rs b/backend/windmill-trigger/src/handler.rs index ad3a6b097a..00d8cbf828 100644 --- a/backend/windmill-trigger/src/handler.rs +++ b/backend/windmill-trigger/src/handler.rs @@ -16,6 +16,7 @@ use windmill_api_auth::{build_scope_path_predicate, check_scopes, ApiAuthed}; use windmill_common::{ db::UserDB, error::{Error, JsonResult, Result}, + trigger_history::{self, TriggerHistoryEvent, TriggerOperation, TriggerSource}, user_drafts::{ delete_all_drafts_for_path, delete_own_draft_for_path, fetch_draft_only_list_rows, overlay_or_draft_only, UserDraftItemKind, WithDraftOverlay, WithDraftQuery, @@ -458,6 +459,45 @@ pub trait TriggerCrud: Send + Sync + 'static { } } +/// Append this mutation to `trigger_history`, diffing the row at `path` against +/// `before`. +/// +/// Call it on the transaction that made the change, after the change: the +/// snapshot it takes is the "after" side of the diff, and the two commit or roll +/// back together. +/// +/// Records nothing when the snapshots say no row was written. `TriggerCrud::update_trigger` +/// returns `Result<()>` and several impls do not check `rows_affected`, so an +/// update aimed at a path that does not exist — or that RLS hides from the +/// caller — reaches here having changed nothing; without this the caller could +/// forge history rows at any path, since the insert policy is `WITH CHECK (true)`. +async fn record_trigger_history( + tx: &mut PgConnection, + authed: &ApiAuthed, + workspace_id: &str, + path: &str, + operation: TriggerOperation, + before: Option, +) -> Result<()> { + let after = trigger_history::snapshot_row(&mut *tx, T::TABLE_NAME, workspace_id, path).await?; + if after.is_none() || (operation == TriggerOperation::Update && before.is_none()) { + return Ok(()); + } + trigger_history::record( + &mut *tx, + TriggerHistoryEvent { + workspace_id, + trigger_kind: T::TRIGGER_TYPE, + path, + operation, + source: TriggerSource::of_request(authed.is_session_token), + username: Some(&authed.username), + changes: trigger_history::summarize_changes(before.as_ref(), after.as_ref()), + }, + ) + .await +} + pub fn trigger_routes() -> Router { let mut router = Router::new() .route("/create", post(create_trigger::)) @@ -556,6 +596,16 @@ async fn create_trigger( .await?; } + record_trigger_history::( + &mut *tx, + &authed, + &workspace_id, + &new_path, + TriggerOperation::Create, + None, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -782,6 +832,9 @@ async fn update_trigger( &authed.username, ); + let before = + trigger_history::snapshot_row(&mut *tx, T::TABLE_NAME, &workspace_id, path).await?; + handler .update_trigger(&db, &mut *tx, &authed, &workspace_id, path, edit_trigger) .await?; @@ -799,6 +852,18 @@ async fn update_trigger( .await?; } + // Recorded at the new path, so a rename reads as one event there with + // `path` among the changed fields rather than a delete plus a create. + record_trigger_history::( + &mut *tx, + &authed, + &workspace_id, + &new_path, + TriggerOperation::Update, + before, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -913,6 +978,22 @@ async fn delete_trigger( .await?; } + // No diff: the row is gone, and the trashbin above already keeps its full + // contents for a restore. + trigger_history::record( + &mut *tx, + TriggerHistoryEvent { + workspace_id: &workspace_id, + trigger_kind: T::TRIGGER_TYPE, + path, + operation: TriggerOperation::Delete, + source: TriggerSource::of_request(authed.is_session_token), + username: Some(&authed.username), + changes: None, + }, + ) + .await?; + audit_log( &mut *tx, &authed, @@ -1052,6 +1133,9 @@ async fn set_trigger_mode( } } + let before = + trigger_history::snapshot_row(&mut *tx, T::TABLE_NAME, &workspace_id, path).await?; + let updated = handler .set_trigger_mode(&authed, &mut *tx, &workspace_id, path, &payload.mode) .await?; @@ -1063,6 +1147,20 @@ async fn set_trigger_mode( ))); } + record_trigger_history::( + &mut *tx, + &authed, + &workspace_id, + path, + match payload.mode { + TriggerMode::Enabled => TriggerOperation::Enable, + TriggerMode::Disabled => TriggerOperation::Disable, + TriggerMode::Suspended => TriggerOperation::Suspend, + }, + before, + ) + .await?; + tx.commit().await?; handle_deployment_metadata( diff --git a/backend/windmill-trigger/src/listener.rs b/backend/windmill-trigger/src/listener.rs index e327ed3155..d2c80b4d65 100644 --- a/backend/windmill-trigger/src/listener.rs +++ b/backend/windmill-trigger/src/listener.rs @@ -386,9 +386,14 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { error: String, ) { if listening_trigger.trigger_mode { - // SAFETY: Self::TABLE_NAME is a compile-time constant. - let report_status = sqlx::query(&format!( - r#" + // Contract on `record_in_disable_tx`: one transaction so the row + // lock spans both writes. + let mut history_err = None; + let report_status = async { + let mut tx = db.begin().await?; + // SAFETY: Self::TABLE_NAME is a compile-time constant. + let rows = sqlx::query(&format!( + r#" UPDATE {} SET @@ -398,18 +403,60 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { last_server_ping = NULL WHERE workspace_id = $2 AND - path = $3 + path = $3 AND + mode <> 'disabled'::TRIGGER_MODE "#, - Self::TABLE_NAME - )) - .bind(&error) - .bind(&listening_trigger.workspace_id) - .bind(&listening_trigger.path) - .execute(db) + Self::TABLE_NAME + )) + .bind(&error) + .bind(&listening_trigger.workspace_id) + .bind(&listening_trigger.path) + .execute(&mut *tx) + .await? + .rows_affected(); + + // Zero rows: deleted, or a user disabled it first — no + // transition of ours to record. + if rows > 0 { + // `to_key`, not `Display`: it is what lines up with the + // `TRIGGER_TYPE` the API records under. + let trigger_kind = Self::TRIGGER_KIND.to_key(); + history_err = windmill_common::trigger_history::record_in_disable_tx( + &mut tx, + windmill_common::trigger_history::TriggerHistoryEvent::server_disable( + &listening_trigger.workspace_id, + &trigger_kind, + &listening_trigger.path, + serde_json::json!({ "mode": { "new": "disabled" } }), + &error, + ), + ) + .await; + } + tx.commit().await?; + Ok::<(), Error>(()) + } .await; + if let Some(history_err) = history_err { + // Spawned: the commit above made the cleared `server_id` visible, + // so the ping branch of the enclosing `select!` is about to + // finish and drop everything left in this future. Awaiting the + // alert here would lose the one signal that the row is missing. + let message = format!( + "Disabled {} trigger {} but could not record it in the trigger history: {}", + Self::TRIGGER_KIND, + listening_trigger.path, + history_err + ); + let (db, workspace_id) = (db.clone(), listening_trigger.workspace_id.clone()); + tokio::spawn(async move { + report_critical_error(message, db, Some(&workspace_id), None).await; + }); + } + match report_status { - Ok(_) => { + Ok(()) => { report_critical_error( format!( "Disabling {} trigger {} because of error: {}", diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 7066b35945..42246c6358 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -2923,15 +2923,39 @@ pub async fn handle_flow( // its own disable write failed. Retry it: rearm_schedule turns // these into NoOp, so without disabling here the schedule would // stay enabled yet never run. - if let Err(disable_err) = sqlx::query!( - "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3", - err.to_string(), - &flow_job.workspace_id, - &schedule.path - ) - .execute(db) - .await - { + // Contract on `record_in_disable_tx`. Worth knowing here: + // this is the last chance to disable, because a flow + // schedule arms its next occurrence when the flow *starts*, + // so once the flow is gone nothing reaches this code again. + let mut history_lost = None; + let disable_result = async { + let mut tx = db.begin().await?; + let rows = sqlx::query!( + "UPDATE schedule SET enabled = false, error = $1 WHERE workspace_id = $2 AND path = $3 AND enabled = true", + err.to_string(), + &flow_job.workspace_id, + &schedule.path + ) + .execute(&mut *tx) + .await? + .rows_affected(); + if rows > 0 { + history_lost = windmill_common::trigger_history::record_in_disable_tx( + &mut tx, + windmill_queue::jobs::schedule_auto_disable_event( + &flow_job.workspace_id, + &schedule.path, + &err.to_string(), + ), + ) + .await; + } + tx.commit().await?; + Ok::<(), Error>(()) + } + .await; + + if let Err(disable_err) = disable_result { report_error_to_workspace_handler_or_critical_side_channel( &mini_job, db, @@ -2942,6 +2966,17 @@ pub async fn handle_flow( ) .await; } + if let Some(history_err) = history_lost { + report_error_to_workspace_handler_or_critical_side_channel( + &mini_job, + db, + format!( + "Disabled schedule {} but could not record it in the trigger history: {history_err}", + schedule.path, + ), + ) + .await; + } } else { // Transient error (DB contention, timeout) after retry exhaustion: // not the schedule's fault. Report it but leave the schedule diff --git a/cli/src/core/client.ts b/cli/src/core/client.ts index f7073fe501..32776bd498 100644 --- a/cli/src/core/client.ts +++ b/cli/src/core/client.ts @@ -14,6 +14,17 @@ export function markRequestsAsSyncOrigin() { OpenAPI.HEADERS = { ...existing, "X-Windmill-Deploy-Origin": "sync" }; } +/** + * Name this process as the CLI on every subsequent request, so a trigger the + * CLI created or disabled is attributed to `cli` rather than to a bare API + * call in `trigger_history`. Attribution only — nothing on the server grants + * anything on the strength of it. + */ +export function markRequestsAsCliClient() { + const existing = typeof OpenAPI.HEADERS === "object" ? OpenAPI.HEADERS : {}; + OpenAPI.HEADERS = { ...existing, "X-Windmill-Client": "cli" }; +} + export function setClient(token?: string, baseUrl?: string) { if (baseUrl === undefined) { baseUrl = process.env["BASE_INTERNAL_URL"] ?? diff --git a/cli/src/main.ts b/cli/src/main.ts index ff8c576ee9..ccce20da30 100755 --- a/cli/src/main.ts +++ b/cli/src/main.ts @@ -32,6 +32,7 @@ import { OpenAPI } from "../gen/index.ts"; import { getHeaders } from "./utils/utils.ts"; import { detectAuthGatewayChallenge } from "./utils/http_guards.ts"; import { setShowDiffs } from "./core/conf.ts"; +import { markRequestsAsCliClient } from "./core/client.ts"; import { NpmProvider } from "./utils/upgrade.ts"; import { pull as hubPull } from "./commands/hub/hub.ts"; import { pull, push } from "./commands/sync/sync.ts"; @@ -300,6 +301,7 @@ async function main() { if (extraHeaders) { OpenAPI.HEADERS = extraHeaders; } + markRequestsAsCliClient(); OpenAPI.interceptors.response.use(async (response) => { await detectAuthGatewayChallenge(response); return response; diff --git a/frontend/src/lib/components/triggers/TriggerEditorToolbar.svelte b/frontend/src/lib/components/triggers/TriggerEditorToolbar.svelte index f3c3600543..72f0d83fa0 100644 --- a/frontend/src/lib/components/triggers/TriggerEditorToolbar.svelte +++ b/frontend/src/lib/components/triggers/TriggerEditorToolbar.svelte @@ -5,7 +5,7 @@ import { Tooltip } from '../meltComponents' import DeleteTriggerButton from './DeleteTriggerButton.svelte' - import { type Trigger } from './utils' + import { type Trigger, type TriggerType } from './utils' import TriggerSuspendedJobsModal from './TriggerSuspendedJobsModal.svelte' import type { TriggerMode } from '$lib/gen' import TriggerModeToggle from './TriggerModeToggle.svelte' @@ -14,6 +14,7 @@ import { pageDrawerSessionSource } from '../sessions/pageDrawerSession' import { page } from '$app/state' import { workspaceStore } from '$lib/stores' + import TriggerHistoryButton from './TriggerHistoryButton.svelte' interface Props { saveDisabled: any @@ -35,6 +36,10 @@ /** Path of the trigger being edited, used to deep-link "Open in AI session" * at this trigger. Empty while creating one. */ triggerPath?: string + /** Kind the modification history is recorded under. Only the schedule + * editor has to pass it: every other editor renders with a `trigger`, + * whose `type` is the same value. */ + triggerKind?: TriggerType } let { @@ -54,7 +59,8 @@ trigger, suspendedJobsModal, disableSuspendedMode = false, - triggerPath + triggerPath, + triggerKind }: Props = $props() const canSave = $derived((permissions === 'write' && edit) || permissions === 'create') @@ -77,10 +83,20 @@ ) : undefined ) + + // Only a deployed trigger has a history: a draft has never been written. + // `triggerKind` is what opts an editor in, so the kinds `trigger_history` + // does not record (native triggers) simply never pass it. + const historyPath = $derived( + triggerKind && edit && !trigger?.isDraft ? triggerPath || trigger?.path : undefined + ) {#if !allowDraft} {@render extra?.()} + {#if triggerKind && historyPath} + + {/if} {#if edit} + {#if triggerKind && historyPath} + + {/if} {#if !trigger?.draftConfig}
diff --git a/frontend/src/lib/components/triggers/TriggerHistoryButton.svelte b/frontend/src/lib/components/triggers/TriggerHistoryButton.svelte new file mode 100644 index 0000000000..f8bf4c7bb4 --- /dev/null +++ b/frontend/src/lib/components/triggers/TriggerHistoryButton.svelte @@ -0,0 +1,122 @@ + + + + drawer?.closeDrawer()}> + {#if loading} +
+ {#each new Array(3) as _, i (i)} + + {/each} +
+ {:else if error} +

{error}

+ {:else if !entries || entries.length === 0} +

No modification recorded yet

+ {:else} +
+ {#each entries as entry (entry.id)} +
+
+ {entry.operation} + {entry.source} + + {entry.username ?? 'Windmill'} + + {displayDate(entry.created_at)} +
+ {#if entry.changes} + + {/if} +
+ {/each} +
+ {#if entries.length === PER_PAGE} +

+ Showing the {PER_PAGE} most recent modifications +

+ {/if} + {/if} +
+
+ + + {:else} + + {#if change.kind === 'changed' || (change.kind === 'removed' && change.prev !== undefined)} + {formatScalar(change.prev)} + + {/if} + {#if change.kind === 'removed'} + unset + {:else} + {formatScalar(change.next)} + {/if} + + {/if} + {#if tree && expanded[change.field]} +
+ {#if change.kind !== 'added' && change.prev !== undefined} +
before
+ + {/if} + {#if change.kind !== 'removed'} +
after
+ + {/if} +
+ {/if} + {/each} +
+{/if} diff --git a/frontend/src/lib/components/triggers/amqp/AmqpTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/amqp/AmqpTriggerEditorInner.svelte index 1ddf2db2b2..9fb44491bb 100644 --- a/frontend/src/lib/components/triggers/amqp/AmqpTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/amqp/AmqpTriggerEditorInner.svelte @@ -445,6 +445,7 @@ {#snippet actionsSnippet()} {#if !drawerLoading} { + it('classifies each side by whether it holds a value', () => { + expect( + parseChanges({ + schedule: { old: '0 0 * * *', new: '0 1 * * *' }, + summary: { new: 'created' }, + on_failure: { old: 'u/admin/handler', new: null }, + description: { old: null, new: 'set now' } + }) + ).toEqual({ + kind: 'fields', + changes: [ + { kind: 'added', field: 'description', next: 'set now' }, + { kind: 'removed', field: 'on_failure', prev: 'u/admin/handler' }, + { kind: 'changed', field: 'schedule', prev: '0 0 * * *', next: '0 1 * * *' }, + { kind: 'added', field: 'summary', next: 'created' } + ] + }) + }) + + it('reads the truncation marker and the empty cases', () => { + expect(parseChanges({ truncated_fields: ['args', 'schedule'] })).toEqual({ + kind: 'truncated', + fields: ['args', 'schedule'] + }) + expect(parseChanges(null)).toEqual({ kind: 'none' }) + expect(parseChanges({})).toEqual({ kind: 'none' }) + }) +}) diff --git a/frontend/src/lib/components/triggers/triggerHistoryChanges.ts b/frontend/src/lib/components/triggers/triggerHistoryChanges.ts new file mode 100644 index 0000000000..b9cf98932b --- /dev/null +++ b/frontend/src/lib/components/triggers/triggerHistoryChanges.ts @@ -0,0 +1,70 @@ +/** + * The `changes` column of `trigger_history`, in the three shapes the backend + * writes it (see `windmill-common/src/trigger_history.rs`): + * + * null no field-level diff — a delete + * { truncated_fields: [...] } the diff exceeded the 32 KiB cap + * { : { old?, new } } the fields that changed + * + * Parsing it into this closed set up front is what lets the viewer lay every + * entry out the same way, instead of rendering whatever tree came back. + */ + +/** + * How one field moved. A missing or null side is treated as the field having no + * value: the writer omits `old` where it never read one, and the backend writes + * `new: null` for a cleared column, so both collapse to added/removed. + */ +export type FieldChange = + | { kind: 'added'; field: string; next: unknown } + | { kind: 'removed'; field: string; prev: unknown | undefined } + | { kind: 'changed'; field: string; prev: unknown; next: unknown } + +export type ParsedChanges = + | { kind: 'none' } + | { kind: 'truncated'; fields: string[] } + | { kind: 'fields'; changes: FieldChange[] } + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +export function parseChanges(changes: unknown): ParsedChanges { + if (!isRecord(changes)) return { kind: 'none' } + + const truncated = changes['truncated_fields'] + if (Array.isArray(truncated)) { + return { kind: 'truncated', fields: truncated.map(String) } + } + + const parsed: FieldChange[] = [] + for (const [field, entry] of Object.entries(changes)) { + if (!isRecord(entry)) continue + // A null on either side is the field having no value, which reads as + // added/removed rather than as a change to or from `null`. + const prev = entry['old'] ?? undefined + const next = entry['new'] ?? undefined + if (next === undefined) { + parsed.push({ kind: 'removed', field, prev }) + } else if (prev === undefined) { + parsed.push({ kind: 'added', field, next }) + } else { + parsed.push({ kind: 'changed', field, prev, next }) + } + } + // Alphabetical: jsonb hands back its own key order, which is by length then + // bytes and would shuffle as values change. + parsed.sort((a, b) => a.field.localeCompare(b.field)) + + return parsed.length === 0 ? { kind: 'none' } : { kind: 'fields', changes: parsed } +} + +/** True when a value needs the tree viewer rather than a one-line rendering. */ +export function isComplex(value: unknown): boolean { + return typeof value === 'object' && value !== null +} + +/** One-line rendering of a scalar, quoted so `""` and `"0"` stay visible. */ +export function formatScalar(value: unknown): string { + return typeof value === 'string' ? JSON.stringify(value) : String(value) +} diff --git a/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte index 9910666fea..e255b7b552 100644 --- a/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte @@ -506,6 +506,7 @@ {#snippet actionsButtons()} {#if !drawerLoading}