| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161 |
- from __future__ import annotations
- import hashlib
- import io
- import re
- from dataclasses import dataclass
- @dataclass(frozen=True)
- class ArtifactRecord:
- uid: str
- source_uid: str
- filename: str
- media_type: str
- size_bytes: int
- content_hash: str
- storage_ref: str
- parser_version: str
- class ArtifactService:
- def __init__(self, repository, storage, file_policy, *, uid_factory):
- self.repository = repository
- self.storage = storage
- self.file_policy = file_policy
- self.uid_factory = uid_factory
- def store(self, source_uid, filename, media_type, content, parser_version):
- validated = self.file_policy.validate(filename, media_type, content)
- digest = hashlib.sha256(content).hexdigest()
- existing = self.repository.find(source_uid, digest, parser_version)
- if existing is not None:
- return existing, False
- object_key = f"data-research/{source_uid}/{digest}"
- storage_ref = self.storage.put(object_key, content, validated.media_type)
- record = ArtifactRecord(
- uid=self.uid_factory(),
- source_uid=str(source_uid),
- filename=validated.filename,
- media_type=validated.media_type,
- size_bytes=validated.size_bytes,
- content_hash=digest,
- storage_ref=storage_ref,
- parser_version=str(parser_version),
- )
- return self.repository.save(record), True
- def redact_excerpt(value):
- value = re.sub(r"(?<!\d)(1\d{2})\d{4}(\d{4})(?!\d)", r"\1****\2", str(value or ""))
- value = re.sub(
- r"(?i)\b(token|password|secret|api[_-]?key)\s*[=:]\s*[^\s,;]+",
- lambda match: f"{match.group(1)}=[redacted]",
- value,
- )
- return value
- class EvidenceService:
- def __init__(self, session, storage):
- self.session = session
- self.storage = storage
- def preview(self, uid):
- from app.models.data_research import EvidenceFragment
- model = self.session.get(EvidenceFragment, str(uid))
- if model is None:
- raise LookupError("evidence was not found")
- data = model.to_dict()
- data["excerpt"] = redact_excerpt(data.get("excerpt"))
- return data
- def list_for_job(self, job_uid):
- from app.models.data_research import EvidenceFragment
- records = self.session.query(EvidenceFragment).filter_by(
- job_uid=str(job_uid)
- ).order_by(EvidenceFragment.created_at.asc()).all()
- result = []
- for model in records:
- data = model.to_dict()
- data["excerpt"] = redact_excerpt(data.get("excerpt"))
- result.append(data)
- return result
- def download(self, uid):
- from app.models.data_research import EvidenceFragment, SourceArtifact
- evidence = self.session.get(EvidenceFragment, str(uid))
- artifact = (
- self.session.get(SourceArtifact, str(evidence.artifact_uid))
- if evidence is not None and evidence.artifact_uid
- else None
- )
- if artifact is None:
- raise LookupError("evidence artifact was not found")
- return self.storage.get(artifact.storage_ref), artifact.filename, artifact.media_type
- class SqlAlchemyArtifactRepository:
- def __init__(self, session):
- self.session = session
- @staticmethod
- def _record(model):
- return ArtifactRecord(
- uid=str(model.uid),
- source_uid=str(model.source_uid),
- filename=model.filename,
- media_type=model.media_type,
- size_bytes=int(model.size_bytes),
- content_hash=model.content_hash,
- storage_ref=model.storage_ref,
- parser_version=model.parser_version,
- )
- def find(self, source_uid, content_hash, parser_version):
- from app.models.data_research import SourceArtifact
- model = self.session.query(SourceArtifact).filter_by(
- source_uid=str(source_uid),
- content_hash=content_hash,
- parser_version=str(parser_version),
- ).first()
- return self._record(model) if model is not None else None
- def save(self, record):
- from app.models.data_research import SourceArtifact
- model = SourceArtifact(**record.__dict__)
- self.session.add(model)
- self.session.flush()
- return self._record(model)
- class MinioArtifactStorage:
- def __init__(self, client, bucket):
- self.client = client
- self.bucket = bucket
- def put(self, object_key, content, media_type):
- self.client.put_object(
- self.bucket,
- object_key,
- io.BytesIO(content),
- len(content),
- content_type=media_type,
- )
- return f"minio://{self.bucket}/{object_key}"
- def get(self, storage_ref):
- prefix = f"minio://{self.bucket}/"
- if not str(storage_ref).startswith(prefix):
- raise ValueError("invalid storage reference")
- response = self.client.get_object(self.bucket, str(storage_ref)[len(prefix):])
- try:
- return response.read()
- finally:
- response.close()
- response.release_conn()
|