diff --git a/backend/.sqlx/query-895a4feb0b3cc01ad711ce0eff8eef4f68cc3ecaf680c55a94175f6dfc2e2912.json b/backend/.sqlx/query-895a4feb0b3cc01ad711ce0eff8eef4f68cc3ecaf680c55a94175f6dfc2e2912.json new file mode 100644 index 0000000000..a498ecccd0 --- /dev/null +++ b/backend/.sqlx/query-895a4feb0b3cc01ad711ce0eff8eef4f68cc3ecaf680c55a94175f6dfc2e2912.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO dbt_column_edge (workspace_id, script_path, script_hash, job_id,\n parent_unique_id, parent_column, child_unique_id,\n child_column, lineage_kind)\n VALUES ($1, $2, NULL, $3, 'model.p.draft_src', 'raw', 'model.p.draft', 'clean', 'mod')", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "895a4feb0b3cc01ad711ce0eff8eef4f68cc3ecaf680c55a94175f6dfc2e2912" +} diff --git a/backend/.sqlx/query-1fb5590cb1fe706b0d2dffcd04a41bb2fce6012b812d16d37c895b2eb9aeca4c.json b/backend/.sqlx/query-c793b147014cb1b6e138aed1c2eb8b93b4a8383d8cf835a7d350056565669e7f.json similarity index 61% rename from backend/.sqlx/query-1fb5590cb1fe706b0d2dffcd04a41bb2fce6012b812d16d37c895b2eb9aeca4c.json rename to backend/.sqlx/query-c793b147014cb1b6e138aed1c2eb8b93b4a8383d8cf835a7d350056565669e7f.json index d58ac0c4e4..fee6801a6a 100644 --- a/backend/.sqlx/query-1fb5590cb1fe706b0d2dffcd04a41bb2fce6012b812d16d37c895b2eb9aeca4c.json +++ b/backend/.sqlx/query-c793b147014cb1b6e138aed1c2eb8b93b4a8383d8cf835a7d350056565669e7f.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "INSERT INTO dbt_node (workspace_id, script_path, script_hash, job_id, unique_id,\n resource_type, name, asset_path, raw_code, tags)\n VALUES ($1, $2, NULL, $3, 'model.p.draft', 'model', 'draft',\n 'u/a/wh/analytics/draft', 'select 3', '{}')", + "query": "INSERT INTO dbt_node (workspace_id, script_path, script_hash, job_id, unique_id,\n resource_type, name, asset_path, raw_code, tags)\n VALUES ($1, $2, NULL, $3, 'model.p.draft', 'model', 'draft',\n 'u/a/wh/analytics/draft', 'select 3', '{}'),\n ($1, $2, NULL, $3, 'model.p.draft_src', 'model', 'draft_src',\n 'u/a/wh/analytics/draft_src', 'select 4', '{}')", "describe": { "columns": [], "parameters": { @@ -12,5 +12,5 @@ }, "nullable": [] }, - "hash": "1fb5590cb1fe706b0d2dffcd04a41bb2fce6012b812d16d37c895b2eb9aeca4c" + "hash": "c793b147014cb1b6e138aed1c2eb8b93b4a8383d8cf835a7d350056565669e7f" } diff --git a/backend/.sqlx/query-754484aa4e2d676ddf6c0890fed39f096c5df4317068b31bcb9d4cae0d6bf233.json b/backend/.sqlx/query-f50c16f3a6070642f2c87b735974704643d51daa0080178c4fe49eebcc650ed9.json similarity index 50% rename from backend/.sqlx/query-754484aa4e2d676ddf6c0890fed39f096c5df4317068b31bcb9d4cae0d6bf233.json rename to backend/.sqlx/query-f50c16f3a6070642f2c87b735974704643d51daa0080178c4fe49eebcc650ed9.json index 84c27ca4a8..bdd6d618c2 100644 --- a/backend/.sqlx/query-754484aa4e2d676ddf6c0890fed39f096c5df4317068b31bcb9d4cae0d6bf233.json +++ b/backend/.sqlx/query-f50c16f3a6070642f2c87b735974704643d51daa0080178c4fe49eebcc650ed9.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "WITH live AS (\n SELECT * FROM (\n SELECT DISTINCT ON (s.path) s.path, s.hash\n FROM script s\n WHERE $5::text IS NULL AND s.workspace_id = $1 AND s.language = 'dbt'\n AND ($3::bigint IS NULL OR s.hash = $3)\n AND ($3::bigint IS NOT NULL OR (s.deleted = false AND s.archived = false))\n ORDER BY s.path, s.created_at DESC\n ) cur\n UNION ALL\n SELECT $5::text, $3::bigint WHERE $5::text IS NOT NULL\n ),\n chosen AS (\n SELECT CASE WHEN $4::uuid IS NOT NULL AND EXISTS (\n SELECT 1 FROM dbt_graph_snapshot g\n WHERE g.workspace_id = $1 AND g.job_id = $4)\n THEN $4::uuid\n ELSE '00000000-0000-0000-0000-000000000000'::uuid END AS job_id\n )\n SELECT p.asset_path AS \"from_path!\", e.parent_column AS \"from_column!\",\n c.asset_path AS \"to_path!\", e.child_column AS \"to_column!\",\n e.lineage_kind AS \"lineage_kind!\", e.script_path AS \"script_path!\",\n -- Same gate as the model's own SQL, applied in Rust beside it:\n -- a column graph is the shape of what the author WROTE, one\n -- level finer than the `ref()` graph, which is ungated only\n -- because it draws relations the caller already sees in\n -- `asset`.\n EXISTS (SELECT 1 FROM script sc\n WHERE sc.workspace_id = e.workspace_id\n AND sc.path = e.script_path\n AND sc.hash = e.script_hash) AS \"script_visible!\"\n FROM dbt_column_edge e\n JOIN live l ON l.path = e.script_path\n AND (e.script_hash = l.hash\n OR ($5::text IS NOT NULL AND l.hash IS NULL\n AND e.script_hash IS NULL))\n JOIN chosen ch ON ch.job_id = e.job_id\n JOIN dbt_node p ON p.workspace_id = e.workspace_id\n AND p.script_path = e.script_path\n AND (p.script_hash = e.script_hash\n OR ($5::text IS NOT NULL AND e.script_hash IS NULL\n AND p.script_hash IS NULL))\n AND p.job_id = ch.job_id\n AND p.unique_id = e.parent_unique_id\n JOIN dbt_node c ON c.workspace_id = e.workspace_id\n AND c.script_path = e.script_path\n AND (c.script_hash = e.script_hash\n OR ($5::text IS NOT NULL AND e.script_hash IS NULL\n AND c.script_hash IS NULL))\n AND c.job_id = ch.job_id\n AND c.unique_id = e.child_unique_id\n WHERE e.workspace_id = $1\n AND e.lineage_kind IN ('copy', 'mod')\n AND p.asset_path IS NOT NULL AND c.asset_path IS NOT NULL\n AND ($3::bigint IS NOT NULL OR $5::text IS NOT NULL OR EXISTS (\n SELECT 1 FROM asset a\n WHERE a.workspace_id = $1 AND a.kind = 'dbt'\n AND a.path = c.asset_path\n AND ($2::text IS NULL OR a.usage_path LIKE $2)))\n LIMIT $6", + "query": "WITH live AS (\n SELECT * FROM (\n SELECT DISTINCT ON (s.path) s.path, s.hash\n FROM script s\n WHERE $5::text IS NULL AND s.workspace_id = $1 AND s.language = 'dbt'\n AND ($3::bigint IS NULL OR s.hash = $3)\n AND ($3::bigint IS NOT NULL OR (s.deleted = false AND s.archived = false))\n ORDER BY s.path, s.created_at DESC\n ) cur\n UNION ALL\n SELECT $5::text, $3::bigint WHERE $5::text IS NOT NULL\n ),\n chosen AS (\n SELECT CASE WHEN $4::uuid IS NOT NULL AND EXISTS (\n SELECT 1 FROM dbt_graph_snapshot g\n WHERE g.workspace_id = $1 AND g.job_id = $4)\n THEN $4::uuid\n ELSE '00000000-0000-0000-0000-000000000000'::uuid END AS job_id\n )\n SELECT p.asset_path AS \"from_path!\", e.parent_column AS \"from_column!\",\n c.asset_path AS \"to_path!\", e.child_column AS \"to_column!\",\n e.lineage_kind AS \"lineage_kind!\", e.script_path AS \"script_path!\",\n -- Same gate as the model's own SQL, applied in Rust beside it:\n -- a column graph is the shape of what the author WROTE, one\n -- level finer than the `ref()` graph, which is ungated only\n -- because it draws relations the caller already sees in\n -- `asset`.\n --\n -- The NULL arm is not optional, for the reason the node query\n -- gives: a version-less row has no `script` row to ask and\n -- needs none, and `sc.hash = NULL` is never true — so without\n -- it an editor buffer's parse renders its columns and none of\n -- their lineage.\n (e.script_hash IS NULL OR EXISTS (\n SELECT 1 FROM script sc\n WHERE sc.workspace_id = e.workspace_id\n AND sc.path = e.script_path\n AND sc.hash = e.script_hash)) AS \"script_visible!\"\n FROM dbt_column_edge e\n JOIN live l ON l.path = e.script_path\n AND (e.script_hash = l.hash\n OR ($5::text IS NOT NULL AND l.hash IS NULL\n AND e.script_hash IS NULL))\n JOIN chosen ch ON ch.job_id = e.job_id\n JOIN dbt_node p ON p.workspace_id = e.workspace_id\n AND p.script_path = e.script_path\n AND (p.script_hash = e.script_hash\n OR ($5::text IS NOT NULL AND e.script_hash IS NULL\n AND p.script_hash IS NULL))\n AND p.job_id = ch.job_id\n AND p.unique_id = e.parent_unique_id\n JOIN dbt_node c ON c.workspace_id = e.workspace_id\n AND c.script_path = e.script_path\n AND (c.script_hash = e.script_hash\n OR ($5::text IS NOT NULL AND e.script_hash IS NULL\n AND c.script_hash IS NULL))\n AND c.job_id = ch.job_id\n AND c.unique_id = e.child_unique_id\n WHERE e.workspace_id = $1\n AND e.lineage_kind IN ('copy', 'mod')\n AND p.asset_path IS NOT NULL AND c.asset_path IS NOT NULL\n AND ($3::bigint IS NOT NULL OR $5::text IS NOT NULL OR EXISTS (\n SELECT 1 FROM asset a\n WHERE a.workspace_id = $1 AND a.kind = 'dbt'\n AND a.path = c.asset_path\n AND ($2::text IS NULL OR a.usage_path LIKE $2)))\n LIMIT $6", "describe": { "columns": [ { @@ -59,5 +59,5 @@ null ] }, - "hash": "754484aa4e2d676ddf6c0890fed39f096c5df4317068b31bcb9d4cae0d6bf233" + "hash": "f50c16f3a6070642f2c87b735974704643d51daa0080178c4fe49eebcc650ed9" } diff --git a/backend/windmill-api-assets/src/lib.rs b/backend/windmill-api-assets/src/lib.rs index edce5dc10c..5e4290115d 100644 --- a/backend/windmill-api-assets/src/lib.rs +++ b/backend/windmill-api-assets/src/lib.rs @@ -1541,10 +1541,17 @@ pub async fn asset_graph_for( -- level finer than the `ref()` graph, which is ungated only -- because it draws relations the caller already sees in -- `asset`. - EXISTS (SELECT 1 FROM script sc - WHERE sc.workspace_id = e.workspace_id - AND sc.path = e.script_path - AND sc.hash = e.script_hash) AS "script_visible!" + -- + -- The NULL arm is not optional, for the reason the node query + -- gives: a version-less row has no `script` row to ask and + -- needs none, and `sc.hash = NULL` is never true — so without + -- it an editor buffer's parse renders its columns and none of + -- their lineage. + (e.script_hash IS NULL OR EXISTS ( + SELECT 1 FROM script sc + WHERE sc.workspace_id = e.workspace_id + AND sc.path = e.script_path + AND sc.hash = e.script_hash)) AS "script_visible!" FROM dbt_column_edge e JOIN live l ON l.path = e.script_path AND (e.script_hash = l.hash diff --git a/backend/windmill-api-assets/tests/dbt_pinned_graph.rs b/backend/windmill-api-assets/tests/dbt_pinned_graph.rs index 7f4979a96f..e91ed08517 100644 --- a/backend/windmill-api-assets/tests/dbt_pinned_graph.rs +++ b/backend/windmill-api-assets/tests/dbt_pinned_graph.rs @@ -364,7 +364,25 @@ async fn seed_editor_graph(db: &Pool, job: uuid::Uuid) { r#"INSERT INTO dbt_node (workspace_id, script_path, script_hash, job_id, unique_id, resource_type, name, asset_path, raw_code, tags) VALUES ($1, $2, NULL, $3, 'model.p.draft', 'model', 'draft', - 'u/a/wh/analytics/draft', 'select 3', '{}')"#, + 'u/a/wh/analytics/draft', 'select 3', '{}'), + ($1, $2, NULL, $3, 'model.p.draft_src', 'model', 'draft_src', + 'u/a/wh/analytics/draft_src', 'select 4', '{}')"#, + WS, + PATH, + job + ) + .execute(db) + .await + .unwrap(); + // A version-less row's `script_hash` is NULL on both sides of every join and + // every visibility check, and `= NULL` is never true — so the column edges + // need the same NULL arm the node query has, or a buffer parse renders its + // columns and none of their lineage. + sqlx::query!( + "INSERT INTO dbt_column_edge (workspace_id, script_path, script_hash, job_id, + parent_unique_id, parent_column, child_unique_id, + child_column, lineage_kind) + VALUES ($1, $2, NULL, $3, 'model.p.draft_src', 'raw', 'model.p.draft', 'clean', 'mod')", WS, PATH, job @@ -410,6 +428,17 @@ async fn an_editor_graph_renders_only_through_its_own_job(db: Pool) { serde_json::json!(parse), "labelled as a graph of its own, so the editor can say where it came from" ); + assert_eq!( + body["dbt_column_edges"], + serde_json::json!([{ + "from_asset_path": "u/a/wh/analytics/draft_src", + "from_column": "raw", + "to_asset_path": "u/a/wh/analytics/draft", + "to_column": "clean", + "kind": "mod", + }]), + "and its column lineage, which the buffer parse is the whole point of: {body}" + ); assert!( !body.to_string().contains("select 1"), "and the deployed version's models are not mixed into it: {body}" diff --git a/backend/windmill-common/src/dbt_manifest.rs b/backend/windmill-common/src/dbt_manifest.rs index 616faa2305..a9f9213a4a 100644 --- a/backend/windmill-common/src/dbt_manifest.rs +++ b/backend/windmill-common/src/dbt_manifest.rs @@ -406,6 +406,13 @@ pub struct IngestedNode { /// from the engine's static analysis. `None` when the project did not ask /// for it or the engine wrote none. Beside `columns` rather than merged into /// it: that one is what the author DECLARED, and stays that. + /// + /// Skipped when absent, unlike its neighbours, because `graph_digest` + /// serializes these nodes: emitting `"column_schema":null` would change + /// every stored digest at once, and every dynamic run of a project that + /// never asked for the pass would store a full snapshot until its script is + /// redeployed. + #[serde(skip_serializing_if = "Option::is_none")] pub column_schema: Option, pub freshness: Option, /// The transform itself, for the graph to render. The copy taken at @@ -510,10 +517,11 @@ impl IngestedManifest { // stores a snapshot at all: parquet row order is the engine's and two // passes over one project must not read as two different graphs. // - // Direct kinds first, so the cap below spends the budget on the edges a - // column trace actually draws. `scan` is both the bulk of a wide - // project's lineage and the kind nothing renders, so a plain sort would - // let it evict the lineage this exists for. + // Direct kinds first, so what the truncation below gives up is `scan` — + // the bulk of a wide project's lineage and the kind nothing renders. + // The worker's own reader already prioritizes them while decoding, since + // the memory bound has to apply there; this is the same order for a + // caller that did not come through it, which is the agent-worker wire. edges.sort_by(|a, b| { is_direct(&b.lineage_kind) .cmp(&is_direct(&a.lineage_kind)) diff --git a/backend/windmill-worker/src/dbt_column_index.rs b/backend/windmill-worker/src/dbt_column_index.rs index bcba36df3e..3a3390b6b7 100644 --- a/backend/windmill-worker/src/dbt_column_index.rs +++ b/backend/windmill-worker/src/dbt_column_index.rs @@ -22,13 +22,15 @@ //! asks which engine it is beyond "has the flag" — a release that starts //! writing them is picked up with no change. +use std::collections::HashSet; use std::path::Path; +use std::time::Duration; use parquet::file::reader::{FileReader, SerializedFileReader}; use parquet::record::{Field, Row}; use uuid::Uuid; use windmill_common::dbt_manifest::{ - ColumnIndex, IndexedColumn, IngestedColumnEdge, MAX_COLUMN_EDGES, + is_direct, ColumnIndex, IndexedColumn, IngestedColumnEdge, MAX_COLUMN_EDGES, }; use windmill_common::error; use windmill_common::worker::Connection; @@ -50,9 +52,16 @@ const NODE_COLUMNS_PARQUET: &str = "dbt.node_columns.parquet"; /// Run the lineage pass and read what it produced. /// -/// Best-effort throughout: every failure returns `None` with a line in the job -/// log saying which one, because the graph without column lineage is exactly the -/// graph this project had before it asked for any. +/// Best-effort about the COMPILE and about the artifact: a wrong engine, a +/// failed analysis, a missing or unreadable parquet, or this phase outrunning +/// its budget all return `None` with a line in the job log saying which, because +/// the graph without column lineage is exactly the graph this project had before +/// it asked for any. +/// +/// NOT best-effort about the job: a cancellation, the job's own deadline or the +/// output ceiling are returned as `Err` and fail it. Swallowing those would let +/// a run that blew its timeout inside an optional annotation go on to publish a +/// graph and report success. pub(crate) async fn collect( p: &PreparedProject, descriptor: &DbtDescriptor, @@ -61,6 +70,44 @@ pub(crate) async fn collect( job_id: &Uuid, w_id: &str, conn: &Connection, + kept: &HashSet<&str>, +) -> error::Result> { + let Some(budget) = phase_budget(ctx) else { + return run_pass(p, descriptor, inv, ctx, job_id, w_id, conn, kept).await; + }; + match tokio::time::timeout( + budget, + run_pass(p, descriptor, inv, ctx, job_id, w_id, conn, kept), + ) + .await + { + Ok(r) => r, + Err(_) => { + append_logs( + job_id, + w_id, + format!( + "\nNo column lineage: the analysis pass did not finish within {}s, half of \ + what was left of this job's time. The build below gets the rest.\n", + budget.as_secs() + ), + conn, + ) + .await; + Ok(None) + } + } +} + +async fn run_pass( + p: &PreparedProject, + descriptor: &DbtDescriptor, + inv: &Invocation, + ctx: &mut JobCtx<'_>, + job_id: &Uuid, + w_id: &str, + conn: &Connection, + kept: &HashSet<&str>, ) -> error::Result> { if !descriptor.column_lineage { return Ok(None); @@ -126,11 +173,10 @@ pub(crate) async fn collect( w_id, conn, CLL_MAX_OUTPUT_BYTES, - ctx_share(ctx), ) .await?; - let index = read_index(&index_dir).await; + let index = read_index(&index_dir, kept).await; // What the caller cannot say for itself. The COUNTS are logged where the // index is folded into the graph, since the graph is what decides how much // of it is kept; this is the part only the pass knows. @@ -180,8 +226,16 @@ const CLL_MAX_OUTPUT_BYTES: usize = 1 << 20; /// so an unbounded pass on a slow project would hand `dbt build` an expired /// budget and fail the run it exists only to annotate. Half leaves the build at /// least as long as the annotation was allowed to take. -fn ctx_share(ctx: &JobCtx<'_>) -> Option { - ctx.timeout().map(|left| (left / 2).max(1)) +/// +/// Spent as a race around the whole pass rather than as a shortened deadline +/// handed to the runner: the runner reports its expiry as an `Err`, which is +/// indistinguishable from a cancellation or the job's own deadline, and those +/// two MUST fail the job. Expiring here is this budget and nothing else, so it +/// degrades to no lineage. Dropping the future kills the child, which +/// `run_captured` spawns with `kill_on_drop`. +fn phase_budget(ctx: &JobCtx<'_>) -> Option { + ctx.timeout() + .map(|left| Duration::from_secs((left.max(0) as u64 / 2).max(1))) } /// The tail of what the engine said, bounded. The whole of it is every rendered @@ -202,15 +256,20 @@ fn diagnostics(out: &str) -> String { /// The column schemas alone are not worth a graph: they arrive with the lineage /// or not at all, and a node's declared columns already answer for the case /// where the pass never ran. -async fn read_index(index_dir: &Path) -> Option { +async fn read_index(index_dir: &Path, kept: &HashSet<&str>) -> Option { let lineage = index_dir.join(COLUMN_LINEAGE_PARQUET); if !tokio::fs::try_exists(&lineage).await.unwrap_or(false) { return None; } let columns = index_dir.join(NODE_COLUMNS_PARQUET); + // Owned, because the decode moves to a blocking thread. The index describes + // the whole project while this graph describes one selection of it, so + // scoping HERE is what keeps the cap below from being spent on rows the + // graph would discard anyway. + let kept: HashSet = kept.iter().map(|s| (*s).to_string()).collect(); // Decompressing and decoding a parquet is CPU work on a file the engine just // wrote, so it does not belong on the runtime's poll thread. - tokio::task::spawn_blocking(move || read_index_blocking(&lineage, &columns)) + tokio::task::spawn_blocking(move || read_index_blocking(&lineage, &columns, &kept)) .await .map_err(|e| tracing::warn!("reading the dbt column index: {e:#}")) .ok()? @@ -218,37 +277,58 @@ async fn read_index(index_dir: &Path) -> Option { .ok() } -fn read_index_blocking(lineage: &Path, columns: &Path) -> error::Result { +fn read_index_blocking( + lineage: &Path, + columns: &Path, + kept: &HashSet, +) -> error::Result { let mut out = ColumnIndex::default(); - for_each_row(lineage, MAX_COLUMN_EDGES, |row| { - let parent_unique_id = string(row, "from_node_unique_id"); - let child_unique_id = string(row, "to_node_unique_id"); - let parent_column = string(row, "from_column_name"); - let child_column = string(row, "to_column_name"); - // A column of a node the analysis could not name is not an endpoint the - // graph can draw. - if parent_unique_id.is_empty() - || child_unique_id.is_empty() - || parent_column.is_empty() - || child_column.is_empty() - { - return; + // Two passes, direct kinds first. The cap is a memory bound, so it has to + // apply while decoding — but applied to the file's own row order it would + // let `scan` edges, which are the bulk of a wide project's index and which + // nothing renders, fill the budget before a single `copy` edge is read. + // Reading a row and dropping it costs no memory, so the second pass is only + // time, on a file the engine just wrote. + for direct in [true, false] { + let remaining = MAX_COLUMN_EDGES - out.edges.len(); + if remaining == 0 { + break; } - out.edges.push(IngestedColumnEdge { - parent_unique_id, - parent_column, - child_unique_id, - child_column, - lineage_kind: string(row, "lineage_kind"), - }); - })?; + for_each_row(lineage, remaining, |row| { + let lineage_kind = string(row, "lineage_kind"); + if is_direct(&lineage_kind) != direct { + return Kept::No; + } + let parent_unique_id = string(row, "from_node_unique_id"); + let child_unique_id = string(row, "to_node_unique_id"); + let parent_column = string(row, "from_column_name"); + let child_column = string(row, "to_column_name"); + // A column of a node the analysis could not name is not an endpoint + // the graph can draw, and neither is one outside this graph's nodes. + if parent_column.is_empty() + || child_column.is_empty() + || !kept.contains(&parent_unique_id) + || !kept.contains(&child_unique_id) + { + return Kept::No; + } + out.edges.push(IngestedColumnEdge { + parent_unique_id, + parent_column, + child_unique_id, + child_column, + lineage_kind, + }); + Kept::Yes + })?; + } // Absent is normal — an engine can write the lineage table and not this one — // and unreadable is not worth losing the lineage over. let _ = for_each_row(columns, MAX_INDEXED_COLUMNS, |row| { let unique_id = string(row, "unique_id"); let name = string(row, "column_name"); - if unique_id.is_empty() || name.is_empty() { - return; + if name.is_empty() || !kept.contains(&unique_id) { + return Kept::No; } // The author's `data_type` where `schema.yml` gives one, since that is // what the project calls the column; the analysis's own inference @@ -265,6 +345,7 @@ fn read_index_blocking(lineage: &Path, columns: &Path) -> error::Result error::Result` first — each row holding -/// its own copy of every column NAME — is what would take the worker process -/// down, and this module's whole contract is that it cannot fail a deploy or a -/// run. -fn for_each_row(path: &Path, limit: usize, mut f: impl FnMut(&Row)) -> error::Result<()> { +/// Counted on what the caller keeps, not on what the file holds, because the +/// input this defends against is the one that cannot be collected: `scan` +/// lineage is emitted from every predicate and join column to every output +/// column, so a project shaped that way writes an index whose row count is +/// quadratic in its widest model. Materializing that into a `Vec` first — +/// each row holding its own copy of every column NAME — is what would take the +/// worker process down, and this module's whole contract is that it cannot fail +/// a deploy or a run. A row the caller skips costs nothing, so skipping is free +/// and only keeping is budgeted. +fn for_each_row(path: &Path, limit: usize, mut f: impl FnMut(&Row) -> Kept) -> error::Result<()> { let fail = |e: parquet::errors::ParquetError| { error::Error::internal_err(format!("reading {}: {e}", path.display())) }; let file = std::fs::File::open(path) .map_err(|e| error::Error::internal_err(format!("opening {}: {e}", path.display())))?; let reader = SerializedFileReader::new(file).map_err(fail)?; - for (kept, row) in reader.get_row_iter(None).map_err(fail)?.enumerate() { - if kept >= limit { + let mut budget = limit; + for row in reader.get_row_iter(None).map_err(fail)? { + if budget == 0 { tracing::warn!( - "dbt column index: {} holds more than {limit} rows; the rest is dropped", + "dbt column index: {} yielded more than {limit} usable rows; the rest is dropped", path.display() ); break; } - f(&row.map_err(fail)?); + if f(&row.map_err(fail)?) == Kept::Yes { + budget -= 1; + } } Ok(()) } +/// Whether the row the closure just saw was retained, which is what the budget +/// counts. +#[derive(PartialEq, Eq)] +enum Kept { + Yes, + No, +} + /// By NAME, not by position: these tables are the engine's own schema and it /// adds columns to them between releases. fn field<'a>(row: &'a Row, name: &str) -> Option<&'a Field> { diff --git a/backend/windmill-worker/src/dbt_executor.rs b/backend/windmill-worker/src/dbt_executor.rs index ec404c70c6..5efa8a1db6 100644 --- a/backend/windmill-worker/src/dbt_executor.rs +++ b/backend/windmill-worker/src/dbt_executor.rs @@ -3182,9 +3182,14 @@ async fn attach_column_index( w_id: &str, conn: &Connection, ) -> error::Result<()> { - let Some(index) = - crate::dbt_column_index::collect(p, descriptor, inv, ctx, job_id, w_id, conn).await? - else { + // The nodes this graph kept, so the pass reads only rows it could store: the + // index describes the whole project, this graph one selection of it. + let kept: std::collections::HashSet<&str> = + ingested.nodes.iter().map(|n| n.unique_id.as_str()).collect(); + let index = + crate::dbt_column_index::collect(p, descriptor, inv, ctx, job_id, w_id, conn, &kept).await?; + drop(kept); + let Some(index) = index else { return Ok(()); }; let found = index.edges.len(); @@ -3574,11 +3579,6 @@ pub(crate) async fn run_captured( w_id: &str, conn: &Connection, max_stdout_bytes: usize, - // How much of the job's wall clock this phase may spend. `None` is what is - // left of it, which is what a phase the job exists to run wants; a phase - // that only ANNOTATES the job passes less, so it cannot starve the one that - // does the work. - timeout_secs: Option, ) -> error::Result { use tokio::io::AsyncReadExt; @@ -3600,7 +3600,7 @@ pub(crate) async fn run_captured( let out = run_future_with_polling_update_job_poller( *job_id, - timeout_secs.or_else(|| ctx.timeout()), + ctx.timeout(), conn, ctx.mem_peak, ctx.canceled_by, @@ -3680,8 +3680,7 @@ pub(crate) async fn run_capturing( conn: &Connection, max_stdout_bytes: usize, ) -> error::Result { - let captured = - run_captured(cmd, name, ctx, job_id, w_id, conn, max_stdout_bytes, None).await?; + let captured = run_captured(cmd, name, ctx, job_id, w_id, conn, max_stdout_bytes).await?; if !captured.success { return Err(Error::ExecutionErr(format!( "{name} failed: {}",