Files
govoplan-files/src/govoplan_files/backend/storage/campaign_usage.py
2026-07-02 14:59:52 +02:00

303 lines
11 KiB
Python

from __future__ import annotations
from collections import defaultdict
from pathlib import PurePosixPath
from typing import TYPE_CHECKING, 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
if TYPE_CHECKING:
from govoplan_campaign.backend.db.models import CampaignJob
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]
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[CampaignJob], *, 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: CampaignJob, *, stage: str) -> set[AttachmentUseKey]:
return _known_use_keys_for_jobs(session, [job], stage=stage)
def _add_use(
session: Session,
job: CampaignJob,
*,
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[CampaignJob],
*,
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
managed_by_job: dict[str, list[dict[str, object]]] = defaultdict(list)
fallback_attachments_by_job: dict[str, list[dict[str, object]]] = defaultdict(list)
job_by_id = {job.id: job for job in job_list}
asset_ids: set[str] = set()
version_ids: set[str] = set()
blob_ids: set[str] = set()
for job in job_list:
attachments = job.resolved_attachments or []
if not isinstance(attachments, list):
continue
for attachment in attachments:
if not isinstance(attachment, dict):
continue
managed_matches = attachment.get("managed_matches")
if isinstance(managed_matches, list) and managed_matches:
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
managed_by_job[job.id].append(item)
asset_ids.add(asset_id)
version_ids.add(version_id)
blob_ids.add(blob_id)
else:
fallback_attachments_by_job[job.id].append(attachment)
assets_by_id = {
item.id: item
for item in session.query(FileAsset).filter(FileAsset.id.in_(asset_ids)).all()
} if asset_ids else {}
versions_by_id = {
item.id: item
for item in session.query(FileVersion).filter(FileVersion.id.in_(version_ids)).all()
} if version_ids else {}
blobs_by_id = {
item.id: item
for item in session.query(FileBlob).filter(FileBlob.id.in_(blob_ids)).all()
} if blob_ids else {}
known_keys = _known_use_keys_for_jobs(session, job_list, stage=stage)
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 asset or not version or not blob:
continue
if asset.tenant_id != job.tenant_id or version.tenant_id != job.tenant_id or blob.tenant_id != job.tenant_id:
continue
if version.file_asset_id != asset.id or version.blob_id != blob.id:
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,
)
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 = 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 = {}
for asset in assets:
by_key[asset.display_path.strip("/")] = asset
by_key.setdefault(asset.filename, asset)
assets_by_campaign[campaign_key] = by_key
version_blobs_by_campaign[campaign_key] = current_versions_and_blobs(session, assets)
version_blobs = version_blobs_by_campaign[campaign_key]
for attachment in attachments:
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: CampaignJob, *, 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: CampaignJob) -> 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,
)
)