From 2e9232787498827b3b3477ca575b3aa087826a93 Mon Sep 17 00:00:00 2001 From: Linghong Date: Thu, 16 Jul 2026 07:36:20 +0000 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E8=B6=85=E6=97=B6=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=9C=A8=E9=80=94=E5=86=99=E5=9B=9E=E7=AB=9E=E4=BA=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 超时/续租失败时不再 drop 在途 work 并抢先写失败态(已发出的 procedure 无法撤销,可能造成业务已提交但任务被标记失败并退款)。 改为脱管运行至租约仲裁窗口结束,由服务端 lease fencing 裁决写回, 任务重派与耗尽结算走队列侧原子事务。 Co-Authored-By: Claude Fable 5 --- .../src/external_generation_worker.rs | 158 ++++++++++++++++-- 1 file changed, 147 insertions(+), 11 deletions(-) diff --git a/server-rs/crates/api-server/src/external_generation_worker.rs b/server-rs/crates/api-server/src/external_generation_worker.rs index 56d87938f..2bedd96a3 100644 --- a/server-rs/crates/api-server/src/external_generation_worker.rs +++ b/server-rs/crates/api-server/src/external_generation_worker.rs @@ -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> { + 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>, + 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::>().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 并释放其持有的资源" ); }