Files
govoplan-files/src/govoplan_files/backend/storage/campaign_usage.py
2026-07-14 13:22:11 +02:00

391 lines
14 KiB
Python

from __future__ import annotations
from collections import defaultdict
from dataclasses import dataclass, field
from pathlib import PurePosixPath
from typing import Any, Iterable
from sqlalchemy.orm import Session
from govoplan_files.backend.db.models import CampaignAttachmentUse, FileAsset, FileBlob, FileVersion
from govoplan_files.backend.storage.common import utcnow
from govoplan_files.backend.storage.files import current_versions_and_blobs, list_assets_for_user
CampaignJobLike = Any
def _candidate_match_keys(raw_match: str) -> set[str]:
cleaned = raw_match.replace("\\", "/").strip().strip("/")
result = {cleaned}
if cleaned:
result.add(PurePosixPath(cleaned).name)
return {item for item in result if item}
AttachmentUseKey = tuple[str, str, str, str]
@dataclass(slots=True)
class _AttachmentBatchRefs:
managed_by_job: dict[str, list[dict[str, object]]] = field(default_factory=lambda: defaultdict(list))
fallback_attachments_by_job: dict[str, list[dict[str, object]]] = field(default_factory=lambda: defaultdict(list))
asset_ids: set[str] = field(default_factory=set)
version_ids: set[str] = field(default_factory=set)
blob_ids: set[str] = field(default_factory=set)
def _attachment_use_key(*, job_id: str, file_version_id: str, filename_used: str, stage: str) -> AttachmentUseKey:
return job_id, file_version_id, filename_used, stage
def _known_use_keys_for_jobs(session: Session, jobs: Iterable[CampaignJobLike], *, stage: str) -> set[AttachmentUseKey]:
job_ids = {job.id for job in jobs if job.id}
if not job_ids:
return set()
keys = {
_attachment_use_key(
job_id=job_id,
file_version_id=file_version_id,
filename_used=filename_used,
stage=use_stage,
)
for job_id, file_version_id, filename_used, use_stage in (
session.query(
CampaignAttachmentUse.campaign_job_id,
CampaignAttachmentUse.file_version_id,
CampaignAttachmentUse.filename_used,
CampaignAttachmentUse.use_stage,
)
.filter(
CampaignAttachmentUse.campaign_job_id.in_(job_ids),
CampaignAttachmentUse.use_stage == stage,
)
.all()
)
if job_id is not None
}
for pending in session.new:
if not isinstance(pending, CampaignAttachmentUse):
continue
if pending.campaign_job_id not in job_ids or pending.use_stage != stage:
continue
keys.add(
_attachment_use_key(
job_id=pending.campaign_job_id,
file_version_id=pending.file_version_id,
filename_used=pending.filename_used,
stage=pending.use_stage,
)
)
return keys
def _known_use_keys(session: Session, job: CampaignJobLike, *, stage: str) -> set[AttachmentUseKey]:
return _known_use_keys_for_jobs(session, [job], stage=stage)
def _add_use(
session: Session,
job: CampaignJobLike,
*,
asset: FileAsset,
version: FileVersion,
blob: FileBlob,
filename_used: str,
stage: str,
known_keys: set[AttachmentUseKey],
) -> None:
key = _attachment_use_key(
job_id=job.id,
file_version_id=version.id,
filename_used=filename_used,
stage=stage,
)
if key in known_keys:
return
known_keys.add(key)
session.add(
CampaignAttachmentUse(
tenant_id=job.tenant_id,
campaign_id=job.campaign_id,
campaign_version_id=job.campaign_version_id,
campaign_job_id=job.id,
entry_index=job.entry_index,
entry_id=job.entry_id,
file_asset_id=asset.id,
file_version_id=version.id,
file_blob_id=blob.id,
filename_used=filename_used,
checksum_sha256=blob.checksum_sha256,
size_bytes=blob.size_bytes,
content_type=blob.content_type,
use_stage=stage,
)
)
def record_campaign_attachment_uses_for_jobs(
session: Session,
jobs: Iterable[CampaignJobLike],
*,
stage: str = "built",
) -> None:
"""Record immutable attachment evidence for multiple jobs efficiently.
Build operations may create thousands of jobs in one transaction. Loading
every campaign-shared asset and flushing once per recipient turns that into
an avoidable query/flush multiplier. This batch path resolves the exact
managed IDs once, fetches each referenced table in bulk, and reuses legacy
filename maps per campaign only when an older snapshot needs them.
"""
job_list = [job for job in jobs if job.id]
if not job_list:
return
refs = _collect_attachment_batch_refs(job_list)
job_by_id = {job.id: job for job in job_list}
known_keys = _known_use_keys_for_jobs(session, job_list, stage=stage)
assets_by_id, versions_by_id, blobs_by_id = _load_managed_attachment_entities(session, refs)
_record_managed_attachment_uses(
session,
job_by_id=job_by_id,
managed_by_job=refs.managed_by_job,
assets_by_id=assets_by_id,
versions_by_id=versions_by_id,
blobs_by_id=blobs_by_id,
stage=stage,
known_keys=known_keys,
)
_record_fallback_attachment_uses(
session,
job_by_id=job_by_id,
fallback_attachments_by_job=refs.fallback_attachments_by_job,
stage=stage,
known_keys=known_keys,
)
def _collect_attachment_batch_refs(jobs: list[CampaignJobLike]) -> _AttachmentBatchRefs:
refs = _AttachmentBatchRefs()
for job in jobs:
attachments = job.resolved_attachments or []
if not isinstance(attachments, list):
continue
for attachment in attachments:
if not isinstance(attachment, dict):
continue
if not _collect_managed_attachment_refs(refs, job.id, attachment):
refs.fallback_attachments_by_job[job.id].append(attachment)
return refs
def _collect_managed_attachment_refs(refs: _AttachmentBatchRefs, job_id: str, attachment: dict[str, object]) -> bool:
managed_matches = attachment.get("managed_matches")
if not isinstance(managed_matches, list) or not managed_matches:
return False
for item in managed_matches:
if not isinstance(item, dict):
continue
asset_id = str(item.get("asset_id") or "")
version_id = str(item.get("version_id") or "")
blob_id = str(item.get("blob_id") or "")
if not asset_id or not version_id or not blob_id:
continue
refs.managed_by_job[job_id].append(item)
refs.asset_ids.add(asset_id)
refs.version_ids.add(version_id)
refs.blob_ids.add(blob_id)
return True
def _load_managed_attachment_entities(
session: Session,
refs: _AttachmentBatchRefs,
) -> tuple[dict[str, FileAsset], dict[str, FileVersion], dict[str, FileBlob]]:
assets_by_id = {item.id: item for item in session.query(FileAsset).filter(FileAsset.id.in_(refs.asset_ids)).all()} if refs.asset_ids else {}
versions_by_id = {item.id: item for item in session.query(FileVersion).filter(FileVersion.id.in_(refs.version_ids)).all()} if refs.version_ids else {}
blobs_by_id = {item.id: item for item in session.query(FileBlob).filter(FileBlob.id.in_(refs.blob_ids)).all()} if refs.blob_ids else {}
return assets_by_id, versions_by_id, blobs_by_id
def _record_managed_attachment_uses(
session: Session,
*,
job_by_id: dict[str, CampaignJobLike],
managed_by_job: dict[str, list[dict[str, object]]],
assets_by_id: dict[str, FileAsset],
versions_by_id: dict[str, FileVersion],
blobs_by_id: dict[str, FileBlob],
stage: str,
known_keys: set[AttachmentUseKey],
) -> None:
for job_id, items in managed_by_job.items():
job = job_by_id[job_id]
for item in items:
asset = assets_by_id.get(str(item.get("asset_id") or ""))
version = versions_by_id.get(str(item.get("version_id") or ""))
blob = blobs_by_id.get(str(item.get("blob_id") or ""))
if not _managed_attachment_entities_match_job(job, asset, version, blob):
continue
_add_use(
session,
job,
asset=asset,
version=version,
blob=blob,
filename_used=str(item.get("filename") or asset.filename),
stage=stage,
known_keys=known_keys,
)
def _managed_attachment_entities_match_job(
job: CampaignJobLike,
asset: FileAsset | None,
version: FileVersion | None,
blob: FileBlob | None,
) -> bool:
if not asset or not version or not blob:
return False
if asset.tenant_id != job.tenant_id or version.tenant_id != job.tenant_id or blob.tenant_id != job.tenant_id:
return False
return version.file_asset_id == asset.id and version.blob_id == blob.id
def _record_fallback_attachment_uses(
session: Session,
*,
job_by_id: dict[str, CampaignJobLike],
fallback_attachments_by_job: dict[str, list[dict[str, object]]],
stage: str,
known_keys: set[AttachmentUseKey],
) -> None:
assets_by_campaign: dict[tuple[str, str], dict[str, FileAsset]] = {}
version_blobs_by_campaign: dict[tuple[str, str], dict[str, tuple[FileVersion, FileBlob]]] = {}
for job_id, attachments in fallback_attachments_by_job.items():
job = job_by_id[job_id]
campaign_key = (job.tenant_id, job.campaign_id)
by_key, version_blobs = _fallback_campaign_assets(
session,
job,
campaign_key=campaign_key,
assets_by_campaign=assets_by_campaign,
version_blobs_by_campaign=version_blobs_by_campaign,
)
for attachment in attachments:
_record_fallback_attachment(session, job, attachment, by_key=by_key, version_blobs=version_blobs, stage=stage, known_keys=known_keys)
def _fallback_campaign_assets(
session: Session,
job: CampaignJobLike,
*,
campaign_key: tuple[str, str],
assets_by_campaign: dict[tuple[str, str], dict[str, FileAsset]],
version_blobs_by_campaign: dict[tuple[str, str], dict[str, tuple[FileVersion, FileBlob]]],
) -> tuple[dict[str, FileAsset], dict[str, tuple[FileVersion, FileBlob]]]:
by_key = assets_by_campaign.get(campaign_key)
if by_key is None:
assets = list_assets_for_user(session, tenant_id=job.tenant_id, user_id="", campaign_id=job.campaign_id, is_admin=True)
by_key = _asset_lookup_by_path_and_filename(assets)
assets_by_campaign[campaign_key] = by_key
version_blobs_by_campaign[campaign_key] = current_versions_and_blobs(session, assets)
return by_key, version_blobs_by_campaign[campaign_key]
def _asset_lookup_by_path_and_filename(assets: list[FileAsset]) -> dict[str, FileAsset]:
by_key: dict[str, FileAsset] = {}
for asset in assets:
by_key[asset.display_path.strip("/")] = asset
by_key.setdefault(asset.filename, asset)
return by_key
def _record_fallback_attachment(
session: Session,
job: CampaignJobLike,
attachment: dict[str, object],
*,
by_key: dict[str, FileAsset],
version_blobs: dict[str, tuple[FileVersion, FileBlob]],
stage: str,
known_keys: set[AttachmentUseKey],
) -> None:
matches = attachment.get("matches") if isinstance(attachment.get("matches"), list) else []
for raw in matches:
if not isinstance(raw, str):
continue
asset = next((by_key[key] for key in _candidate_match_keys(raw) if key in by_key), None)
if not asset:
continue
version_blob = version_blobs.get(asset.id)
if not version_blob:
continue
version, blob = version_blob
_add_use(
session,
job,
asset=asset,
version=version,
blob=blob,
filename_used=asset.filename,
stage=stage,
known_keys=known_keys,
)
def record_campaign_attachment_uses_for_job(session: Session, job: CampaignJobLike, *, stage: str = "built") -> None:
"""Record immutable managed file versions used by one built/sent job."""
record_campaign_attachment_uses_for_jobs(session, [job], stage=stage)
def mark_job_attachment_uses_sent(session: Session, job: CampaignJobLike) -> None:
record_campaign_attachment_uses_for_job(session, job, stage="built")
# Sessions use autoflush=False. Flush any compatibility-built evidence so
# the following query can copy it to the sent stage in the same call.
session.flush()
now = utcnow()
uses = (
session.query(CampaignAttachmentUse)
.filter(
CampaignAttachmentUse.tenant_id == job.tenant_id,
CampaignAttachmentUse.campaign_job_id == job.id,
CampaignAttachmentUse.use_stage == "built",
)
.all()
)
sent_keys = _known_use_keys(session, job, stage="sent")
for use in uses:
key = _attachment_use_key(
job_id=job.id,
file_version_id=use.file_version_id,
filename_used=use.filename_used,
stage="sent",
)
if key in sent_keys:
continue
sent_keys.add(key)
session.add(
CampaignAttachmentUse(
tenant_id=use.tenant_id,
campaign_id=use.campaign_id,
campaign_version_id=use.campaign_version_id,
campaign_job_id=use.campaign_job_id,
entry_index=use.entry_index,
entry_id=use.entry_id,
file_asset_id=use.file_asset_id,
file_version_id=use.file_version_id,
file_blob_id=use.file_blob_id,
filename_used=use.filename_used,
checksum_sha256=use.checksum_sha256,
size_bytes=use.size_bytes,
content_type=use.content_type,
use_stage="sent",
used_at=now,
)
)