mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-07 16:03:21 +00:00
feat: use sse for flow status updates
This commit is contained in:
+10
-4
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"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 CASE\n -- flow step:\n WHEN flow_step_id IS NOT NULL THEN NULL\n -- completed:\n WHEN c.id IS NOT NULL THEN COALESCE(\n c.workflow_as_code_status || c.flow_status,\n c.workflow_as_code_status,\n c.flow_status\n )\n -- not completed:\n ELSE COALESCE(\n f.workflow_as_code_status || f.flow_status,\n f.workflow_as_code_status,\n f.flow_status\n )\n END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",\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))",
|
||||
"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<Box<RawValue>>\",\n COALESCE(c.workflow_as_code_status, f.workflow_as_code_status) AS \"workflow_as_code_status: sqlx::types::Json<Box<RawValue>>\",\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": [
|
||||
{
|
||||
@@ -30,16 +30,21 @@
|
||||
},
|
||||
{
|
||||
"ordinal": 5,
|
||||
"name": "workflow_as_code_status: sqlx::types::Json<Box<RawValue>>",
|
||||
"type_info": "Jsonb"
|
||||
},
|
||||
{
|
||||
"ordinal": 6,
|
||||
"name": "log_offset",
|
||||
"type_info": "Int4"
|
||||
},
|
||||
{
|
||||
"ordinal": 6,
|
||||
"ordinal": 7,
|
||||
"name": "created_by!",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 7,
|
||||
"ordinal": 8,
|
||||
"name": "progress",
|
||||
"type_info": "Int4"
|
||||
}
|
||||
@@ -61,9 +66,10 @@
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
false,
|
||||
null
|
||||
]
|
||||
},
|
||||
"hash": "28258e41e95c7c86af0bc3b7469df614a9b71bc372b7994802edac93e6c2ada6"
|
||||
"hash": "278bc6b4f149f824b5db32dacfaa714ee3852dc2ddf2d661dfdd5a986a9bb62b"
|
||||
}
|
||||
+1
-1
@@ -15,7 +15,7 @@
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
null
|
||||
true
|
||||
]
|
||||
},
|
||||
"hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55"
|
||||
|
||||
+2
-1
@@ -13,7 +13,8 @@
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"s3object",
|
||||
"resource"
|
||||
"resource",
|
||||
"variable"
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
+17
-14
@@ -1835,12 +1835,14 @@ async fn handle_zombie_jobs(db: &Pool<Postgres>, base_internal_url: &str, worker
|
||||
);
|
||||
}
|
||||
|
||||
let jobs = sqlx::query_as::<_, QueuedJob>("SELECT * FROM v2_as_queue WHERE id = ANY($1)")
|
||||
.bind(&timeouts[..])
|
||||
.fetch_all(db)
|
||||
.await
|
||||
.map_err(|e| tracing::error!("Error fetching same worker jobs: {:?}", e))
|
||||
.unwrap_or_default();
|
||||
let jobs = sqlx::query_as::<_, QueuedJob>(
|
||||
"SELECT *, null as workflow_as_code_status FROM v2_as_queue WHERE id = ANY($1)",
|
||||
)
|
||||
.bind(&timeouts[..])
|
||||
.fetch_all(db)
|
||||
.await
|
||||
.map_err(|e| tracing::error!("Error fetching same worker jobs: {:?}", e))
|
||||
.unwrap_or_default();
|
||||
|
||||
jobs
|
||||
};
|
||||
@@ -1848,7 +1850,7 @@ async fn handle_zombie_jobs(db: &Pool<Postgres>, base_internal_url: &str, worker
|
||||
let non_restartable_jobs = if *RESTART_ZOMBIE_JOBS {
|
||||
vec![]
|
||||
} else {
|
||||
sqlx::query_as::<_, QueuedJob>("SELECT * FROM v2_as_queue WHERE last_ping < now() - ($1 || ' seconds')::interval
|
||||
sqlx::query_as::<_, QueuedJob>("SELECT *, null as workflow_as_code_status FROM v2_as_queue WHERE last_ping < now() - ($1 || ' seconds')::interval
|
||||
AND running = true AND job_kind NOT IN ('flow', 'flowpreview', 'flownode', 'singlescriptflow') AND same_worker = false")
|
||||
.bind(ZOMBIE_JOB_TIMEOUT.as_str())
|
||||
.fetch_all(db)
|
||||
@@ -1873,13 +1875,14 @@ async fn handle_zombie_jobs(db: &Pool<Postgres>, base_internal_url: &str, worker
|
||||
}
|
||||
}
|
||||
|
||||
let zombie_jobs_restart_limit_reached =
|
||||
sqlx::query_as::<_, QueuedJob>("SELECT * FROM v2_as_queue WHERE id = ANY($1)")
|
||||
.bind(&zombie_jobs_uuid_restart_limit_reached[..])
|
||||
.fetch_all(db)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_else(|| vec![]);
|
||||
let zombie_jobs_restart_limit_reached = sqlx::query_as::<_, QueuedJob>(
|
||||
"SELECT *, null as workflow_as_code_status FROM v2_as_queue WHERE id = ANY($1)",
|
||||
)
|
||||
.bind(&zombie_jobs_uuid_restart_limit_reached[..])
|
||||
.fetch_all(db)
|
||||
.await
|
||||
.ok()
|
||||
.unwrap_or_else(|| vec![]);
|
||||
|
||||
let timeouts = non_restartable_jobs
|
||||
.into_iter()
|
||||
|
||||
@@ -7664,7 +7664,9 @@ paths:
|
||||
progress:
|
||||
type: integer
|
||||
flow_status:
|
||||
$ref: "#/components/schemas/WorkflowStatusRecord"
|
||||
$ref: "../../openflow.openapi.yaml#/components/schemas/FlowStatus"
|
||||
workflow_as_code_status:
|
||||
$ref: "#/components/schemas/WorkflowStatus"
|
||||
|
||||
/w/{workspace}/jobs_u/getupdate_sse/{id}:
|
||||
get:
|
||||
@@ -14056,6 +14058,8 @@ components:
|
||||
the execution of this script will be permissioned_as and by extension its DT_TOKEN.
|
||||
flow_status:
|
||||
$ref: "../../openflow.openapi.yaml#/components/schemas/FlowStatus"
|
||||
workflow_as_code_status:
|
||||
$ref: "#/components/schemas/WorkflowStatus"
|
||||
raw_flow:
|
||||
$ref: "../../openflow.openapi.yaml#/components/schemas/FlowValue"
|
||||
is_flow_step:
|
||||
@@ -14163,6 +14167,8 @@ components:
|
||||
the execution of this script will be permissioned_as and by extension its DT_TOKEN.
|
||||
flow_status:
|
||||
$ref: "../../openflow.openapi.yaml#/components/schemas/FlowStatus"
|
||||
workflow_as_code_status:
|
||||
$ref: "#/components/schemas/WorkflowStatus"
|
||||
raw_flow:
|
||||
$ref: "../../openflow.openapi.yaml#/components/schemas/FlowValue"
|
||||
is_flow_step:
|
||||
|
||||
@@ -18,6 +18,7 @@ use quick_cache::sync::Cache;
|
||||
use serde_json::value::RawValue;
|
||||
use sqlx::Pool;
|
||||
use std::collections::HashMap;
|
||||
use std::hash::{DefaultHasher, Hash, Hasher};
|
||||
use std::ops::{Deref, DerefMut};
|
||||
use std::str::FromStr;
|
||||
use std::time::Instant;
|
||||
@@ -723,7 +724,7 @@ macro_rules! get_job_query {
|
||||
CASE WHEN jsonb_typeof(args) = 'object' THEN args
|
||||
ELSE jsonb_build_object('value', args)
|
||||
END
|
||||
ELSE '{{\"reason\": \"WINDMILL_TOO_BIG\"}}'::jsonb END as args, COALESCE(flow_status, workflow_as_code_status) AS flow_status, \
|
||||
ELSE '{{\"reason\": \"WINDMILL_TOO_BIG\"}}'::jsonb END as args, flow_status, workflow_as_code_status, \
|
||||
{logs} as logs, {code} as raw_code, canceled_by is not null as canceled, canceled_by, canceled_reason, kind as job_kind, \
|
||||
CASE WHEN trigger_kind = 'schedule'::job_trigger_kind THEN trigger END AS schedule_path, permissioned_as, \
|
||||
{flow} as raw_flow, flow_step_id IS NOT NULL AS is_flow_step, script_lang as language, \
|
||||
@@ -3003,6 +3004,7 @@ impl<'a> From<UnifiedJob> for Job {
|
||||
result_columns: None,
|
||||
logs: None,
|
||||
flow_status: None,
|
||||
workflow_as_code_status: None,
|
||||
deleted: uj.deleted,
|
||||
canceled: uj.canceled,
|
||||
canceled_by: uj.canceled_by,
|
||||
@@ -3039,6 +3041,7 @@ impl<'a> From<UnifiedJob> for Job {
|
||||
scheduled_for: uj.scheduled_for.unwrap(),
|
||||
logs: None,
|
||||
flow_status: None,
|
||||
workflow_as_code_status: None,
|
||||
canceled: uj.canceled,
|
||||
canceled_by: uj.canceled_by,
|
||||
canceled_reason: None,
|
||||
@@ -5686,25 +5689,31 @@ pub struct JobUpdate {
|
||||
pub mem_peak: Option<i32>,
|
||||
pub progress: Option<i32>,
|
||||
pub flow_status: Option<Box<serde_json::value::RawValue>>,
|
||||
pub workflow_as_code_status: Option<Box<serde_json::value::RawValue>>,
|
||||
pub job: Option<Job>,
|
||||
pub only_result: Option<Box<serde_json::value::RawValue>>,
|
||||
}
|
||||
|
||||
#[derive(PartialEq)]
|
||||
pub struct JobUpdateLastStatus {
|
||||
pub running: Option<bool>,
|
||||
pub completed: Option<bool>,
|
||||
pub log_offset: Option<i32>,
|
||||
pub mem_peak: Option<i32>,
|
||||
impl JobUpdate {
|
||||
pub fn hash_str(&self) -> String {
|
||||
let mut hasher = DefaultHasher::new();
|
||||
self.hash(&mut hasher);
|
||||
format!("{:x}", hasher.finish())
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&JobUpdate> for JobUpdateLastStatus {
|
||||
fn from(update: &JobUpdate) -> Self {
|
||||
Self {
|
||||
running: update.running,
|
||||
completed: update.completed,
|
||||
log_offset: update.log_offset,
|
||||
mem_peak: update.mem_peak,
|
||||
impl Hash for JobUpdate {
|
||||
fn hash<H: Hasher>(&self, state: &mut H) {
|
||||
self.running.hash(state);
|
||||
self.completed.hash(state);
|
||||
self.log_offset.hash(state);
|
||||
self.mem_peak.hash(state);
|
||||
self.progress.hash(state);
|
||||
if !self.completed.unwrap_or(false) {
|
||||
self.flow_status.as_ref().map(|x| x.get().hash(state));
|
||||
self.workflow_as_code_status
|
||||
.as_ref()
|
||||
.map(|x| x.get().hash(state));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5852,7 +5861,7 @@ fn get_job_update_sse_stream(
|
||||
|
||||
tokio::spawn(async move {
|
||||
let mut log_offset = initial_log_offset;
|
||||
let mut last_update: Option<JobUpdateLastStatus> = None;
|
||||
let mut last_update_hash: Option<String> = None;
|
||||
|
||||
// Send initial update immediately
|
||||
let mut running = running;
|
||||
@@ -5872,7 +5881,7 @@ fn get_job_update_sse_stream(
|
||||
.await
|
||||
{
|
||||
Ok(update) => {
|
||||
last_update = Some((&update).into());
|
||||
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);
|
||||
@@ -5951,9 +5960,9 @@ fn get_job_update_sse_stream(
|
||||
// if !only_result.unwrap_or(false) {
|
||||
// tracing::error!("update {:?}", update);
|
||||
// }
|
||||
let update_last_status = (&update).into();
|
||||
let update_last_status = update.hash_str();
|
||||
// Only send if the update has changed
|
||||
if last_update.as_ref() != Some(&update_last_status) {
|
||||
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);
|
||||
@@ -5966,7 +5975,7 @@ fn get_job_update_sse_stream(
|
||||
break;
|
||||
}
|
||||
|
||||
last_update = Some(update_last_status);
|
||||
last_update_hash = Some(update_last_status);
|
||||
}
|
||||
}
|
||||
Err(_) => {
|
||||
@@ -6068,6 +6077,7 @@ async fn get_job_update_data(
|
||||
progress: None,
|
||||
job: None,
|
||||
flow_status: None,
|
||||
workflow_as_code_status: None,
|
||||
only_result: result.0,
|
||||
})
|
||||
} else {
|
||||
@@ -6080,22 +6090,8 @@ async fn get_job_update_data(
|
||||
END AS running,
|
||||
SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs,
|
||||
COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,
|
||||
CASE
|
||||
-- flow step:
|
||||
WHEN flow_step_id IS NOT NULL THEN NULL
|
||||
-- completed:
|
||||
WHEN c.id IS NOT NULL THEN COALESCE(
|
||||
c.workflow_as_code_status || c.flow_status,
|
||||
c.workflow_as_code_status,
|
||||
c.flow_status
|
||||
)
|
||||
-- not completed:
|
||||
ELSE COALESCE(
|
||||
f.workflow_as_code_status || f.flow_status,
|
||||
f.workflow_as_code_status,
|
||||
f.flow_status
|
||||
)
|
||||
END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
COALESCE(c.flow_status, f.flow_status) AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
COALESCE(c.workflow_as_code_status, f.workflow_as_code_status) AS \"workflow_as_code_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset,
|
||||
created_by AS \"created_by!\",
|
||||
CASE WHEN $4::BOOLEAN THEN (
|
||||
@@ -6140,6 +6136,9 @@ async fn get_job_update_data(
|
||||
new_logs: record.logs,
|
||||
mem_peak: record.mem_peak,
|
||||
progress: record.progress,
|
||||
workflow_as_code_status: record
|
||||
.workflow_as_code_status
|
||||
.map(|x: sqlx::types::Json<Box<RawValue>>| x.0),
|
||||
job,
|
||||
flow_status: record
|
||||
.flow_status
|
||||
|
||||
@@ -97,6 +97,8 @@ pub struct QueuedJob {
|
||||
pub permissioned_as: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub flow_status: Option<Json<Box<RawValue>>>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub workflow_as_code_status: Option<Json<Box<RawValue>>>,
|
||||
pub is_flow_step: bool,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub language: Option<ScriptLang>,
|
||||
@@ -179,6 +181,7 @@ impl Default for QueuedJob {
|
||||
job_kind: JobKind::Identity,
|
||||
schedule_path: None,
|
||||
permissioned_as: "".to_string(),
|
||||
workflow_as_code_status: None,
|
||||
flow_status: None,
|
||||
is_flow_step: false,
|
||||
language: None,
|
||||
@@ -236,6 +239,8 @@ pub struct CompletedJob {
|
||||
pub permissioned_as: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub flow_status: Option<sqlx::types::Json<Box<RawValue>>>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub workflow_as_code_status: Option<sqlx::types::Json<Box<RawValue>>>,
|
||||
pub is_flow_step: bool,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub language: Option<ScriptLang>,
|
||||
|
||||
@@ -3145,7 +3145,7 @@ async fn get_queued_job_tx<'c>(
|
||||
tx: &mut Transaction<'c, Postgres>,
|
||||
) -> error::Result<Option<QueuedJob>> {
|
||||
sqlx::query_as::<_, QueuedJob>(
|
||||
"SELECT *
|
||||
"SELECT *, null as workflow_as_code_status
|
||||
FROM v2_as_queue WHERE id = $1 AND workspace_id = $2",
|
||||
)
|
||||
.bind(id)
|
||||
@@ -3157,7 +3157,7 @@ async fn get_queued_job_tx<'c>(
|
||||
|
||||
pub async fn get_queued_job(id: &Uuid, w_id: &str, db: &DB) -> error::Result<Option<QueuedJob>> {
|
||||
sqlx::query_as::<_, QueuedJob>(
|
||||
"SELECT *
|
||||
"SELECT *, null as workflow_as_code_status
|
||||
FROM v2_as_queue WHERE id = $1 AND workspace_id = $2",
|
||||
)
|
||||
.bind(id)
|
||||
|
||||
@@ -559,7 +559,7 @@
|
||||
wideResults
|
||||
{flowStateStore}
|
||||
{jobId}
|
||||
on:done={() => {
|
||||
on:done={(x) => {
|
||||
isRunning = false
|
||||
$executionCount = $executionCount + 1
|
||||
onJobDone?.()
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
import { emptyString } from '$lib/utils'
|
||||
import type { DurationStatus } from './graph'
|
||||
import type { Writable } from 'svelte/store'
|
||||
import Badge from './Badge.svelte'
|
||||
import Badge from './common/badge/Badge.svelte'
|
||||
|
||||
interface Props {
|
||||
job: Job
|
||||
@@ -84,9 +84,14 @@
|
||||
|
||||
<div>
|
||||
<span class="inline-flex gap-1">
|
||||
<div class="text-xs flex items-center px-2 py-2">
|
||||
<Badge>{mod.id}</Badge>
|
||||
</div>
|
||||
<Badge
|
||||
color="indigo"
|
||||
wrapperClass="max-w-full"
|
||||
baseClass="max-w-full truncate !px-1"
|
||||
title={mod.id}
|
||||
>
|
||||
<span class="max-w-full text-2xs truncate">{mod.id}</span></Badge
|
||||
>
|
||||
<span class="font-medium text-primary mt-0.5">
|
||||
{#if !emptyString(rawMod?.summary)}
|
||||
{rawMod?.summary ?? ''}
|
||||
|
||||
@@ -36,6 +36,7 @@
|
||||
import FlowPreviewResult from './FlowPreviewResult.svelte'
|
||||
import type { FlowGraphAssetContext } from './flows/types'
|
||||
import { createState } from '$lib/svelte5Utils.svelte'
|
||||
import JobLoader from './JobLoader.svelte'
|
||||
|
||||
const dispatch = createEventDispatcher()
|
||||
|
||||
@@ -406,7 +407,8 @@
|
||||
JobService.getJob({
|
||||
workspace: workspaceId ?? $workspaceStore ?? '',
|
||||
id: mod.job ?? '',
|
||||
noLogs: true
|
||||
noLogs: true,
|
||||
noCode: true
|
||||
})
|
||||
.then((job) => {
|
||||
const newState = {
|
||||
@@ -456,6 +458,12 @@
|
||||
true
|
||||
)
|
||||
}
|
||||
if (mod.flow_jobs_success) {
|
||||
setModuleState(mod.id ?? '', {
|
||||
flow_jobs_success: mod.flow_jobs_success
|
||||
})
|
||||
}
|
||||
// console.log('updateInnerModules', mod.id, mod)
|
||||
|
||||
/**
|
||||
* else if (mod.type === 'Failure' || mod.type === 'WaitingForPriorSteps') {
|
||||
@@ -473,23 +481,6 @@
|
||||
}
|
||||
}
|
||||
|
||||
async function getNewJob(jobId: string, initialJob: Job | undefined) {
|
||||
if (
|
||||
jobId == initialJob?.id &&
|
||||
initialJob?.id != undefined &&
|
||||
initialJob?.type === 'CompletedJob'
|
||||
) {
|
||||
return initialJob
|
||||
} else {
|
||||
let r = await JobService.getJob({
|
||||
workspace: workspaceId ?? $workspaceStore ?? '',
|
||||
id: jobId ?? '',
|
||||
noLogs: true
|
||||
})
|
||||
return r
|
||||
}
|
||||
}
|
||||
|
||||
let debounceJobId: string | undefined = undefined
|
||||
let lastRefreshed: Date | undefined = undefined
|
||||
function debounceLoadJobInProgress() {
|
||||
@@ -509,9 +500,21 @@
|
||||
}, pollingRate)
|
||||
}
|
||||
|
||||
let errorCount = 0
|
||||
let notAnonynmous = $state(false)
|
||||
let started = false
|
||||
let jobLoader: JobLoader | undefined = undefined
|
||||
|
||||
function setJob(newJob: Job, force: boolean) {
|
||||
if (!deepEqual(job, newJob) || isForloopSelected || force) {
|
||||
job = newJob
|
||||
job?.flow_status && updateStatus(job?.flow_status)
|
||||
dispatch('jobsLoaded', { job, force: false })
|
||||
notAnonynmous = false
|
||||
if (job?.type == 'CompletedJob' && !destroyed) {
|
||||
dispatch('done', job)
|
||||
}
|
||||
}
|
||||
}
|
||||
async function loadJobInProgress() {
|
||||
if (!started) {
|
||||
started = true
|
||||
@@ -519,30 +522,29 @@
|
||||
}
|
||||
if (jobId != '00000000-0000-0000-0000-000000000000') {
|
||||
try {
|
||||
const newJob = await getNewJob(jobId, initialJob)
|
||||
if (!deepEqual(job, newJob) || isForloopSelected) {
|
||||
job = newJob
|
||||
job?.flow_status && updateStatus(job?.flow_status)
|
||||
dispatch('jobsLoaded', { job, force: false })
|
||||
if (
|
||||
jobId == initialJob?.id &&
|
||||
initialJob?.id != undefined &&
|
||||
initialJob?.type === 'CompletedJob'
|
||||
) {
|
||||
setJob(initialJob, false)
|
||||
} else {
|
||||
jobLoader?.watchJob(jobId, {
|
||||
change(newJob) {
|
||||
setJob(newJob, true)
|
||||
}
|
||||
})
|
||||
}
|
||||
errorCount = 0
|
||||
notAnonynmous = false
|
||||
} catch (e) {
|
||||
if (
|
||||
e?.body?.includes('As a non logged in user, you can only see jobs ran by anonymous users')
|
||||
) {
|
||||
notAnonynmous = true
|
||||
} else {
|
||||
errorCount += 1
|
||||
console.error(e)
|
||||
}
|
||||
}
|
||||
}
|
||||
if (job?.type !== 'CompletedJob' && errorCount < 4 && !destroyed) {
|
||||
debounceLoadJobInProgress()
|
||||
} else {
|
||||
dispatch('done', job)
|
||||
}
|
||||
}
|
||||
|
||||
let destroyed = false
|
||||
@@ -930,6 +932,7 @@
|
||||
let selected = $derived(isListJob ? 'sequence' : 'graph')
|
||||
</script>
|
||||
|
||||
<JobLoader noCode bind:this={jobLoader} />
|
||||
{#if notAnonynmous}
|
||||
<Alert type="error" title="Required Auth">
|
||||
As a non logged in user, you can only see jobs ran by anonymous users like you
|
||||
@@ -1030,7 +1033,8 @@
|
||||
storedJob = await JobService.getJob({
|
||||
workspace: workspaceId ?? $workspaceStore ?? '',
|
||||
id: loopJobId,
|
||||
noLogs: true
|
||||
noLogs: true,
|
||||
noCode: true
|
||||
})
|
||||
storedListJobs[j] = storedJob
|
||||
}
|
||||
|
||||
@@ -4,7 +4,8 @@
|
||||
JobService,
|
||||
type FlowStatus,
|
||||
type Preview,
|
||||
type GetJobUpdatesResponse
|
||||
type GetJobUpdatesResponse,
|
||||
type WorkflowStatus
|
||||
} from '$lib/gen'
|
||||
import { workspaceStore } from '$lib/stores'
|
||||
import { onDestroy, tick, untrack } from 'svelte'
|
||||
@@ -18,6 +19,7 @@
|
||||
done?: (x: Job & { result?: any }) => void
|
||||
doneResult?: ({ id, result }: { id: string; result: any }) => void
|
||||
doneError?: ({ id, error }: { id?: string; error: Error }) => void
|
||||
change?: (x: Job) => void
|
||||
cancel?: ({ id }: { id: string }) => void
|
||||
started?: ({ id }: { id: string }) => void
|
||||
running?: ({ id }: { id: string }) => void
|
||||
@@ -27,6 +29,7 @@
|
||||
isLoading?: boolean
|
||||
job?: Job | undefined
|
||||
noCode?: boolean
|
||||
noLogs?: boolean
|
||||
workspaceOverride?: string | undefined
|
||||
notfound?: boolean
|
||||
allowConcurentRequests?: boolean
|
||||
@@ -52,6 +55,7 @@
|
||||
lazyLogs = false,
|
||||
onlyResult = false,
|
||||
scriptProgress = $bindable(undefined),
|
||||
noLogs = false,
|
||||
children
|
||||
}: Props = $props()
|
||||
|
||||
@@ -346,9 +350,23 @@
|
||||
if (previewJobUpdates.flow_status) {
|
||||
job.flow_status = previewJobUpdates.flow_status as FlowStatus
|
||||
}
|
||||
if (previewJobUpdates.workflow_as_code_status) {
|
||||
job.workflow_as_code_status = previewJobUpdates.workflow_as_code_status as WorkflowStatus
|
||||
}
|
||||
if (previewJobUpdates.mem_peak && job) {
|
||||
job.mem_peak = previewJobUpdates.mem_peak
|
||||
}
|
||||
|
||||
if (
|
||||
job &&
|
||||
(previewJobUpdates.running ||
|
||||
previewJobUpdates.progress ||
|
||||
previewJobUpdates.new_logs ||
|
||||
previewJobUpdates.flow_status ||
|
||||
previewJobUpdates.mem_peak)
|
||||
) {
|
||||
callbacks?.change?.(job)
|
||||
}
|
||||
}
|
||||
async function loadTestJob(id: string, callbacks?: Callbacks): Promise<boolean> {
|
||||
let isCompleted = false
|
||||
@@ -369,7 +387,13 @@
|
||||
})
|
||||
|
||||
if ((previewJobUpdates.running ?? false) || (previewJobUpdates.completed ?? false)) {
|
||||
job = await JobService.getJob({ workspace: workspace!, id, noCode, noLogs: onlyResult })
|
||||
job = await JobService.getJob({
|
||||
workspace: workspace!,
|
||||
id,
|
||||
noCode,
|
||||
noLogs: onlyResult || noLogs
|
||||
})
|
||||
callbacks?.change?.(job)
|
||||
}
|
||||
|
||||
updateJobFromProgress(previewJobUpdates, job, callbacks)
|
||||
@@ -377,7 +401,7 @@
|
||||
job = await JobService.getJob({
|
||||
workspace: workspace!,
|
||||
id,
|
||||
noLogs: lazyLogs || onlyResult,
|
||||
noLogs: lazyLogs || onlyResult || noLogs,
|
||||
noCode
|
||||
})
|
||||
}
|
||||
@@ -434,6 +458,7 @@
|
||||
} else {
|
||||
callbacks?.done?.(job)
|
||||
}
|
||||
callbacks?.change?.(job)
|
||||
|
||||
if (!allowConcurentRequests) {
|
||||
currentId = undefined
|
||||
@@ -467,7 +492,12 @@
|
||||
try {
|
||||
// First load the job to get initial state
|
||||
if (!job && !onlyResult) {
|
||||
job = await JobService.getJob({ workspace: workspace!, id, noLogs: lazyLogs, noCode })
|
||||
job = await JobService.getJob({
|
||||
workspace: workspace!,
|
||||
id,
|
||||
noLogs: lazyLogs || noLogs,
|
||||
noCode
|
||||
})
|
||||
}
|
||||
|
||||
// If job is already completed, don't start SSE
|
||||
|
||||
@@ -100,13 +100,15 @@
|
||||
placement: 'bottom',
|
||||
gutter: 0,
|
||||
offset: { mainAxis: 3, crossAxis: 69 * zoom },
|
||||
overflowPadding: historyOpen ? 250 : 8
|
||||
overflowPadding: historyOpen ? 250 : 8,
|
||||
flip: false
|
||||
})
|
||||
popover?.updatePositioning({
|
||||
placement: 'bottom',
|
||||
gutter: 0,
|
||||
offset: { mainAxis: 3, crossAxis: showInput ? -69 * zoom : 0 },
|
||||
overflowPadding: historyOpen ? 250 : 8
|
||||
overflowPadding: historyOpen ? 250 : 8,
|
||||
flip: false
|
||||
})
|
||||
}
|
||||
|
||||
@@ -160,7 +162,8 @@
|
||||
placement: 'bottom',
|
||||
gutter: 0,
|
||||
offset: { mainAxis: 3, crossAxis: 69 },
|
||||
overflowPadding: historyOpen ? 250 : 8
|
||||
overflowPadding: historyOpen ? 250 : 8,
|
||||
flip: false
|
||||
}}
|
||||
usePointerDownOutside
|
||||
closeOnOutsideClick={false}
|
||||
@@ -195,12 +198,14 @@
|
||||
{/snippet}
|
||||
</Popover>
|
||||
{/if}
|
||||
|
||||
<Popover
|
||||
floatingConfig={{
|
||||
placement: 'bottom',
|
||||
gutter: 0,
|
||||
offset: { mainAxis: 3, crossAxis: showInput ? -69 : 0 },
|
||||
overflowPadding: historyOpen ? 250 : 8
|
||||
overflowPadding: historyOpen ? 250 : 8,
|
||||
flip: false
|
||||
}}
|
||||
usePointerDownOutside
|
||||
closeOnOutsideClick={false}
|
||||
|
||||
@@ -380,7 +380,7 @@
|
||||
}
|
||||
let newGraph = graph
|
||||
newGraph.nodes.sort((a, b) => b.id.localeCompare(a.id))
|
||||
console.log('compute')
|
||||
// console.log('compute')
|
||||
;[nodes, edges] = computeAssetNodes(layoutNodes(newGraph.nodes), newGraph.edges)
|
||||
await tick()
|
||||
height = Math.max(...nodes.map((n) => n.position.y + NODE.height + 100), minHeight)
|
||||
|
||||
@@ -189,9 +189,9 @@
|
||||
{/if}
|
||||
|
||||
<div class=" w-full rounded-md min-h-full">
|
||||
{#if job?.is_flow_step == false && job?.flow_status && (isScriptPreview(job?.job_kind) || job?.job_kind == 'script') && !(typeof job.flow_status == 'object' && '_metadata' in job.flow_status)}
|
||||
{#if job?.workflow_as_code_status}
|
||||
<WorkflowTimeline
|
||||
flow_status={asWorkflowStatus(job.flow_status)}
|
||||
flow_status={asWorkflowStatus(job.workflow_as_code_status)}
|
||||
flowDone={job.type == 'CompletedJob'}
|
||||
/>
|
||||
{/if}
|
||||
|
||||
@@ -129,10 +129,10 @@
|
||||
{#if selectedTab === 'logs'}
|
||||
<SplitPanesWrapper>
|
||||
<Splitpanes horizontal>
|
||||
{#if previewJob?.is_flow_step == false && previewJob?.flow_status && !(typeof previewJob.flow_status == 'object' && '_metadata' in previewJob.flow_status)}
|
||||
{#if previewJob?.workflow_as_code_status}
|
||||
<Pane class="relative">
|
||||
<WorkflowTimeline
|
||||
flow_status={asWorkflowStatus(previewJob.flow_status)}
|
||||
flow_status={asWorkflowStatus(previewJob.workflow_as_code_status)}
|
||||
flowDone={previewJob.type == 'CompletedJob'}
|
||||
/>
|
||||
</Pane>
|
||||
|
||||
@@ -921,10 +921,10 @@
|
||||
<ExecutionDuration bind:job bind:longRunning={currentJobIsLongRunning} />
|
||||
{/if}
|
||||
<div class="max-w-7xl mx-auto w-full px-4 mb-10">
|
||||
{#if job?.flow_status && typeof job.flow_status == 'object' && !('_metadata' in job.flow_status)}
|
||||
{#if job?.workflow_as_code_status}
|
||||
<div class="mt-10"></div>
|
||||
<WorkflowTimeline
|
||||
flow_status={asWorkflowStatus(job.flow_status)}
|
||||
flow_status={asWorkflowStatus(job.workflow_as_code_status)}
|
||||
flowDone={job.type == 'CompletedJob'}
|
||||
/>
|
||||
{/if}
|
||||
|
||||
Reference in New Issue
Block a user