use crate::{
backend::PartitionHeadClaim,
commands::enqueue_continuations,
core::BgJobHandler,
models::{FailureTransition, Job, RequeuedStage, Stage},
UtcDateTime,
};
use std::{sync::Arc, time::Duration};
#[cfg(feature = "prometheus")]
use crate::metrics::{HandlerOutcome, HandlerTimer};
pub(crate) struct PartitionAssignment {
topic: String,
partition: crate::topic::PartitionId,
lease_epoch: i64,
lease_lost: tokio::sync::watch::Receiver<()>,
renewal: tokio::task::JoinHandle<()>,
}
impl Drop for PartitionAssignment {
fn drop(&mut self) {
self.renewal.abort();
}
}
pub(crate) async fn poll_partition_head<C, H>(
handler: &Arc<H>,
owner: &str,
live_workers: &[String],
assignment: &mut Option<PartitionAssignment>,
#[cfg(feature = "prometheus")] worker_id: &str,
) -> anyhow::Result<bool>
where
C: Sync + Send,
H: BgJobHandler<C> + Sync + Send + 'static,
{
let publisher = handler.get_publisher();
if let Some(held) = assignment.as_ref() {
let still_owned = crate::topic::rendezvous_owner(&held.topic, held.partition, live_workers)
== Some(owner);
let lease_lost = held.lease_lost.has_changed().unwrap_or(true);
if !still_owned && !lease_lost {
let released = publisher
.release_partition_assignment(owner, &held.topic, held.partition, held.lease_epoch)
.await?;
if !released {
tracing::debug!(
topic = %held.topic,
partition = held.partition.0,
lease_epoch = held.lease_epoch,
"Partition assignment was already released before rebalancing"
);
}
*assignment = None;
} else if lease_lost {
*assignment = None;
}
}
let claim = match assignment.as_ref() {
Some(held) => {
publisher
.peek_assigned_partition_head(owner, &held.topic, held.partition, held.lease_epoch)
.await?
}
None => None,
};
if claim.is_none() {
if let Some(held) = assignment.as_ref() {
let released = publisher
.release_partition_assignment(owner, &held.topic, held.partition, held.lease_epoch)
.await?;
if !released {
tracing::debug!(
topic = %held.topic,
partition = held.partition.0,
lease_epoch = held.lease_epoch,
"Partition assignment was already released before it became idle"
);
}
*assignment = None;
}
}
let claim = match claim {
Some(claim) => claim,
None => {
let claimed = publisher.claim_partition_head(owner, live_workers).await?;
#[cfg(feature = "prometheus")]
crate::metrics::record_partition_claim(
publisher.metrics_queue(),
claimed.as_ref().map(|claim| claim.topic.as_str()),
);
let Some(claimed) = claimed else {
return Ok(false);
};
let (lease_lost_tx, lease_lost_rx) = tokio::sync::watch::channel(());
let renewal = spawn_lease_renewal(handler.clone(), claimed.clone(), lease_lost_tx);
*assignment = Some(PartitionAssignment {
topic: claimed.topic.clone(),
partition: claimed.partition,
lease_epoch: claimed.lease_epoch,
lease_lost: lease_lost_rx,
renewal,
});
claimed
}
};
let Some(job) = publisher.storage.get_job(claim.job_id.clone()).await? else {
tracing::warn!(
job_id = %claim.job_id,
topic = %claim.topic,
partition = claim.partition.0,
"Partition head job is missing; removing orphaned partition row"
);
publisher
.complete_partition_head(&claim, Vec::new())
.await?;
return Ok(true);
};
let runnable = match &job.stage {
Stage::Enqueued(_) | Stage::Running(_) => Some(job),
Stage::Delayed(delayed) if delayed.is_time() => Some(job.transition()),
Stage::Requeued(requeued) if requeued.is_ready() => Some(job.transition()),
Stage::Delayed(delayed) => {
release_not_ready(handler, &claim, delayed.not_before).await?;
None
}
Stage::Requeued(requeued) => {
release_not_ready(
handler,
&claim,
requeued.not_before.unwrap_or_else(chrono::Utc::now),
)
.await?;
None
}
Stage::Waiting(waiting) => {
let parent_id = waiting.parent_id.clone();
let parent_settled = match publisher.storage.get_job(parent_id).await? {
Some(parent) => parent.stage.is_success(),
None => true,
};
if parent_settled {
Some(job.transition())
} else {
let retry_shortly = chrono::Utc::now()
.checked_add_signed(chrono::Duration::seconds(1))
.unwrap_or_else(chrono::Utc::now);
release_not_ready(handler, &claim, retry_shortly).await?;
None
}
}
Stage::Success(_) | Stage::Failed(_) => {
tracing::warn!(
job_id = %claim.job_id,
"Partition head job is already terminal; removing stale partition row"
);
publisher
.complete_partition_head(&claim, Vec::new())
.await?;
None
}
};
let Some(job) = runnable else {
return Ok(true);
};
let lease_lost = &mut assignment
.as_mut()
.expect("assignment set above")
.lease_lost;
handle_partitioned_job(
handler,
&claim,
job,
lease_lost,
#[cfg(feature = "prometheus")]
worker_id,
)
.await?;
Ok(true)
}
async fn release_not_ready<C, H>(
handler: &Arc<H>,
claim: &PartitionHeadClaim,
earliest_run_at: UtcDateTime,
) -> anyhow::Result<()>
where
C: Sync + Send,
H: BgJobHandler<C> + Sync + Send + 'static,
{
let released = handler
.get_publisher()
.release_partition_head(claim, earliest_run_at, Vec::new())
.await?;
if !released {
tracing::debug!(
job_id = %claim.job_id,
topic = %claim.topic,
partition = claim.partition.0,
lease_epoch = claim.lease_epoch,
"Partition lease was lost before an unready head could be released"
);
}
Ok(())
}
fn spawn_lease_renewal<C, H>(
handler: Arc<H>,
claim: PartitionHeadClaim,
lease_lost: tokio::sync::watch::Sender<()>,
) -> tokio::task::JoinHandle<()>
where
C: Sync + Send,
H: BgJobHandler<C> + Sync + Send + 'static,
{
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(10)).await;
let renewed = handler.get_publisher().renew_partition_lease(&claim).await;
#[cfg(feature = "prometheus")]
if let Ok(renewed) = &renewed {
crate::metrics::record_partition_lease_renewal(
handler.get_publisher().metrics_queue(),
&claim.topic,
*renewed,
);
}
match renewed {
Ok(true) => {}
Ok(false) => {
tracing::warn!(
job_id = %claim.job_id,
topic = %claim.topic,
partition = claim.partition.0,
"Partition lease was lost"
);
let _ = lease_lost.send(());
return;
}
Err(error) => {
tracing::warn!(%error, job_id = %claim.job_id, "Failed to renew partition lease");
}
}
}
})
}
#[tracing::instrument(skip_all, fields(job_id = %job.id, job_type = %job.payload_type, topic = %claim.topic, partition = claim.partition.0))]
async fn handle_partitioned_job<C, H>(
handler: &Arc<H>,
claim: &PartitionHeadClaim,
job: Job,
lease_lost: &mut tokio::sync::watch::Receiver<()>,
#[cfg(feature = "prometheus")] worker_id: &str,
) -> anyhow::Result<()>
where
C: Sync + Send,
H: BgJobHandler<C> + Sync + Send + 'static,
{
let ptype = job.payload_type.clone();
let payload = job.payload.clone();
let job_id = job.id.clone();
let publisher = handler.get_publisher();
let running_job = job.transition();
publisher.save(&running_job).await?;
#[cfg(feature = "prometheus")]
let _partition_active_guard =
crate::metrics::PartitionActiveGuard::start(publisher.metrics_queue(), &claim.topic);
#[cfg(feature = "prometheus")]
let handler_timer = HandlerTimer::start(publisher.metrics_queue(), worker_id, &ptype);
let handler_result = tokio::select! {
result = handler.dispatch(ptype.clone(), &payload, job_id) => result,
changed = lease_lost.changed() => match changed {
Ok(()) => Err(anyhow::anyhow!(
"partition lease lost while handler was running"
)),
Err(_) => Err(anyhow::anyhow!(
"partition lease renewal task stopped while handler was running"
)),
},
};
#[cfg(feature = "prometheus")]
handler_timer.finish(match &handler_result {
Ok(()) => HandlerOutcome::Success,
Err(_) => HandlerOutcome::Error,
});
match handler_result {
Ok(_) => {
let success_job = running_job.transition_success()?;
let success_job_id = success_job.id.clone();
let operations = publisher
.job_save_and_expire_operations(&success_job, Duration::from_secs(3600))?;
if publisher.complete_partition_head(claim, operations).await? {
enqueue_continuations(handler.clone(), &success_job, &success_job_id).await?;
} else {
tracing::warn!(
job_id = %success_job_id,
topic = %claim.topic,
partition = claim.partition.0,
lease_epoch = claim.lease_epoch,
"Partition lease was lost before successful completion"
);
}
}
Err(error) => {
tracing::warn!("Failed partitioned job {}: {}", running_job.id, error);
let mut running_job = running_job;
if running_job.config.needs_retry_policy() {
match handler.retry_policy(&ptype, &payload).await {
Ok(retry_policy) => running_job
.config
.resolve_retry_policy(retry_policy.as_ref()),
Err(policy_error) => tracing::warn!(
"Could not resolve retry policy for job {}: {}",
running_job.id,
policy_error
),
}
}
let failed_job_id = running_job.id.clone();
match running_job.transition_failure(error.to_string())? {
FailureTransition::Retry(job) => {
let not_before = requeued_not_before(&job).unwrap_or_else(chrono::Utc::now);
let operations = publisher.job_save_operations(&job)?;
let released = publisher
.release_partition_head(claim, not_before, operations)
.await?;
if !released {
tracing::warn!(
job_id = %failed_job_id,
topic = %claim.topic,
partition = claim.partition.0,
lease_epoch = claim.lease_epoch,
"Partition lease was lost before retry release"
);
}
}
FailureTransition::Failed(job) => {
let operations = publisher
.job_save_and_expire_operations(&job, Duration::from_secs(3600))?;
if !publisher.complete_partition_head(claim, operations).await? {
tracing::warn!(
job_id = %failed_job_id,
topic = %claim.topic,
partition = claim.partition.0,
lease_epoch = claim.lease_epoch,
"Partition lease was lost before final failure completion"
);
}
}
}
}
}
Ok(())
}
fn requeued_not_before(job: &Job) -> Option<UtcDateTime> {
match &job.stage {
Stage::Requeued(RequeuedStage { not_before, .. }) => *not_before,
_ => None,
}
}