diff --git a/cli/src/commands/pipeline/pipeline.ts b/cli/src/commands/pipeline/pipeline.ts new file mode 100644 index 0000000000..27985af3c7 --- /dev/null +++ b/cli/src/commands/pipeline/pipeline.ts @@ -0,0 +1,274 @@ +import { Command } from "@cliffy/command"; +import { Table } from "@cliffy/table"; +import { colors } from "@cliffy/ansi/colors"; + +import { OpenAPI } from "../../../gen/index.ts"; +import * as wmill from "../../../gen/services.gen.ts"; +import { requireLogin } from "../../core/auth.ts"; +import { resolveWorkspace } from "../../core/context.ts"; +import * as log from "../../core/log.ts"; +import { GlobalOptions } from "../../types.ts"; + +// Mirrors the asset-graph endpoint payload (backend/windmill-api-assets). +// Raw-fetched because these routes are newer than the checked-in generated +// client; swap to the generated functions on the next client regen. +type GraphRunnable = { + path: string; + usage_kind: "script" | "flow" | "job"; + in_pipeline?: boolean; +}; +type GraphEdge = { + runnable_kind: string; + runnable_path: string; + asset_kind: string; + asset_path: string; + access_type?: "r" | "w" | "rw"; +}; +type GraphTrigger = + | { + trigger_kind: "asset"; + asset_kind: string; + asset_path: string; + runnable_kind: string; + runnable_path: string; + } + | { + trigger_kind: string; + path?: string; + runnable_kind: string; + runnable_path: string; + missing?: boolean; + }; +type AssetGraph = { + runnables: GraphRunnable[]; + assets: { kind: string; path: string }[]; + edges: GraphEdge[]; + triggers: GraphTrigger[]; +}; + +async function apiGet(path: string): Promise { + const response = await fetch(`${OpenAPI.BASE}${path}`, { + headers: { Authorization: `Bearer ${OpenAPI.TOKEN}` }, + }); + if (!response.ok) { + const body = await response.text(); + throw new Error(`GET ${path} -> ${response.status}: ${body}`); + } + return (await response.json()) as T; +} + +async function list(opts: GlobalOptions & { json?: boolean }) { + if (opts.json) log.setSilent(true); + const workspace = await resolveWorkspace(opts); + await requireLogin(opts); + + const items = await apiGet<{ folder: string; script_count: number }[]>( + `/w/${workspace.workspaceId}/assets/pipelines`, + ); + if (opts.json) { + console.log(JSON.stringify(items)); + } else if (items.length === 0) { + log.info( + "No pipelines in this workspace. Mark scripts with a `// pipeline` comment (plus `// on ` triggers) and push them into a folder.", + ); + } else { + new Table() + .header(["Folder", "Scripts"]) + .padding(2) + .border(true) + .body(items.map((p) => [`f/${p.folder}`, String(p.script_count)])) + .render(); + } +} + +const ASSET_KINDS = "s3object,ducklake,datatable,volume"; + +function assetUri(kind: string, path: string): string { + const prefix = kind === "s3object" ? "s3" : kind; + return `${prefix}://${path}`; +} + +function shortName(scriptPath: string): string { + return scriptPath.split("/").pop() ?? scriptPath; +} + +async function show( + opts: GlobalOptions & { json?: boolean }, + folder: string, +) { + if (opts.json) log.setSilent(true); + const workspace = await resolveWorkspace(opts); + await requireLogin(opts); + + const f = folder.replace(/^f\//, "").replace(/\/$/, ""); + const graph = await apiGet( + `/w/${workspace.workspaceId}/assets/graph?folder=${encodeURIComponent(f)}&asset_kinds=${ASSET_KINDS}`, + ); + if (opts.json) { + console.log(JSON.stringify(graph)); + return; + } + if (graph.runnables.length === 0) { + log.info( + `No pipeline scripts in f/${f}. Mark scripts with a \`// pipeline\` comment and push them.`, + ); + return; + } + + // Index the graph: writes per script, subscribers per asset, native + // trigger markers per script, asset subscriptions per script. + const writesByScript = new Map(); + for (const e of graph.edges) { + if (e.access_type === "w" || e.access_type === "rw") { + const uri = assetUri(e.asset_kind, e.asset_path); + writesByScript.set(e.runnable_path, [ + ...(writesByScript.get(e.runnable_path) ?? []), + uri, + ]); + } + } + const subsByAsset = new Map(); + const subsByScript = new Map(); + const nativeByScript = new Map< + string, + { kind: string; path?: string; missing?: boolean }[] + >(); + for (const t of graph.triggers) { + if (t.trigger_kind === "asset") { + const at = t as Extract; + const uri = assetUri(at.asset_kind, at.asset_path); + subsByAsset.set(uri, [...(subsByAsset.get(uri) ?? []), t.runnable_path]); + subsByScript.set(t.runnable_path, [ + ...(subsByScript.get(t.runnable_path) ?? []), + uri, + ]); + } else { + const nt = t as Exclude; + nativeByScript.set(t.runnable_path, [ + ...(nativeByScript.get(t.runnable_path) ?? []), + { kind: nt.trigger_kind, path: nt.path, missing: nt.missing }, + ]); + } + } + + function triggerBadges(script: string): string { + const out: string[] = []; + for (const t of nativeByScript.get(script) ?? []) { + if (t.kind === "data_upload") { + out.push(colors.magenta("[data upload]")); + } else if (t.missing) { + out.push(colors.red(`[${t.kind} ✗ missing]`)); + } else { + out.push(colors.yellow(`[${t.kind}${t.path ? ` ${t.path}` : ""}]`)); + } + } + return out.length > 0 ? " " + out.join(" ") : ""; + } + + const printed = new Set(); + const lines: string[] = []; + + function printScript(script: string, prefix: string, extraOn?: string[]) { + const alsoOn = + extraOn && extraOn.length > 0 + ? colors.dim(` (also on: ${extraOn.join(", ")})`) + : ""; + if (printed.has(script)) { + lines.push( + `${prefix}${colors.bold(shortName(script))}${colors.dim(" ↻ shown above")}`, + ); + return; + } + printed.add(script); + lines.push(`${prefix}${colors.bold(shortName(script))}${triggerBadges(script)}${alsoOn}`); + const childPrefix = prefix.replace(/├─ $/, "│ ").replace(/└─ $/, " "); + const writes = [...(writesByScript.get(script) ?? [])].sort(); + writes.forEach((uri, i) => { + const lastAsset = i === writes.length - 1; + const assetBranch = lastAsset ? "└─▶ " : "├─▶ "; + lines.push(`${childPrefix}${assetBranch}${colors.cyan(uri)}`); + const assetChildPrefix = childPrefix + (lastAsset ? " " : "│ "); + const subs = [...(subsByAsset.get(uri) ?? [])].sort(); + subs.forEach((sub, j) => { + const branch = j === subs.length - 1 ? "└─ " : "├─ "; + const otherOn = (subsByScript.get(sub) ?? []).filter((u) => u !== uri); + printScript(sub, assetChildPrefix + branch, otherOn); + }); + }); + } + + // Roots: pipeline scripts that aren't subscribed to any asset — sources + // (data upload, schedule, webhook) and manual entries. + const roots = graph.runnables + .map((r) => r.path) + .filter((p) => !(subsByScript.get(p)?.length)) + .sort(); + + // UI-first markers (data_upload, webhook) have no trigger row — the + // graph endpoint can't surface them, they live as `// on ` + // annotations in the body. Roots are where sources matter, so fetch just + // those bodies and lift the marker kinds the canvas would show. + const MARKER_KINDS = ["data_upload", "webhook", "email"]; + await Promise.all( + roots.map(async (p) => { + const r = graph.runnables.find((x) => x.path === p); + if (r?.usage_kind !== "script") return; + try { + const script = await wmill.getScriptByPath({ + workspace: workspace.workspaceId, + path: p, + }); + const existing = nativeByScript.get(p) ?? []; + for (const line of (script.content ?? "").split("\n")) { + const m = line.match(/^\s*(?:\/\/|--|#)\s*on\s+(\w+)\s*$/); + if (!m) continue; + const kind = m[1]; + if (!MARKER_KINDS.includes(kind)) continue; + if (!existing.some((t) => t.kind === kind)) { + existing.push({ kind }); + } + } + if (existing.length > 0) nativeByScript.set(p, existing); + } catch { + // body fetch is best-effort enrichment only + } + }), + ); + + const scriptCount = graph.runnables.length; + const assetCount = graph.assets.length; + log.info( + colors.bold(`Pipeline f/${f}`) + + colors.dim(` — ${scriptCount} script${scriptCount === 1 ? "" : "s"} · ${assetCount} asset${assetCount === 1 ? "" : "s"}`), + ); + lines.push(""); + for (const root of roots) { + printScript(root, ""); + lines.push(""); + } + // Anything unreachable from the roots (e.g. cycles) still gets listed. + for (const r of graph.runnables) { + if (!printed.has(r.path)) { + printScript(r.path, ""); + lines.push(""); + } + } + console.log(lines.join("\n")); +} + +const command = new Command() + .description( + "inspect asset-driven pipelines (scripts marked `// pipeline`, wired by `// on ` annotations)", + ) + .command("list", "list pipeline folders in the workspace") + .option("--json", "Output as JSON (for piping to jq)") + .action(list as any) + .command( + "show", + "render a pipeline folder's DAG (sources, lineage, subscriptions) in the terminal", + ) + .arguments("") + .option("--json", "Output the raw asset graph as JSON") + .action(show as any); + +export default command; diff --git a/cli/src/guidance/skills.gen.ts b/cli/src/guidance/skills.gen.ts index cdb0c19234..aded724843 100644 --- a/cli/src/guidance/skills.gen.ts +++ b/cli/src/guidance/skills.gen.ts @@ -6432,6 +6432,17 @@ Validate Windmill flow, schedule, and trigger YAML files in a directory - \`--csv-separator \` - CSV column separator (default ,) - \`--csv-header\` - Treat the first CSV row as a header +### pipeline + +inspect asset-driven pipelines (scripts marked \`// pipeline\`, wired by \`// on \` annotations) + +**Subcommands:** + +- \`pipeline list\` - list pipeline folders in the workspace + - \`--json\` - Output as JSON (for piping to jq) +- \`pipeline show \` - render a pipeline folder's DAG (sources, lineage, subscriptions) in the terminal + - \`--json\` - Output the raw asset graph as JSON + ### protection-rules **Subcommands:** diff --git a/cli/src/main.ts b/cli/src/main.ts index c9e8d557c4..6d71f45d04 100755 --- a/cli/src/main.ts +++ b/cli/src/main.ts @@ -51,6 +51,7 @@ import generateMetadata from "./commands/generate-metadata/generate-metadata.ts" import docs from "./commands/docs/docs.ts"; import config from "./commands/config/config.ts"; import datatable from "./commands/datatable/datatable.ts"; +import pipeline from "./commands/pipeline/pipeline.ts"; import ducklake from "./commands/ducklake/ducklake.ts"; import objectStorage from "./commands/object-storage/object-storage.ts"; import { fetchVersion } from "./core/context.ts"; @@ -77,6 +78,7 @@ export { docs, config, datatable, + pipeline, ducklake, objectStorage, hubPull, @@ -215,6 +217,7 @@ const command = new Command() .command("docs", docs) .command("config", config) .command("datatable", datatable) + .command("pipeline", pipeline) .command("ducklake", ducklake) .command("object-storage", objectStorage) .command("version --version", "Show version information") diff --git a/system_prompts/auto-generated/cli/cli-commands.md b/system_prompts/auto-generated/cli/cli-commands.md index ac1a5b45f6..9018b8161d 100644 --- a/system_prompts/auto-generated/cli/cli-commands.md +++ b/system_prompts/auto-generated/cli/cli-commands.md @@ -407,6 +407,17 @@ Validate Windmill flow, schedule, and trigger YAML files in a directory - `--csv-separator ` - CSV column separator (default ,) - `--csv-header` - Treat the first CSV row as a header +### pipeline + +inspect asset-driven pipelines (scripts marked `// pipeline`, wired by `// on ` annotations) + +**Subcommands:** + +- `pipeline list` - list pipeline folders in the workspace + - `--json` - Output as JSON (for piping to jq) +- `pipeline show ` - render a pipeline folder's DAG (sources, lineage, subscriptions) in the terminal + - `--json` - Output the raw asset graph as JSON + ### protection-rules **Subcommands:** diff --git a/system_prompts/auto-generated/prompts.ts b/system_prompts/auto-generated/prompts.ts index dc0f3702f7..57d2d70b04 100644 --- a/system_prompts/auto-generated/prompts.ts +++ b/system_prompts/auto-generated/prompts.ts @@ -2957,6 +2957,17 @@ Validate Windmill flow, schedule, and trigger YAML files in a directory - \`--csv-separator \` - CSV column separator (default ,) - \`--csv-header\` - Treat the first CSV row as a header +### pipeline + +inspect asset-driven pipelines (scripts marked \`// pipeline\`, wired by \`// on \` annotations) + +**Subcommands:** + +- \`pipeline list\` - list pipeline folders in the workspace + - \`--json\` - Output as JSON (for piping to jq) +- \`pipeline show \` - render a pipeline folder's DAG (sources, lineage, subscriptions) in the terminal + - \`--json\` - Output the raw asset graph as JSON + ### protection-rules **Subcommands:** diff --git a/system_prompts/auto-generated/skills/cli-commands/SKILL.md b/system_prompts/auto-generated/skills/cli-commands/SKILL.md index f67d475856..3ef96cd135 100644 --- a/system_prompts/auto-generated/skills/cli-commands/SKILL.md +++ b/system_prompts/auto-generated/skills/cli-commands/SKILL.md @@ -412,6 +412,17 @@ Validate Windmill flow, schedule, and trigger YAML files in a directory - `--csv-separator ` - CSV column separator (default ,) - `--csv-header` - Treat the first CSV row as a header +### pipeline + +inspect asset-driven pipelines (scripts marked `// pipeline`, wired by `// on ` annotations) + +**Subcommands:** + +- `pipeline list` - list pipeline folders in the workspace + - `--json` - Output as JSON (for piping to jq) +- `pipeline show ` - render a pipeline folder's DAG (sources, lineage, subscriptions) in the terminal + - `--json` - Output the raw asset graph as JSON + ### protection-rules **Subcommands:**