refactor: add fenced transfer recovery leases

This commit is contained in:
jxxghp
2026-08-27 14:39:38 +08:00
parent 22865cb35b
commit a62c541ec0
32 changed files with 4135 additions and 650 deletions
+78 -11
View File
@@ -627,6 +627,14 @@ class TransferPlanningStateError(RuntimeError):
"""计划检查点无法从当前持久状态推进时抛出的状态错误。"""
class TransferAdmissionProjectionError(RuntimeError):
"""持久登记无法安全恢复为应用层整理准入投影时抛出的错误。"""
class TransferLeaseLostError(TransferPlanningStateError):
"""整理 worker 已失去持久租约、不得继续推进任务时抛出的错误。"""
class TransferTask(OptionalMediaIdentityMixin, BaseModel):
"""
文件整理任务。
@@ -661,6 +669,8 @@ class TransferTask(OptionalMediaIdentityMixin, BaseModel):
_planning_input: Optional[TransferPlanningInput] = PrivateAttr(default=None)
_plan_checkpoint: Optional[TransferPlanCheckpoint] = PrivateAttr(default=None)
_planning_context_restored: bool = PrivateAttr(default=False)
_lease_owner: Optional[str] = PrivateAttr(default=None)
_lease_token: Optional[str] = PrivateAttr(default=None)
@property
def admission_task_id(self) -> Optional[str]:
@@ -698,6 +708,23 @@ class TransferTask(OptionalMediaIdentityMixin, BaseModel):
"""标记领域上下文已离线恢复,禁止旧流程再次在线补充。"""
self._planning_context_restored = True
@property
def lease_owner(self) -> Optional[str]:
"""返回当前任务绑定的进程级 worker owner。"""
return self._lease_owner
@property
def lease_token(self) -> Optional[str]:
"""返回当前任务绑定的持久租约令牌。"""
return self._lease_token
def bind_execution_lease(self, *, owner_id: str, lease_token: str) -> None:
"""绑定持久 claim 结果,且不改变插件可见的旧任务序列化字段。"""
if not owner_id or not lease_token:
raise ValueError("整理执行租约缺少 owner 或 token")
self._lease_owner = owner_id
self._lease_token = lease_token
def to_dict(self):
"""
返回字典。
@@ -743,6 +770,11 @@ class TransferAdmission:
input_fingerprint: Optional[str] = None
planning_input: Optional[TransferPlanningInput] = None
checkpoint: Optional[TransferPlanCheckpoint] = None
lease_owner: Optional[str] = None
lease_token: Optional[str] = None
lease_expires_at: Optional[str] = None
heartbeat_at: Optional[str] = None
attempt_count: int = 0
class TransferAdmissionRepository(Protocol):
@@ -758,10 +790,6 @@ class TransferAdmissionRepository(Protocol):
"""按规划输入幂等登记源文件并返回稳定任务身份。"""
...
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
"""按登记顺序返回等待恢复或执行的任务。"""
...
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
"""记录内存队列接收失败,保留任务供后续恢复。"""
...
@@ -770,22 +798,61 @@ class TransferAdmissionRepository(Protocol):
self,
*,
task_id: str,
lease_token: str,
input_fingerprint: str,
checkpoint: TransferPlanCheckpoint,
) -> TransferAdmission:
"""原子保存完整计划并匹配输入的任务推进到已规划"""
"""按有效租约原子保存完整计划并推进匹配输入的任务。"""
...
def record_planning_failure(self, *, task_id: str, error: str) -> None:
"""记录规划失败但保留接纳状态供下次恢复重试。"""
def record_planning_failure(
self, *, task_id: str, lease_token: str, error: str
) -> None:
"""按有效租约记录规划失败,保留业务状态供恢复重试。"""
...
def list_recoverable(self, limit: int = 5000) -> list[TransferAdmission]:
"""按登记顺序返回接纳或已规划的可恢复任务。"""
def claim_task(
self,
*,
task_id: str,
owner_id: str,
lease_seconds: int,
) -> Optional[TransferAdmission]:
"""为指定任务取得唯一执行租约;已有有效租约时返回 None。"""
...
def discard_task(self, *, task_id: str) -> int:
"""按稳定任务身份删除已经到达终态的登记。"""
def claim_recoverable(
self,
*,
owner_id: str,
limit: int,
lease_seconds: int,
) -> list[TransferAdmission]:
"""按登记顺序原子取得可恢复任务的唯一执行租约。"""
...
def heartbeat(
self,
*,
task_id: str,
lease_token: str,
lease_seconds: int,
) -> Optional[TransferAdmission]:
"""仅为当前且未过期的 token 延长租约;过期 token 不得复活。"""
...
def release_claim(
self,
*,
task_id: str,
lease_token: str,
error: Optional[str] = None,
) -> bool:
"""按 token 释放当前 claim,并可记录可恢复错误。"""
...
def discard_claimed(self, *, task_id: str, lease_token: str) -> int:
"""仅允许当前 lease owner 删除已经到达终态的登记。"""
...
+644 -42
View File
@@ -46,6 +46,7 @@ from app.application.transfer import (
TransferAdmission,
TransferFailureNotification,
TransferFailureNotificationAggregator,
TransferLeaseLostError,
TransferPlanCheckpoint,
TransferPlanningInput,
TransferPlanningStateError,
@@ -146,6 +147,10 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
_WORKER_RESTART_TIMEOUT_SECONDS = 30.0
_WORKER_CLOSE_TIMEOUT_SECONDS = 30.0
_QUEUE_STOP_SENTINEL = object()
_WORKER_LEASE_SECONDS = 120
_LEASE_HEARTBEAT_INTERVAL_SECONDS = 30.0
_RECOVERY_POLL_INTERVAL_SECONDS = 15.0
_RECOVERY_CLAIM_LIMIT = 100
@staticmethod
def _transfer_result_payload(
@@ -254,9 +259,16 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self._worker_lifecycle_lock = threading.RLock()
self._worker_state_lock = threading.RLock()
self._closing = False
# pending 回放同样由整理链持有,关闭时可阻止继续处理下一条登记
# 恢复调度与租约续期均由整理链持有,关闭时与 worker 一并收敛。
self._worker_owner_id = uuid.uuid4().hex
self._owned_leases: Dict[str, Tuple[str, float]] = {}
self._queued_lease_tokens: set[Tuple[str, str]] = set()
self._replay_thread: Optional[threading.Thread] = None
self._replay_stop_event = threading.Event()
self._recovery_wakeup_event = threading.Event()
self._lease_heartbeat_thread: Optional[threading.Thread] = None
self._lease_heartbeat_stop_event = threading.Event()
self._lease_release_thread: Optional[threading.Thread] = None
self._active_tasks = 0
self._processed_num = 0
self._fail_num = 0
@@ -366,6 +378,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
:param timeout_seconds: 生命周期锁、worker 与回放线程共享的最大等待秒数
:return: 全部后台线程均已收敛时返回 True,否则返回 False
"""
self.__ensure_lease_runtime_state()
deadline = time.monotonic() + max(0.0, timeout_seconds)
if not self.__acquire_worker_lifecycle_lock(deadline):
logger.error(
@@ -378,7 +391,9 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self._closing = True
worker_threads = self.__request_worker_stop()
replay_thread = self._replay_thread
heartbeat_thread = self._lease_heartbeat_thread
self._replay_stop_event.set()
self._recovery_wakeup_event.set()
alive_workers = self.__join_threads(worker_threads, deadline)
alive_replays = self.__join_threads(
@@ -387,7 +402,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
with self._worker_state_lock:
self._threads = []
self._retiring_threads = alive_workers
if replay_thread and not alive_replays and self._replay_thread is replay_thread:
if (
replay_thread
and replay_thread not in alive_replays
and self._replay_thread is replay_thread
):
self._replay_thread = None
alive_threads = [*alive_workers, *alive_replays]
@@ -398,6 +417,38 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
", ".join(thread.name for thread in alive_threads),
)
return False
with self._worker_state_lock:
release_thread = self.__start_lease_release_owner_locked(
error="整理宿主关闭,释放未结算任务租约"
)
alive_releases = self.__join_threads(
[release_thread] if release_thread else [], deadline
)
if alive_releases:
logger.error(
"整理租约释放线程未在 %.1f 秒内收敛,heartbeat 保持运行:%s",
max(0.0, timeout_seconds),
", ".join(thread.name for thread in alive_releases),
)
return False
with self._worker_state_lock:
if self._lease_release_thread is release_thread:
self._lease_release_thread = None
self._lease_heartbeat_stop_event.set()
alive_heartbeats = self.__join_threads(
[heartbeat_thread] if heartbeat_thread else [], deadline
)
if heartbeat_thread and not alive_heartbeats:
with self._worker_state_lock:
if self._lease_heartbeat_thread is heartbeat_thread:
self._lease_heartbeat_thread = None
if alive_heartbeats:
logger.error(
"整理租约续期线程未在 %.1f 秒内收敛:%s",
max(0.0, timeout_seconds),
", ".join(thread.name for thread in alive_heartbeats),
)
return False
logger.info("文件整理 worker 与待处理回放线程已关闭")
return True
finally:
@@ -885,6 +936,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
if self._closing:
logger.warning("文件整理链已关闭,拒绝新的队列任务")
return False
if isinstance(task.lease_token, str) and task.lease_token:
return self.__enqueue_claimed_task(task)
return self._transfer_queue_service().put(task, self.__default_callback)
def _transfer_queue_service(self) -> TransferQueueService:
@@ -902,39 +955,369 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
def replay_pending(self) -> None:
"""
回放上次进程退出时仍未整理完的文件
启动唯一恢复调度 owner,并唤醒一次即时恢复扫描
在后台线程执行:回放要 stat 源文件,而启动期挂载可能尚未就绪甚至处于
挂死状态,同步执行会把整个启动流程堵住
启动回放、同进程入队补偿和租约过期接管都经由这个入口。调度线程只负责
claim 和重新入队;实际业务仍由普通整理 worker 执行
"""
self.__ensure_recovery_scheduler(immediate=True)
def __ensure_recovery_scheduler(self, *, immediate: bool) -> None:
"""确保唯一恢复调度 owner 存在,并只为显式请求执行即时扫描。"""
self.__ensure_lease_runtime_state()
with self._worker_state_lock:
if self._closing:
logger.info("文件整理链正在关闭,跳过待处理文件回放")
return
self.__start_lease_heartbeat_owner_locked()
if self._replay_thread and self._replay_thread.is_alive():
logger.info("待处理文件回放已在运行,跳过重复启动")
if immediate:
self._recovery_wakeup_event.set()
return
stop_event = threading.Event()
thread = threading.Thread(
target=self.__run_replay_pending,
args=(stop_event,),
args=(self._replay_stop_event, immediate),
name="MoviePilot-TransferReplay",
daemon=True,
)
self._replay_stop_event = stop_event
self._replay_thread = thread
if immediate:
self._recovery_wakeup_event.set()
# 在状态锁内启动,避免 close_workers 看到尚未 start 的线程后错误 join。
thread.start()
def __run_replay_pending(self, stop_event: threading.Event) -> None:
"""执行一次受控回放,并在自然结束后释放当前线程句柄。"""
def __run_replay_pending(
self,
stop_event: threading.Event,
initial_immediate: bool = True,
) -> None:
"""按初次扫描意图持续恢复,失败兜底新建 owner 时先等待固定轮询。"""
try:
self.__replay_pending(stop_event)
if not initial_immediate:
self._recovery_wakeup_event.wait(
timeout=self._RECOVERY_POLL_INTERVAL_SECONDS
)
while not stop_event.is_set():
self._recovery_wakeup_event.clear()
self.__replay_pending(stop_event)
self._recovery_wakeup_event.wait(
timeout=self._RECOVERY_POLL_INTERVAL_SECONDS
)
finally:
with self._worker_state_lock:
if self._replay_thread is threading.current_thread():
self._replay_thread = None
def __ensure_lease_runtime_state(self) -> None:
"""为绕过构造器的兼容调用补齐进程 owner 与租约线程状态。"""
if not hasattr(self, "_worker_owner_id"):
self._worker_owner_id = uuid.uuid4().hex
if not hasattr(self, "_owned_leases"):
self._owned_leases = {}
if not hasattr(self, "_queued_lease_tokens"):
self._queued_lease_tokens = set()
if not hasattr(self, "_recovery_wakeup_event"):
self._recovery_wakeup_event = threading.Event()
if not hasattr(self, "_lease_heartbeat_thread"):
self._lease_heartbeat_thread = None
if not hasattr(self, "_lease_heartbeat_stop_event"):
self._lease_heartbeat_stop_event = threading.Event()
if not hasattr(self, "_lease_release_thread"):
self._lease_release_thread = None
if not hasattr(self, "_replay_thread"):
self._replay_thread = None
if not hasattr(self, "_replay_stop_event"):
self._replay_stop_event = threading.Event()
if not hasattr(self, "_worker_state_lock"):
self._worker_state_lock = threading.RLock()
if not hasattr(self, "_closing"):
self._closing = False
def __start_lease_heartbeat_owner_locked(self) -> None:
"""在状态锁内确保当前进程只有一个租约续期线程。"""
heartbeat_thread = self._lease_heartbeat_thread
if heartbeat_thread and heartbeat_thread.is_alive():
return
heartbeat_thread = threading.Thread(
target=self.__run_lease_heartbeat,
args=(self._lease_heartbeat_stop_event,),
name="MoviePilot-TransferLeaseHeartbeat",
daemon=True,
)
self._lease_heartbeat_thread = heartbeat_thread
heartbeat_thread.start()
def __ensure_lease_heartbeat_owner(self) -> None:
"""按需启动进程级租约续期 owner,供启动前到达的普通任务使用。"""
self.__ensure_lease_runtime_state()
with self._worker_state_lock:
if not self._closing:
self.__start_lease_heartbeat_owner_locked()
def __run_lease_heartbeat(self, stop_event: threading.Event) -> None:
"""按固定周期续期本进程已 claim 且尚未结算的任务。"""
try:
while not stop_event.wait(self._LEASE_HEARTBEAT_INTERVAL_SECONDS):
self.__heartbeat_owned_leases()
finally:
with self._worker_state_lock:
if self._lease_heartbeat_thread is threading.current_thread():
self._lease_heartbeat_thread = None
def __heartbeat_owned_leases(self) -> None:
"""经 Application Port 续期所有排队中或执行中的任务租约。"""
self.__ensure_lease_runtime_state()
with self._worker_state_lock:
owned_leases = list(self._owned_leases.items())
for task_id, (lease_token, deadline) in owned_leases:
try:
admission = self._transfer_admissions.heartbeat(
task_id=task_id,
lease_token=lease_token,
lease_seconds=self._WORKER_LEASE_SECONDS,
)
except Exception as err:
if time.monotonic() >= deadline:
self.__forget_owned_lease(task_id, lease_token)
logger.error(
"整理任务租约续期持续失败并已超过本地期限:%s - %s",
task_id,
err,
)
else:
logger.error(f"整理任务租约续期失败:{task_id} - {err}")
continue
if (
admission is None
or admission.lease_owner != self._worker_owner_id
or admission.lease_token != lease_token
):
self.__forget_owned_lease(task_id, lease_token)
logger.error(f"整理任务租约已失效或被接管:{task_id}")
continue
with self._worker_state_lock:
current = self._owned_leases.get(task_id)
if current and current[0] == lease_token:
self._owned_leases[task_id] = (
lease_token,
time.monotonic() + self._WORKER_LEASE_SECONDS,
)
def __bind_claimed_admission(
self,
task: TransferTask,
admission: TransferAdmission,
) -> None:
"""校验仓储 claim 投影并把 token 私有绑定到执行任务。"""
if (
admission.task_id != task.admission_task_id
):
raise TransferLeaseLostError(
f"整理任务 claim 投影无效:{admission.task_id}"
)
self.__register_claimed_admission(admission)
assert admission.lease_owner is not None
assert admission.lease_token is not None
task.bind_execution_lease(
owner_id=admission.lease_owner,
lease_token=admission.lease_token,
)
def __register_claimed_admission(
self,
admission: TransferAdmission,
) -> None:
"""把仓储 claim 加入续期集合,覆盖恢复构造阶段可能发生的同步 I/O。"""
if (
admission.lease_owner != self._worker_owner_id
or not admission.lease_token
):
raise TransferLeaseLostError(
f"整理任务 claim 投影无效:{admission.task_id}"
)
with self._worker_state_lock:
self._owned_leases[admission.task_id] = (
admission.lease_token,
time.monotonic() + self._WORKER_LEASE_SECONDS,
)
self.__ensure_lease_heartbeat_owner()
def __forget_owned_lease(self, task_id: str, lease_token: str) -> None:
"""仅在 token 仍匹配时移除本进程租约镜像,避免删掉新接管记录。"""
with self._worker_state_lock:
current = self._owned_leases.get(task_id)
if current and current[0] == lease_token:
self._owned_leases.pop(task_id, None)
self._queued_lease_tokens.discard((task_id, lease_token))
def __is_claimed_task_enqueued(self, task_id: str, lease_token: str) -> bool:
"""返回指定 claim 是否已经成功进入普通 worker 队列。"""
with self._worker_state_lock:
return (task_id, lease_token) in self._queued_lease_tokens
def __owns_lease(self, task_id: str, lease_token: Optional[str]) -> bool:
"""返回本地续期镜像是否仍持有指定 token。"""
if not lease_token:
return False
with self._worker_state_lock:
current = self._owned_leases.get(task_id)
return bool(current and current[0] == lease_token)
def __assert_owned_lease(self, task: TransferTask) -> None:
"""拒绝无 token、已过本地期限或已被 heartbeat 判失效的任务推进。"""
if task.preview:
return
task_id = task.admission_task_id
lease_token = task.lease_token
if (
not task_id
or not lease_token
or task.lease_owner != self._worker_owner_id
):
raise TransferLeaseLostError("整理任务缺少当前进程的有效执行租约")
with self._worker_state_lock:
current = self._owned_leases.get(task_id)
if (
current is None
or current[0] != lease_token
or current[1] <= time.monotonic()
):
if current and current[0] == lease_token:
self._owned_leases.pop(task_id, None)
raise TransferLeaseLostError(f"整理任务租约已经失效:{task_id}")
def __claim_task_for_execution(self, task: TransferTask) -> None:
"""让普通队列任务在业务执行前取得唯一租约,恢复任务复用既有 token。"""
if task.preview:
return
self.__ensure_lease_runtime_state()
if task.lease_token:
self.__assert_owned_lease(task)
return
if not task.admission_task_id:
admitted = self.__admit_transfer(task)
task.bind_admission_task_id(admitted.task_id)
task_id = task.admission_task_id
if task_id is None:
raise TransferLeaseLostError("整理任务准入后仍缺少 durable 身份")
claimed = self._transfer_admissions.claim_task(
task_id=task_id,
owner_id=self._worker_owner_id,
lease_seconds=self._WORKER_LEASE_SECONDS,
)
if claimed is None:
raise TransferLeaseLostError(
f"整理任务已由其他 worker claim{task_id}"
)
try:
self.__bind_claimed_admission(task, claimed)
except Exception as err:
self.__release_admission_claim(claimed, error=str(err))
raise
def __release_task_claim(
self,
task: TransferTask,
*,
error: Optional[str] = None,
) -> bool:
"""按 task token 释放未到终态的 claim,并按固定轮询等待恢复。"""
task_id = task.admission_task_id if task else None
lease_token = task.lease_token if task else None
if not task_id or not lease_token:
return False
try:
return self._transfer_admissions.release_claim(
task_id=task_id,
lease_token=lease_token,
error=error,
)
except Exception as err:
logger.error(f"释放整理任务租约失败:{task_id} - {err}")
return False
finally:
self.__forget_owned_lease(task_id, lease_token)
self.__ensure_recovery_scheduler(immediate=False)
def __release_admission_claim(
self,
admission: TransferAdmission,
*,
error: Optional[str] = None,
) -> None:
"""释放尚未绑定或成功入队的恢复 claim。"""
if not admission.lease_token:
return
try:
self._transfer_admissions.release_claim(
task_id=admission.task_id,
lease_token=admission.lease_token,
error=error,
)
except Exception as err:
logger.error(f"释放恢复任务租约失败:{admission.task_id} - {err}")
finally:
self.__forget_owned_lease(admission.task_id, admission.lease_token)
def __release_all_owned_leases(self, *, error: str) -> None:
"""在执行 owner 全部收敛后释放剩余 claim,避免关停后等待租约自然过期。"""
with self._worker_state_lock:
owned_leases = list(self._owned_leases.items())
for task_id, (lease_token, _deadline) in owned_leases:
try:
self._transfer_admissions.release_claim(
task_id=task_id,
lease_token=lease_token,
error=error,
)
except Exception as err:
logger.error(f"关闭时释放整理任务租约失败:{task_id} - {err}")
finally:
self.__forget_owned_lease(task_id, lease_token)
def __start_lease_release_owner_locked(
self, *, error: str
) -> Optional[threading.Thread]:
"""启动并持有唯一租约释放线程,使同步数据库阻塞不突破关闭预算。"""
release_thread = self._lease_release_thread
if release_thread is not None:
return release_thread
if not self._owned_leases:
return None
release_thread = threading.Thread(
target=self.__release_all_owned_leases,
kwargs={"error": error},
name="MoviePilot-TransferLeaseRelease",
daemon=True,
)
self._lease_release_thread = release_thread
release_thread.start()
return release_thread
def __enqueue_claimed_task(self, task: TransferTask) -> bool:
"""把已 claim 的恢复任务送入普通队列,禁止再次准入或二次 claim。"""
self.__assert_owned_lease(task)
if not self.__put_to_jobview(task):
return False
try:
self._register_scrape_batch_task(task)
assert task.admission_task_id is not None
assert task.lease_token is not None
with self._worker_state_lock:
self._queued_lease_tokens.add(
(task.admission_task_id, task.lease_token)
)
self._queue.put(
TransferQueue(task=task, callback=self.__default_callback)
)
except Exception as err:
try:
self.__record_enqueue_failure(task, err)
finally:
self.jobview.remove_task(task.fileitem)
raise
return True
def __replay_pending(
self, stop_event: Optional[threading.Event] = None
) -> None:
@@ -948,17 +1331,35 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
stop_event = stop_event or threading.Event()
if stop_event.is_set():
return
self.__ensure_lease_runtime_state()
try:
pendings = self._transfer_admissions.list_recoverable()
pendings = self._transfer_admissions.claim_recoverable(
owner_id=self._worker_owner_id,
limit=self._RECOVERY_CLAIM_LIMIT,
lease_seconds=self._WORKER_LEASE_SECONDS,
)
except Exception as err:
logger.error(f"读取待整理文件登记失败:{err}")
return
if not pendings:
return
try:
for admission in pendings:
self.__register_claimed_admission(admission)
except Exception as err:
logger.error(f"登记恢复任务租约失败:{err}")
for claimed in pendings:
self.__release_admission_claim(claimed, error=str(err))
return
logger.info(f"发现 {len(pendings)} 个上次未整理完的文件,正在重新送入整理链 ...")
replayed = 0
for admission in pendings:
for index, admission in enumerate(pendings):
if stop_event.is_set():
for unprocessed in pendings[index:]:
self.__release_admission_claim(
unprocessed,
error="整理宿主关闭,恢复任务尚未入队",
)
break
storage = admission.storage
src_path = admission.src_path
@@ -970,17 +1371,50 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
)
# stat 等同步 I/O 返回后重新检查,关闭期间不得注销尚未完成的登记。
if stop_event.is_set():
self.__release_admission_claim(
admission,
error="整理宿主关闭,恢复任务尚未入队",
)
for unprocessed in pendings[index + 1:]:
self.__release_admission_claim(
unprocessed,
error="整理宿主关闭,恢复任务尚未入队",
)
break
if not fileitem:
if should_discard:
# 源文件确认已消失,注销登记避免每次启动重复回放
self._transfer_admissions.discard_task(
task_id=admission.task_id
lease_token = admission.lease_token
if lease_token is None:
raise TransferLeaseLostError(
f"恢复任务缺少 lease token{admission.task_id}"
)
discarded = self._transfer_admissions.discard_claimed(
task_id=admission.task_id,
lease_token=lease_token,
)
if not discarded:
logger.warning(
f"恢复任务终态注销被 CAS 拒绝:{admission.task_id}"
)
self.__forget_owned_lease(
admission.task_id,
lease_token,
)
else:
self.__release_admission_claim(
admission,
error="恢复源文件暂时不可读取",
)
continue
if admission.checkpoint:
if self.__queue_planned_replay(fileitem, admission):
replayed += 1
else:
self.__release_admission_claim(
admission,
error="恢复任务未进入内存队列",
)
continue
planning_input = admission.planning_input
if planning_input and (
@@ -990,6 +1424,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
):
if self.__queue_accepted_replay(fileitem, admission):
replayed += 1
else:
self.__release_admission_claim(
admission,
error="恢复任务未进入内存队列",
)
continue
replay_kwargs = self.__build_replay_kwargs(planning_input)
self._execute_transfer(
@@ -997,9 +1436,21 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
recovery_admission=admission,
**replay_kwargs,
)
replayed += 1
assert admission.lease_token is not None
if self.__is_claimed_task_enqueued(
admission.task_id,
admission.lease_token,
):
replayed += 1
else:
self.__release_admission_claim(
admission,
error="旧恢复入口未产生可执行队列任务",
)
except Exception as err:
logger.error(f"回放待整理文件失败:{storage}:{src_path} - {err}")
if self.__owns_lease(admission.task_id, admission.lease_token):
self.__release_admission_claim(admission, error=str(err))
if stop_event.is_set():
logger.info(
"待整理文件回放收到关闭请求,已送入 %s 个文件,其余登记保持待处理",
@@ -1116,6 +1567,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
preview=False,
)
task.bind_admission_task_id(admission.task_id)
self.__bind_claimed_admission(task, admission)
task.bind_planning_input(planning_input)
if planning_input.mediainfo:
task.mark_planning_context_restored()
@@ -1168,6 +1620,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
preview=False,
)
task.bind_admission_task_id(admission.task_id)
self.__bind_claimed_admission(task, admission)
task.bind_planning_input(checkpoint.planning_input)
task.bind_plan_checkpoint(checkpoint)
self.__restore_planned_task(task)
@@ -1442,11 +1895,17 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
error: object,
) -> None:
"""记录 checkpoint 前失败,保留 accepted 任务供后续重新规划。"""
if task.preview or task.plan_checkpoint or not task.admission_task_id:
if (
task.preview
or task.plan_checkpoint
or not task.admission_task_id
or not task.lease_token
):
return
try:
self._transfer_admissions.record_planning_failure(
task_id=task.admission_task_id,
lease_token=task.lease_token,
error=str(error),
)
except Exception as record_error:
@@ -1488,9 +1947,9 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
"""先提交冻结 provider 调用,空结果时再提交并执行宿主计划。"""
planning_input = task.planning_input or self.__build_planning_input(task)
task.bind_planning_input(planning_input)
if not task.preview and not task.admission_task_id:
admission = self.__admit_transfer(task)
task.bind_admission_task_id(admission.task_id)
if not task.preview:
self.__claim_task_for_execution(task)
self.__assert_owned_lease(task)
checkpoint = task.plan_checkpoint
if checkpoint is None:
@@ -1540,6 +1999,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
raise
self.__restore_planned_task(task)
self.__assert_owned_lease(task)
legacy_result = self.__execute_legacy_transfer_providers(
task,
@@ -1571,6 +2031,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self.__record_checkpoint_failure(task, error)
raise
self.__assert_owned_lease(task)
result = self.execute_transfer_plan(
checkpoint,
meta=task.meta,
@@ -1623,8 +2084,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
return checkpoint
if not task.admission_task_id:
raise RuntimeError("整理计划提交前缺少持久任务身份")
self.__assert_owned_lease(task)
assert task.lease_token is not None
persisted = self._transfer_admissions.checkpoint_plan(
task_id=task.admission_task_id,
lease_token=task.lease_token,
input_fingerprint=planning_input.fingerprint,
checkpoint=checkpoint,
)
@@ -1638,11 +2102,12 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
error: Exception,
) -> None:
"""记录当前 checkpoint 阶段失败并保留原状态供重启恢复。"""
if task.preview or not task.admission_task_id:
if task.preview or not task.admission_task_id or not task.lease_token:
return
try:
self._transfer_admissions.record_planning_failure(
task_id=task.admission_task_id,
lease_token=task.lease_token,
error=str(error),
)
setattr(error, "_transfer_planning_failure_recorded", True)
@@ -1790,6 +2255,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
)
except Exception as error:
logger.error(f"旧整理兼容命令执行失败:{error}")
self.__release_task_claim(task, error=str(error))
return TransferInfo(
success=False,
message=str(error),
@@ -1797,7 +2263,16 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
fail_list=[fileitem.path],
transfer_type=transfer_type,
)
self.__discard_pending(task)
if not self.__discard_pending(task):
message = "旧整理兼容命令已执行,但 durable 终态结算失去租约"
logger.error(message)
return TransferInfo(
success=False,
message=message,
fileitem=fileitem,
fail_list=[fileitem.path],
transfer_type=transfer_type,
)
return result
def __record_enqueue_failure(
@@ -1805,9 +2280,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
task: TransferTask,
error: Exception,
) -> None:
"""记录内存入队失败并撤销该任务的批次占位。"""
"""按是否已经 claim 选择 fencing 失败记录,并撤销批次占位。"""
try:
if task.admission_task_id:
if task.lease_token:
self.__release_task_claim(task, error=str(error))
elif task.admission_task_id:
self._transfer_admissions.record_enqueue_failure(
task_id=task.admission_task_id,
error=str(error),
@@ -1819,8 +2296,9 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
)
finally:
self._finish_scrape_batch_task(task)
self.__ensure_recovery_scheduler(immediate=False)
def __discard_pending(self, task: TransferTask):
def __discard_pending(self, task: TransferTask) -> bool:
"""
注销一个待整理文件登记,整理到达终态(成功或失败)时调用。
@@ -1829,16 +2307,34 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
:param task: 任务信息
"""
if not task or not task.admission_task_id:
return
try:
self._transfer_admissions.discard_task(
task_id=task.admission_task_id
return True
if not task.lease_token:
logger.error(
f"整理任务缺少 lease token,拒绝伪装 durable 终态:"
f"{task.admission_task_id}"
)
return False
discarded = 0
try:
discarded = self._transfer_admissions.discard_claimed(
task_id=task.admission_task_id,
lease_token=task.lease_token,
)
if not discarded:
logger.error(
f"整理任务终态注销被 CAS 拒绝:{task.admission_task_id}"
)
except Exception as err:
logger.error(
"注销整理任务 durable admission 失败: "
f"{task.admission_task_id} - {err}"
)
finally:
self.__forget_owned_lease(
task.admission_task_id,
task.lease_token,
)
return bool(discarded)
def __put_to_jobview(self, task: TransferTask) -> bool:
"""
@@ -1905,13 +2401,23 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
if marker:
marker(task)
def __finish_job_execution(self, task: TransferTask, *, terminal: bool = True):
"""结束内存执行;仅确定到达业务终态时注销 durable 记录。"""
def __finish_job_execution(
self,
task: TransferTask,
*,
terminal: bool = True,
terminal_settlement: Optional[bool] = None,
) -> bool:
"""结束内存执行,并复用成功回调前已经提交的 durable 终态结果。"""
marker = getattr(self.jobview, "finish_execution", None)
if marker:
marker(task)
if terminal:
self.__discard_pending(task)
if terminal_settlement is not None:
return terminal_settlement
return self.__discard_pending(task)
self.__release_task_claim(task)
return True
def __expire_stale_transfer_tasks(self):
"""清理外部接管后失去状态心跳的运行中整理任务。"""
@@ -2001,6 +2507,32 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self.__settle_transfer_progress_if_idle()
continue
if task.admission_task_id and task.lease_token:
with self._worker_state_lock:
self._queued_lease_tokens.discard(
(task.admission_task_id, task.lease_token)
)
try:
self.__claim_task_for_execution(task)
except TransferLeaseLostError as err:
logger.info(f"跳过未取得执行租约的整理任务:{err}")
self.__release_task_claim(task, error=str(err))
self.jobview.try_remove_job(task)
self._finish_scrape_batch_task(task)
self._queue.task_done()
self.__settle_transfer_progress_if_idle()
continue
except Exception as err:
logger.error(
f"整理任务 claim 失败,保留 durable admission{err}"
)
self.jobview.try_remove_job(task)
self._finish_scrape_batch_task(task)
self._queue.task_done()
self.__settle_transfer_progress_if_idle()
self.__ensure_recovery_scheduler(immediate=False)
continue
# 文件信息
fileitem = task.fileitem
@@ -2030,6 +2562,29 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self._active_tasks += 1
terminal = False
terminal_settlement: Optional[bool] = None
state = False
err_msg = ""
def callback_after_terminal_settlement(
callback_task: TransferTask,
transferinfo: TransferInfo,
) -> Tuple[bool, str]:
"""成功结果先提交 durable 终态,再执行历史、事件和通知回调。"""
nonlocal terminal_settlement
if transferinfo.success and not callback_task.preview:
terminal_settlement = self.__discard_pending(callback_task)
if not terminal_settlement:
self.__fail_transfer_task(callback_task)
return False, "整理任务 durable 终态结算失去租约"
if item.callback:
callback_result: Tuple[bool, str] = item.callback(
callback_task,
transferinfo,
)
return callback_result
return transferinfo.success, transferinfo.message or ""
try:
self.__start_job_execution(task)
# 更新进度
@@ -2044,7 +2599,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
)
# 整理
state, err_msg = self.__handle_transfer(
task=task, callback=item.callback
task=task,
callback=callback_after_terminal_settlement,
)
terminal = task.plan_checkpoint is not None
@@ -2063,6 +2619,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
text=__process_msg,
)
except Exception as e:
if terminal_settlement is not None:
terminal = True
logger.error(
f"{fileitem.name} 整理任务处理出现错误:{e} - {traceback.format_exc()}"
)
@@ -2071,12 +2629,25 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self._processed_num += 1
self._fail_num += 1
finally:
self.__finish_job_execution(task, terminal=terminal)
self._queue.task_done()
with task_lock:
# 减少运行中的任务数
self._active_tasks -= 1
self.__settle_transfer_progress_if_idle()
durable_settled = False
try:
durable_settled = self.__finish_job_execution(
task,
terminal=terminal,
terminal_settlement=terminal_settlement,
)
except Exception as err:
logger.error(
f"整理任务终态结算异常:{task.admission_task_id} - {err}"
)
finally:
self._queue.task_done()
with task_lock:
# 减少运行中的任务数
self._active_tasks -= 1
if terminal and state and not durable_settled:
self._fail_num += 1
self.__settle_transfer_progress_if_idle()
except queue.Empty:
# 即使队列空了,如果还有任务在运行,也不应该结束进度
@@ -3548,6 +4119,10 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
and file_item.path == recovery_admission.src_path
):
transfer_task.bind_admission_task_id(recovery_admission.task_id)
self.__bind_claimed_admission(
transfer_task,
recovery_admission,
)
if recovery_admission.planning_input:
transfer_task.bind_planning_input(
recovery_admission.planning_input
@@ -3651,14 +4226,37 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
},
)
terminal = False
terminal_settlement: Optional[bool] = None
def callback_after_terminal_settlement(
callback_task: TransferTask,
transferinfo: TransferInfo,
) -> Tuple[bool, str]:
"""同步成功结果先结算 durable 行,再执行兼容回调副作用。"""
nonlocal terminal_settlement
if transferinfo.success and not callback_task.preview:
terminal_settlement = self.__discard_pending(callback_task)
if not terminal_settlement:
self.__fail_transfer_task(callback_task)
return False, "整理任务 durable 终态结算失去租约"
callback = (
_preview_callback
if preview
else self.__default_callback
)
return callback(callback_task, transferinfo)
try:
self.__claim_task_for_execution(transfer_task)
self.__start_job_execution(transfer_task)
state, err_msg = self.__handle_transfer(
task=transfer_task,
callback=_preview_callback if preview else self.__default_callback,
callback=callback_after_terminal_settlement,
)
terminal = bool(preview or transfer_task.plan_checkpoint is not None)
except Exception as e:
if terminal_settlement is not None:
terminal = True
logger.error(
f"{transfer_task.fileitem.name} 整理任务处理出现错误:"
f"{e} - {traceback.format_exc()}"
@@ -3667,10 +4265,14 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo
self.__fail_transfer_task(transfer_task)
state, err_msg = False, str(e)
finally:
self.__finish_job_execution(
durable_settled = self.__finish_job_execution(
transfer_task,
terminal=terminal,
terminal_settlement=terminal_settlement,
)
if terminal and not durable_settled:
state = False
err_msg = "整理任务 durable 终态结算失去租约"
if not state:
all_success = False
logger.warn(f"{transfer_task.fileitem.name} {err_msg}")
+290 -26
View File
@@ -1,7 +1,9 @@
"""整理任务持久准入端口的 SQLAlchemy 适配器。"""
import logging
from collections.abc import Callable
from datetime import datetime
from datetime import datetime, timedelta, timezone
from json import JSONDecodeError
from typing import Optional
from uuid import uuid4
@@ -14,6 +16,8 @@ from app.application.transfer import (
TRANSFER_ADMISSION_PROVIDER_PENDING,
TransferAdmission,
TransferAdmissionConflictError,
TransferAdmissionProjectionError,
TransferLeaseLostError,
TransferPlanCheckpoint,
TransferPlanningInput,
TransferPlanningStateError,
@@ -22,10 +26,14 @@ from app.db.models.transferpending import TransferPending
from app.db.oper.transferpending import TransferPendingOper
from app.db.uow import SqlAlchemyUnitOfWork
_diagnostic_logger = logging.getLogger(__name__)
class TransactionalTransferAdmissionRepository:
"""以短生命周期 Session 实现整理任务持久准入端口。"""
_MAX_RECOVERY_SCAN_TASKS = 5000
def __init__(self, session_factory: Callable[[], Session]) -> None:
"""保存由组合根提供的同步会话工厂。"""
self._session_factory = session_factory
@@ -35,6 +43,33 @@ class TransactionalTransferAdmissionRepository:
"""生成与历史登记时间可按字典序比较的当前时间。"""
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
@staticmethod
def _lease_now() -> datetime:
"""生成不受宿主时区影响的当前 UTC 租约时间。"""
return datetime.now(timezone.utc)
@staticmethod
def _format_lease_time(value: datetime) -> str:
"""把 UTC 时间编码为可稳定排序的固定宽度字符串。"""
return value.astimezone(timezone.utc).strftime("%Y-%m-%d %H:%M:%S.%f")
@staticmethod
def _validate_claim_arguments(*, owner_id: str, lease_seconds: int) -> None:
"""拒绝无法建立有效租约身份或正向期限的调用。"""
if not owner_id:
raise ValueError("整理任务 claim 缺少 owner_id")
if lease_seconds <= 0:
raise ValueError("整理任务 lease_seconds 必须大于零")
@staticmethod
def _recoverable_states() -> tuple[str, ...]:
"""返回允许被 worker claim 的稳定业务状态。"""
return (
TRANSFER_ADMISSION_ACCEPTED,
TRANSFER_ADMISSION_PROVIDER_PENDING,
TRANSFER_ADMISSION_PLANNED,
)
@staticmethod
def _project(pending: TransferPending) -> TransferAdmission:
"""在 Session 有效期内把 ORM 行冻结为应用层 DTO。"""
@@ -84,6 +119,11 @@ class TransactionalTransferAdmissionRepository:
input_fingerprint=pending.input_fingerprint,
planning_input=planning_input,
checkpoint=checkpoint,
lease_owner=pending.lease_owner,
lease_token=pending.lease_token,
lease_expires_at=pending.lease_expires_at,
heartbeat_at=pending.heartbeat_at,
attempt_count=pending.attempt_count,
)
@staticmethod
@@ -151,27 +191,225 @@ class TransactionalTransferAdmissionRepository:
self._assert_input_match(pending, effective_input)
return self._project(pending)
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
"""在独立只读会话中投影等待恢复或执行的准入记录。"""
def claim_task(
self,
*,
task_id: str,
owner_id: str,
lease_seconds: int,
) -> Optional[TransferAdmission]:
"""原子 claim 指定任务,任何已存在的有效租约都拒绝重复领取。"""
self._validate_claim_arguments(
owner_id=owner_id,
lease_seconds=lease_seconds,
)
if not task_id:
raise ValueError("整理任务 claim 缺少 task_id")
now = self._lease_now()
now_time = self._format_lease_time(now)
lease_expires_at = self._format_lease_time(
now + timedelta(seconds=lease_seconds)
)
with self._session_factory() as session:
pending_items = TransferPendingOper(db=session).list_by_state(
state=TRANSFER_ADMISSION_ACCEPTED,
limit=limit,
)
return [self._project(pending) for pending in pending_items]
transaction = SqlAlchemyUnitOfWork(session)
try:
oper = TransferPendingOper(db=session)
lease_token = uuid4().hex
updated = oper.stage_claim_task(
task_id=task_id,
states=self._recoverable_states(),
owner_id=owner_id,
lease_token=lease_token,
now_time=now_time,
lease_expires_at=lease_expires_at,
updated_at=self._now(),
)
session.flush()
session.expire_all()
try:
pending = oper.get_by_task_id(task_id=task_id)
except JSONDecodeError as error:
raise TransferAdmissionProjectionError(
f"整理任务持久 JSON 无法解码: {task_id} - {error}"
) from error
if updated:
if pending is None:
raise TransferPlanningStateError(
f"claim 后未找到整理任务: {task_id}"
)
try:
admission = self._project(pending)
except (
TransferAdmissionConflictError,
TransferPlanningStateError,
TypeError,
ValueError,
) as error:
raise TransferAdmissionProjectionError(
f"整理任务持久投影损坏: {task_id} - {error}"
) from error
transaction.commit()
return admission
transaction.commit()
return None
except Exception:
transaction.rollback()
raise
def list_recoverable(self, limit: int = 5000) -> list[TransferAdmission]:
"""投影接纳、provider 待执行或已规划的全部可恢复任务。"""
with self._session_factory() as session:
pending_items = TransferPendingOper(db=session).list_by_states(
states=(
TRANSFER_ADMISSION_ACCEPTED,
TRANSFER_ADMISSION_PROVIDER_PENDING,
TRANSFER_ADMISSION_PLANNED,
),
limit=limit,
def claim_recoverable(
self,
*,
owner_id: str,
limit: int,
lease_seconds: int,
) -> list[TransferAdmission]:
"""按登记顺序逐条 CAS claim 未租用或租约已过期的恢复任务。"""
self._validate_claim_arguments(
owner_id=owner_id,
lease_seconds=lease_seconds,
)
if limit <= 0:
return []
claimed: list[TransferAdmission] = []
after_cursor: Optional[tuple[str, int]] = None
scanned_count = 0
scan_limit = self._MAX_RECOVERY_SCAN_TASKS
while len(claimed) < limit and scanned_count < scan_limit:
candidate_limit = min(
limit - len(claimed),
scan_limit - scanned_count,
)
return [self._project(pending) for pending in pending_items]
now_time = self._format_lease_time(self._lease_now())
with self._session_factory() as session:
candidates = TransferPendingOper(db=session).list_claimable_candidates(
states=self._recoverable_states(),
now_time=now_time,
limit=candidate_limit,
after_cursor=after_cursor,
)
if not candidates:
break
scanned_count += len(candidates)
_, cursor_created_at, cursor_id = candidates[-1]
after_cursor = (cursor_created_at, cursor_id)
for task_id, _, _ in candidates:
try:
admission = self.claim_task(
task_id=task_id,
owner_id=owner_id,
lease_seconds=lease_seconds,
)
except TransferAdmissionProjectionError as error:
self._record_projection_failure(
task_id=task_id,
error=error,
)
continue
if admission is not None:
claimed.append(admission)
if len(claimed) >= limit:
break
return claimed
def _record_projection_failure(
self,
*,
task_id: str,
error: TransferAdmissionProjectionError,
) -> bool:
"""以独立 CAS 留存变化后的投影错误,并仅为新诊断记一次运行日志。"""
diagnostic = f"恢复投影失败: {error}"
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
recorded = TransferPendingOper(
db=session
).stage_record_projection_failure(
task_id=task_id,
states=self._recoverable_states(),
error=diagnostic,
now_time=self._format_lease_time(self._lease_now()),
updated_at=self._now(),
)
transaction.commit()
except Exception:
transaction.rollback()
raise
if recorded:
_diagnostic_logger.error(
f"整理恢复任务投影损坏:task_id={task_id}, error={error}"
)
return bool(recorded)
def heartbeat(
self,
*,
task_id: str,
lease_token: str,
lease_seconds: int,
) -> Optional[TransferAdmission]:
"""仅以当前且未过期的 token 续租,禁止陈旧 worker 复活租约。"""
if not task_id or not lease_token:
raise ValueError("整理任务 heartbeat 缺少任务或租约身份")
if lease_seconds <= 0:
raise ValueError("整理任务 lease_seconds 必须大于零")
now = self._lease_now()
now_time = self._format_lease_time(now)
lease_expires_at = self._format_lease_time(
now + timedelta(seconds=lease_seconds)
)
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
oper = TransferPendingOper(db=session)
updated = oper.stage_heartbeat(
task_id=task_id,
lease_token=lease_token,
now_time=now_time,
lease_expires_at=lease_expires_at,
)
if not updated:
transaction.commit()
return None
session.flush()
session.expire_all()
pending = oper.get_by_task_id(task_id=task_id)
if pending is None:
raise TransferPlanningStateError(
f"heartbeat 后未找到整理任务: {task_id}"
)
admission = self._project(pending)
transaction.commit()
return admission
except Exception:
transaction.rollback()
raise
def release_claim(
self,
*,
task_id: str,
lease_token: str,
error: Optional[str] = None,
) -> bool:
"""仅以当前未过期 token 释放租约,陈旧 worker 不得改变任务。"""
if not task_id or not lease_token:
return False
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
released = TransferPendingOper(db=session).stage_release_claim(
task_id=task_id,
lease_token=lease_token,
error=error,
now_time=self._format_lease_time(self._lease_now()),
updated_at=self._now(),
)
transaction.commit()
return bool(released)
except Exception:
transaction.rollback()
raise
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
"""独立提交最近一次入队失败,保留准入记录供后续恢复。"""
@@ -192,6 +430,7 @@ class TransactionalTransferAdmissionRepository:
self,
*,
task_id: str,
lease_token: str,
input_fingerprint: str,
checkpoint: TransferPlanCheckpoint,
) -> TransferAdmission:
@@ -223,7 +462,9 @@ class TransactionalTransferAdmissionRepository:
checkpoint_payload=checkpoint_payload,
source_states=source_states,
target_state=target_state,
now_time=self._now(),
lease_token=lease_token,
now_time=self._format_lease_time(self._lease_now()),
updated_at=self._now(),
)
session.flush()
session.expire_all()
@@ -232,6 +473,13 @@ class TransactionalTransferAdmissionRepository:
raise TransferPlanningStateError(f"未找到整理任务: {task_id}")
if pending.input_fingerprint != input_fingerprint:
raise TransferAdmissionConflictError("整理任务输入指纹已经改变")
now_time = self._format_lease_time(self._lease_now())
if (
pending.lease_token != lease_token
or not pending.lease_expires_at
or pending.lease_expires_at <= now_time
):
raise TransferLeaseLostError("整理任务租约已过期或已被其他 worker 接管")
if not updated and not (
pending.state == target_state
and pending.checkpoint_payload == checkpoint_payload
@@ -246,28 +494,44 @@ class TransactionalTransferAdmissionRepository:
transaction.rollback()
raise
def record_planning_failure(self, *, task_id: str, error: str) -> None:
def record_planning_failure(
self,
*,
task_id: str,
lease_token: str,
error: str,
) -> None:
"""独立提交规划错误并保持任务处于接纳态供恢复重试。"""
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
TransferPendingOper(db=session).stage_record_planning_failure(
updated = TransferPendingOper(db=session).stage_record_planning_failure(
task_id=task_id,
lease_token=lease_token,
error=error,
now_time=self._now(),
now_time=self._format_lease_time(self._lease_now()),
updated_at=self._now(),
)
if not updated:
raise TransferLeaseLostError(
"整理任务租约已过期或已被其他 worker 接管"
)
transaction.commit()
except Exception:
transaction.rollback()
raise
def discard_task(self, *, task_id: str) -> int:
"""在独立事务中按稳定任务标识删除已到终态的准入记录"""
def discard_claimed(self, *, task_id: str, lease_token: str) -> int:
"""仅以当前未过期 token 删除终态任务,拒绝陈旧 worker 变更"""
if not task_id or not lease_token:
return 0
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
deleted = TransferPendingOper(db=session).stage_discard_task(
deleted = TransferPendingOper(db=session).stage_discard_claimed(
task_id=task_id,
lease_token=lease_token,
now_time=self._format_lease_time(self._lease_now()),
)
transaction.commit()
return deleted
+338 -148
View File
@@ -4,7 +4,20 @@ from datetime import datetime
from typing import Any, List, Optional, cast
from uuid import uuid4
from sqlalchemy import JSON, Index, Integer, String, Text, UniqueConstraint, delete, select, update
from sqlalchemy import (
JSON,
Index,
Integer,
String,
Text,
UniqueConstraint,
and_,
delete,
func,
or_,
select,
update,
)
from sqlalchemy.orm import Mapped, Session, mapped_column
from app.db.base import Base, execute_dml, get_id_column
@@ -72,7 +85,7 @@ class TransferPending(Base):
准入时保存版本化规划输入和指纹;纯规划完成后以同一行原子保存完整有序计划并
推进到 planned。重启恢复可直接消费已规划路径,避免再次触发 rename 等插件事件。
旧路径登记接口仍生成最小 legacy_replan 输入,供插件兼容调用方继续使用
所有执行期 mutation 都以稳定任务身份和租约 token 进行 CAS fencing
"""
id = get_id_column()
@@ -111,6 +124,16 @@ class TransferPending(Base):
checkpoint_payload: Mapped[Optional[dict[str, Any]]] = mapped_column(JSON)
# 规划完成时间
planned_at: Mapped[Optional[str]] = mapped_column(String(40))
# 当前租约拥有者
lease_owner: Mapped[Optional[str]] = mapped_column(String(128))
# 当前租约的唯一防陈旧令牌
lease_token: Mapped[Optional[str]] = mapped_column(String(64))
# 当前租约的 UTC 到期时间
lease_expires_at: Mapped[Optional[str]] = mapped_column(String(40))
# 最近一次成功 claim 或 heartbeat 的 UTC 时间
heartbeat_at: Mapped[Optional[str]] = mapped_column(String(40))
# 真正取得新 token 的累计次数
attempt_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
__table_args__ = (
# 同一个文件重复入队只保留一条,回放时不会重复送入整理链
@@ -122,41 +145,17 @@ class TransferPending(Base):
"created_at",
"id",
),
# 恢复调度按业务状态和租约到期时间筛选可接管任务
Index(
"ix_transferpending_recovery_lease",
"state",
"lease_expires_at",
"created_at",
"id",
),
UniqueConstraint("task_id", name="uq_transferpending_task_id"),
)
@classmethod
def register(cls, db: Session, storage: str, src_path: str,
now_time: str) -> Optional["TransferPending"]:
"""
登记一个待整理文件,已存在时保持原登记时间不变。
:param db: 数据库会话
:param storage: 存储
:param src_path: 源文件路径
:param now_time: 当前时间
:return: 登记记录
"""
if not storage or not src_path:
return None
pending = db.execute(
select(cls).where(cls.storage == storage, cls.src_path == src_path)
).scalars().first()
if pending:
return cast("TransferPending", pending)
planning_input = _legacy_planning_payload(storage, src_path)
pending = cls(
storage=storage,
src_path=src_path,
state="accepted",
created_at=now_time,
updated_at=now_time,
input_version=1,
planning_input=planning_input,
input_fingerprint=_planning_fingerprint(planning_input),
)
db.add(pending)
return pending
@classmethod
def stage_admit(cls, db: Session, *, task_id: str, storage: str,
src_path: str, state: str,
@@ -201,44 +200,6 @@ class TransferPending(Base):
db.add(pending)
return pending
@classmethod
def list_by_state(cls, db: Session, *, state: str,
limit: Optional[int] = 5000) -> List["TransferPending"]:
"""
按登记顺序列出指定持久状态的接纳记录。
:param db: 数据库会话
:param state: 持久状态
:param limit: 单次读取上限
:return: 接纳记录列表
"""
if not state:
return []
return list(db.execute(
select(cls)
.where(cls.state == state)
.order_by(cls.created_at.asc(), cls.id.asc())
.limit(limit)
).scalars().all())
@classmethod
def list_by_states(cls, db: Session, *, states: tuple[str, ...],
limit: Optional[int] = 5000) -> List["TransferPending"]:
"""
按登记顺序列出多个可恢复持久状态的记录。
:param db: 数据库会话
:param states: 允许恢复的状态集合
:param limit: 单次读取上限
:return: 接纳记录列表
"""
if not states:
return []
return list(db.execute(
select(cls)
.where(cls.state.in_(states))
.order_by(cls.created_at.asc(), cls.id.asc())
.limit(limit)
).scalars().all())
@classmethod
def get_by_identity(cls, db: Session, *, storage: str,
src_path: str) -> Optional["TransferPending"]:
@@ -276,12 +237,283 @@ class TransferPending(Base):
db.execute(select(cls).where(cls.task_id == task_id)).scalars().first(),
)
@classmethod
def list_claimable_candidates(
cls,
db: Session,
*,
states: tuple[str, ...],
now_time: str,
limit: int,
after_cursor: Optional[tuple[str, int]] = None,
) -> List[tuple[str, str, int]]:
"""
按稳定游标列出未租用或租约已过期的候选任务。
返回候选不等于取得租约;调用方必须继续执行带相同过期条件的 claim CAS,
并以受影响行数决定竞争结果。
:param db: 数据库会话
:param states: 可恢复业务状态
:param now_time: 当前 UTC 时间
:param limit: 候选数量上限
:param after_cursor: 上一页最后一条的规范登记时间与主键
:return: 任务标识、规范登记时间与主键组成的稳定游标列表
"""
if not states or not now_time or limit <= 0:
return []
cursor_created_at = func.coalesce(cls.created_at, "")
statement = select(cls.task_id, cursor_created_at, cls.id).where(
cls.state.in_(states),
or_(
cls.lease_token.is_(None),
cls.lease_expires_at.is_(None),
cls.lease_expires_at <= now_time,
),
)
if after_cursor is not None:
after_created_at, after_id = after_cursor
statement = statement.where(or_(
cursor_created_at > after_created_at,
and_(
cursor_created_at == after_created_at,
cls.id > after_id,
),
))
rows = db.execute(
statement
.order_by(cursor_created_at.asc(), cls.id.asc())
.limit(limit)
).all()
return [
(task_id, created_at or "", int(row_id))
for task_id, created_at, row_id in rows
]
@classmethod
def claim_task(
cls,
db: Session,
*,
task_id: str,
states: tuple[str, ...],
owner_id: str,
lease_token: str,
now_time: str,
lease_expires_at: str,
updated_at: str,
) -> int:
"""
以未租用或租约已过期为条件原子取得任务租约。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param states: 允许 claim 的业务状态
:param owner_id: 新租约拥有者
:param lease_token: 新租约唯一令牌
:param now_time: 当前 UTC 时间
:param lease_expires_at: 新租约到期时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数,1 表示赢得竞争
"""
if not all((
task_id,
states,
owner_id,
lease_token,
now_time,
lease_expires_at,
updated_at,
)):
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.state.in_(states),
or_(
cls.lease_token.is_(None),
cls.lease_expires_at.is_(None),
cls.lease_expires_at <= now_time,
),
)
.values(
lease_owner=owner_id,
lease_token=lease_token,
lease_expires_at=lease_expires_at,
heartbeat_at=now_time,
attempt_count=cls.attempt_count + 1,
updated_at=updated_at,
),
execution_options={"synchronize_session": False},
)
@classmethod
def record_projection_failure(
cls,
db: Session,
*,
task_id: str,
states: tuple[str, ...],
error: str,
now_time: str,
updated_at: str,
) -> int:
"""
在没有有效租约且诊断发生变化时原子记录恢复投影损坏。
claim 的投影失败会先回滚,因此这里不得重新占用租约。CAS 同时保护
已被其他 worker 领取的任务,并避免周期恢复反复刷新相同错误。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param states: 可恢复业务状态
:param error: 可持久化的稳定诊断文本
:param now_time: 当前 UTC 租约时间
:param updated_at: 宿主本地业务审计时间
:return: 更新的记录数,1 表示首次或变化后的诊断被记录
"""
if not all((task_id, states, error, now_time, updated_at)):
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.state.in_(states),
or_(
cls.lease_token.is_(None),
cls.lease_expires_at.is_(None),
cls.lease_expires_at <= now_time,
),
cls.last_error.is_distinct_from(error),
)
.values(
last_error=error,
updated_at=updated_at,
),
execution_options={"synchronize_session": False},
)
@classmethod
def heartbeat(
cls,
db: Session,
*,
task_id: str,
lease_token: str,
now_time: str,
lease_expires_at: str,
) -> int:
"""
仅以当前且未过期的 token 原子延长任务租约。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param now_time: 当前 UTC 时间
:param lease_expires_at: 新租约到期时间
:return: 更新的记录数
"""
if not all((task_id, lease_token, now_time, lease_expires_at)):
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.lease_token == lease_token,
cls.lease_expires_at.is_not(None),
cls.lease_expires_at > now_time,
)
.values(
lease_expires_at=lease_expires_at,
heartbeat_at=now_time,
),
execution_options={"synchronize_session": False},
)
@classmethod
def release_claim(
cls,
db: Session,
*,
task_id: str,
lease_token: str,
error: Optional[str],
now_time: str,
updated_at: str,
) -> int:
"""
仅以当前且未过期的 token 释放租约并保存本次执行错误。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param error: 本次执行错误,成功释放时为空
:param now_time: 当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
if not task_id or not lease_token or not now_time or not updated_at:
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.lease_token == lease_token,
cls.lease_expires_at.is_not(None),
cls.lease_expires_at > now_time,
)
.values(
lease_owner=None,
lease_token=None,
lease_expires_at=None,
heartbeat_at=None,
last_error=error,
updated_at=updated_at,
),
execution_options={"synchronize_session": False},
)
@classmethod
def discard_claimed(
cls,
db: Session,
*,
task_id: str,
lease_token: str,
now_time: str,
) -> int:
"""
仅以当前且未过期的 token 删除已经到达终态的租约任务。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param now_time: 当前 UTC 时间
:return: 删除的记录数
"""
if not task_id or not lease_token or not now_time:
return 0
return execute_dml(
db,
delete(cls).where(
cls.task_id == task_id,
cls.lease_token == lease_token,
cls.lease_expires_at.is_not(None),
cls.lease_expires_at > now_time,
),
execution_options={"synchronize_session": False},
)
@classmethod
def checkpoint_plan(cls, db: Session, *, task_id: str,
input_fingerprint: str, checkpoint_version: int,
checkpoint_payload: dict[str, Any],
source_states: tuple[str, ...], target_state: str,
now_time: str) -> int:
lease_token: str, now_time: str,
updated_at: str) -> int:
"""
以输入指纹为 CAS 条件原子保存计划并推进到已规划。
:param db: 数据库会话
@@ -291,7 +523,9 @@ class TransferPending(Base):
:param checkpoint_payload: 完整有序计划 JSON
:param source_states: 允许推进检查点的起始状态
:param target_state: 检查点提交后的目标状态
:param now_time: 当前时间
:param lease_token: 当前且未过期的租约令牌
:param now_time: 用于租约 fencing 的当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
if (
@@ -300,8 +534,19 @@ class TransferPending(Base):
or not checkpoint_payload
or not source_states
or not target_state
or not lease_token
or not updated_at
):
return 0
values: dict[str, Any] = {
"state": target_state,
"checkpoint_version": checkpoint_version,
"checkpoint_payload": checkpoint_payload,
"last_error": None,
"updated_at": updated_at,
}
if target_state == "planned":
values["planned_at"] = updated_at
return execute_dml(
db,
update(cls)
@@ -309,30 +554,29 @@ class TransferPending(Base):
cls.task_id == task_id,
cls.state.in_(source_states),
cls.input_fingerprint == input_fingerprint,
cls.lease_token == lease_token,
cls.lease_expires_at.is_not(None),
cls.lease_expires_at > now_time,
)
.values(
state=target_state,
checkpoint_version=checkpoint_version,
checkpoint_payload=checkpoint_payload,
planned_at=now_time,
last_error=None,
updated_at=now_time,
),
.values(**values),
execution_options={"synchronize_session": False},
)
@classmethod
def record_planning_failure(cls, db: Session, *, task_id: str,
error: str, now_time: str) -> int:
lease_token: str, error: str,
now_time: str, updated_at: str) -> int:
"""
为接纳态或 provider 待执行任务记录规划失败,不改变其恢复状态。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param lease_token: 当前且未过期的租约令牌
:param error: 失败原因
:param now_time: 当前时间
:param now_time: 用于租约 fencing 的当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
if not task_id:
if not task_id or not lease_token or not now_time or not updated_at:
return 0
return execute_dml(
db,
@@ -340,8 +584,11 @@ class TransferPending(Base):
.where(
cls.task_id == task_id,
cls.state.in_(("accepted", "provider_pending")),
cls.lease_token == lease_token,
cls.lease_expires_at.is_not(None),
cls.lease_expires_at > now_time,
)
.values(last_error=error, updated_at=now_time),
.values(last_error=error, updated_at=updated_at),
execution_options={"synchronize_session": False},
)
@@ -361,67 +608,10 @@ class TransferPending(Base):
return execute_dml(
db,
update(cls)
.where(cls.task_id == task_id)
.where(
cls.task_id == task_id,
cls.lease_token.is_(None),
)
.values(last_error=error, updated_at=now_time),
execution_options={"synchronize_session": False},
)
@classmethod
def discard_task(cls, db: Session, *, task_id: str) -> int:
"""
在调用方会话中按任务标识删除接纳记录。
:param db: 数据库会话
:param task_id: 任务标识
:return: 删除的记录数
"""
if not task_id:
return 0
return execute_dml(
db, delete(cls).where(cls.task_id == task_id),
execution_options={"synchronize_session": False},
)
@classmethod
def discard(cls, db: Session, storage: str, src_path: str) -> int:
"""
注销一个待整理文件登记,整理到达终态(成功或失败)时调用。
:param db: 数据库会话
:param storage: 存储
:param src_path: 源文件路径
:return: 删除的记录数
"""
if not storage or not src_path:
return 0
return execute_dml(
db, delete(cls).where(cls.storage == storage, cls.src_path == src_path),
execution_options={"synchronize_session": False},
)
@classmethod
def list_all(cls, db: Session, limit: Optional[int] = 5000) -> List["TransferPending"]:
"""
列出全部待整理登记,供启动回放使用。
按登记时间升序回放,保持与原入队顺序一致;上限避免异常积压时
一次性把整理链压垮。
:param db: 数据库会话
:param limit: 单次回放上限
:return: 待整理登记列表
"""
return list(db.execute(
select(cls)
.order_by(cls.created_at.asc(), cls.id.asc())
.limit(limit)
).scalars().all())
@classmethod
def clear(cls, db: Session) -> int:
"""
清空全部待整理登记。
:param db: 数据库会话
:return: 删除的记录数
"""
return execute_dml(
db, delete(cls),
execution_options={"synchronize_session": False},
execution_options={"synchronize_session": "fetch"},
)
+191 -110
View File
@@ -1,5 +1,4 @@
from datetime import datetime
from typing import Any, List, Optional, Tuple
from typing import Any, Optional
from app.db.base import DbOper
from app.db.models.transferpending import TransferPending
@@ -10,27 +9,9 @@ class TransferPendingOper(DbOper):
待整理文件登记管理
保存稳定任务身份存储源文件路径和准入状态用于在进程重启后把没走完
整理链的文件重新送回去避免挂载故障重启后永久漏件旧版路径登记接口继续
保留供插件和兼容调用方使用
整理链的文件重新送回去避免挂载故障重启后永久漏件
"""
def register(self, storage: str, src_path: str) -> Optional[TransferPending]:
"""
登记一个待整理文件
:param storage: 存储
:param src_path: 源文件路径
:return: 登记记录
"""
now_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
return self._execute_sync_write(
lambda session: TransferPending.register(
session,
storage=storage,
src_path=src_path,
now_time=now_time,
)
)
def stage_admit(self, *, task_id: str, storage: str, src_path: str,
state: str, now_time: str, input_version: int = 1,
planning_input: Optional[dict[str, Any]] = None,
@@ -63,38 +44,6 @@ class TransferPendingOper(DbOper):
)
)
def list_by_state(self, *, state: str,
limit: Optional[int] = 5000) -> List[TransferPending]:
"""
使用当前会话列出指定状态记录
:param state: 持久状态
:param limit: 单次读取上限
:return: ORM 接纳记录列表
"""
return self._execute_sync_query(
lambda session: TransferPending.list_by_state(
session,
state=state,
limit=limit,
)
) or []
def list_by_states(self, *, states: tuple[str, ...],
limit: Optional[int] = 5000) -> List[TransferPending]:
"""
使用当前会话列出多个可恢复状态的记录
:param states: 允许恢复的状态集合
:param limit: 单次读取上限
:return: ORM 接纳记录列表
"""
return self._execute_sync_query(
lambda session: TransferPending.list_by_states(
session,
states=states,
limit=limit,
)
) or []
def get_by_identity(self, *, storage: str,
src_path: str) -> Optional[TransferPending]:
"""
@@ -133,7 +82,9 @@ class TransferPendingOper(DbOper):
checkpoint_payload: dict[str, Any],
source_states: tuple[str, ...],
target_state: str,
lease_token: str,
now_time: str,
updated_at: str,
) -> int:
"""
在当前会话中以输入指纹为条件暂存完整计划检查点
@@ -143,7 +94,9 @@ class TransferPendingOper(DbOper):
:param checkpoint_payload: 完整有序计划 JSON
:param source_states: 允许执行 CAS 的起始状态
:param target_state: 检查点提交后的目标状态
:param now_time: 当前时间
:param lease_token: 当前且未过期的租约令牌
:param now_time: 用于租约 fencing 的当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
return self._execute_sync_write(
@@ -155,25 +108,206 @@ class TransferPendingOper(DbOper):
checkpoint_payload=checkpoint_payload,
source_states=source_states,
target_state=target_state,
lease_token=lease_token,
now_time=now_time,
updated_at=updated_at,
)
)
def stage_record_planning_failure(self, *, task_id: str, error: str,
now_time: str) -> int:
def stage_record_planning_failure(self, *, task_id: str, lease_token: str,
error: str, now_time: str,
updated_at: str) -> int:
"""
在当前会话中记录规划失败并保持任务处于接纳态
:param task_id: 稳定任务标识
:param lease_token: 当前且未过期的租约令牌
:param error: 失败原因
:param now_time: 当前时间
:param now_time: 用于租约 fencing 的当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.record_planning_failure(
session,
task_id=task_id,
lease_token=lease_token,
error=error,
now_time=now_time,
updated_at=updated_at,
)
)
def list_claimable_candidates(
self,
*,
states: tuple[str, ...],
now_time: str,
limit: int,
after_cursor: Optional[tuple[str, int]] = None,
) -> list[tuple[str, str, int]]:
"""
使用当前会话按稳定游标读取未租用或已过期的恢复候选
:param states: 可恢复业务状态
:param now_time: 当前 UTC 时间
:param limit: 候选数量上限
:param after_cursor: 上一页最后一条的规范登记时间与主键
:return: 任务标识规范登记时间与主键组成的稳定游标列表
"""
return self._execute_sync_query(
lambda session: TransferPending.list_claimable_candidates(
session,
states=states,
now_time=now_time,
limit=limit,
after_cursor=after_cursor,
)
) or []
def stage_claim_task(
self,
*,
task_id: str,
states: tuple[str, ...],
owner_id: str,
lease_token: str,
now_time: str,
lease_expires_at: str,
updated_at: str,
) -> int:
"""
使用当前会话以未租或过期条件竞争一个新租约
:param task_id: 稳定任务标识
:param states: 可恢复业务状态
:param owner_id: 新租约拥有者
:param lease_token: 新租约唯一令牌
:param now_time: 当前 UTC 时间
:param lease_expires_at: 新租约到期时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.claim_task(
session,
task_id=task_id,
states=states,
owner_id=owner_id,
lease_token=lease_token,
now_time=now_time,
lease_expires_at=lease_expires_at,
updated_at=updated_at,
)
)
def stage_record_projection_failure(
self,
*,
task_id: str,
states: tuple[str, ...],
error: str,
now_time: str,
updated_at: str,
) -> int:
"""
使用当前会话按无有效租约和诊断变化条件记录投影损坏
:param task_id: 稳定任务标识
:param states: 可恢复业务状态
:param error: 可持久化的稳定诊断文本
:param now_time: 当前 UTC 租约时间
:param updated_at: 宿主本地业务审计时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.record_projection_failure(
session,
task_id=task_id,
states=states,
error=error,
now_time=now_time,
updated_at=updated_at,
)
)
def stage_heartbeat(
self,
*,
task_id: str,
lease_token: str,
now_time: str,
lease_expires_at: str,
) -> int:
"""
使用当前会话以当前未过期 token 延长租约
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param now_time: 当前 UTC 时间
:param lease_expires_at: 新租约到期时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.heartbeat(
session,
task_id=task_id,
lease_token=lease_token,
now_time=now_time,
lease_expires_at=lease_expires_at,
)
)
def stage_release_claim(
self,
*,
task_id: str,
lease_token: str,
error: Optional[str],
now_time: str,
updated_at: str,
) -> int:
"""
使用当前会话按未过期 token 释放租约并记录本次错误
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param error: 本次执行错误成功释放时为空
:param now_time: 当前 UTC 时间
:param updated_at: 与既有业务审计字段一致的宿主本地时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.release_claim(
session,
task_id=task_id,
lease_token=lease_token,
error=error,
now_time=now_time,
updated_at=updated_at,
)
)
def stage_discard_claimed(
self,
*,
task_id: str,
lease_token: str,
now_time: str,
) -> int:
"""
使用当前会话按当前未过期 token 删除终态任务
:param task_id: 稳定任务标识
:param lease_token: 当前租约令牌
:param now_time: 当前 UTC 时间
:return: 删除的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.discard_claimed(
session,
task_id=task_id,
lease_token=lease_token,
now_time=now_time,
)
)
@@ -194,56 +328,3 @@ class TransferPendingOper(DbOper):
now_time=now_time,
)
)
def stage_discard_task(self, *, task_id: str) -> int:
"""
在当前会话中暂存按任务标识删除接纳记录
:param task_id: 任务标识
:return: 删除的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.discard_task(
session,
task_id=task_id,
)
)
def discard(self, storage: str, src_path: str) -> int:
"""
注销一个待整理文件登记
:param storage: 存储
:param src_path: 源文件路径
:return: 删除的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.discard(
session,
storage=storage,
src_path=src_path,
)
)
def list_all(self, limit: Optional[int] = 5000) -> List[Tuple[str, str]]:
"""
列出全部待整理登记供启动回放使用
返回纯元组而不是 ORM 实例回放发生在会话之外ORM 实例脱离 session
后访问属性会触发 DetachedInstanceError
:param limit: 单次回放上限
:return: (存储, 源文件路径) 列表
"""
items = self._execute_sync_query(
lambda session: TransferPending.list_all(session, limit=limit)
)
return [
(item.storage, item.src_path)
for item in items or []
if item and item.storage and item.src_path
]
def clear(self) -> int:
"""
清空全部待整理登记
:return: 删除的记录数
"""
return self._execute_sync_write(TransferPending.clear)
+3 -3
View File
@@ -135,10 +135,10 @@ MODULE_ALIASES: Dict[str, ModuleAlias] = {
owner="sdk",
),
"app.db.transferpending_oper": ModuleAlias(
target="app.db.oper.transferpending",
replacement="app.db.oper.transferpending",
target="app.sdk._legacy.transferpending",
replacement="app.application.transfer",
introduced="v3.0.0",
owner="db",
owner="sdk",
),
"app.db.user_oper": ModuleAlias(
target="app.sdk._legacy.user",
+185
View File
@@ -0,0 +1,185 @@
"""兼容旧 ``app.db.transferpending_oper`` 的无 Session 数据访问接口。"""
from datetime import datetime
from typing import List, Optional, Tuple
from uuid import uuid4
from sqlalchemy import delete, select
from app.db.base import DbOper, execute_dml
from app.db.models.transferpending import TransferPending as _TransferPending
class TransferPendingOper(DbOper):
"""
保留旧待整理登记 ABI并对执行租约实施兼容写 fencing
本类只由精确旧导入映射加载宿主整理链仍使用 Application Port 和显式
Session canonical Oper旧删除入口只能处理从未取得租约的记录
"""
def register(self, storage: str, src_path: str) -> Optional[_TransferPending]:
"""
按旧签名登记待整理文件重复登记保持原任务和租约不变
:param storage: 存储
:param src_path: 源文件路径
:return: 登记记录
"""
now_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
return self._execute_sync_write(
lambda session: _TransferPending.stage_admit(
session,
task_id=uuid4().hex,
storage=storage,
src_path=src_path,
state="accepted",
now_time=now_time,
)
)
def list_by_state(
self,
*,
state: str,
limit: Optional[int] = 5000,
) -> List[_TransferPending]:
"""
按旧签名和登记顺序列出指定状态记录
:param state: 持久状态
:param limit: 单次读取上限
:return: ORM 接纳记录列表
"""
if not state:
return []
return self._execute_sync_query(
lambda session: list(session.execute(
select(_TransferPending)
.where(_TransferPending.state == state)
.order_by(_TransferPending.created_at.asc(), _TransferPending.id.asc())
.limit(limit)
).scalars().all())
)
def list_by_states(
self,
*,
states: tuple[str, ...],
limit: Optional[int] = 5000,
) -> List[_TransferPending]:
"""
按旧签名和登记顺序列出多个状态记录
:param states: 持久状态集合
:param limit: 单次读取上限
:return: ORM 接纳记录列表
"""
if not states:
return []
return self._execute_sync_query(
lambda session: list(session.execute(
select(_TransferPending)
.where(_TransferPending.state.in_(states))
.order_by(_TransferPending.created_at.asc(), _TransferPending.id.asc())
.limit(limit)
).scalars().all())
)
def get_by_identity(
self,
*,
storage: str,
src_path: str,
) -> Optional[_TransferPending]:
"""
按旧签名查询指定存储与源路径的登记
:param storage: 存储
:param src_path: 源文件路径
:return: 登记记录
"""
return self._execute_sync_query(
lambda session: _TransferPending.get_by_identity(
session,
storage=storage,
src_path=src_path,
)
)
def get_by_task_id(self, *, task_id: str) -> Optional[_TransferPending]:
"""
按旧签名查询稳定任务标识对应的登记
:param task_id: 稳定任务标识
:return: 登记记录
"""
return self._execute_sync_query(
lambda session: _TransferPending.get_by_task_id(
session,
task_id=task_id,
)
)
def discard(self, storage: str, src_path: str) -> int:
"""
按旧签名删除未 claim 的指定登记
任何带 token 的记录都由当前租约拥有者通过 fenced canonical API 收口
即使租约已经过期旧插件也不得越权代替恢复调度器删除
:param storage: 存储
:param src_path: 源文件路径
:return: 删除的记录数
"""
if not storage or not src_path:
return 0
return self._execute_sync_write(
lambda session: execute_dml(
session,
delete(_TransferPending).where(
_TransferPending.storage == storage,
_TransferPending.src_path == src_path,
_TransferPending.lease_token.is_(None),
),
execution_options={"synchronize_session": False},
)
)
def list_all(self, limit: Optional[int] = 5000) -> List[Tuple[str, str]]:
"""
按旧返回形态列出全部待整理路径
:param limit: 单次读取上限
:return: ``(存储, 源文件路径)`` 列表
"""
rows = self._execute_sync_query(
lambda session: session.execute(
select(_TransferPending.storage, _TransferPending.src_path)
.order_by(_TransferPending.created_at.asc(), _TransferPending.id.asc())
.limit(limit)
).all()
)
return [
(storage, src_path)
for storage, src_path in rows
if storage and src_path
]
def clear(self) -> int:
"""
清空全部未 claim 登记保留任何带租约 token 的任务
:return: 删除的记录数
"""
return self._execute_sync_write(
lambda session: execute_dml(
session,
delete(_TransferPending).where(
_TransferPending.lease_token.is_(None),
),
execution_options={"synchronize_session": False},
)
)
__all__ = ["TransferPendingOper"]
+2 -2
View File
@@ -3,12 +3,12 @@ from app.chain.transfer import TransferChain
def replay_pending_transfers():
"""
回放上次进程退出时仍未整理完的文件
启动唯一整理恢复调度器
整理队列是纯内存的挂载挂死后的人工重启版本升级OOM宿主重启都会让
队列连同这些文件还没整理这个事实一起蒸发而已稳定落地的文件不会再产生
任何监控事件也不会有新的补偿扫描起点结果就是永久漏件
回放本身在后台线程执行不阻塞启动流程
启动回放同进程补偿和过期租约接管共用同一个后台调度入口不阻塞启动流程
"""
TransferChain().replay_pending()