from sqlalchemy.orm import Session from packages.adapters.sqlalchemy_impl.models import IngestJobModel from packages.domain import IngestJob, IngestJobStatus class SQLAlchemyIngestJobRepository: def __init__(self, session: Session): self.session = session def create(self, job: IngestJob) -> IngestJob: model = IngestJobModel( id=job.id, project_id=job.project_id, library_id=job.library_id, storage_key=job.storage_key, status=job.status.value, error_message=job.error_message, result_asset_id=job.result_asset_id, file_hash=job.file_hash, created_at=job.created_at, updated_at=job.updated_at, ) self.session.add(model) self.session.commit() return job def get(self, job_id: str) -> IngestJob | None: model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job_id).first() if model is None: return None return IngestJob( id=model.id, project_id=model.project_id, library_id=model.library_id, storage_key=model.storage_key, status=IngestJobStatus(model.status), error_message=model.error_message, result_asset_id=model.result_asset_id, file_hash=model.file_hash or "", created_at=model.created_at, updated_at=model.updated_at, ) def list_by_project(self, project_id: str) -> list[IngestJob]: models = self.session.query(IngestJobModel).filter(IngestJobModel.project_id == project_id).all() return [self.get(model.id) for model in models if self.get(model.id) is not None] def update(self, job: IngestJob) -> IngestJob: model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job.id).first() if model is None: raise ValueError(f"IngestJob {job.id} not found") model.status = job.status.value model.error_message = job.error_message model.result_asset_id = job.result_asset_id model.file_hash = job.file_hash model.updated_at = job.updated_at self.session.commit() return job