Files
ARR-2.0-0918/arr_web/processing_runtime.py

102 lines
3.3 KiB
Python

"""Production composition for ARR-owned programmatic XML processing."""
from __future__ import annotations
import os
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable, Optional
from arr_ingestion.postgres import DatabaseConfig, PostgresIngestionRepository
from arr_ingestion.service import IngestionService
from arr_ingestion.validation import DeliveryValidator
from arr_processing.local import LocalDailyProcessor
from arr_processing.policy import ProcessorPolicy, load_processor_policy
from arr_storage.aliyun_oss_v2 import AliyunOssConfig, AliyunOssV2Client
from arr_storage.contracts import ObjectKeyPolicy
from arr_storage.remote import CloudObjectBackend
from arr_storage.store import ManagedObjectStore
from arr_web.programmatic import ProgrammaticUploadCoordinator
@dataclass
class ProcessingInputRuntime:
coordinator: ProgrammaticUploadCoordinator
oss_client: AliyunOssV2Client
object_store: ManagedObjectStore
repository: PostgresIngestionRepository
ingestion: IngestionService
processor: LocalDailyProcessor
policy: ProcessorPolicy
def close(self) -> None:
self.oss_client.close()
@dataclass
class ObjectStoreRuntime:
"""Shared OSS/object-store runtime for report publishers and downloads."""
oss_client: AliyunOssV2Client
object_store: ManagedObjectStore
def close(self) -> None:
self.oss_client.close()
def compose_object_store() -> ObjectStoreRuntime:
oss_client = AliyunOssV2Client(AliyunOssConfig.from_environment())
try:
oss_client.assert_immutable_writes_supported()
object_store = ManagedObjectStore(
CloudObjectBackend(oss_client),
ObjectKeyPolicy(os.environ.get("ARR_OBJECT_PREFIX", "arr")),
)
return ObjectStoreRuntime(oss_client, object_store)
except Exception:
try:
oss_client.close()
except Exception:
pass
raise
def compose_programmatic_processing(
*,
project_root: Path,
connect: Optional[Callable[[str], Any]] = None,
) -> ProcessingInputRuntime:
storage = compose_object_store()
try:
oss_client = storage.oss_client
object_store = storage.object_store
database_config = (
DatabaseConfig("controlled")
if connect is not None
else DatabaseConfig.from_environment()
)
repository = PostgresIngestionRepository(
database_config,
connect=connect,
)
policy = load_processor_policy(project_root)
ingestion = IngestionService(DeliveryValidator(object_store, policy), repository)
processor = LocalDailyProcessor(policy)
coordinator = ProgrammaticUploadCoordinator(
object_store=object_store,
ingestion_repository=repository,
ingestion_service=ingestion,
processor=processor,
processor_version=policy.processor_version,
rule_set_sha256=policy.rule_set_sha256,
)
return ProcessingInputRuntime(coordinator, oss_client, object_store,
repository, ingestion, processor, policy)
except Exception:
storage.close()
raise
# Compatibility name for older local launch wrappers.
compose_oss_processing_input = compose_programmatic_processing