from __future__ import annotations from dataclasses import dataclass from typing import Any, Protocol @dataclass(frozen=True) class ProjectionJob: projection_id: str external_document_id: str content: str metadata: dict[str, Any] track_id: str | None = None class ProjectionRepository(Protocol): def mark( self, projection_id: str, status: str, error: str | None = None ) -> None: ... class ProjectionWorker: def __init__(self, *, client, repository: ProjectionRepository) -> None: self._client = client self._repository = repository def project(self, job: ProjectionJob) -> str: self._repository.mark(job.projection_id, "processing") try: track_id = job.track_id if track_id is None: receipt = self._client.insert( external_document_id=job.external_document_id, content=job.content, metadata=job.metadata, ) track_id = receipt.track_id if hasattr(self._repository, "set_track"): self._repository.set_track(job.projection_id, track_id) status = self._client.track_status(track_id) if status == "processed": self._repository.mark(job.projection_id, "ready") return "ready" self._repository.mark(job.projection_id, "processing") return "processing" except Exception as exc: self._repository.mark(job.projection_id, "failed", str(exc)[:1000]) return "failed" def delete(self, job: ProjectionJob) -> str: self._repository.mark(job.projection_id, "deleting") try: receipt = self._client.delete(job.external_document_id) status = "deleted" if receipt.verified else "unverified" self._repository.mark(job.projection_id, status) return status except Exception as exc: self._repository.mark(job.projection_id, "failed", str(exc)[:1000]) return "failed"