artifacts.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161
  1. from __future__ import annotations
  2. import hashlib
  3. import io
  4. import re
  5. from dataclasses import dataclass
  6. @dataclass(frozen=True)
  7. class ArtifactRecord:
  8. uid: str
  9. source_uid: str
  10. filename: str
  11. media_type: str
  12. size_bytes: int
  13. content_hash: str
  14. storage_ref: str
  15. parser_version: str
  16. class ArtifactService:
  17. def __init__(self, repository, storage, file_policy, *, uid_factory):
  18. self.repository = repository
  19. self.storage = storage
  20. self.file_policy = file_policy
  21. self.uid_factory = uid_factory
  22. def store(self, source_uid, filename, media_type, content, parser_version):
  23. validated = self.file_policy.validate(filename, media_type, content)
  24. digest = hashlib.sha256(content).hexdigest()
  25. existing = self.repository.find(source_uid, digest, parser_version)
  26. if existing is not None:
  27. return existing, False
  28. object_key = f"data-research/{source_uid}/{digest}"
  29. storage_ref = self.storage.put(object_key, content, validated.media_type)
  30. record = ArtifactRecord(
  31. uid=self.uid_factory(),
  32. source_uid=str(source_uid),
  33. filename=validated.filename,
  34. media_type=validated.media_type,
  35. size_bytes=validated.size_bytes,
  36. content_hash=digest,
  37. storage_ref=storage_ref,
  38. parser_version=str(parser_version),
  39. )
  40. return self.repository.save(record), True
  41. def redact_excerpt(value):
  42. value = re.sub(r"(?<!\d)(1\d{2})\d{4}(\d{4})(?!\d)", r"\1****\2", str(value or ""))
  43. value = re.sub(
  44. r"(?i)\b(token|password|secret|api[_-]?key)\s*[=:]\s*[^\s,;]+",
  45. lambda match: f"{match.group(1)}=[redacted]",
  46. value,
  47. )
  48. return value
  49. class EvidenceService:
  50. def __init__(self, session, storage):
  51. self.session = session
  52. self.storage = storage
  53. def preview(self, uid):
  54. from app.models.data_research import EvidenceFragment
  55. model = self.session.get(EvidenceFragment, str(uid))
  56. if model is None:
  57. raise LookupError("evidence was not found")
  58. data = model.to_dict()
  59. data["excerpt"] = redact_excerpt(data.get("excerpt"))
  60. return data
  61. def list_for_job(self, job_uid):
  62. from app.models.data_research import EvidenceFragment
  63. records = self.session.query(EvidenceFragment).filter_by(
  64. job_uid=str(job_uid)
  65. ).order_by(EvidenceFragment.created_at.asc()).all()
  66. result = []
  67. for model in records:
  68. data = model.to_dict()
  69. data["excerpt"] = redact_excerpt(data.get("excerpt"))
  70. result.append(data)
  71. return result
  72. def download(self, uid):
  73. from app.models.data_research import EvidenceFragment, SourceArtifact
  74. evidence = self.session.get(EvidenceFragment, str(uid))
  75. artifact = (
  76. self.session.get(SourceArtifact, str(evidence.artifact_uid))
  77. if evidence is not None and evidence.artifact_uid
  78. else None
  79. )
  80. if artifact is None:
  81. raise LookupError("evidence artifact was not found")
  82. return self.storage.get(artifact.storage_ref), artifact.filename, artifact.media_type
  83. class SqlAlchemyArtifactRepository:
  84. def __init__(self, session):
  85. self.session = session
  86. @staticmethod
  87. def _record(model):
  88. return ArtifactRecord(
  89. uid=str(model.uid),
  90. source_uid=str(model.source_uid),
  91. filename=model.filename,
  92. media_type=model.media_type,
  93. size_bytes=int(model.size_bytes),
  94. content_hash=model.content_hash,
  95. storage_ref=model.storage_ref,
  96. parser_version=model.parser_version,
  97. )
  98. def find(self, source_uid, content_hash, parser_version):
  99. from app.models.data_research import SourceArtifact
  100. model = self.session.query(SourceArtifact).filter_by(
  101. source_uid=str(source_uid),
  102. content_hash=content_hash,
  103. parser_version=str(parser_version),
  104. ).first()
  105. return self._record(model) if model is not None else None
  106. def save(self, record):
  107. from app.models.data_research import SourceArtifact
  108. model = SourceArtifact(**record.__dict__)
  109. self.session.add(model)
  110. self.session.flush()
  111. return self._record(model)
  112. class MinioArtifactStorage:
  113. def __init__(self, client, bucket):
  114. self.client = client
  115. self.bucket = bucket
  116. def put(self, object_key, content, media_type):
  117. self.client.put_object(
  118. self.bucket,
  119. object_key,
  120. io.BytesIO(content),
  121. len(content),
  122. content_type=media_type,
  123. )
  124. return f"minio://{self.bucket}/{object_key}"
  125. def get(self, storage_ref):
  126. prefix = f"minio://{self.bucket}/"
  127. if not str(storage_ref).startswith(prefix):
  128. raise ValueError("invalid storage reference")
  129. response = self.client.get_object(self.bucket, str(storage_ref)[len(prefix):])
  130. try:
  131. return response.read()
  132. finally:
  133. response.close()
  134. response.release_conn()