修复超时任务在途写回竞争
超时/续租失败时不再 drop 在途 work 并抢先写失败态(已发出的 procedure 无法撤销,可能造成业务已提交但任务被标记失败并退款)。 改为脱管运行至租约仲裁窗口结束,由服务端 lease fencing 裁决写回, 任务重派与耗尽结算走队列侧原子事务。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -8,7 +8,10 @@ use spacetime_client::{
|
||||
ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord,
|
||||
ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription,
|
||||
};
|
||||
use tokio::{task::JoinSet, time::sleep};
|
||||
use tokio::{
|
||||
task::{JoinHandle, JoinSet},
|
||||
time::sleep,
|
||||
};
|
||||
use tracing::{error, info, warn};
|
||||
|
||||
const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048;
|
||||
@@ -322,33 +325,113 @@ async fn process_external_generation_job(
|
||||
return Err(message);
|
||||
}
|
||||
};
|
||||
let work = with_external_generation_billing_attempt_context(
|
||||
// work 单独 spawn:已发出的 SpacetimeDB procedure 无法通过 drop 客户端 future 撤销,
|
||||
// 若超时时 drop 在途写回再抢先写失败态,可能出现“业务写回已提交、任务却被标记失败
|
||||
// 并退款”的不一致。因此超时/续租失败时不取消 work,也不在客户端写失败态,交由服务端
|
||||
// lease fencing 仲裁:写回在租约有效期内到达则任务照常完成,否则被拒绝;租约到期后
|
||||
// 任务被重新认领,attempt 耗尽时由认领事务原子地标记失败并结算退款。
|
||||
let mut work_handle = tokio::spawn(with_external_generation_billing_attempt_context(
|
||||
job.job_id.clone(),
|
||||
billing_claim_attempt,
|
||||
billing_price_mud_points,
|
||||
process_external_generation_job_once(state.clone(), worker_id.clone(), job.clone()),
|
||||
);
|
||||
));
|
||||
let heartbeat =
|
||||
maintain_external_generation_job_lease(&state, &worker_id, &job, lease, heartbeat_interval);
|
||||
match await_external_generation_job_execution(work, heartbeat, job_timeout).await {
|
||||
match await_external_generation_job_execution(
|
||||
join_external_generation_work(&mut work_handle),
|
||||
heartbeat,
|
||||
job_timeout,
|
||||
)
|
||||
.await
|
||||
{
|
||||
ExternalGenerationJobExecutionOutcome::Finished(result) => result,
|
||||
ExternalGenerationJobExecutionOutcome::TimedOut => {
|
||||
// work 与 heartbeat future 已在调度函数返回前被 drop,它们持有的
|
||||
// SpacetimeDB 连接租约也已释放;此时才写失败态,避免单连接池自等待。
|
||||
let message = external_generation_worker_timeout_message(&job, job_timeout);
|
||||
warn!(
|
||||
job_id = %job.job_id,
|
||||
job_kind = %job.job_kind,
|
||||
timeout_seconds = job_timeout.as_secs(),
|
||||
"external generation worker 任务超过执行预算,停止当前尝试并释放 worker 槽位"
|
||||
"external generation worker 任务超过执行预算,停止续租并释放 worker 槽位,在途执行交由租约仲裁"
|
||||
);
|
||||
detach_external_generation_work_until_lease_expiry(
|
||||
work_handle,
|
||||
&job,
|
||||
lease,
|
||||
"任务超过执行预算",
|
||||
);
|
||||
fail_job(&state, &worker_id, &job, message.clone()).await?;
|
||||
Err(message)
|
||||
}
|
||||
ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error) => Err(error),
|
||||
ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error) => {
|
||||
detach_external_generation_work_until_lease_expiry(
|
||||
work_handle,
|
||||
&job,
|
||||
lease,
|
||||
"任务租约续期失败",
|
||||
);
|
||||
Err(error)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn join_external_generation_work(
|
||||
work_handle: &mut JoinHandle<Result<(), String>>,
|
||||
) -> Result<(), String> {
|
||||
match work_handle.await {
|
||||
Ok(result) => result,
|
||||
Err(join_error) => Err(format!("external generation work 任务 panic:{join_error}")),
|
||||
}
|
||||
}
|
||||
|
||||
/// 停止等待 work 后不能直接取消它:在途 procedure 可能已在服务端提交。让它继续跑完当前
|
||||
/// 尝试,写回由 lease fencing 裁决是否生效。等到租约必然过期(2 倍租约时长,此后一切
|
||||
/// 围栏写回都会被拒绝)仍未结束的任务才安全取消,取消触发的计费补偿退款按 attempt
|
||||
/// 账本幂等,与队列侧耗尽结算不会重复。
|
||||
fn detach_external_generation_work_until_lease_expiry(
|
||||
mut work_handle: JoinHandle<Result<(), String>>,
|
||||
job: &ExternalGenerationJobRecord,
|
||||
lease: Duration,
|
||||
reason: &'static str,
|
||||
) {
|
||||
let job_id = job.job_id.clone();
|
||||
let job_kind = job.job_kind.clone();
|
||||
let grace = lease.saturating_mul(2);
|
||||
tokio::spawn(async move {
|
||||
match tokio::time::timeout(grace, &mut work_handle).await {
|
||||
Ok(Ok(Ok(()))) => info!(
|
||||
job_id = %job_id,
|
||||
job_kind = %job_kind,
|
||||
reason,
|
||||
"external generation worker 脱管任务已在租约仲裁窗口内完成写回"
|
||||
),
|
||||
Ok(Ok(Err(error))) => info!(
|
||||
job_id = %job_id,
|
||||
job_kind = %job_kind,
|
||||
reason,
|
||||
error = %error,
|
||||
"external generation worker 脱管任务在租约仲裁窗口内以失败结束"
|
||||
),
|
||||
Ok(Err(join_error)) => error!(
|
||||
job_id = %job_id,
|
||||
job_kind = %job_kind,
|
||||
reason,
|
||||
error = %join_error,
|
||||
"external generation worker 脱管任务 panic"
|
||||
),
|
||||
Err(_) => {
|
||||
work_handle.abort();
|
||||
warn!(
|
||||
job_id = %job_id,
|
||||
job_kind = %job_kind,
|
||||
reason,
|
||||
grace_ms = grace.as_millis() as u64,
|
||||
"external generation worker 脱管任务超过租约仲裁窗口仍未结束,已取消;其后续写回将被 lease fencing 拒绝"
|
||||
);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
enum ExternalGenerationJobExecutionOutcome {
|
||||
Finished(Result<(), String>),
|
||||
@@ -1443,7 +1526,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_deadline_drops_work_connection_before_failure_writeback() {
|
||||
async fn worker_deadline_drops_local_work_future_before_returning() {
|
||||
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
|
||||
let work_connection = connection.clone();
|
||||
let work = async move {
|
||||
@@ -1462,7 +1545,60 @@ mod tests {
|
||||
assert_eq!(outcome, ExternalGenerationJobExecutionOutcome::TimedOut);
|
||||
assert!(
|
||||
connection.try_acquire().is_ok(),
|
||||
"deadline 返回前应先 drop work 并释放连接"
|
||||
"deadline 返回前应先 drop 传入的 work future 并释放其持有的资源"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_detached_work_keeps_running_after_deadline() {
|
||||
let finished = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
|
||||
let finished_flag = finished.clone();
|
||||
let work_handle = tokio::spawn(async move {
|
||||
tokio::time::sleep(Duration::from_millis(20)).await;
|
||||
finished_flag.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||
Ok(())
|
||||
});
|
||||
let job = external_generation_job_record_fixture(Some("lease-1"));
|
||||
|
||||
detach_external_generation_work_until_lease_expiry(
|
||||
work_handle,
|
||||
&job,
|
||||
Duration::from_millis(200),
|
||||
"任务超过执行预算",
|
||||
);
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
assert!(
|
||||
finished.load(std::sync::atomic::Ordering::SeqCst),
|
||||
"超时后在途 work 应继续执行完成,而不是被取消"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_detached_work_is_aborted_after_lease_arbitration_window() {
|
||||
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
|
||||
let work_connection = connection.clone();
|
||||
let work_handle = tokio::spawn(async move {
|
||||
let _permit = work_connection
|
||||
.acquire_owned()
|
||||
.await
|
||||
.expect("work should acquire the only connection");
|
||||
std::future::pending::<Result<(), String>>().await
|
||||
});
|
||||
let job = external_generation_job_record_fixture(Some("lease-1"));
|
||||
|
||||
detach_external_generation_work_until_lease_expiry(
|
||||
work_handle,
|
||||
&job,
|
||||
Duration::from_millis(10),
|
||||
"任务超过执行预算",
|
||||
);
|
||||
|
||||
let reacquired =
|
||||
tokio::time::timeout(Duration::from_secs(2), connection.acquire_owned()).await;
|
||||
assert!(
|
||||
reacquired.is_ok(),
|
||||
"超过租约仲裁窗口后应取消仍未结束的 work 并释放其持有的资源"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user