generated from john/python-template
+12
-11
@@ -56,18 +56,19 @@ async def _lifespan(app: FastAPI):
|
||||
|
||||
async with AsyncExitStack() as stack:
|
||||
stack.push_async_callback(dispose_database_runtime)
|
||||
stop_event, worker_notifier, worker_health = await stack.enter_async_context(
|
||||
worker_consumer_lifespan(
|
||||
session_factory=app.state.runtime.session_factory,
|
||||
poll_interval_seconds=settings.worker_poll_interval_seconds,
|
||||
shutdown_timeout_seconds=(
|
||||
settings.worker_provider_timeout_seconds + settings.worker_shutdown_grace_seconds
|
||||
),
|
||||
if settings.run_embedded_worker:
|
||||
stop_event, worker_notifier, worker_health = await stack.enter_async_context(
|
||||
worker_consumer_lifespan(
|
||||
session_factory=app.state.runtime.session_factory,
|
||||
poll_interval_seconds=settings.worker_poll_interval_seconds,
|
||||
shutdown_timeout_seconds=(
|
||||
settings.worker_provider_timeout_seconds + settings.worker_shutdown_grace_seconds
|
||||
),
|
||||
)
|
||||
)
|
||||
)
|
||||
app.state.worker_stop_event = stop_event
|
||||
app.state.worker_notifier = worker_notifier
|
||||
app.state.worker_health = worker_health
|
||||
app.state.worker_stop_event = stop_event
|
||||
app.state.worker_notifier = worker_notifier
|
||||
app.state.worker_health = worker_health
|
||||
yield
|
||||
|
||||
|
||||
|
||||
@@ -98,6 +98,7 @@ class Settings(BaseSettings):
|
||||
# --- runtime environment ---
|
||||
environment: Literal["development", "test", "production"] = "development"
|
||||
transcription_commit: NonEmptyStr | None = None
|
||||
run_embedded_worker: bool = True
|
||||
|
||||
# --- persistence ---
|
||||
database: DatabaseSettings = Field(default_factory=SqliteSettings)
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
"""Standalone worker process entrypoint for production deployments."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
from .config import configure_logging
|
||||
from .config import parse_cli_settings
|
||||
from .db import create_all
|
||||
from .db import dispose_database_runtime
|
||||
from .db import initialize_database_runtime
|
||||
from .worker import run_worker_loop
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def _run() -> None:
|
||||
settings = parse_cli_settings()
|
||||
configure_logging(settings)
|
||||
runtime = initialize_database_runtime(settings=settings)
|
||||
|
||||
if settings.should_bootstrap_schema:
|
||||
await create_all(engine=runtime.engine)
|
||||
|
||||
try:
|
||||
await run_worker_loop(
|
||||
session_factory=runtime.session_factory,
|
||||
poll_interval_seconds=settings.worker_poll_interval_seconds,
|
||||
)
|
||||
finally:
|
||||
await dispose_database_runtime()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
try:
|
||||
asyncio.run(_run())
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Worker service received shutdown signal")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user