feat: archive logs export (#38207)

Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
非法操作
2026-07-14 08:25:37 +00:00
committed by GitHub
co-authored by autofix-ci[bot]
parent 82ff93cbdd
commit f09d2b191f
67 changed files with 4779 additions and 18 deletions
@@ -0,0 +1,347 @@
"""
Workflow-run archive bundle index helpers.
Archive manifests in object storage remain the recoverable source of truth. This module mirrors their small query
surface into `workflow_run_archive_bundles` so console listing and download jobs can avoid listing R2 on request.
The backfill path is intentionally idempotent: every manifest is decoded, checked against the V2 schema markers, and
upserted by immutable bundle identity.
"""
import datetime
import json
import logging
import time
from collections.abc import Sequence
from dataclasses import dataclass, field
from typing import TypedDict, cast
from sqlalchemy import select
from sqlalchemy.orm import Session, sessionmaker
from extensions.ext_database import db
from libs.archive_storage import ArchiveStorage, get_archive_storage
from models.workflow import WorkflowRunArchiveBundle
from services.retention.workflow_run.constants import (
ARCHIVE_BUNDLE_FORMAT,
ARCHIVE_BUNDLE_MANIFEST_NAME,
ARCHIVE_BUNDLE_SCHEMA_VERSION,
)
logger = logging.getLogger(__name__)
ARCHIVE_BUNDLE_ROOT_PREFIX = "workflow-runs/v2/"
class ArchiveBundleTableManifestEntry(TypedDict):
row_count: int
checksum: str
size_bytes: int
object_key: str
class ArchiveBundleManifest(TypedDict):
schema_version: str
archive_format: str
tenant_id: str
tenant_prefix: str
year: int
month: int
shard: str
bundle_id: str
object_prefix: str
workflow_run_count: int
workflow_node_execution_count: int
min_created_at: str
max_created_at: str
min_run_id: str
max_run_id: str
archived_at: str
tables: dict[str, ArchiveBundleTableManifestEntry]
run_ids: list[str]
@dataclass(frozen=True)
class ArchiveBundleIndexValues:
"""Computed DB-index values derived from one manifest."""
row_count: int
archive_bytes: int
archived_at: datetime.datetime
@dataclass
class ArchiveBundleIndexBackfillSummary:
"""Aggregate result for a manifest-to-DB-index reconciliation run."""
manifests_found: int = 0
bundles_processed: int = 0
bundles_upserted: int = 0
bundles_failed: int = 0
workflow_run_count: int = 0
row_count: int = 0
archive_bytes: int = 0
elapsed_time: float = 0.0
errors: list[str] = field(default_factory=list)
def decode_archive_bundle_manifest(manifest_data: bytes) -> ArchiveBundleManifest:
"""Decode raw manifest bytes into the V2 archive manifest shape."""
return cast(ArchiveBundleManifest, json.loads(manifest_data.decode("utf-8")))
def parse_archive_manifest_datetime(value: str) -> datetime.datetime:
"""Parse manifest datetimes and normalize timezone-aware values to naive UTC for DB storage."""
parsed = datetime.datetime.fromisoformat(value)
if parsed.tzinfo is None:
return parsed
return parsed.astimezone(datetime.UTC).replace(tzinfo=None)
def calculate_archive_bundle_index_values(
manifest: ArchiveBundleManifest,
manifest_size_bytes: int,
) -> ArchiveBundleIndexValues:
"""Calculate row count, stored bytes, and archived timestamp for the DB index."""
_validate_archive_bundle_manifest(manifest)
row_count = sum(entry["row_count"] for entry in manifest["tables"].values())
archive_bytes = manifest_size_bytes + sum(entry["size_bytes"] for entry in manifest["tables"].values())
return ArchiveBundleIndexValues(
row_count=row_count,
archive_bytes=archive_bytes,
archived_at=parse_archive_manifest_datetime(manifest["archived_at"]),
)
def upsert_archive_bundle_index_from_manifest(
session: Session,
manifest: ArchiveBundleManifest,
manifest_size_bytes: int,
) -> WorkflowRunArchiveBundle:
"""
Persist one archive manifest into `workflow_run_archive_bundles`.
The caller owns transaction boundaries. Re-running this function for the same manifest is safe and refreshes the
mutable metrics derived from object sizes and row counts.
"""
values = calculate_archive_bundle_index_values(manifest, manifest_size_bytes)
existing = session.scalar(
select(WorkflowRunArchiveBundle).where(
WorkflowRunArchiveBundle.tenant_id == manifest["tenant_id"],
WorkflowRunArchiveBundle.year == manifest["year"],
WorkflowRunArchiveBundle.month == manifest["month"],
WorkflowRunArchiveBundle.shard == manifest["shard"],
WorkflowRunArchiveBundle.bundle_id == manifest["bundle_id"],
)
)
if existing is None:
bundle = WorkflowRunArchiveBundle(
tenant_id=manifest["tenant_id"],
year=manifest["year"],
month=manifest["month"],
shard=manifest["shard"],
bundle_id=manifest["bundle_id"],
workflow_run_count=manifest["workflow_run_count"],
row_count=values.row_count,
archive_bytes=values.archive_bytes,
archived_at=values.archived_at,
)
session.add(bundle)
return bundle
existing.workflow_run_count = manifest["workflow_run_count"]
existing.row_count = values.row_count
existing.archive_bytes = values.archive_bytes
existing.archived_at = values.archived_at
return existing
class WorkflowRunArchiveBundleIndexBackfill:
"""
Rebuild the DB bundle index by scanning object-store manifests.
Tenant IDs are the cheapest scope because they map directly to the object prefix. Tenant prefixes are supported for
rollout reconciliation, but they still require listing all tenants under that prefix and filtering keys locally.
"""
storage: ArchiveStorage | None
session_factory: sessionmaker[Session]
def __init__(
self,
*,
storage: ArchiveStorage | None = None,
session_factory: sessionmaker[Session] | None = None,
) -> None:
self.storage = storage
self.session_factory = session_factory or sessionmaker(bind=db.engine, expire_on_commit=False)
def run(
self,
*,
tenant_ids: Sequence[str] | None = None,
tenant_prefixes: Sequence[str] | None = None,
year: int | None = None,
month: int | None = None,
limit: int | None = None,
dry_run: bool = False,
) -> ArchiveBundleIndexBackfillSummary:
"""Scan matching manifest objects and idempotently upsert their DB index rows."""
start_time = time.time()
summary = ArchiveBundleIndexBackfillSummary()
storage = self.storage or get_archive_storage()
manifest_keys = self._list_manifest_keys(
storage,
tenant_ids=tenant_ids,
tenant_prefixes=tenant_prefixes,
year=year,
month=month,
)
summary.manifests_found = len(manifest_keys)
if limit is not None:
manifest_keys = manifest_keys[:limit]
for manifest_key in manifest_keys:
try:
manifest_data = storage.get_object(manifest_key)
manifest = decode_archive_bundle_manifest(manifest_data)
self._validate_manifest_scope(
manifest,
manifest_key=manifest_key,
tenant_ids=tenant_ids,
tenant_prefixes=tenant_prefixes,
year=year,
month=month,
)
values = calculate_archive_bundle_index_values(manifest, len(manifest_data))
summary.bundles_processed += 1
summary.workflow_run_count += manifest["workflow_run_count"]
summary.row_count += values.row_count
summary.archive_bytes += values.archive_bytes
if dry_run:
continue
with self.session_factory() as session:
upsert_archive_bundle_index_from_manifest(session, manifest, len(manifest_data))
session.commit()
summary.bundles_upserted += 1
except Exception as exc:
logger.warning("Failed to backfill workflow archive bundle index from %s", manifest_key, exc_info=True)
summary.bundles_failed += 1
summary.errors.append(f"{manifest_key}: {exc}")
summary.elapsed_time = time.time() - start_time
return summary
@classmethod
def _list_manifest_keys(
cls,
storage: ArchiveStorage,
*,
tenant_ids: Sequence[str] | None,
tenant_prefixes: Sequence[str] | None,
year: int | None,
month: int | None,
) -> list[str]:
prefixes = cls._list_prefixes(tenant_ids=tenant_ids, tenant_prefixes=tenant_prefixes, year=year, month=month)
keys: list[str] = []
for prefix in prefixes:
keys.extend(storage.list_objects(prefix))
return sorted(
key
for key in keys
if key.endswith(f"/{ARCHIVE_BUNDLE_MANIFEST_NAME}")
and cls._manifest_key_matches_scope(
key,
tenant_ids=tenant_ids,
tenant_prefixes=tenant_prefixes,
year=year,
month=month,
)
)
@staticmethod
def _list_prefixes(
*,
tenant_ids: Sequence[str] | None,
tenant_prefixes: Sequence[str] | None,
year: int | None,
month: int | None,
) -> list[str]:
if tenant_ids:
prefixes = []
for tenant_id in sorted(set(tenant_ids)):
prefix = f"{ARCHIVE_BUNDLE_ROOT_PREFIX}tenant_prefix={tenant_id[0].lower()}/tenant_id={tenant_id}/"
if year is not None:
prefix += f"year={year:04d}/"
if month is not None:
prefix += f"month={month:02d}/"
prefixes.append(prefix)
return prefixes
if tenant_prefixes:
return [
f"{ARCHIVE_BUNDLE_ROOT_PREFIX}tenant_prefix={tenant_prefix}/"
for tenant_prefix in sorted(set(tenant_prefixes))
]
return [ARCHIVE_BUNDLE_ROOT_PREFIX]
@staticmethod
def _manifest_key_matches_scope(
key: str,
*,
tenant_ids: Sequence[str] | None,
tenant_prefixes: Sequence[str] | None,
year: int | None,
month: int | None,
) -> bool:
if tenant_ids and _extract_key_part(key, "tenant_id") not in set(tenant_ids):
return False
if tenant_prefixes and _extract_key_part(key, "tenant_prefix") not in set(tenant_prefixes):
return False
if year is not None and _extract_key_part(key, "year") != f"{year:04d}":
return False
if month is not None and _extract_key_part(key, "month") != f"{month:02d}":
return False
return True
@staticmethod
def _validate_manifest_scope(
manifest: ArchiveBundleManifest,
*,
manifest_key: str,
tenant_ids: Sequence[str] | None,
tenant_prefixes: Sequence[str] | None,
year: int | None,
month: int | None,
) -> None:
expected_object_prefix = manifest_key.removesuffix(f"/{ARCHIVE_BUNDLE_MANIFEST_NAME}")
if manifest["object_prefix"] != expected_object_prefix:
raise ValueError(
f"manifest object_prefix mismatch: expected={expected_object_prefix}, "
f"actual={manifest['object_prefix']}"
)
if tenant_ids and manifest["tenant_id"] not in tenant_ids:
raise ValueError(f"manifest tenant_id is outside requested scope: {manifest['tenant_id']}")
if tenant_prefixes and manifest["tenant_prefix"] not in tenant_prefixes:
raise ValueError(f"manifest tenant_prefix is outside requested scope: {manifest['tenant_prefix']}")
if year is not None and manifest["year"] != year:
raise ValueError(f"manifest year is outside requested scope: {manifest['year']}")
if month is not None and manifest["month"] != month:
raise ValueError(f"manifest month is outside requested scope: {manifest['month']}")
def _validate_archive_bundle_manifest(manifest: ArchiveBundleManifest) -> None:
if manifest["schema_version"] != ARCHIVE_BUNDLE_SCHEMA_VERSION:
raise ValueError(f"unsupported archive bundle schema version: {manifest['schema_version']}")
if manifest["archive_format"] != ARCHIVE_BUNDLE_FORMAT:
raise ValueError(f"unsupported archive bundle format: {manifest['archive_format']}")
def _extract_key_part(key: str, name: str) -> str | None:
prefix = f"{name}="
for part in key.split("/"):
if part.startswith(prefix):
return part[len(prefix) :]
return None
@@ -0,0 +1,354 @@
"""
Prepare monthly workflow-run archive downloads.
Console requests create a short-lived Redis task and Celery runs this module in the background. The DB bundle index is
the online lookup source: this preparer never lists archive storage, and it validates the indexed bundle set against the
stable download id before packaging archive Parquet objects into one user-facing CSV ZIP file.
"""
import datetime
import hashlib
import io
import logging
import re
import zipfile
from collections.abc import Sequence
from typing import cast
import pyarrow as pa
import pyarrow.compute as pc
import pyarrow.csv as pa_csv
import pyarrow.parquet as pq
from sqlalchemy import select
from sqlalchemy.orm import Session, sessionmaker
from core.helper.csv_sanitizer import CSVSanitizer
from extensions.ext_database import db
from libs.archive_storage import ArchiveStorage, get_archive_storage, get_export_storage
from models.workflow import WorkflowRunArchiveBundle
from services.retention.workflow_run.archive_bundle_index import (
ARCHIVE_BUNDLE_ROOT_PREFIX,
ArchiveBundleManifest,
ArchiveBundleTableManifestEntry,
decode_archive_bundle_manifest,
)
from services.retention.workflow_run.archive_download_task_cache import (
WorkflowRunArchiveDownloadStatus,
WorkflowRunArchiveDownloadTask,
WorkflowRunArchiveDownloadTaskCache,
build_archive_download_id,
)
from services.retention.workflow_run.constants import (
ARCHIVE_BUNDLE_FORMAT,
ARCHIVE_BUNDLE_MANIFEST_NAME,
ARCHIVE_BUNDLE_SCHEMA_VERSION,
)
logger = logging.getLogger(__name__)
ARCHIVE_DOWNLOAD_ROOT_PREFIX = "workflow-runs/downloads/v1/"
ARCHIVE_DOWNLOAD_MIME_TYPE = "application/zip"
_CSV_FORMULA_PREFIX_PATTERN = f"^[{re.escape(''.join(CSVSanitizer.FORMULA_CHARS))}]"
class WorkflowRunArchiveDownloadPreparer:
"""
Build one ready-to-download CSV ZIP for a Redis archive download task.
The output object is deterministic for a given `download_id`, so retrying a failed task overwrites the same
temporary object instead of creating unbounded duplicate files. Source archive bundles are read from the archive
bucket, while the prepared ZIP is written to the export bucket so object lifecycle policies can expire downloads
without touching long-lived archives.
"""
archive_storage: ArchiveStorage | None
download_storage: ArchiveStorage | None
cache: WorkflowRunArchiveDownloadTaskCache
session_factory: sessionmaker[Session]
def __init__(
self,
*,
storage: ArchiveStorage | None = None,
archive_storage: ArchiveStorage | None = None,
download_storage: ArchiveStorage | None = None,
cache: WorkflowRunArchiveDownloadTaskCache | None = None,
session_factory: sessionmaker[Session] | None = None,
) -> None:
self.archive_storage = archive_storage or storage
self.download_storage = download_storage or storage
self.cache = cache or WorkflowRunArchiveDownloadTaskCache()
self.session_factory = session_factory or sessionmaker(bind=db.engine, expire_on_commit=False)
def prepare(self, *, tenant_id: str, download_id: str) -> WorkflowRunArchiveDownloadTask | None:
"""Prepare a ZIP for an existing Redis task and persist terminal task state."""
with self.cache.lock(tenant_id=tenant_id, download_id=download_id):
task = self.cache.get(tenant_id=tenant_id, download_id=download_id)
if task is None:
logger.info("Workflow run archive download task expired before preparation: %s", download_id)
return None
if task.status != WorkflowRunArchiveDownloadStatus.PENDING:
logger.info("Skipping workflow run archive download task in %s state: %s", task.status, download_id)
return task
processing_task = self._mark_processing(task)
try:
archive_storage = self.archive_storage or get_archive_storage()
download_storage = self.download_storage or get_export_storage()
bundles = self._get_task_bundles(processing_task)
payload = self._build_zip_payload(archive_storage, processing_task, bundles)
storage_key = build_archive_download_storage_key(processing_task)
download_storage.put_object(storage_key, payload)
return self._mark_ready(processing_task, storage_key=storage_key, file_size_bytes=len(payload))
except Exception as exc:
logger.exception("Failed to prepare workflow run archive download %s", download_id)
return self._mark_failed(processing_task, error=str(exc))
def _get_task_bundles(self, task: WorkflowRunArchiveDownloadTask) -> list[WorkflowRunArchiveBundle]:
with self.session_factory() as session:
return _list_task_bundles(session, task)
def _build_zip_payload(
self,
storage: ArchiveStorage,
task: WorkflowRunArchiveDownloadTask,
bundles: Sequence[WorkflowRunArchiveBundle],
) -> bytes:
zip_root = f"workflow-run-logs-{task.year:04d}-{task.month:02d}"
csv_buffers: dict[str, io.BytesIO] = {}
csv_headers_written: set[str] = set()
for bundle in bundles:
object_prefix = _build_archive_bundle_object_prefix(task, bundle)
_, manifest = _load_and_validate_manifest(storage, task, bundle, object_prefix)
for table_name in sorted(manifest["tables"]):
entry = manifest["tables"][table_name]
object_key = entry["object_key"]
table_payload = storage.get_object(object_key)
_validate_table_payload(object_key=object_key, entry=entry, payload=table_payload)
csv_payload = _parquet_payload_to_csv(
table_payload,
include_header=table_name not in csv_headers_written,
)
if not csv_payload:
continue
csv_buffers.setdefault(table_name, io.BytesIO()).write(csv_payload)
csv_headers_written.add(table_name)
buffer = io.BytesIO()
with zipfile.ZipFile(buffer, mode="w", compression=zipfile.ZIP_DEFLATED) as archive:
for table_name, csv_buffer in sorted(csv_buffers.items()):
archive.writestr(f"{zip_root}/{table_name}.csv", csv_buffer.getvalue())
return buffer.getvalue()
def _mark_processing(self, task: WorkflowRunArchiveDownloadTask) -> WorkflowRunArchiveDownloadTask:
now = datetime.datetime.now(datetime.UTC)
processing_task = task.model_copy(
update={
"status": WorkflowRunArchiveDownloadStatus.PROCESSING,
"error": None,
"updated_at": now,
"started_at": task.started_at or now,
}
)
self.cache.save(processing_task)
return processing_task
def _mark_ready(
self,
task: WorkflowRunArchiveDownloadTask,
*,
storage_key: str,
file_size_bytes: int,
) -> WorkflowRunArchiveDownloadTask:
now = datetime.datetime.now(datetime.UTC)
ready_task = task.model_copy(
update={
"status": WorkflowRunArchiveDownloadStatus.READY,
"file_name": build_archive_download_file_name(task),
"storage_key": storage_key,
"file_size_bytes": file_size_bytes,
"error": None,
"updated_at": now,
"finished_at": now,
}
)
return self._save_terminal_task(task, ready_task)
def _mark_failed(self, task: WorkflowRunArchiveDownloadTask, *, error: str) -> WorkflowRunArchiveDownloadTask:
now = datetime.datetime.now(datetime.UTC)
failed_task = task.model_copy(
update={
"status": WorkflowRunArchiveDownloadStatus.FAILED,
"error": error,
"updated_at": now,
"finished_at": now,
}
)
return self._save_terminal_task(task, failed_task)
def _save_terminal_task(
self,
processing_task: WorkflowRunArchiveDownloadTask,
terminal_task: WorkflowRunArchiveDownloadTask,
) -> WorkflowRunArchiveDownloadTask:
with self.cache.lock(tenant_id=processing_task.tenant_id, download_id=processing_task.download_id):
current = self.cache.get(
tenant_id=processing_task.tenant_id,
download_id=processing_task.download_id,
)
if (
current is None
or current.status != WorkflowRunArchiveDownloadStatus.PROCESSING
or current.celery_task_id != processing_task.celery_task_id
):
return current or terminal_task
self.cache.save(terminal_task)
return terminal_task
def build_archive_download_file_name(task: WorkflowRunArchiveDownloadTask) -> str:
"""Return the browser download filename for one monthly archive."""
return f"workflow-run-logs-{task.year:04d}-{task.month:02d}.zip"
def build_archive_download_storage_key(task: WorkflowRunArchiveDownloadTask) -> str:
"""Return the deterministic object-store key for a prepared download ZIP."""
return (
f"{ARCHIVE_DOWNLOAD_ROOT_PREFIX}tenant_prefix={task.tenant_id[0].lower()}/tenant_id={task.tenant_id}/"
f"year={task.year:04d}/month={task.month:02d}/{task.download_id}.zip"
)
def _list_task_bundles(session: Session, task: WorkflowRunArchiveDownloadTask) -> list[WorkflowRunArchiveBundle]:
stmt = (
select(WorkflowRunArchiveBundle)
.where(
WorkflowRunArchiveBundle.tenant_id == task.tenant_id,
WorkflowRunArchiveBundle.year == task.year,
WorkflowRunArchiveBundle.month == task.month,
)
.order_by(WorkflowRunArchiveBundle.shard, WorkflowRunArchiveBundle.bundle_id)
)
indexed_bundles = list(session.scalars(stmt))
if task.bundle_refs:
requested_refs = [(ref.shard, ref.bundle_id) for ref in task.bundle_refs]
else:
requested_bundle_ids = set(task.bundle_ids)
requested_refs = [
(bundle.shard, bundle.bundle_id) for bundle in indexed_bundles if bundle.bundle_id in requested_bundle_ids
]
bundle_by_ref = {(bundle.shard, bundle.bundle_id): bundle for bundle in indexed_bundles}
missing_refs = [ref for ref in requested_refs if ref not in bundle_by_ref]
if missing_refs:
raise ValueError(f"archive bundle index is missing requested bundles: {missing_refs}")
bundles = [bundle_by_ref[ref] for ref in requested_refs]
if len(bundles) != task.bundle_count:
raise ValueError(f"archive bundle count changed: expected={task.bundle_count}, actual={len(bundles)}")
download_id = build_archive_download_id(
tenant_id=task.tenant_id,
year=task.year,
month=task.month,
bundle_refs=requested_refs,
)
if download_id != task.download_id:
raise ValueError("archive download id no longer matches indexed bundle set")
return bundles
def _build_archive_bundle_object_prefix(
task: WorkflowRunArchiveDownloadTask,
bundle: WorkflowRunArchiveBundle,
) -> str:
return (
f"{ARCHIVE_BUNDLE_ROOT_PREFIX}tenant_prefix={task.tenant_id[0].lower()}/tenant_id={task.tenant_id}/"
f"year={task.year:04d}/month={task.month:02d}/shard={bundle.shard}/bundle={bundle.bundle_id}"
)
def _load_and_validate_manifest(
storage: ArchiveStorage,
task: WorkflowRunArchiveDownloadTask,
bundle: WorkflowRunArchiveBundle,
object_prefix: str,
) -> tuple[bytes, ArchiveBundleManifest]:
manifest_key = f"{object_prefix}/{ARCHIVE_BUNDLE_MANIFEST_NAME}"
manifest_data = storage.get_object(manifest_key)
manifest = decode_archive_bundle_manifest(manifest_data)
_validate_manifest(task=task, bundle=bundle, manifest=manifest, object_prefix=object_prefix)
return manifest_data, manifest
def _validate_manifest(
*,
task: WorkflowRunArchiveDownloadTask,
bundle: WorkflowRunArchiveBundle,
manifest: ArchiveBundleManifest,
object_prefix: str,
) -> None:
if manifest["schema_version"] != ARCHIVE_BUNDLE_SCHEMA_VERSION:
raise ValueError(f"unsupported archive bundle schema version: {manifest['schema_version']}")
if manifest["archive_format"] != ARCHIVE_BUNDLE_FORMAT:
raise ValueError(f"unsupported archive bundle format: {manifest['archive_format']}")
if manifest["tenant_id"] != task.tenant_id:
raise ValueError(f"manifest tenant_id mismatch: expected={task.tenant_id}, actual={manifest['tenant_id']}")
if manifest["year"] != task.year:
raise ValueError(f"manifest year mismatch: expected={task.year}, actual={manifest['year']}")
if manifest["month"] != task.month:
raise ValueError(f"manifest month mismatch: expected={task.month}, actual={manifest['month']}")
if manifest["shard"] != bundle.shard:
raise ValueError(f"manifest shard mismatch: expected={bundle.shard}, actual={manifest['shard']}")
if manifest["bundle_id"] != bundle.bundle_id:
raise ValueError(f"manifest bundle_id mismatch: expected={bundle.bundle_id}, actual={manifest['bundle_id']}")
if manifest["object_prefix"] != object_prefix:
raise ValueError(
f"manifest object_prefix mismatch: expected={object_prefix}, actual={manifest['object_prefix']}"
)
if not manifest["tables"]:
raise ValueError("manifest tables must not be empty")
for table_name, raw_entry in manifest["tables"].items():
entry = cast(ArchiveBundleTableManifestEntry, raw_entry)
expected_object_key = f"{object_prefix}/{table_name}.parquet"
if entry["object_key"] != expected_object_key:
raise ValueError(
f"manifest object_key mismatch for {table_name}: "
f"expected={expected_object_key}, actual={entry['object_key']}"
)
def _validate_table_payload(
*,
object_key: str,
entry: ArchiveBundleTableManifestEntry,
payload: bytes,
) -> None:
if len(payload) != entry["size_bytes"]:
raise ValueError(f"archive object size mismatch for {object_key}")
checksum = hashlib.md5(payload).hexdigest()
if checksum != entry["checksum"]:
raise ValueError(f"archive object checksum mismatch for {object_key}")
def _parquet_payload_to_csv(payload: bytes, *, include_header: bool) -> bytes:
table = pq.read_table(io.BytesIO(payload))
if table.num_columns == 0:
return b""
for index, field in enumerate(table.schema):
if pa.types.is_string(field.type) or pa.types.is_large_string(field.type):
sanitized_column = pc.replace_substring_regex( # pyrefly: ignore[missing-attribute]
table.column(index),
pattern=_CSV_FORMULA_PREFIX_PATTERN,
replacement=r"'\0",
max_replacements=1,
)
table = table.set_column(index, field, sanitized_column)
buffer = io.BytesIO()
pa_csv.write_csv(
table,
buffer,
write_options=pa_csv.WriteOptions(include_header=include_header),
)
return buffer.getvalue()
@@ -0,0 +1,177 @@
"""Redis-backed temporary state for workflow-run archive downloads."""
import datetime
import hashlib
import json
import logging
from collections.abc import Sequence
from enum import StrEnum
from typing import Any
from pydantic import BaseModel, ConfigDict, Field
from extensions.ext_redis import RedisClientWrapper, redis_client
logger = logging.getLogger(__name__)
ARCHIVE_DOWNLOAD_FORMAT_VERSION = "v1"
DEFAULT_ARCHIVE_DOWNLOAD_TASK_TTL_SECONDS = 24 * 60 * 60
ARCHIVE_DOWNLOAD_TASK_LOCK_TIMEOUT_SECONDS = 30
_CACHE_KEY_PREFIX = "workflow_run_archive_download"
class WorkflowRunArchiveDownloadStatus(StrEnum):
"""Lifecycle state for an asynchronous archive download request."""
PENDING = "pending"
PROCESSING = "processing"
READY = "ready"
FAILED = "failed"
class WorkflowRunArchiveBundleRef(BaseModel):
"""Immutable object-store identity for one bundle included in a download task."""
model_config = ConfigDict(extra="forbid")
shard: str
bundle_id: str
class WorkflowRunArchiveDownloadTask(BaseModel):
"""Temporary Redis payload for a monthly archive download request."""
model_config = ConfigDict(extra="forbid")
download_id: str
tenant_id: str
requested_by: str
year: int = Field(ge=1)
month: int = Field(ge=1, le=12)
bundle_ids: list[str]
bundle_refs: list[WorkflowRunArchiveBundleRef] = Field(default_factory=list)
bundle_count: int = Field(ge=0)
archive_bytes: int = Field(ge=0)
status: WorkflowRunArchiveDownloadStatus
file_name: str | None = None
storage_key: str | None = None
file_size_bytes: int | None = Field(default=None, ge=0)
celery_task_id: str | None = None
error: str | None = None
created_at: datetime.datetime
updated_at: datetime.datetime
expires_at: datetime.datetime
started_at: datetime.datetime | None = None
finished_at: datetime.datetime | None = None
class WorkflowRunArchiveDownloadTaskCache:
"""Store ephemeral archive download task state in Redis with a TTL."""
_redis: RedisClientWrapper
def __init__(self, redis: RedisClientWrapper = redis_client) -> None:
self._redis = redis
def get(self, *, tenant_id: str, download_id: str) -> WorkflowRunArchiveDownloadTask | None:
raw = self._redis.get(self._cache_key(tenant_id=tenant_id, download_id=download_id))
if raw is None:
return None
data = raw.decode("utf-8") if isinstance(raw, bytes | bytearray) else raw
try:
return WorkflowRunArchiveDownloadTask.model_validate_json(data)
except ValueError:
logger.warning("Malformed workflow run archive download task cache entry: %s", download_id)
return None
def save(self, task: WorkflowRunArchiveDownloadTask) -> None:
ttl_seconds = self._ttl_seconds(task.expires_at)
self._redis.setex(
self._cache_key(tenant_id=task.tenant_id, download_id=task.download_id),
ttl_seconds,
task.model_dump_json(),
)
def lock(self, *, tenant_id: str, download_id: str) -> Any:
return self._redis.lock(
f"{self._cache_key(tenant_id=tenant_id, download_id=download_id)}:lock",
timeout=ARCHIVE_DOWNLOAD_TASK_LOCK_TIMEOUT_SECONDS,
blocking_timeout=ARCHIVE_DOWNLOAD_TASK_LOCK_TIMEOUT_SECONDS,
)
def delete(self, *, tenant_id: str, download_id: str) -> None:
self._redis.delete(self._cache_key(tenant_id=tenant_id, download_id=download_id))
@staticmethod
def _cache_key(*, tenant_id: str, download_id: str) -> str:
return f"{_CACHE_KEY_PREFIX}:{tenant_id}:{download_id}"
@staticmethod
def _ttl_seconds(expires_at: datetime.datetime) -> int:
expires_at_utc = expires_at if expires_at.tzinfo else expires_at.replace(tzinfo=datetime.UTC)
remaining = expires_at_utc - datetime.datetime.now(datetime.UTC)
return max(int(remaining.total_seconds()), 1)
def build_pending_archive_download_task(
*,
tenant_id: str,
requested_by: str,
year: int,
month: int,
bundle_ids: Sequence[str],
bundle_refs: Sequence[tuple[str, str]] = (),
archive_bytes: int,
download_id: str,
ttl_seconds: int = DEFAULT_ARCHIVE_DOWNLOAD_TASK_TTL_SECONDS,
now: datetime.datetime | None = None,
) -> WorkflowRunArchiveDownloadTask:
"""Create the Redis payload stored when the console starts an archive download."""
created_at = now or datetime.datetime.now(datetime.UTC)
if created_at.tzinfo is None:
created_at = created_at.replace(tzinfo=datetime.UTC)
normalized_bundle_ids = list(bundle_ids)
normalized_bundle_refs = [
WorkflowRunArchiveBundleRef(shard=shard, bundle_id=bundle_id) for shard, bundle_id in bundle_refs
]
return WorkflowRunArchiveDownloadTask(
download_id=download_id,
tenant_id=tenant_id,
requested_by=requested_by,
year=year,
month=month,
bundle_ids=normalized_bundle_ids,
bundle_refs=normalized_bundle_refs,
bundle_count=len(normalized_bundle_ids),
archive_bytes=archive_bytes,
status=WorkflowRunArchiveDownloadStatus.PENDING,
created_at=created_at,
updated_at=created_at,
expires_at=created_at + datetime.timedelta(seconds=ttl_seconds),
)
def build_archive_download_id(
*,
tenant_id: str,
year: int,
month: int,
bundle_refs: Sequence[tuple[str, str]],
download_format_version: str = ARCHIVE_DOWNLOAD_FORMAT_VERSION,
) -> str:
"""Build a stable id for the exact archive download content."""
if not bundle_refs:
raise ValueError("bundle_refs must not be empty")
normalized_refs = sorted(f"{shard}:{bundle_id}" for shard, bundle_id in bundle_refs)
payload = json.dumps(
{
"tenant_id": tenant_id,
"year": year,
"month": month,
"bundle_refs": normalized_refs,
"download_format_version": download_format_version,
},
sort_keys=True,
separators=(",", ":"),
)
return hashlib.sha256(payload.encode("utf-8")).hexdigest()[:32]
@@ -0,0 +1,298 @@
"""
Console-facing workflow-run archive queries.
The object store remains the recoverable archive source of truth. This module only reads the DB bundle index and writes
temporary Redis download-task state, so console requests never list R2 online.
"""
import datetime
import logging
import uuid
from collections.abc import Callable
from dataclasses import dataclass
from sqlalchemy import select
from sqlalchemy.orm import Session
from models.workflow import WorkflowRunArchiveBundle
from services.retention.workflow_run.archive_download_task_cache import (
WorkflowRunArchiveDownloadStatus,
WorkflowRunArchiveDownloadTask,
WorkflowRunArchiveDownloadTaskCache,
build_archive_download_id,
build_pending_archive_download_task,
)
logger = logging.getLogger(__name__)
ArchiveDownloadTaskDispatcher = Callable[
[WorkflowRunArchiveDownloadTask, WorkflowRunArchiveDownloadTaskCache],
WorkflowRunArchiveDownloadTask,
]
@dataclass(frozen=True)
class WorkflowRunArchiveMonth:
"""Aggregated archive metadata for one tenant/month."""
year: int
month: int
bundle_count: int
workflow_run_count: int
row_count: int
archive_bytes: int
latest_archived_at: datetime.datetime
download_task: WorkflowRunArchiveDownloadTask | None
@dataclass(frozen=True)
class WorkflowRunArchiveSummary:
"""Top-level archive totals shown on the console page."""
archived_month_count: int
workflow_run_count: int
archive_bytes: int
latest_archived_at: datetime.datetime | None
@dataclass(frozen=True)
class WorkflowRunArchiveList:
"""Console response model before controller serialization."""
summary: WorkflowRunArchiveSummary
months: list[WorkflowRunArchiveMonth]
class WorkflowRunArchiveNotFoundError(Exception):
"""Raised when no archive bundles exist for a requested tenant/month."""
class WorkflowRunArchiveDownloadTaskNotFoundError(Exception):
"""Raised when the temporary Redis task has expired or never existed."""
class WorkflowRunArchiveDownloadNotReadyError(Exception):
"""Raised when a cached download task has not produced a file yet."""
def list_workflow_run_archives(
session: Session,
tenant_id: str,
*,
cache: WorkflowRunArchiveDownloadTaskCache | None = None,
) -> WorkflowRunArchiveList:
"""Return monthly archive metadata for one tenant from the DB bundle index."""
stmt = (
select(WorkflowRunArchiveBundle)
.where(WorkflowRunArchiveBundle.tenant_id == tenant_id)
.order_by(
WorkflowRunArchiveBundle.year.desc(),
WorkflowRunArchiveBundle.month.desc(),
WorkflowRunArchiveBundle.shard,
WorkflowRunArchiveBundle.bundle_id,
)
)
month_bundles: dict[tuple[int, int], list[WorkflowRunArchiveBundle]] = {}
for bundle in session.scalars(stmt):
month_bundles.setdefault((bundle.year, bundle.month), []).append(bundle)
task_cache = cache or WorkflowRunArchiveDownloadTaskCache()
months: list[WorkflowRunArchiveMonth] = []
for (year, month), bundles in month_bundles.items():
bundle_refs = [(bundle.shard, bundle.bundle_id) for bundle in bundles]
months.append(
WorkflowRunArchiveMonth(
year=year,
month=month,
bundle_count=len(bundles),
workflow_run_count=sum(bundle.workflow_run_count for bundle in bundles),
row_count=sum(bundle.row_count for bundle in bundles),
archive_bytes=sum(bundle.archive_bytes for bundle in bundles),
latest_archived_at=max(bundle.archived_at for bundle in bundles),
download_task=_get_cached_month_download_task(
task_cache,
tenant_id=tenant_id,
year=year,
month=month,
bundle_refs=bundle_refs,
),
)
)
latest_archived_at = max((month.latest_archived_at for month in months), default=None)
return WorkflowRunArchiveList(
summary=WorkflowRunArchiveSummary(
archived_month_count=len(months),
workflow_run_count=sum(month.workflow_run_count for month in months),
archive_bytes=sum(month.archive_bytes for month in months),
latest_archived_at=latest_archived_at,
),
months=months,
)
def _get_cached_month_download_task(
cache: WorkflowRunArchiveDownloadTaskCache,
*,
tenant_id: str,
year: int,
month: int,
bundle_refs: list[tuple[str, str]],
) -> WorkflowRunArchiveDownloadTask | None:
if not bundle_refs:
return None
download_id = build_archive_download_id(
tenant_id=tenant_id,
year=year,
month=month,
bundle_refs=bundle_refs,
)
try:
return cache.get(tenant_id=tenant_id, download_id=download_id)
except Exception:
logger.warning("Failed to read cached workflow run archive download task: %s", download_id, exc_info=True)
return None
def create_workflow_run_archive_download_task(
session: Session,
*,
tenant_id: str,
requested_by: str,
year: int,
month: int,
cache: WorkflowRunArchiveDownloadTaskCache | None = None,
dispatcher: ArchiveDownloadTaskDispatcher | None = None,
) -> WorkflowRunArchiveDownloadTask:
"""
Create or return the idempotent Redis task for downloading one tenant/month archive.
The task identity is based on the exact ordered bundle set currently indexed for the month. If the month receives a
new bundle later, the next request gets a different download id and prepares a fresh file.
"""
bundles = _list_archive_bundles(session, tenant_id=tenant_id, year=year, month=month)
if not bundles:
raise WorkflowRunArchiveNotFoundError(f"Workflow run archive not found: {year:04d}-{month:02d}")
bundle_refs = [(bundle.shard, bundle.bundle_id) for bundle in bundles]
download_id = build_archive_download_id(
tenant_id=tenant_id,
year=year,
month=month,
bundle_refs=bundle_refs,
)
task = build_pending_archive_download_task(
tenant_id=tenant_id,
requested_by=requested_by,
year=year,
month=month,
bundle_ids=[bundle.bundle_id for bundle in bundles],
bundle_refs=bundle_refs,
archive_bytes=sum(bundle.archive_bytes for bundle in bundles),
download_id=download_id,
)
task_cache = cache or WorkflowRunArchiveDownloadTaskCache()
dispatch = dispatcher or _dispatch_workflow_run_archive_download_task
with task_cache.lock(tenant_id=tenant_id, download_id=download_id):
existing = task_cache.get(tenant_id=tenant_id, download_id=download_id)
if existing is None or existing.status == WorkflowRunArchiveDownloadStatus.FAILED:
task_to_queue = task
elif existing.status == WorkflowRunArchiveDownloadStatus.PENDING and not existing.celery_task_id:
task_to_queue = existing
else:
return existing
queued_task = task_to_queue.model_copy(
update={"celery_task_id": uuid.uuid4().hex, "updated_at": datetime.datetime.now(datetime.UTC)}
)
task_cache.save(queued_task)
return dispatch(queued_task, task_cache)
def get_workflow_run_archive_download_task(
*,
tenant_id: str,
download_id: str,
cache: WorkflowRunArchiveDownloadTaskCache | None = None,
) -> WorkflowRunArchiveDownloadTask:
"""Return a cached archive download task or raise when the TTL has expired."""
task_cache = cache or WorkflowRunArchiveDownloadTaskCache()
task = task_cache.get(tenant_id=tenant_id, download_id=download_id)
if task is None:
raise WorkflowRunArchiveDownloadTaskNotFoundError(f"Workflow run archive download not found: {download_id}")
return task
def get_ready_workflow_run_archive_download_task(
*,
tenant_id: str,
download_id: str,
cache: WorkflowRunArchiveDownloadTaskCache | None = None,
) -> WorkflowRunArchiveDownloadTask:
"""Return a ready cached archive download task or raise when the file is not available."""
task = get_workflow_run_archive_download_task(tenant_id=tenant_id, download_id=download_id, cache=cache)
if task.status != WorkflowRunArchiveDownloadStatus.READY or not task.storage_key or not task.file_name:
raise WorkflowRunArchiveDownloadNotReadyError(f"Workflow run archive download is not ready: {download_id}")
return task
def _list_archive_bundles(
session: Session,
*,
tenant_id: str,
year: int,
month: int,
) -> list[WorkflowRunArchiveBundle]:
stmt = (
select(WorkflowRunArchiveBundle)
.where(
WorkflowRunArchiveBundle.tenant_id == tenant_id,
WorkflowRunArchiveBundle.year == year,
WorkflowRunArchiveBundle.month == month,
)
.order_by(WorkflowRunArchiveBundle.shard, WorkflowRunArchiveBundle.bundle_id)
)
return list(session.scalars(stmt))
def _dispatch_workflow_run_archive_download_task(
task: WorkflowRunArchiveDownloadTask,
cache: WorkflowRunArchiveDownloadTaskCache,
) -> WorkflowRunArchiveDownloadTask:
"""
Enqueue background ZIP preparation after the caller atomically claimed and cached the Celery id.
"""
from tasks.workflow_run_archive_download_tasks import prepare_workflow_run_archive_download_task
celery_task_id = task.celery_task_id
if celery_task_id is None:
raise ValueError("celery_task_id is required before dispatch")
try:
prepare_workflow_run_archive_download_task.apply_async(
args=(task.tenant_id, task.download_id),
task_id=celery_task_id,
)
except Exception:
failure_time = datetime.datetime.now(datetime.UTC)
failed_task = task.model_copy(
update={
"status": WorkflowRunArchiveDownloadStatus.FAILED,
"error": "Failed to enqueue archive download task.",
"updated_at": failure_time,
"finished_at": failure_time,
}
)
with cache.lock(tenant_id=task.tenant_id, download_id=task.download_id):
current = cache.get(tenant_id=task.tenant_id, download_id=task.download_id)
if (
current is not None
and current.status == WorkflowRunArchiveDownloadStatus.PENDING
and current.celery_task_id == celery_task_id
):
cache.save(failed_task)
current = failed_task
logger.exception("Failed to enqueue workflow run archive download task %s", task.download_id)
return current or failed_task
return task
@@ -5,8 +5,8 @@ This service archives workflow run logs for paid plan users older than the confi
90 days) to S3-compatible storage.
Archive V2 writes bundle-level Parquet objects. A bundle contains many workflow runs and their related table rows.
Bundle metadata lives in the object-store manifest instead of a database table, so archive/delete/restore does not move
the large-table retention problem into another OLTP table.
Bundle metadata lives in the object-store manifest as the recoverable source of truth. Completed bundles are also
mirrored into a small database index so console listing and download jobs do not list object storage online.
Archive campaigns should use fixed absolute UTC windows for every tenant-prefix/shard execution. Relative windows are
evaluated at process start and are not safe for multi-day rollout because each command would scan a different window.
@@ -64,6 +64,12 @@ from repositories.api_workflow_node_execution_repository import DifyAPIWorkflowN
from repositories.api_workflow_run_repository import APIWorkflowRunRepository
from repositories.sqlalchemy_workflow_trigger_log_repository import SQLAlchemyWorkflowTriggerLogRepository
from services.billing_service import BillingService
from services.retention.workflow_run.archive_bundle_index import (
ArchiveBundleManifest,
ArchiveBundleTableManifestEntry,
decode_archive_bundle_manifest,
upsert_archive_bundle_index_from_manifest,
)
from services.retention.workflow_run.constants import (
ARCHIVE_BUNDLE_FORMAT,
ARCHIVE_BUNDLE_INDEX_NAME,
@@ -634,6 +640,7 @@ class WorkflowRunArchiver:
raise ArchiveStorageNotConfiguredError("Archive storage not configured")
if storage.object_exists(self._get_manifest_object_key(identity)):
self._write_bundle_index(storage, identity)
self._sync_existing_bundle_index(session, storage, identity)
result.success = True
result.skipped = True
result.error = "bundle already archived"
@@ -657,6 +664,7 @@ class WorkflowRunArchiver:
result.run_count = len(runs)
if storage.object_exists(self._get_manifest_object_key(identity)):
self._write_bundle_index(storage, identity)
self._sync_existing_bundle_index(session, storage, identity)
result.success = True
result.skipped = True
result.error = "filtered bundle already archived"
@@ -688,6 +696,8 @@ class WorkflowRunArchiver:
storage.put_object(self._get_table_object_key(identity, table_name), payload)
storage.put_object(self._get_manifest_object_key(identity), manifest_data)
self._merge_bundle_manifest_into_index(storage, identity, [run.id for run in runs])
manifest = decode_archive_bundle_manifest(manifest_data)
upsert_archive_bundle_index_from_manifest(session, manifest, len(manifest_data))
session.commit()
logger.info(
@@ -714,6 +724,23 @@ class WorkflowRunArchiver:
result.elapsed_time = time.time() - start_time
return result
def _sync_existing_bundle_index(
self,
session: Session,
storage: ArchiveStorage,
identity: ArchiveBundleIdentity,
) -> None:
"""Best-effort DB index sync for a bundle whose manifest already exists in archive storage."""
manifest_key = self._get_manifest_object_key(identity)
try:
manifest_data = storage.get_object(manifest_key)
manifest = decode_archive_bundle_manifest(manifest_data)
upsert_archive_bundle_index_from_manifest(session, manifest, len(manifest_data))
session.commit()
except Exception:
session.rollback()
logger.warning("Failed to sync workflow archive bundle index for %s", manifest_key, exc_info=True)
def _lock_runs_for_archive(
self,
session: Session,
@@ -841,9 +868,9 @@ class WorkflowRunArchiver:
identity: ArchiveBundleIdentity,
runs: Sequence[WorkflowRun],
table_stats: list[TableStats],
) -> ArchiveManifestDict:
) -> ArchiveBundleManifest:
"""Generate a manifest for the archived workflow run bundle."""
tables: dict[str, TableStatsManifestEntry] = {
tables: dict[str, ArchiveBundleTableManifestEntry] = {
stat.table_name: {
"row_count": stat.row_count,
"checksum": stat.checksum,
@@ -1,9 +1,10 @@
"""
Maintain V2 workflow-run archive bundles.
Archive V2 keeps bundle metadata in object-store manifests, not in a database table. This module discovers bundles by
listing `manifest.json` objects, uses object-store marker files for delete/restore state, and only touches the database
for source-table validation, deletion, and restoration.
Archive V2 keeps object-store manifests as the recoverable bundle source of truth. This maintenance module still
discovers delete/restore targets by listing `manifest.json` objects and uses object-store marker files for
delete/restore state. The separate database bundle index is intended for console listing and download jobs, not as the
source of truth for destructive maintenance.
Each bundle is processed in its own database transaction. A failed bundle leaves source rows unchanged unless the
transaction has already committed; marker handling makes the next run able to reconcile the common committed-but-marker