use super::{HeartbeatResult, RunnerClaimResult, StaleClaim, TaskExecutionDAL};
use crate::dal::unified::models::{NewUnifiedExecutionEvent, UnifiedTaskExecution};
use crate::database::schema::unified::{execution_events, task_executions};
use crate::database::universal_types::{UniversalTimestamp, UniversalUuid};
use crate::error::ValidationError;
use crate::models::execution_event::ExecutionEventType;
use crate::models::task_execution::TaskExecution;
use diesel::prelude::*;
impl<'a> TaskExecutionDAL<'a> {
pub async fn schedule_retry(
&self,
task_id: UniversalUuid,
retry_at: UniversalTimestamp,
new_attempt: i32,
) -> Result<(), ValidationError> {
use diesel::connection::Connection;
crate::interact_on_backend!(self.dal, |conn| {
conn.transaction::<_, diesel::result::Error, _>(|conn| {
let now = UniversalTimestamp::now();
let task: UnifiedTaskExecution =
task_executions::table.find(task_id).first(conn)?;
diesel::update(task_executions::table.find(task_id))
.set((
task_executions::status.eq("Ready"),
task_executions::attempt.eq(new_attempt),
task_executions::retry_at.eq(Some(retry_at)),
task_executions::started_at.eq(None::<UniversalTimestamp>),
task_executions::completed_at.eq(None::<UniversalTimestamp>),
task_executions::updated_at.eq(now),
))
.execute(conn)?;
let event_data = serde_json::json!({
"attempt": new_attempt,
"retry_at": retry_at.to_string()
})
.to_string();
let event = NewUnifiedExecutionEvent {
id: UniversalUuid::new_v4(),
workflow_execution_id: task.workflow_execution_id,
task_execution_id: Some(task_id),
event_type: ExecutionEventType::TaskRetryScheduled.as_str().to_string(),
event_data: Some(event_data),
worker_id: None,
created_at: now,
request_id: None,
runner_id: None,
tenant_id: None,
};
diesel::insert_into(execution_events::table)
.values(&event)
.execute(conn)?;
Ok(())
})
})?;
Ok(())
}
pub async fn claim_for_runner(
&self,
task_id: UniversalUuid,
runner_id: UniversalUuid,
) -> Result<RunnerClaimResult, ValidationError> {
let now = UniversalTimestamp::now();
let rows_updated: usize = crate::interact_on_backend!(self.dal, |conn| {
diesel::update(
task_executions::table
.find(task_id)
.filter(task_executions::claimed_by.is_null()),
)
.set((
task_executions::claimed_by.eq(Some(runner_id)),
task_executions::heartbeat_at.eq(Some(now)),
task_executions::updated_at.eq(now),
))
.execute(conn)
})?;
Ok(if rows_updated > 0 {
RunnerClaimResult::Claimed
} else {
RunnerClaimResult::AlreadyClaimed
})
}
pub async fn heartbeat(
&self,
task_id: UniversalUuid,
runner_id: UniversalUuid,
) -> Result<HeartbeatResult, ValidationError> {
let now = UniversalTimestamp::now();
let rows_updated: usize = crate::interact_on_backend!(self.dal, |conn| {
diesel::update(
task_executions::table
.find(task_id)
.filter(task_executions::claimed_by.eq(Some(runner_id))),
)
.set((
task_executions::heartbeat_at.eq(Some(now)),
task_executions::updated_at.eq(now),
))
.execute(conn)
})?;
Ok(if rows_updated > 0 {
HeartbeatResult::Ok
} else {
HeartbeatResult::ClaimLost
})
}
pub async fn release_runner_claim(
&self,
task_id: UniversalUuid,
) -> Result<(), ValidationError> {
let now = UniversalTimestamp::now();
crate::interact_on_backend!(self.dal, |conn| {
diesel::update(task_executions::table.find(task_id))
.set((
task_executions::claimed_by.eq(None::<UniversalUuid>),
task_executions::heartbeat_at.eq(None::<UniversalTimestamp>),
task_executions::updated_at.eq(now),
))
.execute(conn)
})?;
Ok(())
}
pub async fn find_stale_claims(
&self,
threshold: std::time::Duration,
) -> Result<Vec<StaleClaim>, ValidationError> {
let cutoff = UniversalTimestamp(
chrono::Utc::now()
- chrono::Duration::from_std(threshold).unwrap_or(chrono::Duration::seconds(60)),
);
let stale: Vec<UnifiedTaskExecution> = crate::interact_on_backend!(self.dal, |conn| {
task_executions::table
.filter(task_executions::claimed_by.is_not_null())
.filter(task_executions::heartbeat_at.lt(Some(cutoff)))
.filter(task_executions::status.eq_any(vec!["Running", "Ready"]))
.load(conn)
})?;
Ok(stale
.into_iter()
.filter_map(|t| {
Some(StaleClaim {
task_id: t.id,
claimed_by: t.claimed_by?,
heartbeat_at: t.heartbeat_at?.0,
recovery_attempts: t.recovery_attempts,
})
})
.collect())
}
pub async fn get_ready_for_retry(&self) -> Result<Vec<TaskExecution>, ValidationError> {
let now = UniversalTimestamp::now();
let ready_tasks: Vec<UnifiedTaskExecution> =
crate::interact_on_backend!(self.dal, |conn| {
task_executions::table
.filter(task_executions::status.eq("Ready"))
.filter(
task_executions::retry_at
.is_null()
.or(task_executions::retry_at.le(now)),
)
.filter(task_executions::claimed_by.is_null())
.load(conn)
})?;
Ok(ready_tasks.into_iter().map(Into::into).collect())
}
}