| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859 |
- 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"
|