use std::any::Any;
use std::panic::AssertUnwindSafe;
use std::sync::Arc;
use futures_util::FutureExt;
use runledger_core::jobs::{JobCompletion, JobContext, JobFailure};
use runledger_postgres::QueryErrorKind;
use runledger_postgres::jobs::{self, JobLeaseIdentity, JobProgressUpdate};
use tokio::time::{Duration, Instant, MissedTickBehavior, sleep_until};
use tracing::{Instrument, info, info_span, warn};
use super::completion::{
CompletionContext, CompletionObservation, complete_job_after_handler,
complete_job_failure_after_handler,
};
use super::observers::{JobRunningNotification, TerminalJobObserverEvent, TerminalObserverTasks};
use crate::WorkerError;
use crate::observer::{JobLeaseLostEvent, JobLifecycleObservers, ObservedJob};
use crate::registry::JobRegistry;
const UNKNOWN_WORKER_ID: &str = "unknown-worker";
const LEASE_OWNER_MISMATCH_CODE: &str = "job.lease_owner_mismatch";
const LEASE_MAINTENANCE_FAILED_CODE: &str = "job.lease_maintenance_failed";
const HANDLER_PANIC_CODE: &str = "job.handler_panic";
const RUNNING_PROGRESS_PERSIST_FAILED_REASON: &str = "RUNNING_PROGRESS_PERSIST_FAILED";
const UNSTARTED_CLAIM_RETRY_DELAY_MS: i32 = 1_000;
enum JobExecutionFailure {
Handler(JobFailure),
LeaseMaintenance(JobFailure),
}
pub(super) struct ClaimedJobExecution {
pool: runledger_postgres::DbPool,
registry: Arc<JobRegistry>,
job: jobs::JobQueueRecord,
lease_ttl_seconds: i32,
observers: JobLifecycleObservers,
terminal_observer_tasks: TerminalObserverTasks,
worker_id: String,
}
impl ClaimedJobExecution {
pub(super) fn new(
pool: runledger_postgres::DbPool,
registry: Arc<JobRegistry>,
job: jobs::JobQueueRecord,
lease_ttl_seconds: i32,
observers: JobLifecycleObservers,
terminal_observer_tasks: TerminalObserverTasks,
) -> Self {
let worker_id = job
.worker_id
.clone()
.unwrap_or_else(|| UNKNOWN_WORKER_ID.to_owned());
Self {
pool,
registry,
job,
lease_ttl_seconds,
observers,
terminal_observer_tasks,
worker_id,
}
}
pub(super) async fn execute(self) {
let job_span = info_span!(
"job",
sentry.name = %self.job.job_type,
sentry.op = "runledger.job",
job_id = %self.job.id,
job_type = %self.job.job_type,
run_number = self.job.run_number,
attempt = self.job.attempt,
organization_id = ?self.job.organization_id,
worker_id = %self.worker_id,
);
async move {
let start = Instant::now();
let context = self.context();
let observed_job = self.observed_job();
if !self.mark_job_running_or_abort(&context).await {
return;
}
let mut running_notification =
JobRunningNotification::spawn(self.observers.clone(), observed_job.clone());
match self.execute_job_handler_with_heartbeats(&context).await {
Ok(completion) => {
complete_job_after_handler(
self.completion_context(&context),
completion,
CompletionObservation::new(
&self.observers,
observed_job.clone(),
start.elapsed(),
&mut running_notification,
&self.terminal_observer_tasks,
),
)
.await;
}
Err(JobExecutionFailure::Handler(failure)) => {
complete_job_failure_after_handler(
self.completion_context(&context),
failure,
CompletionObservation::new(
&self.observers,
observed_job.clone(),
start.elapsed(),
&mut running_notification,
&self.terminal_observer_tasks,
),
)
.await;
}
Err(JobExecutionFailure::LeaseMaintenance(failure)) => {
self.log_lease_maintenance_abort(&failure);
running_notification
.spawn_terminal_observer(
&self.terminal_observer_tasks,
&self.job,
self.observers.clone(),
TerminalJobObserverEvent::LeaseLost(JobLeaseLostEvent {
job: observed_job.clone(),
duration: start.elapsed(),
failure,
}),
)
.await;
}
}
info!(
job_id = %self.job.id,
attempt = self.job.attempt,
run_number = self.job.run_number,
elapsed_ms = start.elapsed().as_millis(),
"job processed"
);
}
.instrument(job_span)
.await;
}
fn context(&self) -> JobContext {
JobContext {
job_id: self.job.id,
run_number: self.job.run_number,
attempt: self.job.attempt,
organization_id: self.job.organization_id,
worker_id: self.worker_id.clone(),
checkpoint: self.job.checkpoint.clone(),
}
}
fn observed_job(&self) -> ObservedJob {
ObservedJob {
job_id: self.job.id,
job_type: self.job.job_type.clone(),
organization_id: self.job.organization_id,
run_number: self.job.run_number,
attempt: self.job.attempt,
max_attempts: self.job.max_attempts,
worker_id: self.worker_id.to_owned(),
}
}
fn lease_identity(&self) -> JobLeaseIdentity<'_> {
JobLeaseIdentity::new(
self.job.id,
self.job.run_number,
self.job.attempt,
&self.worker_id,
)
}
fn completion_context<'execution, 'context>(
&'execution self,
context: &'context JobContext,
) -> CompletionContext<'execution, 'context> {
CompletionContext::new(
&self.pool,
self.registry.as_ref(),
context,
&self.job,
self.lease_identity(),
)
}
fn log_lease_maintenance_abort(&self, failure: &JobFailure) {
warn!(
job_id = %self.job.id,
attempt = self.job.attempt,
failure_code = failure.code,
"job processing aborted because durable lease maintenance was lost"
);
}
async fn mark_job_running_or_abort(&self, context: &JobContext) -> bool {
let running_progress = JobProgressUpdate {
stage: Some(runledger_core::jobs::JobStage::Running),
progress_done: None,
progress_total: None,
checkpoint: None,
};
let Err(source) = jobs::update_job_progress_for_lease(
&self.pool,
self.lease_identity(),
&running_progress,
)
.await
else {
return true;
};
self.handle_running_progress_persist_failure(context, source)
.await;
false
}
async fn handle_running_progress_persist_failure(
&self,
context: &JobContext,
source: runledger_postgres::Error,
) {
let lease_owner_mismatch = is_lease_owner_mismatch_error(&source);
let error = WorkerError::SetRunningProgress {
job_id: self.job.id,
attempt: self.job.attempt,
source,
};
if lease_owner_mismatch {
warn!(
%error,
job_id = %self.job.id,
attempt = self.job.attempt,
"aborting job before execution because lease ownership was already lost"
);
return;
}
match jobs::release_unstarted_job_claim(
&self.pool,
self.job.id,
self.job.run_number,
self.job.attempt,
&context.worker_id,
RUNNING_PROGRESS_PERSIST_FAILED_REASON,
UNSTARTED_CLAIM_RETRY_DELAY_MS,
)
.await
{
Ok(()) => {
warn!(
%error,
job_id = %self.job.id,
attempt = self.job.attempt,
"running progress could not be persisted; released unstarted claim back to pending"
);
}
Err(release_error) => {
let no_longer_releasable =
is_unstarted_claim_release_not_applicable_error(&release_error);
let release_error = WorkerError::ReleaseUnstartedClaim {
job_id: self.job.id,
attempt: self.job.attempt,
source: release_error,
};
if no_longer_releasable {
warn!(
%error,
%release_error,
job_id = %self.job.id,
attempt = self.job.attempt,
"running progress could not be persisted; unstarted release no longer applies and the job will continue under the current lease owner"
);
return;
}
warn!(
%error,
%release_error,
job_id = %self.job.id,
attempt = self.job.attempt,
"running progress could not be persisted; leaving claim for reaper recovery"
);
}
}
}
async fn execute_job_handler_with_heartbeats(
&self,
context: &JobContext,
) -> Result<JobCompletion, JobExecutionFailure> {
let registry = Arc::clone(&self.registry);
let mut execution = Box::pin(
AssertUnwindSafe(execute_job_handler(registry, context, &self.job)).catch_unwind(),
);
let timeout_deadline =
Instant::now() + Duration::from_secs(self.job.timeout_seconds.max(1) as u64);
let mut timeout = Box::pin(sleep_until(timeout_deadline));
let mut ticker = tokio::time::interval(heartbeat_interval(self.lease_ttl_seconds));
ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
ticker.tick().await;
loop {
tokio::select! {
result = &mut execution => {
return match result {
Ok(result) => result.map_err(JobExecutionFailure::Handler),
Err(panic_payload) => {
Err(JobExecutionFailure::Handler(handler_panic_failure(panic_payload)))
}
};
}
_ = &mut timeout => {
return Err(JobExecutionFailure::Handler(JobFailure::timeout(
"job.timeout_exceeded",
"Job exceeded the configured timeout.",
)));
}
_ = ticker.tick() => {
if let Err(error) = jobs::heartbeat_job_for_lease(
&self.pool,
self.lease_identity(),
self.lease_ttl_seconds,
)
.await
{
let lease_owner_mismatch = is_lease_owner_mismatch_error(&error);
let error = WorkerError::Heartbeat {
job_id: self.job.id,
attempt: self.job.attempt,
source: error,
};
if lease_owner_mismatch {
warn!(%error, job_id = %self.job.id, "job heartbeat lost lease ownership");
return Err(JobExecutionFailure::LeaseMaintenance(
lease_owner_mismatch_failure(),
));
}
warn!(
%error,
job_id = %self.job.id,
"aborting job because lease heartbeat could not be persisted"
);
return Err(JobExecutionFailure::LeaseMaintenance(
lease_maintenance_failure(),
));
}
}
}
}
}
}
async fn execute_job_handler(
registry: Arc<JobRegistry>,
context: &JobContext,
job: &jobs::JobQueueRecord,
) -> Result<JobCompletion, JobFailure> {
let Some(handler) = registry.get(job.job_type.as_borrowed()) else {
return Err(JobFailure::terminal(
"job.handler_not_registered",
"No handler is registered for this job type.",
));
};
handler.execute(context.clone(), job.payload.clone()).await
}
pub(super) fn lease_owner_mismatch_failure() -> JobFailure {
JobFailure::lease_expired(
LEASE_OWNER_MISMATCH_CODE,
"Job lease ownership was lost during processing.",
)
}
fn lease_maintenance_failure() -> JobFailure {
JobFailure::lease_expired(
LEASE_MAINTENANCE_FAILED_CODE,
"Job lease could not be durably maintained during processing.",
)
}
fn handler_panic_failure(panic_payload: Box<dyn Any + Send>) -> JobFailure {
JobFailure::panicked(
HANDLER_PANIC_CODE,
format!(
"Job handler panicked: {}",
panic_payload_message(&*panic_payload)
),
)
}
fn panic_payload_message(panic_payload: &(dyn Any + Send)) -> String {
if let Some(message) = panic_payload.downcast_ref::<String>() {
return message.clone();
}
if let Some(message) = panic_payload.downcast_ref::<&'static str>() {
return (*message).to_string();
}
"non-string panic payload".to_string()
}
fn has_query_error_kind(error: &runledger_postgres::Error, expected_kind: QueryErrorKind) -> bool {
matches!(
error,
runledger_postgres::Error::QueryError(query_error)
if query_error.kind() == Some(expected_kind)
)
}
pub(super) fn is_lease_owner_mismatch_error(error: &runledger_postgres::Error) -> bool {
has_query_error_kind(error, QueryErrorKind::JobLeaseOwnerMismatch)
}
fn is_unstarted_claim_release_not_applicable_error(error: &runledger_postgres::Error) -> bool {
has_query_error_kind(error, QueryErrorKind::JobUnstartedClaimReleaseNotApplicable)
}
fn heartbeat_interval(lease_ttl_seconds: i32) -> Duration {
let seconds = (lease_ttl_seconds.max(1) / 3).max(1) as u64;
Duration::from_secs(seconds)
}