merge: preserve existing business features on new foundation
This commit is contained in:
@@ -4,6 +4,7 @@ import logging
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.infrastructure.db import engine
|
||||
from app.worker.offsite_mail_worker import OffsiteMailWorker
|
||||
from app.worker.runtime import WorkerRuntime
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -12,10 +13,12 @@ logger = logging.getLogger(__name__)
|
||||
async def serve(*, once: bool = False) -> None:
|
||||
settings = get_settings()
|
||||
runtime = WorkerRuntime(settings=settings)
|
||||
offsite_worker = OffsiteMailWorker(settings)
|
||||
try:
|
||||
while True:
|
||||
try:
|
||||
worked = await runtime.run_once()
|
||||
worked = await offsite_worker.run_once() or worked
|
||||
except Exception:
|
||||
# 常驻 Worker 不能因为"某一轮"的异常就整体退出:数据库抖动、
|
||||
# 迁移期间锁表、外部依赖瞬断都会命中这里,而 `run_once` 里的
|
||||
@@ -31,6 +34,7 @@ async def serve(*, once: bool = False) -> None:
|
||||
if not worked:
|
||||
await asyncio.sleep(settings.worker_poll_seconds)
|
||||
finally:
|
||||
await offsite_worker.close()
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user