diff --git a/backend/.sqlx/query-278bc6b4f149f824b5db32dacfaa714ee3852dc2ddf2d661dfdd5a986a9bb62b.json b/backend/.sqlx/query-278bc6b4f149f824b5db32dacfaa714ee3852dc2ddf2d661dfdd5a986a9bb62b.json deleted file mode 100644 index 1a18c37a2e..0000000000 --- a/backend/.sqlx/query-278bc6b4f149f824b5db32dacfaa714ee3852dc2ddf2d661dfdd5a986a9bb62b.json +++ /dev/null @@ -1,75 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT\n c.id IS NOT NULL AS completed,\n CASE \n WHEN q.id IS NOT NULL THEN (CASE WHEN NOT $5 AND q.running THEN true ELSE null END)\n ELSE false\n END AS running,\n SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs,\n COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,\n COALESCE(c.flow_status, f.flow_status) AS \"flow_status: sqlx::types::Json>\",\n COALESCE(c.workflow_as_code_status, f.workflow_as_code_status) AS \"workflow_as_code_status: sqlx::types::Json>\",\n job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset,\n created_by AS \"created_by!\",\n CASE WHEN $4::BOOLEAN THEN (\n SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc'\n ) END AS progress\n FROM v2_job j\n LEFT JOIN v2_job_queue q USING (id)\n LEFT JOIN v2_job_runtime r USING (id)\n LEFT JOIN v2_job_status f USING (id)\n LEFT JOIN v2_job_completed c USING (id)\n LEFT JOIN job_logs ON job_logs.job_id = $3\n WHERE j.workspace_id = $2 AND j.id = $3\n AND ($6::text[] IS NULL OR j.tag = ANY($6))", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "completed", - "type_info": "Bool" - }, - { - "ordinal": 1, - "name": "running", - "type_info": "Bool" - }, - { - "ordinal": 2, - "name": "logs", - "type_info": "Text" - }, - { - "ordinal": 3, - "name": "mem_peak", - "type_info": "Int4" - }, - { - "ordinal": 4, - "name": "flow_status: sqlx::types::Json>", - "type_info": "Jsonb" - }, - { - "ordinal": 5, - "name": "workflow_as_code_status: sqlx::types::Json>", - "type_info": "Jsonb" - }, - { - "ordinal": 6, - "name": "log_offset", - "type_info": "Int4" - }, - { - "ordinal": 7, - "name": "created_by!", - "type_info": "Varchar" - }, - { - "ordinal": 8, - "name": "progress", - "type_info": "Int4" - } - ], - "parameters": { - "Left": [ - "Int4", - "Text", - "Uuid", - "Bool", - "Bool", - "TextArray" - ] - }, - "nullable": [ - null, - null, - null, - null, - null, - null, - null, - false, - null - ] - }, - "hash": "278bc6b4f149f824b5db32dacfaa714ee3852dc2ddf2d661dfdd5a986a9bb62b" -} diff --git a/backend/.sqlx/query-4bc533074c720820cebff8d97a203df52520b7606378ecca267e88383a45b49b.json b/backend/.sqlx/query-4bc533074c720820cebff8d97a203df52520b7606378ecca267e88383a45b49b.json new file mode 100644 index 0000000000..8c9de521d2 --- /dev/null +++ b/backend/.sqlx/query-4bc533074c720820cebff8d97a203df52520b7606378ecca267e88383a45b49b.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO job_result_stream (workspace_id, job_id, stream)\n VALUES ($1, $2, $3)\n ON CONFLICT (job_id) DO UPDATE SET stream = job_result_stream.stream || $3\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Uuid", + "Text" + ] + }, + "nullable": [] + }, + "hash": "4bc533074c720820cebff8d97a203df52520b7606378ecca267e88383a45b49b" +} diff --git a/backend/.sqlx/query-4f372d047c78532907adf2d2dc114352aa7b5b28dccc50a1231a7f6539397da7.json b/backend/.sqlx/query-4f372d047c78532907adf2d2dc114352aa7b5b28dccc50a1231a7f6539397da7.json new file mode 100644 index 0000000000..e8377c5ff4 --- /dev/null +++ b/backend/.sqlx/query-4f372d047c78532907adf2d2dc114352aa7b5b28dccc50a1231a7f6539397da7.json @@ -0,0 +1,95 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT\n c.id IS NOT NULL AS completed,\n CASE \n WHEN q.id IS NOT NULL THEN (CASE WHEN NOT $5 AND q.running THEN true ELSE null END)\n ELSE false\n END AS running,\n CASE WHEN $7::BOOLEAN THEN NULL ELSE SUBSTR(logs, GREATEST($1 - log_offset, 0)) END AS logs,\n SUBSTR(rs.stream, $8) AS new_result_stream,\n COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,\n COALESCE(c.flow_status, f.flow_status) AS \"flow_status: sqlx::types::Json>\",\n COALESCE(c.workflow_as_code_status, f.workflow_as_code_status) AS \"workflow_as_code_status: sqlx::types::Json>\",\n CASE WHEN $7::BOOLEAN THEN NULL ELSE job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 END AS log_offset,\n CHAR_LENGTH(rs.stream) + 1 AS stream_offset,\n created_by AS \"created_by!\",\n CASE WHEN $4::BOOLEAN THEN (\n SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc'\n ) END AS progress,\n rs.stream AS \"result_stream: Option\"\n FROM v2_job j\n LEFT JOIN v2_job_queue q USING (id)\n LEFT JOIN v2_job_runtime r USING (id)\n LEFT JOIN v2_job_status f USING (id)\n LEFT JOIN v2_job_completed c USING (id)\n LEFT JOIN job_result_stream rs ON rs.job_id = $3\n LEFT JOIN job_logs ON job_logs.job_id = $3\n WHERE j.workspace_id = $2 AND j.id = $3\n AND ($6::text[] IS NULL OR j.tag = ANY($6))", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "completed", + "type_info": "Bool" + }, + { + "ordinal": 1, + "name": "running", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "logs", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "new_result_stream", + "type_info": "Text" + }, + { + "ordinal": 4, + "name": "mem_peak", + "type_info": "Int4" + }, + { + "ordinal": 5, + "name": "flow_status: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 6, + "name": "workflow_as_code_status: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 7, + "name": "log_offset", + "type_info": "Int4" + }, + { + "ordinal": 8, + "name": "stream_offset", + "type_info": "Int4" + }, + { + "ordinal": 9, + "name": "created_by!", + "type_info": "Varchar" + }, + { + "ordinal": 10, + "name": "progress", + "type_info": "Int4" + }, + { + "ordinal": 11, + "name": "result_stream: Option", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Int4", + "Text", + "Uuid", + "Bool", + "Bool", + "TextArray", + "Bool", + "Int4" + ] + }, + "nullable": [ + null, + null, + null, + null, + null, + null, + null, + null, + null, + false, + null, + false + ] + }, + "hash": "4f372d047c78532907adf2d2dc114352aa7b5b28dccc50a1231a7f6539397da7" +} diff --git a/backend/.sqlx/query-69924462c788dbc8f31aacc7f8ae588d76bf1f25d631833ce4b194818a7d1437.json b/backend/.sqlx/query-69924462c788dbc8f31aacc7f8ae588d76bf1f25d631833ce4b194818a7d1437.json new file mode 100644 index 0000000000..f25cfd9eba --- /dev/null +++ b/backend/.sqlx/query-69924462c788dbc8f31aacc7f8ae588d76bf1f25d631833ce4b194818a7d1437.json @@ -0,0 +1,48 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT result as \"result: sqlx::types::Json>\", v2_job.tag,\n v2_job_queue.running as \"running: Option\", SUBSTR(rs.stream, $3) AS \"result_stream: Option\", CHAR_LENGTH(rs.stream) AS stream_offset\n FROM v2_job\n LEFT JOIN v2_job_queue USING (id)\n LEFT JOIN v2_job_completed USING (id)\n LEFT JOIN job_result_stream rs ON rs.job_id = $2\n WHERE v2_job.id = $2 AND v2_job.workspace_id = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "result: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 1, + "name": "tag", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "running: Option", + "type_info": "Bool" + }, + { + "ordinal": 3, + "name": "result_stream: Option", + "type_info": "Text" + }, + { + "ordinal": 4, + "name": "stream_offset", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Text", + "Uuid", + "Int4" + ] + }, + "nullable": [ + true, + false, + false, + null, + null + ] + }, + "hash": "69924462c788dbc8f31aacc7f8ae588d76bf1f25d631833ce4b194818a7d1437" +} diff --git a/backend/.sqlx/query-71a040866adbd192080da165eb120abf7531b2542a2da225152e19144400f950.json b/backend/.sqlx/query-71a040866adbd192080da165eb120abf7531b2542a2da225152e19144400f950.json deleted file mode 100644 index d17cf3c524..0000000000 --- a/backend/.sqlx/query-71a040866adbd192080da165eb120abf7531b2542a2da225152e19144400f950.json +++ /dev/null @@ -1,184 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT workspace_id, slack_team_id, teams_team_id, teams_team_name, slack_name, slack_command_script, teams_command_script, slack_email, auto_invite_domain, auto_invite_operator, auto_add, customer_id, plan, webhook, deploy_to, ai_config, error_handler, error_handler_extra_args, error_handler_muted_on_cancel, large_file_storage, git_sync, deploy_ui, default_app, default_scripts, mute_critical_alerts, color, operator_settings, git_app_installations FROM workspace_settings WHERE workspace_id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "workspace_id", - "type_info": "Varchar" - }, - { - "ordinal": 1, - "name": "slack_team_id", - "type_info": "Varchar" - }, - { - "ordinal": 2, - "name": "teams_team_id", - "type_info": "Text" - }, - { - "ordinal": 3, - "name": "teams_team_name", - "type_info": "Text" - }, - { - "ordinal": 4, - "name": "slack_name", - "type_info": "Varchar" - }, - { - "ordinal": 5, - "name": "slack_command_script", - "type_info": "Varchar" - }, - { - "ordinal": 6, - "name": "teams_command_script", - "type_info": "Text" - }, - { - "ordinal": 7, - "name": "slack_email", - "type_info": "Varchar" - }, - { - "ordinal": 8, - "name": "auto_invite_domain", - "type_info": "Varchar" - }, - { - "ordinal": 9, - "name": "auto_invite_operator", - "type_info": "Bool" - }, - { - "ordinal": 10, - "name": "auto_add", - "type_info": "Bool" - }, - { - "ordinal": 11, - "name": "customer_id", - "type_info": "Varchar" - }, - { - "ordinal": 12, - "name": "plan", - "type_info": "Varchar" - }, - { - "ordinal": 13, - "name": "webhook", - "type_info": "Text" - }, - { - "ordinal": 14, - "name": "deploy_to", - "type_info": "Varchar" - }, - { - "ordinal": 15, - "name": "ai_config", - "type_info": "Jsonb" - }, - { - "ordinal": 16, - "name": "error_handler", - "type_info": "Varchar" - }, - { - "ordinal": 17, - "name": "error_handler_extra_args", - "type_info": "Json" - }, - { - "ordinal": 18, - "name": "error_handler_muted_on_cancel", - "type_info": "Bool" - }, - { - "ordinal": 19, - "name": "large_file_storage", - "type_info": "Jsonb" - }, - { - "ordinal": 20, - "name": "git_sync", - "type_info": "Jsonb" - }, - { - "ordinal": 21, - "name": "deploy_ui", - "type_info": "Jsonb" - }, - { - "ordinal": 22, - "name": "default_app", - "type_info": "Varchar" - }, - { - "ordinal": 23, - "name": "default_scripts", - "type_info": "Jsonb" - }, - { - "ordinal": 24, - "name": "mute_critical_alerts", - "type_info": "Bool" - }, - { - "ordinal": 25, - "name": "color", - "type_info": "Varchar" - }, - { - "ordinal": 26, - "name": "operator_settings", - "type_info": "Jsonb" - }, - { - "ordinal": 27, - "name": "git_app_installations", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [ - false, - true, - true, - true, - true, - true, - true, - false, - true, - true, - true, - true, - true, - true, - true, - true, - true, - true, - false, - true, - true, - true, - true, - true, - true, - true, - true, - false - ] - }, - "hash": "71a040866adbd192080da165eb120abf7531b2542a2da225152e19144400f950" -} diff --git a/backend/.sqlx/query-e653d36b607a16c0dfc0324690942ab25883b53a81ebb581fe019af2ec5eb567.json b/backend/.sqlx/query-72956f508f66312807738399b57aa01a048fa4f9281327cf5b78111178424b43.json similarity index 91% rename from backend/.sqlx/query-e653d36b607a16c0dfc0324690942ab25883b53a81ebb581fe019af2ec5eb567.json rename to backend/.sqlx/query-72956f508f66312807738399b57aa01a048fa4f9281327cf5b78111178424b43.json index 284cf3338f..2ffdf141b1 100644 --- a/backend/.sqlx/query-e653d36b607a16c0dfc0324690942ab25883b53a81ebb581fe019af2ec5eb567.json +++ b/backend/.sqlx/query-72956f508f66312807738399b57aa01a048fa4f9281327cf5b78111178424b43.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n id AS \"id!\", workspace_id AS \"workspace_id!\", parent_job, is_flow_step,\n flow_status AS \"flow_status: Box\", last_ping, same_worker\n FROM v2_as_queue\n WHERE running = true AND suspend = 0 AND suspend_until IS null AND scheduled_for <= now()\n AND (job_kind = 'flow' OR job_kind = 'flowpreview' OR job_kind = 'flownode')\n AND last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval\n AND canceled = false\n ", + "query": "\n SELECT\n id AS \"id!\", workspace_id AS \"workspace_id!\", parent_job, is_flow_step,\n flow_status AS \"flow_status: Box\", last_ping, same_worker\n FROM v2_as_queue\n WHERE running = true AND suspend = 0 AND suspend_until IS null AND scheduled_for <= now()\n AND (job_kind = 'flow' OR job_kind = 'flowpreview' OR job_kind = 'flownode')\n AND last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval\n AND canceled = false\n \n ", "describe": { "columns": [ { @@ -54,5 +54,5 @@ true ] }, - "hash": "e653d36b607a16c0dfc0324690942ab25883b53a81ebb581fe019af2ec5eb567" + "hash": "72956f508f66312807738399b57aa01a048fa4f9281327cf5b78111178424b43" } diff --git a/backend/.sqlx/query-74a2871aba7e35527dcefb2538b8cf41c35618f7c15ea047c5082f8feb7a8464.json b/backend/.sqlx/query-74a2871aba7e35527dcefb2538b8cf41c35618f7c15ea047c5082f8feb7a8464.json index 6a45170165..9aa5d8ccbd 100644 --- a/backend/.sqlx/query-74a2871aba7e35527dcefb2538b8cf41c35618f7c15ea047c5082f8feb7a8464.json +++ b/backend/.sqlx/query-74a2871aba7e35527dcefb2538b8cf41c35618f7c15ea047c5082f8feb7a8464.json @@ -14,7 +14,8 @@ "Enum": [ "s3object", "resource", - "variable" + "variable", + "ducklake" ] } } diff --git a/backend/.sqlx/query-9f5b677a02690d3e4b4a5f5e141c7107077bbe90423102b5469e219f2a8b9293.json b/backend/.sqlx/query-9f5b677a02690d3e4b4a5f5e141c7107077bbe90423102b5469e219f2a8b9293.json new file mode 100644 index 0000000000..f9a1de974a --- /dev/null +++ b/backend/.sqlx/query-9f5b677a02690d3e4b4a5f5e141c7107077bbe90423102b5469e219f2a8b9293.json @@ -0,0 +1,36 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT result as \"result: sqlx::types::Json>\", SUBSTR(rs.stream, $3) AS \"result_stream: Option\", CHAR_LENGTH(rs.stream) + 1 AS stream_offset\n FROM v2_job_completed FULL OUTER JOIN job_result_stream rs ON rs.job_id = v2_job_completed.id WHERE (v2_job_completed.id = $2 AND v2_job_completed.workspace_id = $1 OR rs.workspace_id = $1)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "result: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 1, + "name": "result_stream: Option", + "type_info": "Text" + }, + { + "ordinal": 2, + "name": "stream_offset", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Text", + "Uuid", + "Int4" + ] + }, + "nullable": [ + true, + null, + null + ] + }, + "hash": "9f5b677a02690d3e4b4a5f5e141c7107077bbe90423102b5469e219f2a8b9293" +} diff --git a/backend/.sqlx/query-b9b3c341fe452da916ee29637e14b5c1ad75462eba17083c6f81ff6ef35af77f.json b/backend/.sqlx/query-b9b3c341fe452da916ee29637e14b5c1ad75462eba17083c6f81ff6ef35af77f.json deleted file mode 100644 index 746d347214..0000000000 --- a/backend/.sqlx/query-b9b3c341fe452da916ee29637e14b5c1ad75462eba17083c6f81ff6ef35af77f.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT result as \"result: sqlx::types::Json>\", v2_job_queue.running as \"running: Option\" FROM v2_job_completed FULL OUTER JOIN v2_job_queue USING (id) WHERE (v2_job_queue.id = $1 AND v2_job_queue.workspace_id = $2) OR (v2_job_completed.id = $1 AND v2_job_completed.workspace_id = $2)", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "result: sqlx::types::Json>", - "type_info": "Jsonb" - }, - { - "ordinal": 1, - "name": "running: Option", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [ - true, - false - ] - }, - "hash": "b9b3c341fe452da916ee29637e14b5c1ad75462eba17083c6f81ff6ef35af77f" -} diff --git a/backend/.sqlx/query-ceb8c2607023883e1eebd4b9539e36ed202a6ecd12e3f90cb070341c38886de4.json b/backend/.sqlx/query-ceb8c2607023883e1eebd4b9539e36ed202a6ecd12e3f90cb070341c38886de4.json deleted file mode 100644 index 8800f8c6cc..0000000000 --- a/backend/.sqlx/query-ceb8c2607023883e1eebd4b9539e36ed202a6ecd12e3f90cb070341c38886de4.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT result as \"result: sqlx::types::Json>\" FROM v2_job_completed WHERE id = $2 AND workspace_id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "result: sqlx::types::Json>", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Text", - "Uuid" - ] - }, - "nullable": [ - true - ] - }, - "hash": "ceb8c2607023883e1eebd4b9539e36ed202a6ecd12e3f90cb070341c38886de4" -} diff --git a/backend/.sqlx/query-ec0f8fa36328507e51c1974dbef884b755504a6cefa4af34fa4659fb95a7ee9a.json b/backend/.sqlx/query-ec0f8fa36328507e51c1974dbef884b755504a6cefa4af34fa4659fb95a7ee9a.json new file mode 100644 index 0000000000..431160d7be --- /dev/null +++ b/backend/.sqlx/query-ec0f8fa36328507e51c1974dbef884b755504a6cefa4af34fa4659fb95a7ee9a.json @@ -0,0 +1,42 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT \n result as \"result: sqlx::types::Json>\",\n v2_job_queue.running as \"running: Option\",\n SUBSTR(rs.stream, $3) AS \"result_stream: Option\",\n CHAR_LENGTH(rs.stream) + 1 AS stream_offset\n FROM v2_job_completed FULL OUTER JOIN v2_job_queue USING (id) \n LEFT JOIN job_result_stream rs ON rs.job_id = $1\n WHERE (v2_job_queue.id = $1 AND v2_job_queue.workspace_id = $2) OR (v2_job_completed.id = $1 AND v2_job_completed.workspace_id = $2)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "result: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 1, + "name": "running: Option", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "result_stream: Option", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "stream_offset", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text", + "Int4" + ] + }, + "nullable": [ + true, + false, + null, + null + ] + }, + "hash": "ec0f8fa36328507e51c1974dbef884b755504a6cefa4af34fa4659fb95a7ee9a" +} diff --git a/backend/.sqlx/query-fab257c4e20aa51b8f785b1882aa0b16fde33b246cbf0749ffa0e4ed63504451.json b/backend/.sqlx/query-fab257c4e20aa51b8f785b1882aa0b16fde33b246cbf0749ffa0e4ed63504451.json deleted file mode 100644 index 14762cc32c..0000000000 --- a/backend/.sqlx/query-fab257c4e20aa51b8f785b1882aa0b16fde33b246cbf0749ffa0e4ed63504451.json +++ /dev/null @@ -1,35 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT result as \"result: sqlx::types::Json>\", v2_job.tag,\n v2_job_queue.running as \"running: Option\"\n FROM v2_job\n LEFT JOIN v2_job_queue USING (id)\n LEFT JOIN v2_job_completed USING (id)\n WHERE v2_job.id = $2 AND v2_job.workspace_id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "result: sqlx::types::Json>", - "type_info": "Jsonb" - }, - { - "ordinal": 1, - "name": "tag", - "type_info": "Varchar" - }, - { - "ordinal": 2, - "name": "running: Option", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Text", - "Uuid" - ] - }, - "nullable": [ - true, - false, - false - ] - }, - "hash": "fab257c4e20aa51b8f785b1882aa0b16fde33b246cbf0749ffa0e4ed63504451" -} diff --git a/backend/migrations/20250804155709_add_stream_result.down.sql b/backend/migrations/20250804155709_add_stream_result.down.sql new file mode 100644 index 0000000000..f8720d6bd6 --- /dev/null +++ b/backend/migrations/20250804155709_add_stream_result.down.sql @@ -0,0 +1,2 @@ +-- Add down migration script here +DROP TABLE job_result_stream; \ No newline at end of file diff --git a/backend/migrations/20250804155709_add_stream_result.up.sql b/backend/migrations/20250804155709_add_stream_result.up.sql new file mode 100644 index 0000000000..70e7ebff25 --- /dev/null +++ b/backend/migrations/20250804155709_add_stream_result.up.sql @@ -0,0 +1,11 @@ +-- Add up migration script here +CREATE TABLE job_result_stream ( + job_id UUID NOT NULL PRIMARY KEY, + workspace_id TEXT NOT NULL, + stream TEXT NOT NULL +); + +ALTER TABLE job_result_stream ADD CONSTRAINT fk_job_result_stream_job_id FOREIGN KEY (job_id) REFERENCES v2_job_queue(id) ON DELETE CASCADE; + +GRANT ALL ON TABLE job_result_stream TO windmill_admin; +GRANT ALL ON TABLE job_result_stream TO windmill_user; \ No newline at end of file diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 4492afe692..f31ef52c97 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -2129,6 +2129,7 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { AND (job_kind = 'flow' OR job_kind = 'flowpreview' OR job_kind = 'flownode') AND last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval AND canceled = false + "#, FLOW_ZOMBIE_TRANSITION_TIMEOUT.as_str() ) diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 5fa2670c5d..cce3f338dc 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -27,7 +27,7 @@ use serde::Serialize; use windmill_common::flows::InputTransform; use windmill_common::worker::WORKER_CONFIG; -#[cfg(feature = "python")] +#[cfg(any(feature = "python", feature = "deno_core"))] use windmill_common::flow_status::{FlowStatus, FlowStatusModule, RestartedFrom}; use windmill_common::{ diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index e7c3ef37b9..c5be91f428 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -7798,10 +7798,18 @@ paths: in: query schema: type: integer + - name: stream_offset + in: query + schema: + type: integer - name: get_progress in: query schema: type: boolean + - name: no_logs + in: query + schema: + type: boolean responses: "200": @@ -7823,6 +7831,10 @@ paths: type: integer progress: type: integer + stream_offset: + type: integer + new_result_stream: + type: string flow_status: $ref: "../../openflow.openapi.yaml#/components/schemas/FlowStatus" workflow_as_code_status: @@ -7853,6 +7865,10 @@ paths: in: query schema: type: boolean + - name: no_logs + in: query + schema: + type: boolean responses: "200": @@ -14840,7 +14856,7 @@ components: - key - typ required: - - object + - object - type: object properties: list: diff --git a/backend/windmill-api/src/approvals.rs b/backend/windmill-api/src/approvals.rs index f16ed00abe..e03b587f10 100644 --- a/backend/windmill-api/src/approvals.rs +++ b/backend/windmill-api/src/approvals.rs @@ -1,18 +1,24 @@ -use serde::{Deserialize, Serialize}; -use std::collections::HashMap; -use uuid::Uuid; -use std::str::FromStr; -use regex::Regex; -use serde_json::Value; use crate::auth::OptTokened; use crate::db::{ApiAuthed, DB}; -use crate::jobs::{cancel_suspended_job, resume_suspended_job, QueryApprover, QueryOrBody, ResumeUrls, get_resume_urls_internal}; -use axum::{extract::{Path, Query}, Extension}; -use windmill_common::error::Error; +use crate::jobs::{ + cancel_suspended_job, get_resume_urls_internal, resume_suspended_job, QueryApprover, + QueryOrBody, ResumeUrls, +}; +use axum::{ + extract::{Path, Query}, + Extension, +}; +use regex::Regex; +use serde::{Deserialize, Serialize}; +use serde_json::value::RawValue; +use serde_json::Value; +use std::collections::HashMap; +use std::str::FromStr; +use uuid::Uuid; use windmill_common::cache; +use windmill_common::error::Error; use windmill_common::jobs::JobKind; use windmill_common::scripts::ScriptHash; -use serde_json::value::RawValue; #[derive(Debug, Deserialize, Serialize)] pub struct ResumeSchema { @@ -234,7 +240,7 @@ pub async fn get_approval_form_details( .ok_or_else(|| Error::BadRequest("This workflow is no longer running and has either already timed out or been cancelled or completed.".to_string())) .map(|r| (r.job_kind, r.script_hash, r.raw_flow, r.parent_job, r.created_at, r.created_by, r.script_path, r.args))?; - let flow_data = match cache::job::fetch_flow(&db, job_kind, script_hash).await { + let flow_data = match cache::job::fetch_flow(&db, &job_kind, script_hash).await { Ok(data) => data, Err(_) => { if let Some(parent_job_id) = parent_job_id.as_ref() { diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 8f247aaf75..9a0f918978 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -900,11 +900,10 @@ impl<'a> GetQuery<'a> { /// when pushed from an un-updated workers. /// This function is used to make the above change transparent for the API, as the returned jobs /// will have the raw values as if they were still in the tables. - async fn resolve_raw_values( + async fn resolve_raw_values( &self, db: &DB, id: Uuid, - kind: JobKind, hash: Option, job: &mut JobExtended, ) { @@ -917,18 +916,18 @@ impl<'a> GetQuery<'a> { // Try to fetch the flow from the cache, fallback to the preview flow. // NOTE: This could check for the job kinds instead of the `or_else` but it's not // necessary as `fetch_flow` return early if the job kind is not a preview one. - cache::job::fetch_flow(db, kind, hash) + cache::job::fetch_flow(db, job.job_kind(), hash) .or_else(|_| cache::job::fetch_preview_flow(db, &id, raw_flow)) .await .ok() .inspect(|data| job.raw_flow = Some(sqlx::types::Json(data.raw_flow.clone()))); } - if self.with_code { + if self.with_code && job.job_kind() == &JobKind::Preview { // Try to fetch the code from the cache, fallback to the preview code. // NOTE: This could check for the job kinds instead of the `or_else` but it's not // necessary as `fetch_script` return early if the job kind is not a preview one. let conn = Connection::from(db.clone()); - cache::job::fetch_script(db.clone(), kind, hash) + cache::job::fetch_script(db.clone(), job.job_kind(), hash) .or_else(|_| cache::job::fetch_preview_script(&conn, &id, raw_lock, raw_code)) .await .ok() @@ -958,7 +957,7 @@ impl<'a> GetQuery<'a> { self.check_auth(job.as_ref().map(|job| job.created_by.as_str()))?; if let Some(job) = job.as_mut() { - self.resolve_raw_values(&db, job.id, job.job_kind, job.script_hash, job) + self.resolve_raw_values(&db, job.id, job.script_hash, job) .await; } if self.with_flow { @@ -992,7 +991,7 @@ impl<'a> GetQuery<'a> { self.check_auth(cjob.as_ref().map(|job| job.created_by.as_str()))?; if let Some(job) = cjob.as_mut() { - self.resolve_raw_values(db, job.id, job.job_kind, job.script_hash, job) + self.resolve_raw_values(db, job.id, job.script_hash, job) .await; } @@ -2642,7 +2641,7 @@ pub async fn get_resume_urls_internal( } #[derive(sqlx::FromRow, Debug, Serialize)] -pub struct JobExtended { +pub struct JobExtended { #[sqlx(flatten)] #[serde(flatten)] inner: T, @@ -2665,7 +2664,23 @@ pub struct JobExtended { pub aggregate_wait_time_ms: Option, } -impl JobExtended { +pub trait JobCommon { + fn job_kind(&self) -> &JobKind; +} + +impl JobCommon for QueuedJob { + fn job_kind(&self) -> &JobKind { + &self.job_kind + } +} + +impl JobCommon for CompletedJob { + fn job_kind(&self) -> &JobKind { + &self.job_kind + } +} + +impl JobExtended { pub fn new( self_wait_time_ms: Option, aggregate_wait_time_ms: Option, @@ -2683,7 +2698,7 @@ impl JobExtended { } } -impl Deref for JobExtended { +impl Deref for JobExtended { type Target = T; fn deref(&self) -> &Self::Target { @@ -2691,7 +2706,7 @@ impl Deref for JobExtended { } } -impl DerefMut for JobExtended { +impl DerefMut for JobExtended { fn deref_mut(&mut self) -> &mut Self::Target { &mut self.inner } @@ -4538,8 +4553,13 @@ pub async fn run_wait_result_script_by_path_internal( check_queue_too_long(&db, QUEUE_LIMIT_WAIT_RESULT.or(run_query.queue_limit)).await?; let mut tx = user_db.clone().begin(&authed).await?; - let (job_payload, tag, delete_after_use, timeout, on_behalf_of) = - script_path_to_payload(script_path.to_path(), &mut *tx, &w_id, run_query.skip_preprocessor).await?; + let (job_payload, tag, delete_after_use, timeout, on_behalf_of) = script_path_to_payload( + script_path.to_path(), + &mut *tx, + &w_id, + run_query.skip_preprocessor, + ) + .await?; drop(tx); let tag = run_query.tag.clone().or(tag); @@ -5680,22 +5700,38 @@ pub async fn run_job_by_hash_inner( pub struct JobUpdateQuery { pub running: Option, pub log_offset: Option, + pub stream_offset: Option, pub get_progress: Option, + pub no_logs: Option, pub only_result: Option, pub fast: Option, } #[derive(Serialize, Debug)] pub struct JobUpdate { + #[serde(skip_serializing_if = "Option::is_none")] pub running: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub completed: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub new_logs: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub new_result_stream: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub log_offset: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub stream_offset: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub mem_peak: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub progress: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub flow_status: Option>, + #[serde(skip_serializing_if = "Option::is_none")] pub workflow_as_code_status: Option>, + #[serde(skip_serializing_if = "Option::is_none")] pub job: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub only_result: Option>, } @@ -5714,6 +5750,7 @@ impl Hash for JobUpdate { self.log_offset.hash(state); self.mem_peak.hash(state); self.progress.hash(state); + self.stream_offset.hash(state); if !self.completed.unwrap_or(false) { self.flow_status.as_ref().map(|x| x.get().hash(state)); self.workflow_as_code_status @@ -5779,7 +5816,7 @@ async fn get_job_update( opt_tokened: OptTokened, Extension(db): Extension, Path((w_id, job_id)): Path<(String, Uuid)>, - Query(JobUpdateQuery { log_offset, get_progress, running, only_result, .. }): Query< + Query(JobUpdateQuery { log_offset, stream_offset, get_progress, running, only_result, no_logs, .. }): Query< JobUpdateQuery, >, ) -> JsonResult { @@ -5791,11 +5828,13 @@ async fn get_job_update( &w_id, &job_id, log_offset, + stream_offset, get_progress, running, true, false, only_result, + no_logs, ) .await?, )) @@ -5806,7 +5845,7 @@ async fn get_job_update_sse( opt_tokened: OptTokened, Extension(db): Extension, Path((w_id, job_id)): Path<(String, Uuid)>, - Query(JobUpdateQuery { log_offset, get_progress, running, only_result, fast }): Query< + Query(JobUpdateQuery { log_offset, stream_offset, get_progress, running, no_logs, only_result, fast }): Query< JobUpdateQuery, >, ) -> Response { @@ -5817,10 +5856,12 @@ async fn get_job_update_sse( w_id, job_id, log_offset, + stream_offset, get_progress, running, only_result, fast, + no_logs, ) .map(|x| { format!( @@ -5857,19 +5898,24 @@ fn get_job_update_sse_stream( w_id: String, job_id: Uuid, initial_log_offset: Option, + initial_stream_offset: Option, get_progress: Option, running: Option, only_result: Option, fast: Option, + no_logs: Option, ) -> impl futures::Stream { let (tx, rx) = tokio::sync::mpsc::channel(32); tokio::spawn(async move { let mut log_offset = initial_log_offset; + let mut stream_offset = initial_stream_offset; let mut last_update_hash: Option = None; // Send initial update immediately let mut running = running; + let mut mem_peak = 0; + match get_job_update_data( &opt_authed, &opt_tokened, @@ -5877,22 +5923,38 @@ fn get_job_update_sse_stream( &w_id, &job_id, log_offset, + stream_offset, get_progress, running, true, true, only_result, + no_logs, ) .await { - Ok(update) => { + Ok(mut update) => { last_update_hash = Some(update.hash_str()); let completion_sent = update.completed.unwrap_or(false); if running.is_some() && update.running.is_some_and(|x| x) { running = Some(true); } + if let Some(new_mem_peak) = update.mem_peak { + mem_peak = new_mem_peak; + } if let Some(new_offset) = update.log_offset { - log_offset = Some(new_offset); + if new_offset != log_offset.unwrap_or(0) { + log_offset = Some(new_offset); + } else { + update.log_offset = None; + } + } + if let Some(new_stream_offset) = update.stream_offset { + if new_stream_offset != stream_offset.unwrap_or(0) { + stream_offset = Some(new_stream_offset); + } else { + update.stream_offset = None; + } } if tx.send(JobUpdateSSEStream::Update(update)).await.is_err() { tracing::warn!("Failed to send initial job update for job {job_id}"); @@ -5918,6 +5980,7 @@ fn get_job_update_sse_stream( let mut i = 0; let start = Instant::now(); let mut last_ping = Instant::now(); + loop { i += 1; let ms_duration = if i > 100 || !fast.unwrap_or(false) { @@ -5950,18 +6013,27 @@ fn get_job_update_sse_stream( &w_id, &job_id, log_offset, + stream_offset, get_progress, running, false, true, only_result, + no_logs, ) .await { - Ok(update) => { + Ok(mut update) => { if running.is_some() && update.running.is_some_and(|x| x) { running = Some(true); } + if update.completed.is_some_and(|x| !x) { + update.completed = None; + } + if update.new_logs.as_ref().is_some_and(|x| x.is_empty()) { + update.new_logs = None; + } + // if !only_result.unwrap_or(false) { // tracing::error!("update {:?}", update); // } @@ -5970,7 +6042,25 @@ fn get_job_update_sse_stream( if last_update_hash.as_ref() != Some(&update_last_status) { // Update log offset if available if let Some(new_offset) = update.log_offset { - log_offset = Some(new_offset); + if new_offset != log_offset.unwrap_or(0) { + log_offset = Some(new_offset); + } else { + update.log_offset = None; + } + } + if let Some(new_stream_offset) = update.stream_offset { + if new_stream_offset != stream_offset.unwrap_or(0) { + stream_offset = Some(new_stream_offset); + } else { + update.stream_offset = None; + } + } + if let Some(new_mem_peak) = update.mem_peak { + if new_mem_peak != mem_peak { + mem_peak = new_mem_peak; + } else { + update.mem_peak = None; + } } let completed = update.completed.unwrap_or(false); if tx.send(JobUpdateSSEStream::Update(update)).await.is_err() { @@ -5983,7 +6073,8 @@ fn get_job_update_sse_stream( last_update_hash = Some(update_last_status); } } - Err(_) => { + Err(e) => { + tracing::error!("Error getting job update: {:?}", e); if tx.send(JobUpdateSSEStream::NotFound).await.is_err() { tracing::warn!("Failed to send job not found for job {job_id}"); } @@ -6003,11 +6094,13 @@ async fn get_job_update_data( w_id: &str, job_id: &Uuid, log_offset: Option, + stream_offset: Option, get_progress: Option, running: Option, log_view: bool, get_full_job_on_completion: bool, only_result: Option, + no_logs: Option, ) -> error::Result { let tags = if log_view { log_job_view( @@ -6028,19 +6121,22 @@ async fn get_job_update_data( if only_result.unwrap_or(false) { let result = if let Some(tags) = tags { - let r = sqlx::query!( - "SELECT result as \"result: sqlx::types::Json>\", v2_job.tag, - v2_job_queue.running as \"running: Option\" + let r = + sqlx::query!( + "SELECT result as \"result: sqlx::types::Json>\", v2_job.tag, + v2_job_queue.running as \"running: Option\", SUBSTR(rs.stream, $3) AS \"result_stream: Option\", CHAR_LENGTH(rs.stream) AS stream_offset FROM v2_job LEFT JOIN v2_job_queue USING (id) LEFT JOIN v2_job_completed USING (id) + LEFT JOIN job_result_stream rs ON rs.job_id = $2 WHERE v2_job.id = $2 AND v2_job.workspace_id = $1", - w_id, - job_id, - ) - .fetch_optional(db) - .await? - .ok_or_else(|| Error::NotFound(format!("Job not found: {}", job_id)))?; + w_id, + job_id, + stream_offset.unwrap_or(0), + ) + .fetch_optional(db) + .await? + .ok_or_else(|| Error::NotFound(format!("Job not found: {}", job_id)))?; if !tags.contains(&r.tag.as_str()) { return Err(Error::NotAuthorized(format!( @@ -6050,27 +6146,44 @@ async fn get_job_update_data( ))); } let running = r.running.as_ref().map(|x| *x); - (r.result.map(|x| x.0), running) + (r.result.map(|x| x.0), running, r.result_stream.flatten(), r.stream_offset) } else { if running.is_some_and(|x| !x) { let r = sqlx::query!( - "SELECT result as \"result: sqlx::types::Json>\", v2_job_queue.running as \"running: Option\" FROM v2_job_completed FULL OUTER JOIN v2_job_queue USING (id) WHERE (v2_job_queue.id = $1 AND v2_job_queue.workspace_id = $2) OR (v2_job_completed.id = $1 AND v2_job_completed.workspace_id = $2)", + "SELECT + result as \"result: sqlx::types::Json>\", + v2_job_queue.running as \"running: Option\", + SUBSTR(rs.stream, $3) AS \"result_stream: Option\", + CHAR_LENGTH(rs.stream) + 1 AS stream_offset + FROM v2_job_completed FULL OUTER JOIN v2_job_queue USING (id) + LEFT JOIN job_result_stream rs ON rs.job_id = $1 + WHERE (v2_job_queue.id = $1 AND v2_job_queue.workspace_id = $2) OR (v2_job_completed.id = $1 AND v2_job_completed.workspace_id = $2)", job_id, w_id, + stream_offset.unwrap_or(0), ).fetch_optional(db).await?; if let Some(r) = r { let running = r.running.as_ref().map(|x| *x); - (r.result.map(|x| x.0), running) + (r.result.map(|x| x.0), running, r.result_stream.flatten(), r.stream_offset) } else { - (None, None) + (None, None, None, None) } } else { - (sqlx::query_scalar!( - "SELECT result as \"result: sqlx::types::Json>\" FROM v2_job_completed WHERE id = $2 AND workspace_id = $1", + let q = sqlx::query!( + "SELECT result as \"result: sqlx::types::Json>\", SUBSTR(rs.stream, $3) AS \"result_stream: Option\", CHAR_LENGTH(rs.stream) + 1 AS stream_offset + FROM v2_job_completed FULL OUTER JOIN job_result_stream rs ON rs.job_id = v2_job_completed.id WHERE (v2_job_completed.id = $2 AND v2_job_completed.workspace_id = $1 OR rs.workspace_id = $1)", w_id, job_id, - ).fetch_optional(db).await?.flatten() - .map(|x| x.0), running) + stream_offset.unwrap_or(0), + ) + .fetch_optional(db) + .await?; + tracing::error!("q {:?}", q); + if let Some(r) = q { + (r.result.map(|x| x.0), running, r.result_stream.flatten(), r.stream_offset) + } else { + (None, None, None, None) + } } }; Ok(JobUpdate { @@ -6078,6 +6191,8 @@ async fn get_job_update_data( completed: if result.0.is_some() { Some(true) } else { None }, log_offset: None, new_logs: None, + new_result_stream: result.2, + stream_offset: result.3, mem_peak: None, progress: None, job: None, @@ -6093,20 +6208,24 @@ async fn get_job_update_data( WHEN q.id IS NOT NULL THEN (CASE WHEN NOT $5 AND q.running THEN true ELSE null END) ELSE false END AS running, - SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs, + CASE WHEN $7::BOOLEAN THEN NULL ELSE SUBSTR(logs, GREATEST($1 - log_offset, 0)) END AS logs, + SUBSTR(rs.stream, $8) AS new_result_stream, COALESCE(r.memory_peak, c.memory_peak) AS mem_peak, COALESCE(c.flow_status, f.flow_status) AS \"flow_status: sqlx::types::Json>\", COALESCE(c.workflow_as_code_status, f.workflow_as_code_status) AS \"workflow_as_code_status: sqlx::types::Json>\", - job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset, + CASE WHEN $7::BOOLEAN THEN NULL ELSE job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 END AS log_offset, + CHAR_LENGTH(rs.stream) + 1 AS stream_offset, created_by AS \"created_by!\", CASE WHEN $4::BOOLEAN THEN ( SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc' - ) END AS progress + ) END AS progress, + rs.stream AS \"result_stream: Option\" FROM v2_job j LEFT JOIN v2_job_queue q USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status f USING (id) LEFT JOIN v2_job_completed c USING (id) + LEFT JOIN job_result_stream rs ON rs.job_id = $3 LEFT JOIN job_logs ON job_logs.job_id = $3 WHERE j.workspace_id = $2 AND j.id = $3 AND ($6::text[] IS NULL OR j.tag = ANY($6))", @@ -6116,6 +6235,8 @@ async fn get_job_update_data( get_progress.unwrap_or(false), running, tags.as_ref().map(|v| v.as_slice()) as Option<&[&str]>, + no_logs.unwrap_or(false), + stream_offset.unwrap_or(0), ) .fetch_optional(db) .await? @@ -6139,6 +6260,8 @@ async fn get_job_update_data( completed: record.completed, log_offset: record.log_offset, new_logs: record.logs, + new_result_stream: record.new_result_stream, + stream_offset: record.stream_offset, mem_peak: record.mem_peak, progress: record.progress, workflow_as_code_status: record diff --git a/backend/windmill-common/src/cache.rs b/backend/windmill-common/src/cache.rs index 03d8ec7f03..4129948268 100644 --- a/backend/windmill-common/src/cache.rs +++ b/backend/windmill-common/src/cache.rs @@ -797,11 +797,12 @@ pub mod job { #[track_caller] pub fn fetch_script( db: DB, - kind: JobKind, + kind: &JobKind, hash: Option, ) -> impl Future>> { use JobKind::*; let loc = Location::caller(); + let kind = kind.clone(); async move { match (kind, hash.map(|ScriptHash(id)| id)) { (FlowScript, Some(id)) => { @@ -825,11 +826,12 @@ pub mod job { #[track_caller] pub fn fetch_flow<'c>( db: &'c DB, - kind: JobKind, + kind: &JobKind, hash: Option, ) -> impl Future>> + 'c { use JobKind::*; let loc = Location::caller(); + let kind = kind.clone(); async move { match (kind, hash.map(|ScriptHash(id)| id)) { (FlowDependencies, Some(id)) => flow::fetch_version(db, id).await, diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index a1bb0babf5..3076157c02 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -82,6 +82,7 @@ pub mod variables; pub mod worker; pub mod workspaces; pub mod triggers; +pub mod result_stream; pub mod stream; pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50; diff --git a/backend/windmill-common/src/result_stream.rs b/backend/windmill-common/src/result_stream.rs new file mode 100644 index 0000000000..4daa0cbd40 --- /dev/null +++ b/backend/windmill-common/src/result_stream.rs @@ -0,0 +1,33 @@ +use uuid::Uuid; +use crate::{error, DB}; + +pub const STREAM_PREFIX: &str = "WM_STREAM: "; + +pub fn extract_stream_from_logs(line: &str) -> Option { + if line.starts_with(STREAM_PREFIX) { + // Extract the content after "WM_STREAM:" prefix + let stream_content = line.strip_prefix(STREAM_PREFIX).unwrap_or(""); + if !stream_content.is_empty() { + return Some(stream_content.to_string().replace("\\n", "\n")); + } + } + None +} + + + +pub async fn append_result_stream_db(db: &DB, workspace_id: &str, job_id: &Uuid, nstream: &str) -> error::Result<()> { + if !nstream.is_empty() { + sqlx::query!( + r#" + INSERT INTO job_result_stream (workspace_id, job_id, stream) + VALUES ($1, $2, $3) + ON CONFLICT (job_id) DO UPDATE SET stream = job_result_stream.stream || $3 + "#, + workspace_id, + job_id, + nstream, + ).execute(db).await?; + } + Ok(()) +} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 5334a1ac90..76a763f463 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1281,7 +1281,7 @@ async fn restart_job_if_perpetual_inner( #[cfg(feature = "enterprise")] async fn has_failure_module(db: &Pool, job: &MiniPulledJob) -> bool { - if let Ok(flow) = cache::job::fetch_flow(db, job.kind, job.runnable_id).await { + if let Ok(flow) = cache::job::fetch_flow(db, &job.kind, job.runnable_id).await { return flow.value().failure_module.is_some(); } sqlx::query_scalar!( @@ -4785,7 +4785,7 @@ async fn restarted_flows_resolution( )) })?; - let flow_data = cache::job::fetch_flow(db, row.job_kind, row.script_hash) + let flow_data = cache::job::fetch_flow(db, &row.job_kind, row.script_hash) .or_else(|_| cache::job::fetch_preview_flow(db.into(), &completed_flow_id, row.raw_flow)) .await?; let flow_value = flow_data.value(); diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index a508f3084c..49819a41c5 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -364,7 +364,7 @@ pub async fn install_bun_lockfile( occupancy_metrics, None, ) - .await? + .await?; } else { Box::into_pin(child_process.wait()).await?; } @@ -1081,6 +1081,10 @@ function argsObjToArr({{ {spread} }}) {{ return [ {spread} ]; }} +function isAsyncIterable(obj) {{ + return obj != null && typeof obj[Symbol.asyncIterator] === 'function'; +}} + BigInt.prototype.toJSON = function () {{ return this.toString(); }}; @@ -1093,6 +1097,12 @@ async function run() {{ throw new Error("{main_name} function is missing"); }} let res = await Main.{main_name}(...argsArr); + if (isAsyncIterable(res)) {{ + for await (const chunk of res) {{ + console.log("WM_STREAM: " + chunk.replace('\n', '\\n')); + }} + res = null; + }} const res_json = JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value); await fs.writeFile("result.json", res_json); process.exit(0); @@ -1441,7 +1451,7 @@ try {{ .await? }; - handle_child( + let handle_result = handle_child( &job.id, conn, mem_peak, @@ -1474,7 +1484,7 @@ try {{ })?; *new_args = Some(args.clone()); } - read_result(job_dir).await + read_result(job_dir, handle_result.result_stream).await } pub async fn get_common_bun_proc_envs(base_internal_url: Option<&str>) -> HashMap { diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index c29f00c963..ea7f4a85f2 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -359,11 +359,41 @@ pub async fn read_file(path: &str) -> error::Result> { return Ok(r); } +pub async fn merge_result_stream( + result: error::Result>, + result_stream: Option, +) -> error::Result> { + if let Some(result_stream) = result_stream { + result.and_then(|x| { + let mut value: Value = serde_json::from_str(x.get())?; + + // Insert the string at the "wm_stream" field + if let Value::Object(ref mut map) = value { + map.insert("wm_stream".to_string(), Value::String(result_stream)); + } else if value.is_null() { + // return Ok(unsafe_raw(json)) + return Ok(to_raw_value(&json!(result_stream))); + } else { + return Ok(x); + } + + // Convert back to RawValue + let json_string = serde_json::to_string(&value)?; + Ok(RawValue::from_string(json_string)?) + }) + } else { + result + } +} /// Read the `result.json` file. This function assumes that the file contains valid json and will /// result in undefined behaviour if it isn't. If the result.json is user generated or otherwise /// not guaranteed to be valid, use `read_and_check_result` -pub async fn read_result(job_dir: &str) -> error::Result> { - return read_file(&format!("{job_dir}/result.json")).await; +pub async fn read_result( + job_dir: &str, + result_stream: Option, +) -> error::Result> { + let rf = read_file(&format!("{job_dir}/result.json")).await; + merge_result_stream(rf, result_stream).await } pub async fn read_and_check_file(path: &str) -> error::Result> { diff --git a/backend/windmill-worker/src/csharp_executor.rs b/backend/windmill-worker/src/csharp_executor.rs index a8ed06a2e2..4bb554d15f 100644 --- a/backend/windmill-worker/src/csharp_executor.rs +++ b/backend/windmill-worker/src/csharp_executor.rs @@ -641,5 +641,5 @@ pub async fn handle_csharp_job( None, ) .await?; - read_result(job_dir).await + read_result(job_dir, None).await } diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index 75a80e39d8..04f8c591fe 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -126,7 +126,7 @@ pub async fn handle_dedicated_process( let status = Box::into_pin(child.wait()) .await .expect("child process encountered an error"); - if let Err(e) = process_status(&cmd_name, status) { + if let Err(e) = process_status(&cmd_name, status, vec![]) { tracing::error!("child exit status was not success: {e:#}"); } else { tracing::info!("child exit status was success"); diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index e2bf579b53..5c2c22c540 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -283,6 +283,10 @@ BigInt.prototype.toJSON = function () {{ return this.toString(); }}; +function isAsyncIterable(obj) {{ + return obj != null && typeof obj[Symbol.asyncIterator] === 'function'; +}} + async function run() {{ {dates} {preprocessor} @@ -291,6 +295,12 @@ async function run() {{ throw new Error("{main_name} function is missing"); }} let res: any = await {main_name}(...argsArr); + if (isAsyncIterable(res)) {{ + for await (const chunk of res) {{ + console.log("WM_STREAM: " + chunk.replace('\n', '\\n')); + }} + res = null; + }} const res_json = JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value); await Deno.writeTextFile("result.json", res_json); Deno.exit(0); @@ -408,7 +418,7 @@ try {{ }; // logs.push_str(format!("prepare: {:?}\n", start.elapsed().as_micros()).as_str()); // start = Instant::now(); - handle_child( + let handle_result = handle_child( &job.id, conn, mem_peak, @@ -445,7 +455,7 @@ try {{ })?; *new_args = Some(args.clone()); } - read_result(job_dir).await + read_result(job_dir, handle_result.result_stream).await } async fn build_import_map( diff --git a/backend/windmill-worker/src/go_executor.rs b/backend/windmill-worker/src/go_executor.rs index 88b8863884..97a2d15d03 100644 --- a/backend/windmill-worker/src/go_executor.rs +++ b/backend/windmill-worker/src/go_executor.rs @@ -408,7 +408,7 @@ func Run(req Req) (interface{{}}, error){{ run_go.stdout(Stdio::piped()).stderr(Stdio::piped()); start_child_process(run_go, &compiled_executable_name).await? }; - handle_child( + let handle_result = handle_child( &job.id, conn, mem_peak, @@ -425,7 +425,7 @@ func Run(req Req) (interface{{}}, error){{ ) .await?; - read_result(job_dir).await + read_result(job_dir, handle_result.result_stream).await } async fn gen_go_mod( diff --git a/backend/windmill-worker/src/handle_child.rs b/backend/windmill-worker/src/handle_child.rs index 82f708be49..11078c7892 100644 --- a/backend/windmill-worker/src/handle_child.rs +++ b/backend/windmill-worker/src/handle_child.rs @@ -7,6 +7,7 @@ use nix::unistd::Pid; use process_wrap::tokio::TokioChildWrapper; use windmill_common::agent_workers::PingJobStatusResponse; use windmill_common::jobs::LARGE_LOG_THRESHOLD_SIZE; +use windmill_common::result_stream::extract_stream_from_logs; #[cfg(windows)] use std::process::Stdio; @@ -52,7 +53,7 @@ use futures::{ }; use crate::common::{resolve_job_timeout, OccupancyMetrics}; -use crate::job_logger::{append_job_logs, append_with_limit}; +use crate::job_logger::{append_job_logs, append_result_stream, append_with_limit}; use crate::job_logger_oss::process_streaming_log_lines; use crate::worker_utils::{ping_job_status, update_worker_ping_from_job}; use crate::{MAX_RESULT_SIZE, MAX_WAIT_FOR_SIGINT, MAX_WAIT_FOR_SIGTERM}; @@ -87,6 +88,10 @@ async fn kill_process_tree(pid: Option) -> Result<(), String> { } } +pub struct HandleChildResult { + pub result_stream: Option, +} + /// - wait until child exits and return with exit status /// - read lines from stdout and stderr and append them to the "queue"."logs" /// quitting early if output exceedes MAX_LOG_SIZE characters (not bytes) @@ -109,7 +114,7 @@ pub async fn handle_child( occupancy_metrics: &mut Option<&mut OccupancyMetrics>, // Do not print logs to output, but instead save to string. pipe_stdout: Option<&mut String>, -) -> error::Result<()> { +) -> error::Result { let start = Instant::now(); let pid = child.id(); @@ -296,6 +301,7 @@ pub async fn handle_child( } }; + let mut stream_result = Vec::new(); /* a future that reads output from the child and appends to the database */ let lines = write_lines( output, @@ -308,6 +314,7 @@ pub async fn handle_child( pipe_stdout, &mut rx2, child_name, + &mut stream_result, ) .instrument(trace_span!("child_lines")); @@ -322,7 +329,7 @@ pub async fn handle_child( _ if *too_many_logs.borrow() => Err(Error::ExecutionErr(format!( "logs or result reached limit. (current max size: {MAX_RESULT_SIZE} characters)" ))), - Ok(Ok(status)) => process_status(&child_name, status), + Ok(Ok(status)) => process_status(&child_name, status, stream_result), Ok(Err(kill_reason)) => match kill_reason { KillReason::AlreadyCompleted => { Err(Error::AlreadyCompleted("Job already completed".to_string())) @@ -346,6 +353,7 @@ pub async fn write_lines( pipe_stdout: Option<&mut String>, rx2: &mut broadcast::Receiver<()>, child_name: &str, + stream_result: &mut Vec, ) { let max_log_size = if *CLOUD_HOSTED { MAX_RESULT_SIZE @@ -401,13 +409,25 @@ pub async fn write_lines( let mut joined = String::new(); let job_id = job_id.clone(); + let mut nstream = String::new(); while let Some(line) = read_lines.next().await { match line { Ok(line) => { if line.is_empty() { continue; } - append_with_limit(&mut joined, &line, &mut log_remaining); + if let Some(stream) = extract_stream_from_logs(&line) { + let len = stream.len(); + if log_remaining >= len { + log_remaining -= len; + nstream.push_str(&stream); + stream_result.push(stream); + } else { + log_remaining = 0; + } + } else { + append_with_limit(&mut joined, &line, &mut log_remaining); + } if log_remaining == 0 { tracing::info!(%job_id, "Too many logs lines for job {job_id}"); let _ = set_too_many_logs.send(true); @@ -460,6 +480,14 @@ pub async fn write_lines( let job_id = job_id.clone(); let pg_log_total_size = pg_log_total_size.clone(); (do_write, write_result) = tokio::spawn(async move { + if !nstream.is_empty() { + if let Err(err) = append_result_stream(&conn, &w_id, &job_id, &nstream).await { + tracing::error!( + "Unable to send result stream for job {job_id}. Error was: {:?}", + err + ); + } + } append_job_logs( &job_id, &w_id, @@ -762,9 +790,19 @@ pub fn lines_to_stream( }) } -pub fn process_status(program: &str, status: ExitStatus) -> error::Result<()> { +pub fn process_status( + program: &str, + status: ExitStatus, + stream_result: Vec, +) -> error::Result { if status.success() { - Ok(()) + Ok(HandleChildResult { + result_stream: if stream_result.is_empty() { + None + } else { + Some(stream_result.join("")) + }, + }) } else if let Some(code) = status.code() { Err(error::Error::ExitStatus(program.to_string(), code)) } else { diff --git a/backend/windmill-worker/src/java_executor.rs b/backend/windmill-worker/src/java_executor.rs index f62de150a9..48de5ecbe8 100644 --- a/backend/windmill-worker/src/java_executor.rs +++ b/backend/windmill-worker/src/java_executor.rs @@ -88,7 +88,7 @@ pub async fn handle_java_job<'a>(mut args: JobHandlerInput<'a>) -> Result( &mut Some(occupancy_metrics), None, ) - .await + .await?; + Ok(()) } #[derive(Default, Debug)] diff --git a/backend/windmill-worker/src/job_logger.rs b/backend/windmill-worker/src/job_logger.rs index 626dfbce44..263234cee9 100644 --- a/backend/windmill-worker/src/job_logger.rs +++ b/backend/windmill-worker/src/job_logger.rs @@ -1,10 +1,11 @@ use regex::Regex; pub use windmill_common::jobs::LARGE_LOG_THRESHOLD_SIZE; +use windmill_common::result_stream::append_result_stream_db; use windmill_common::utils::WarnAfterExt; use windmill_common::worker::{Connection, CLOUD_HOSTED}; -use windmill_common::DB; +use windmill_common::{error, DB}; use windmill_queue::append_logs; use std::sync::atomic::AtomicU32; @@ -61,6 +62,32 @@ pub async fn append_job_logs( } } +pub async fn append_result_stream( + conn: &Connection, + workspace_id: &str, + job_id: &Uuid, + nstream: &str, +) -> error::Result<()> { + match conn { + Connection::Sql(db) => { + append_result_stream_db(db, workspace_id, job_id, nstream).await?; + } + Connection::Http(client) => { + if let Err(e) = client + .post::<_, String>( + &format!("/api/w/{}/agent_workers/push_logs/{}", workspace_id, job_id), + None, + &nstream, + ) + .await + { + tracing::error!(%job_id, %e, "error sending result stream for job {job_id}: {e}"); + }; + } + } + Ok(()) +} + pub async fn append_logs_with_compaction( job_id: &Uuid, w_id: &str, diff --git a/backend/windmill-worker/src/js_eval.rs b/backend/windmill-worker/src/js_eval.rs index 77ccd9b15f..d5bf2add3b 100644 --- a/backend/windmill-worker/src/js_eval.rs +++ b/backend/windmill-worker/src/js_eval.rs @@ -701,7 +701,7 @@ pub struct MainArgs { #[cfg(feature = "deno_core")] pub struct LogString { - pub s: String, + pub s: mpsc::UnboundedSender, } #[cfg(feature = "deno_core")] @@ -894,6 +894,8 @@ pub async fn eval_fetch_timeout( return y*2; }); + let (log_sender, mut log_receiver) = mpsc::unbounded_channel::(); + { let op_state = js_runtime.op_state(); let mut op_state = op_state.borrow_mut(); @@ -901,7 +903,7 @@ pub async fn eval_fetch_timeout( //reqwest client seems to not be sharable between runtimes unfortunately // op_state.put(HTTP_CLIENT.clone()); op_state.put(MainArgs { args: spread }); - op_state.put(LogString { s: String::new() }); + op_state.put(LogString { s: log_sender }); } sender @@ -913,23 +915,51 @@ pub async fn eval_fetch_timeout( .build()?; let future = async { + use crate::common::merge_result_stream; + + if !extra_logs.is_empty() { + append_logs(&job_id, w_id_.as_str(), format!("{extra_logs}"), &conn_).await; + } + let w_id = w_id_.clone(); + let handle = tokio::spawn(async move { + let mut result_stream = String::new(); + while let Some(log) = log_receiver.recv().await { + use windmill_common::result_stream::extract_stream_from_logs; + + if let Some(stream) = extract_stream_from_logs(&log.trim_end_matches("\n")) { + use crate::job_logger::append_result_stream; + + result_stream.push_str(&stream); + if let Err(e) = append_result_stream(&conn_, &w_id, &job_id, &stream).await + { + tracing::error!("failed to append result stream for job {job_id}: {e}"); + } + } else { + append_logs(&job_id, w_id_.as_str(), log, &conn_).await; + } + } + if !result_stream.is_empty() { + Some(result_stream) + } else { + None + } + }); + let r = tokio::select! { r = eval_fetch(&mut js_runtime, &js_expr, Some(env_code), script_entrypoint_override, load_client, &job_id) => Ok(r), _ = memory_limit_rx.recv() => Err(Error::ExecutionErr("Memory limit reached, killing isolate".to_string())) }; - - append_logs( - &job_id, - w_id_.as_str(), - format!( - "{extra_logs}{}", - js_runtime.op_state().borrow().borrow::().s - ), - &conn_, - ) - .await; - - r + drop(js_runtime); + if let Ok(r) = r { + match handle.await { + Ok(Some(logs)) => Ok(merge_result_stream(r, Some(logs)).await), + Ok(None) => Ok(r), + Err(e) => Err(Error::ExecutionErr(e.to_string())), + } + } else { + r + } + // r }; let r = runtime.block_on(future)?; // tracing::info!("total: {:?}", instant.elapsed()); @@ -1039,8 +1069,46 @@ async fn eval_fetch( "", format!( r#" +function isAsyncIterable(obj) {{ + // return true; // TODO: remove this + return obj != null && typeof obj[Symbol.asyncIterator] === 'function'; +}} + +function processStreamIterative(res) {{ + const iterator = res[Symbol.asyncIterator](); + + function processLoop() {{ + return new Promise(function(resolve) {{ + function step() {{ + iterator.next().then(function(result) {{ + if (!result.done) {{ + const chunk = result.value; + console.log("WM_STREAM: " + chunk.replace('\n', '\\n')); + // Continue the loop + step(); + }} else {{ + resolve("null"); + }} + }}).catch(function(error) {{ + resolve("null"); + }}); + }} + step(); + }}); + }} + + return processLoop(); +}} + let args = Deno.core.ops.op_get_static_args().map(JSON.parse) -import("file:///eval.ts").then((module) => module.{main_override}(...args)).then(JSON.stringify) +import("file:///eval.ts").then((module) => module.{main_override}(...args)) + .then(res => {{ + if (isAsyncIterable(res)) {{ + return processStreamIterative(res) + }} else {{ + return JSON.stringify(res ?? null); + }} + }}) "# ), ) @@ -1120,11 +1188,14 @@ fn op_get_static_args(op_state: Rc>) -> Vec> { #[op2(fast)] fn op_log(op_state: Rc>, #[string] log: &str) { // tracing::error!("log: |{}|", log); - op_state + if let Err(e) = op_state .borrow_mut() .borrow_mut::() .s - .push_str(log); + .send(log.to_string()) + { + tracing::error!("failed to send log: {e}"); + } } #[cfg(feature = "deno_core")] diff --git a/backend/windmill-worker/src/nu_executor.rs b/backend/windmill-worker/src/nu_executor.rs index c31d979c27..3e7d221163 100644 --- a/backend/windmill-worker/src/nu_executor.rs +++ b/backend/windmill-worker/src/nu_executor.rs @@ -20,7 +20,6 @@ use crate::{ }; use windmill_common::client::AuthedClient; - const NSJAIL_CONFIG_RUN_NU_CONTENT: &str = include_str!("../nsjail/run.nu.config.proto"); lazy_static::lazy_static! { static ref NU_PATH: String = std::env::var("NU_PATH").unwrap_or_else(|_| "/usr/bin/nu".to_string()); @@ -69,7 +68,7 @@ pub async fn handle_nu_job<'a>(mut args: JobHandlerInput<'a>) -> Result( &mut Some(occupancy_metrics), None, ) - .await + .await?; + Ok(()) } // #[cfg(test)] // mod test { diff --git a/backend/windmill-worker/src/php_executor.rs b/backend/windmill-worker/src/php_executor.rs index ec8478a9a5..f3e041073f 100644 --- a/backend/windmill-worker/src/php_executor.rs +++ b/backend/windmill-worker/src/php_executor.rs @@ -20,8 +20,7 @@ use crate::{ read_result, start_child_process, OccupancyMetrics, }, handle_child::handle_child, - COMPOSER_CACHE_DIR, COMPOSER_PATH, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, - PHP_PATH, + COMPOSER_CACHE_DIR, COMPOSER_PATH, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PHP_PATH, }; use windmill_common::client::AuthedClient; @@ -345,5 +344,5 @@ try {{ None, ) .await?; - read_result(job_dir).await + read_result(job_dir, None).await } diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index fb760fb168..a7f2e608ad 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -679,8 +679,7 @@ replace_invalid_fields = re.compile(r'(?:\bNaN\b|\\*\\u0000|Infinity|\-Infinity) result_json = os.path.join(os.path.abspath(os.path.dirname(__file__)), "result.json") -def res_to_json(res): - typ = type(res) +def res_to_json(res, typ): if typ.__name__ == 'DataFrame': if typ.__module__ == 'pandas.core.frame': res = res.values.tolist() @@ -704,7 +703,12 @@ try: if inner_script.{main_override} is None or not callable(inner_script.{main_override}): raise ValueError("{main_override} function is missing") res = inner_script.{main_override}(**args) - res_json = res_to_json(res) + typ = type(res) + if hasattr(res, '__iter__') and not isinstance(res, (str, dict, list, bytes, tuple, set, frozenset, range, memoryview, bytearray)) and typ.__name__ != 'DataFrame': + for chunk in res: + print("WM_STREAM: " + chunk.replace('\n', '\\n')) + res = None + res_json = res_to_json(res, typ) with open(result_json, 'w') as f: f.write(res_json) except BaseException as e: @@ -858,7 +862,7 @@ mount {{ start_child_process(python_cmd, &python_path).await? }; - handle_child( + let handle_result = handle_child( &job.id, conn, mem_peak, @@ -892,7 +896,7 @@ mount {{ *new_args = Some(args.clone()); } - read_result(job_dir).await + read_result(job_dir, handle_result.result_stream).await } async fn prepare_wrapper( diff --git a/backend/windmill-worker/src/python_versions.rs b/backend/windmill-worker/src/python_versions.rs index 99fa67f1fe..8b61ca532d 100644 --- a/backend/windmill-worker/src/python_versions.rs +++ b/backend/windmill-worker/src/python_versions.rs @@ -656,7 +656,8 @@ impl PyV { occupancy_metrics, None, ) - .await + .await?; + Ok(()) } async fn find_python(&self) -> error::Result> { #[cfg(windows)] diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index a5de05c9e9..cf455e3f35 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -382,7 +382,7 @@ pub async fn process_result( Err(e) => { let error_value = match e { Error::ExitStatus(program, i) => { - let res = read_result(job_dir).await.ok(); + let res = read_result(job_dir, None).await.ok(); if res.as_ref().is_some_and(|x| !x.get().is_empty()) { res.unwrap() diff --git a/backend/windmill-worker/src/rust_executor.rs b/backend/windmill-worker/src/rust_executor.rs index 045158511e..0475260d86 100644 --- a/backend/windmill-worker/src/rust_executor.rs +++ b/backend/windmill-worker/src/rust_executor.rs @@ -600,5 +600,5 @@ pub async fn handle_rust_job( None, ) .await?; - read_result(job_dir).await + read_result(job_dir, None).await } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 33fb3f5213..7ca2d0d269 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -2401,7 +2401,7 @@ pub async fn handle_queued_job( let flow_data = match preview_data { Some(RawData::Flow(data)) => data, // Not a preview: fetch from the cache or the database. - _ => cache::job::fetch_flow(db, job.kind, job.runnable_id).await?, + _ => cache::job::fetch_flow(db, &job.kind, job.runnable_id).await?, }; handle_flow( job, diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 00477e7510..6abe8e57ad 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -287,7 +287,7 @@ pub async fn update_flow_status_after_job_completion_internal( )) })?; - let flow_data = cache::job::fetch_flow(db, job_kind, script_hash) + let flow_data = cache::job::fetch_flow(db, &job_kind, script_hash) .or_else(|_| cache::job::fetch_preview_flow(db, &flow, raw_flow)) .await?; let flow_value = flow_data.value(); diff --git a/frontend/src/lib/components/DisplayResult.svelte b/frontend/src/lib/components/DisplayResult.svelte index 89ce56aa91..fbb5641b95 100644 --- a/frontend/src/lib/components/DisplayResult.svelte +++ b/frontend/src/lib/components/DisplayResult.svelte @@ -39,6 +39,7 @@ import { getContext, hasContext, createEventDispatcher, onDestroy } from 'svelte' import { toJsonStr } from '$lib/utils' import { userStore } from '$lib/stores' + import ResultStreamDisplay from './ResultStreamDisplay.svelte' const IMG_MAX_SIZE = 10000000 const TABLE_MAX_SIZE = 5000000 @@ -83,12 +84,14 @@ noControls?: boolean drawerOpen?: boolean nodeId?: string | undefined + loading?: boolean | undefined language?: string | undefined appPath?: string | undefined customUi?: DisplayResultUi | undefined isTest?: boolean externalToolbarAvailable?: boolean forceJson?: boolean + result_stream?: string | undefined fixTableSizingToParent?: boolean copilot_fix?: import('svelte').Snippet children?: import('svelte').Snippet @@ -111,9 +114,11 @@ isTest = true, externalToolbarAvailable = false, forceJson = $bindable(false), + result_stream = undefined, fixTableSizingToParent = false, copilot_fix, - children + children, + loading = false }: Props = $props() let enableHtml = $state(false) let s3FileDisplayRawMode = $state(false) @@ -487,7 +492,15 @@ -{#if is_render_all} + +{#if result_stream && result == undefined} +
+
+ Streaming result +
+ +
+{:else if is_render_all}
{#if !noControls}
@@ -690,7 +703,7 @@ {:else if !forceJson && resultKind === 'plain'}
{typeof result === 'string' ? result : result?.['result']}
{#if !noControls} + >{#if !noControls && !loading}
diff --git a/frontend/src/lib/components/FlowStatusViewerInner.svelte b/frontend/src/lib/components/FlowStatusViewerInner.svelte index 2f4f6e1e17..4846956038 100644 --- a/frontend/src/lib/components/FlowStatusViewerInner.svelte +++ b/frontend/src/lib/components/FlowStatusViewerInner.svelte @@ -83,10 +83,17 @@ subflowParentsDurationStatuses?: Writable>[] isForloopSelected?: boolean parentRecursiveRefresh?: Record Promise> - job?: Job | undefined + job?: (Job & { result_stream?: string }) | undefined rightColumnSelect?: 'timeline' | 'node_status' | 'node_definition' | 'user_states' localModuleStates?: Writable> localDurationStatuses?: Writable> + onResultStreamUpdate?: ({ + jobId, + result_stream + }: { + jobId: string + result_stream?: string + }) => void customUi?: { tagLabel?: string | undefined } @@ -119,8 +126,24 @@ rightColumnSelect = $bindable('timeline'), localModuleStates = writable({}), localDurationStatuses = writable({}), - customUi + customUi, + onResultStreamUpdate = undefined }: Props = $props() + + let resultStreams: Record = $state({}) + + if (onResultStreamUpdate == undefined) { + onResultStreamUpdate = ({ + jobId, + result_stream + }: { + jobId: string + result_stream?: string + }) => { + resultStreams[jobId] = result_stream + } + } + let recursiveRefresh: Record Promise> = $state({}) // Add support for the input args assets shown as an asset node @@ -516,6 +539,9 @@ jobLoader?.watchJob(jobId, { change(newJob) { setJob(newJob, true) + }, + resultStreamUpdate({ id, result_stream }: { id: string; result_stream?: string }) { + onResultStreamUpdate?.({ jobId: id, result_stream }) } }) } @@ -959,6 +985,7 @@ @@ -976,6 +1003,7 @@ {innerModules} {suspendStatus} {hideJobId} + result_streams={resultStreams} />
{/if} @@ -1067,6 +1095,7 @@ storedListJobs[j] = job innerJobLoaded(job, j, false, force) }} + {onResultStreamUpdate} />
{/if} @@ -1142,6 +1171,7 @@ {reducedPolling} {workspaceId} jobId={failedRetry} + {onResultStreamUpdate} />
{/each} @@ -1174,6 +1204,7 @@ let { force, job } = e.detail onJobsLoaded(mod, job, force) }} + {onResultStreamUpdate} /> {:else if mod.flow_jobs?.length == 0 && mod.job == '00000000-0000-0000-0000-000000000000'}
no subflow (empty loop?)
@@ -1205,6 +1236,7 @@ let { job, force } = e.detail onJobsLoaded(mod, job, force) }} + {onResultStreamUpdate} /> {/if} {:else} @@ -1428,6 +1460,7 @@ waitingForExecutor={node.type == 'WaitingForExecutor'} refreshLog={node.type == 'InProgress'} col + result_stream={resultStreams[node.job_id ?? '']} result={node.result} tag={node.tag} logs={node.logs} diff --git a/frontend/src/lib/components/JobLoader.svelte b/frontend/src/lib/components/JobLoader.svelte index f6e6730b09..2d4f30303f 100644 --- a/frontend/src/lib/components/JobLoader.svelte +++ b/frontend/src/lib/components/JobLoader.svelte @@ -23,11 +23,12 @@ cancel?: ({ id }: { id: string }) => void started?: ({ id }: { id: string }) => void running?: ({ id }: { id: string }) => void + resultStreamUpdate?: ({ id, result_stream }: { id: string; result_stream?: string }) => void } interface Props { isLoading?: boolean - job?: Job | undefined + job?: (Job & { result_stream?: string }) | undefined noCode?: boolean noLogs?: boolean workspaceOverride?: string | undefined @@ -35,7 +36,6 @@ allowConcurentRequests?: boolean jobUpdateLastFetch?: Date | undefined toastError?: boolean - lazyLogs?: boolean onlyResult?: boolean // If you want to find out progress of subjobs of a flow, check job.flow_status.progress scriptProgress?: number | undefined @@ -52,7 +52,6 @@ notfound = $bindable(false), jobUpdateLastFetch = $bindable(undefined), toastError = false, - lazyLogs = false, onlyResult = false, scriptProgress = $bindable(undefined), noLogs = false, @@ -74,6 +73,7 @@ let errorIteration = 0 let logOffset = 0 + let resultStreamOffset = 0 let lastCallbacks: Callbacks | undefined = undefined let finished: string[] = [] @@ -194,6 +194,9 @@ if (logOffset == 0) { logOffset = job?.logs?.length ? job.logs?.length + 1 : 0 } + if (resultStreamOffset == 0) { + resultStreamOffset = job?.result_stream?.length ? job.result_stream?.length + 1 : 0 + } } export async function getLogs() { if (job) { @@ -285,6 +288,7 @@ let startedWatchingJob: number | undefined = undefined export async function watchJob(testId: string, callbacks?: Callbacks) { logOffset = 0 + resultStreamOffset = 0 syncIteration = 0 errorIteration = 0 currentId = testId @@ -334,7 +338,7 @@ function updateJobFromProgress( previewJobUpdates: GetJobUpdatesResponse, - job: Job, + job: Job & { result_stream?: string }, callbacks: Callbacks | undefined ) { // Clamp number between two values with the following line: @@ -357,10 +361,26 @@ } } + if (previewJobUpdates.new_result_stream) { + if (!job.result_stream) { + job.result_stream = previewJobUpdates.new_result_stream + } else { + job.result_stream = job.result_stream.concat(previewJobUpdates.new_result_stream) + } + callbacks?.resultStreamUpdate?.({ + id: job.id, + result_stream: job.result_stream + }) + } + if (previewJobUpdates.log_offset) { logOffset = previewJobUpdates.log_offset ?? 0 } + if (previewJobUpdates.stream_offset) { + resultStreamOffset = previewJobUpdates.stream_offset ?? 0 + } + if (previewJobUpdates.flow_status) { job.flow_status = previewJobUpdates.flow_status as FlowStatus } @@ -415,7 +435,7 @@ job = await JobService.getJob({ workspace: workspace!, id, - noLogs: lazyLogs || onlyResult || noLogs, + noLogs: onlyResult || noLogs, noCode }) } @@ -504,6 +524,7 @@ callbacks?: Callbacks ): Promise { let isCompleted = false + let resultOnlyResultStream: string = '' if (isCurrentJob(id)) { try { // First load the job to get initial state @@ -511,11 +532,18 @@ job = await JobService.getJob({ workspace: workspace!, id, - noLogs: lazyLogs || noLogs, + noLogs: noLogs, noCode }) } + if (!onlyResult) { + callbacks?.resultStreamUpdate?.({ + id, + result_stream: undefined + }) + } + // If job is already completed, don't start SSE if (job?.type === 'CompletedJob') { isCompleted = true @@ -555,6 +583,12 @@ if (startedWatchingJob && startedWatchingJob > Date.now() - 5000) { params.set('fast', 'true') } + if (noLogs) { + params.set('no_logs', 'true') + } + if (resultStreamOffset) { + params.set('stream_offset', resultStreamOffset.toString()) + } const sseUrl = `/api/w/${workspace}/jobs_u/getupdate_sse/${id}?${params.toString()}` @@ -598,6 +632,17 @@ callbacks?.running?.({ id }) } + if (onlyResult && previewJobUpdates.new_result_stream) { + resultOnlyResultStream = resultOnlyResultStream.concat( + previewJobUpdates.new_result_stream + ) + // console.log('resultOnlyResultStream', resultOnlyResultStream) + callbacks?.resultStreamUpdate?.({ + id, + result_stream: resultOnlyResultStream + }) + } + // Check if job is completed if (previewJobUpdates.completed) { currentEventSource?.close() @@ -611,8 +656,9 @@ }) clearCurrentId() } else { - const njob = previewJobUpdates.job as Job + const njob = previewJobUpdates.job as Job & { result_stream?: string } njob.logs = job?.logs ?? '' + njob.result_stream = job?.result_stream ?? '' job = njob onJobCompleted(id, job, callbacks) } diff --git a/frontend/src/lib/components/ResultStreamDisplay.svelte b/frontend/src/lib/components/ResultStreamDisplay.svelte new file mode 100644 index 0000000000..852a67992e --- /dev/null +++ b/frontend/src/lib/components/ResultStreamDisplay.svelte @@ -0,0 +1,5 @@ + + +
{result_stream}
diff --git a/frontend/src/lib/components/apps/components/display/AppDisplayComponent.svelte b/frontend/src/lib/components/apps/components/display/AppDisplayComponent.svelte index e334e2ab2d..6a32d62e7d 100644 --- a/frontend/src/lib/components/apps/components/display/AppDisplayComponent.svelte +++ b/frontend/src/lib/components/apps/components/display/AppDisplayComponent.svelte @@ -26,6 +26,7 @@ configuration: RichConfigurations } + let result_stream: string | undefined = $state(undefined) let { id, componentInput, @@ -57,6 +58,7 @@ }) let css = $state(initCss($app.css?.displaycomponent, customCss)) + let loading = $state(false) {#each Object.keys(components['displaycomponent'].initialData.configuration) as key (key)} @@ -78,7 +80,15 @@ /> {/each} - +
('AppViewerContext') let result: any = $state(noBackend ? runnable.noBackendValue : undefined) - export function onSuccess() { if (runnable.recomputeIds) { runnable.recomputeIds.forEach((id) => $runnableComponents?.[id]?.cb?.map((cb) => cb())) diff --git a/frontend/src/lib/components/apps/components/helpers/RunnableComponent.svelte b/frontend/src/lib/components/apps/components/helpers/RunnableComponent.svelte index 387065b197..6f3dec9d84 100644 --- a/frontend/src/lib/components/apps/components/helpers/RunnableComponent.svelte +++ b/frontend/src/lib/components/apps/components/helpers/RunnableComponent.svelte @@ -38,6 +38,7 @@ extraQueryParams?: Record autoRefresh?: boolean result?: any + result_stream?: string forceSchemaDisplay?: boolean wrapperClass?: string wrapperStyle?: string @@ -72,6 +73,7 @@ extraQueryParams = {}, autoRefresh = true, result = $bindable(undefined), + result_stream = $bindable(undefined), forceSchemaDisplay = false, wrapperClass = '', wrapperStyle = '', @@ -226,6 +228,15 @@ loading = false dispatch('done', { id, result }) }, + resultStreamUpdate({ + id, + result_stream: nresult_stream + }: { + id: string + result_stream?: string + }) { + setResult(nresult_stream, id) + }, cancel({ id }: { id: string }) { onCancel?.() let jobId = id diff --git a/frontend/src/lib/components/apps/editor/RunnableJobPanelInner.svelte b/frontend/src/lib/components/apps/editor/RunnableJobPanelInner.svelte index f1c0d4aba7..9d9f36b0a7 100644 --- a/frontend/src/lib/components/apps/editor/RunnableJobPanelInner.svelte +++ b/frontend/src/lib/components/apps/editor/RunnableJobPanelInner.svelte @@ -42,11 +42,12 @@
- {:else if testJob != undefined && 'result' in testJob && testJob.result != undefined} + {:else if testJob != undefined && (testJob.type == 'CompletedJob' || testJob.result_stream)}
content?: import('svelte').Snippet onSelectedChange?: (value: string) => void + onTabClick?: (value: string) => void } let { @@ -27,7 +28,8 @@ values = undefined, children, content, - onSelectedChange + onSelectedChange, + onTabClick }: Props = $props() const selectedStore = writable(selected) @@ -37,6 +39,7 @@ update: (value: string) => { selectedStore.set(value) selected = value + onTabClick?.(value) }, hashNavigation }) @@ -55,9 +58,11 @@ } } } + $effect(() => { selected && untrack(() => updateSelected()) }) + $effect(() => { $selectedStore && untrack(() => onSelectedChange?.($selectedStore)) }) diff --git a/frontend/src/lib/components/details/DetailPageHeader.svelte b/frontend/src/lib/components/details/DetailPageHeader.svelte index 42eda23ff5..25e55a2c8a 100644 --- a/frontend/src/lib/components/details/DetailPageHeader.svelte +++ b/frontend/src/lib/components/details/DetailPageHeader.svelte @@ -111,7 +111,6 @@