RICS / worker.py
StormShadow308's picture
feat: async pipeline, job queue, generation hardening, and docs
732b14f
Raw
History Blame Contribute Delete
1.85 kB
"""Temporal worker entry point for report generation workflows."""
from __future__ import annotations
import asyncio
import logging
from app.config import settings
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
async def main() -> None:
try:
from temporalio.client import Client
from temporalio.worker import Worker
except ImportError as exc: # pragma: no cover
raise SystemExit(
"temporalio is required for the worker. Install: pip install 'report-genius-ai[temporal]'"
) from exc
from app.db.database import init_db
from app.services.generation_stale import generation_stale_sweeper_loop
await init_db()
if int(getattr(settings, "generation_stale_sweep_seconds", 120)) > 0:
asyncio.create_task(generation_stale_sweeper_loop())
from app.workflows.report_activities import (
assemble_report,
execute_report_generation,
fetch_sources,
generate_section_activity,
retrieve_context,
validate_output,
)
from app.workflows.report_workflow import ReportGenerationWorkflow
client = await Client.connect(
settings.temporal_host,
namespace=settings.temporal_namespace,
)
worker = Worker(
client,
task_queue=settings.temporal_task_queue,
workflows=[ReportGenerationWorkflow],
activities=[
fetch_sources,
retrieve_context,
generate_section_activity,
execute_report_generation,
validate_output,
assemble_report,
],
)
logger.info(
"Temporal worker started queue=%s host=%s",
settings.temporal_task_queue,
settings.temporal_host,
)
await worker.run()
if __name__ == "__main__":
asyncio.run(main())