refactor: remove legacy database views v2_as_queue and v2_as_completed_job (#6689)

* refactor: remove legacy database views v2_as_queue and v2_as_completed_job

Signed-off-by: Ramtin Mesgari <26694963+iamramtin@users.noreply.github.com>

* fix tests

* fix jobs.rs

* end

* fix

* improvement

* improvement

---------

Signed-off-by: Ramtin Mesgari <26694963+iamramtin@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
This commit is contained in:
Ramtin Mesgari
2025-11-04 13:08:02 +00:00
committed by GitHub
co-authored by Ruben Fiszel
parent 59a79e94ed
commit cedfd183d9
35 changed files with 933 additions and 471 deletions
+144 -61
View File
@@ -6,6 +6,7 @@
* LICENSE-AGPL for a copy of the license.
*/
use std::future::Future;
use std::{collections::HashMap, sync::Arc, vec};
use anyhow::Context;
@@ -154,18 +155,19 @@ pub struct JobCompleted {
pub async fn cancel_single_job<'c>(
username: &str,
reason: Option<String>,
job_running: Arc<QueuedJob>,
job_running: QueuedJobV2,
w_id: &str,
mut tx: Transaction<'c, Postgres>,
db: &Pool<Postgres>,
force_cancel: bool,
) -> error::Result<(Transaction<'c, Postgres>, Option<Uuid>)> {
let id = job_running.id;
if force_cancel || (job_running.parent_job.is_none() && !job_running.running) {
let username = username.to_string();
let w_id = w_id.to_string();
let db = db.clone();
tracing::info!("cancelling job {:?}", job_running.id);
let job_running = job_running.clone();
tokio::task::spawn(async move {
let reason: String = reason
.clone()
@@ -179,10 +181,11 @@ pub async fn cancel_single_job<'c>(
&Connection::from(db.clone()),
)
.await;
let memory_peak = job_running.memory_peak.unwrap_or(0);
let add_job = add_completed_job_error(
&db,
&MiniCompletedJob::from(MiniPulledJob::from(&job_running)),
job_running.mem_peak.unwrap_or(0),
&MiniCompletedJob::from(job_running),
memory_peak,
Some(CanceledBy { username: Some(username.to_string()), reason: Some(reason) }),
e,
"server",
@@ -210,7 +213,7 @@ pub async fn cancel_single_job<'c>(
}
}
Ok((tx, Some(job_running.id)))
Ok((tx, Some(id)))
}
pub async fn cancel_job<'c>(
@@ -224,26 +227,35 @@ pub async fn cancel_job<'c>(
require_anonymous: bool,
) -> error::Result<(Transaction<'c, Postgres>, Option<Uuid>)> {
//TODO fetch mini completed job instead of QueuedJob
let job = get_queued_job_tx(id, &w_id, &mut tx).await?;
let job = get_queued_job_v2(&mut *tx, &id).await?;
if job.is_none() {
return Ok((tx, None));
}
if require_anonymous && job.as_ref().unwrap().created_by != "anonymous" {
let mut job = job.unwrap();
if require_anonymous && job.created_by != "anonymous" {
return Err(Error::BadRequest(
"You are not logged in and this job was not created by an anonymous user like you so you cannot cancel it".to_string(),
));
}
let mut job = job.unwrap();
if job.workspace_id != w_id {
return Err(Error::BadRequest(
"You are not authorized to cancel this job belonging to another workspace".to_string(),
));
}
if force_cancel {
// if force canceling a flow step, make sure we force cancel from the highest parent
loop {
if job.parent_job.is_none() {
break;
}
match get_queued_job_tx(job.parent_job.unwrap(), &w_id, &mut tx).await? {
match get_queued_job_v2(&mut *tx, &job.parent_job.unwrap()).await? {
Some(j) => {
job = j;
}
@@ -253,7 +265,7 @@ pub async fn cancel_job<'c>(
}
// prevent cancelling a future tick of a schedule
if let Some(schedule_path) = job.schedule_path.as_ref() {
if let Some(schedule_path) = job.schedule_path().as_ref() {
let now = now_from_db(&mut *tx).await?;
if job.scheduled_for > now {
return Err(Error::BadRequest(
@@ -266,7 +278,7 @@ pub async fn cancel_job<'c>(
}
}
let job = Arc::new(job);
let job = job;
// get all children using recursive CTE
let mut jobs_to_cancel = sqlx::query!(
@@ -308,7 +320,7 @@ ORDER BY depth, id
let (ntx, _) = cancel_single_job(
username,
reason.clone(),
job.clone(),
job,
w_id,
tx,
db,
@@ -335,13 +347,13 @@ ORDER BY depth, id
}
}
for job_id in jobs_to_cancel {
let job = get_queued_job_tx(job_id, &w_id, &mut tx).await?;
let job = get_queued_job_v2(&mut *tx, &job_id).await?;
if let Some(job) = job {
let (ntx, _) = cancel_single_job(
username,
reason.clone(),
Arc::new(job),
job,
w_id,
tx,
db,
@@ -554,7 +566,7 @@ async fn cancel_persistent_script_jobs_internal<'c>(
// we could have retrieved the job IDs in the first query where we retrieve the hashes, but just in case a job was inserted in the queue right in-between the two above query, we re-do the fetch here
let jobs_to_cancel = sqlx::query_scalar::<_, Uuid>(
"SELECT id FROM v2_as_queue WHERE workspace_id = $1 AND script_path = $2 AND canceled = false",
"SELECT j.id FROM v2_job_queue q JOIN v2_job j USING (id) WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND q.canceled_by IS NULL",
)
.bind(w_id)
.bind(script_path)
@@ -1948,6 +1960,34 @@ pub struct MiniCompletedJob {
pub cache_ttl: Option<i32>,
}
impl From<QueuedJobV2> for MiniCompletedJob {
fn from(job: QueuedJobV2) -> Self {
MiniCompletedJob {
id: job.id,
workspace_id: job.workspace_id,
runnable_id: job.runnable_id,
scheduled_for: job.scheduled_for,
parent_job: job.parent_job,
flow_innermost_root_job: job.flow_innermost_root_job,
runnable_path: job.runnable_path,
kind: job.kind,
started_at: job.started_at,
permissioned_as: job.permissioned_as,
created_by: job.created_by,
script_lang: job.script_lang,
permissioned_as_email: job.permissioned_as_email,
flow_step_id: job.flow_step_id,
trigger_kind: job.trigger_kind,
trigger: job.trigger,
priority: job.priority,
concurrent_limit: job.concurrent_limit,
tag: job.tag,
cache_ttl: job.cache_ttl,
}
}
}
impl From<MiniPulledJob> for MiniCompletedJob {
fn from(job: MiniPulledJob) -> Self {
MiniCompletedJob {
@@ -2008,11 +2048,7 @@ impl MiniCompletedJob {
self.flow_step_id.is_some()
}
pub fn schedule_path(&self) -> Option<String> {
if self.trigger_kind.as_ref().is_some_and(|t| matches!(t, JobTriggerKind::Schedule)) {
self.trigger.clone()
} else {
None
}
schedule_path(&self.trigger_kind, &self.trigger)
}
pub fn is_flow(&self) -> bool {
@@ -2025,6 +2061,14 @@ impl MiniCompletedJob {
}
fn schedule_path(trigger_kind: &Option<JobTriggerKind>, trigger: &Option<String>) -> Option<String> {
if trigger_kind.as_ref().is_some_and(|t| matches!(t, JobTriggerKind::Schedule)) {
trigger.clone()
} else {
None
}
}
#[derive(Serialize, Deserialize, Debug, Clone)]
struct FlowStatusChatInputEnabled {
chat_input_enabled: Option<bool>,
@@ -2111,15 +2155,7 @@ impl MiniPulledJob {
}
pub fn schedule_path(&self) -> Option<String> {
if self
.trigger_kind
.as_ref()
.is_some_and(|t| matches!(t, JobTriggerKind::Schedule))
{
self.trigger.clone()
} else {
None
}
schedule_path(&self.trigger_kind, &self.trigger)
}
pub async fn mark_as_started_if_step(&self, db: &DB) -> Result<(), Error> {
@@ -2316,6 +2352,57 @@ pub async fn get_mini_pulled_job<'c>(
Ok(job)
}
pub struct QueuedJobV2 {
pub id: Uuid,
pub workspace_id: String,
pub runnable_id: Option<ScriptHash>,
pub scheduled_for: chrono::DateTime<chrono::Utc>,
pub parent_job: Option<Uuid>,
// pub root_job: Option<Uuid>,
pub flow_innermost_root_job: Option<Uuid>,
pub runnable_path: Option<String>,
pub kind: JobKind,
pub started_at: Option<chrono::DateTime<chrono::Utc>>,
pub permissioned_as: String,
pub created_by: String,
pub script_lang: Option<ScriptLang>,
pub permissioned_as_email: String,
pub flow_step_id: Option<String>,
pub trigger_kind: Option<JobTriggerKind>,
pub trigger: Option<String>,
pub priority: Option<i16>,
pub concurrent_limit: Option<i32>,
pub tag: String,
pub cache_ttl: Option<i32>,
pub last_ping: Option<chrono::DateTime<chrono::Utc>>,
pub worker: Option<String>,
pub memory_peak: Option<i32>,
pub running: bool,
}
impl QueuedJobV2 {
pub fn schedule_path(&self) -> Option<String> {
schedule_path(&self.trigger_kind, &self.trigger)
}
}
pub async fn get_queued_job_v2<'c>(
e: impl PgExecutor<'c>, job_id: &Uuid) -> error::Result<Option<QueuedJobV2>> {
let job = sqlx::query_as!(
QueuedJobV2,
"SELECT id, q.workspace_id, j.runnable_id as \"runnable_id: ScriptHash\", scheduled_for, parent_job, flow_innermost_root_job, runnable_path, kind as \"kind: JobKind\", started_at, permissioned_as, created_by, script_lang as \"script_lang: ScriptLang\",
permissioned_as_email, flow_step_id, trigger_kind as \"trigger_kind: JobTriggerKind\", trigger, q.priority, concurrent_limit, q.tag, cache_ttl, r.ping as last_ping, worker, memory_peak, running
FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id)
WHERE j.id = $1",
job_id,
)
.fetch_optional(e)
.await?;
Ok(job)
}
#[derive(Serialize, Deserialize, Debug)]
pub struct PulledJobResult {
pub job: Option<PulledJob>,
@@ -3170,33 +3257,29 @@ pub async fn job_is_complete(db: &DB, id: Uuid, w_id: &str) -> error::Result<boo
.unwrap_or(false))
}
async fn get_queued_job_tx<'c>(
id: Uuid,
w_id: &str,
tx: &mut Transaction<'c, Postgres>,
) -> error::Result<Option<QueuedJob>> {
sqlx::query_as::<_, QueuedJob>(
"SELECT *, null as workflow_as_code_status
FROM v2_as_queue WHERE id = $1 AND workspace_id = $2",
)
.bind(id)
.bind(w_id)
.fetch_optional(&mut **tx)
.await
.map_err(Into::into)
pub fn get_mini_completed_job<
'a,
'e,
A: sqlx::Acquire<'e, Database = Postgres> + Send + 'a,
>(id: &'a Uuid, w_id: &'a str, db: A) -> impl Future<Output = error::Result<Option<MiniCompletedJob>>> + Send + 'a {
async move {
let mut conn = db.acquire().await?;
sqlx::query_as!(
MiniCompletedJob,
"SELECT
j.id, j.workspace_id, j.runnable_id AS \"runnable_id!: ScriptHash\", q.scheduled_for, q.started_at, j.parent_job, j.flow_innermost_root_job, j.runnable_path, j.kind as \"kind!: JobKind\", j.permissioned_as,
j.created_by, j.script_lang AS \"script_lang!: ScriptLang\", j.permissioned_as_email, j.flow_step_id, j.trigger_kind AS \"trigger_kind!: JobTriggerKind\", j.trigger, j.priority, j.concurrent_limit, j.tag, j.cache_ttl
FROM v2_job j LEFT JOIN v2_job_queue q ON j.id = q.id
WHERE j.id = $1 AND j.workspace_id = $2",
id,
w_id
)
.fetch_optional(&mut *conn)
.await
.map_err(Into::into)
}
}
pub async fn get_queued_job(id: &Uuid, w_id: &str, db: &DB) -> error::Result<Option<QueuedJob>> {
sqlx::query_as::<_, QueuedJob>(
"SELECT *, null as workflow_as_code_status
FROM v2_as_queue WHERE id = $1 AND workspace_id = $2",
)
.bind(id)
.bind(w_id)
.fetch_optional(db)
.await
.map_err(Into::into)
}
pub enum PushIsolationLevel<'c> {
IsolatedRoot(DB),
@@ -3557,7 +3640,7 @@ pub async fn push<'c, 'd>(
}
let in_queue = sqlx::query_scalar!(
"SELECT COUNT(id) FROM v2_as_queue WHERE email = $1",
"SELECT COUNT(id) FROM v2_job WHERE permissioned_as_email = $1",
email
)
.fetch_one(_db)
@@ -3571,7 +3654,7 @@ pub async fn push<'c, 'd>(
}
let concurrent_runs = sqlx::query_scalar!(
"SELECT COUNT(id) FROM v2_as_queue WHERE running = true AND email = $1",
"SELECT COUNT(j.id) FROM v2_job_queue q JOIN v2_job j USING (id) WHERE q.running = true AND j.permissioned_as_email = $1",
email
)
.fetch_one(_db)
@@ -5155,11 +5238,11 @@ async fn restarted_flows_resolution(
> {
let row = sqlx::query!(
"SELECT
script_path, script_hash AS \"script_hash: ScriptHash\",
job_kind AS \"job_kind!: JobKind\",
flow_status AS \"flow_status: Json<Box<RawValue>>\",
raw_flow AS \"raw_flow: Json<Box<RawValue>>\"
FROM v2_as_completed_job WHERE id = $1 and workspace_id = $2",
j.runnable_path as script_path, j.runnable_id AS \"script_hash: ScriptHash\",
j.kind AS \"job_kind!: JobKind\",
COALESCE(c.flow_status, c.workflow_as_code_status) AS \"flow_status: Json<Box<RawValue>>\",
j.raw_flow AS \"raw_flow: Json<Box<RawValue>>\"
FROM v2_job_completed c JOIN v2_job j USING (id) WHERE j.id = $1 and j.workspace_id = $2",
completed_flow_id,
workspace_id,
)