artifacts.py 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148
  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 download(self, uid):
  62. from app.models.data_research import EvidenceFragment, SourceArtifact
  63. evidence = self.session.get(EvidenceFragment, str(uid))
  64. artifact = (
  65. self.session.get(SourceArtifact, str(evidence.artifact_uid))
  66. if evidence is not None and evidence.artifact_uid
  67. else None
  68. )
  69. if artifact is None:
  70. raise LookupError("evidence artifact was not found")
  71. return self.storage.get(artifact.storage_ref), artifact.filename, artifact.media_type
  72. class SqlAlchemyArtifactRepository:
  73. def __init__(self, session):
  74. self.session = session
  75. @staticmethod
  76. def _record(model):
  77. return ArtifactRecord(
  78. uid=str(model.uid),
  79. source_uid=str(model.source_uid),
  80. filename=model.filename,
  81. media_type=model.media_type,
  82. size_bytes=int(model.size_bytes),
  83. content_hash=model.content_hash,
  84. storage_ref=model.storage_ref,
  85. parser_version=model.parser_version,
  86. )
  87. def find(self, source_uid, content_hash, parser_version):
  88. from app.models.data_research import SourceArtifact
  89. model = self.session.query(SourceArtifact).filter_by(
  90. source_uid=str(source_uid),
  91. content_hash=content_hash,
  92. parser_version=str(parser_version),
  93. ).first()
  94. return self._record(model) if model is not None else None
  95. def save(self, record):
  96. from app.models.data_research import SourceArtifact
  97. model = SourceArtifact(**record.__dict__)
  98. self.session.add(model)
  99. self.session.flush()
  100. return self._record(model)
  101. class MinioArtifactStorage:
  102. def __init__(self, client, bucket):
  103. self.client = client
  104. self.bucket = bucket
  105. def put(self, object_key, content, media_type):
  106. self.client.put_object(
  107. self.bucket,
  108. object_key,
  109. io.BytesIO(content),
  110. len(content),
  111. content_type=media_type,
  112. )
  113. return f"minio://{self.bucket}/{object_key}"
  114. def get(self, storage_ref):
  115. prefix = f"minio://{self.bucket}/"
  116. if not str(storage_ref).startswith(prefix):
  117. raise ValueError("invalid storage reference")
  118. response = self.client.get_object(self.bucket, str(storage_ref)[len(prefix):])
  119. try:
  120. return response.read()
  121. finally:
  122. response.close()
  123. response.release_conn()