Source code for ccat_data_transfer.data_integrity_manager

import os
import json
import sqlalchemy
from ccat_ops_db import models
from .database import DatabaseConnection
from sqlalchemy.orm import Session
from typing import Tuple, Optional

from .setup_celery_app import app, make_celery_task
from .utils import (
    create_local_folder,
    calculate_checksum,
    unpack_local,
    get_redis_connection,
    safe_join,
)
from .logging_utils import get_structured_logger
from .exceptions import (
    UnpackError,
    ChecksumVerificationError,
    ArchiveCorruptionError,
    OperationNotFoundError,
)
from .queue_discovery import route_task_by_location
from .operation_types import OperationType

# Use only task loggers
logger = get_structured_logger(__name__)

redis_ = get_redis_connection()


# The Celery task arg is the UnpackOperation id (the sibling Operation row that
# is the source of truth for the unpack lifecycle, ADR-0003), NOT the legacy
# DataTransfer id. The base task (#153) keys mark_in_progress / get_retry_count /
# recovery off this id. Unpack consumes the transfer's destination copy and
# produces the unpacked raw-data-package copies (copy-anchored lineage).


def _data_transfer_for_unpack(
    session: Session, unpack_operation: models.UnpackOperation
) -> Optional[models.DataTransfer]:
    """Resolve the DataTransfer that delivered this package to THIS unpack
    operation's destination.

    The UnpackOperation is anchored per destination
    (data_transfer_package_id, destination_location_id); the rich unpack
    execution context (raw packages, paths) is read from the matching
    DataTransfer during the additive transition. Resolution is DETERMINISTIC on
    the destination (not ``.first()`` over all of a package's transfers) so the
    correct destination's copy is verified and its produced copy registered at
    the right location. There is one DataTransfer intent per (package,
    destination) slot, so the destination anchor alone identifies it — the
    transfer-completed gate lives on the sibling TransferOperation, and
    DataTransfer.status was dropped in ops-db#95."""
    return (
        session.query(models.DataTransfer)
        .filter_by(
            data_transfer_package_id=unpack_operation.data_transfer_package_id,
            destination_location_id=unpack_operation.destination_location_id,
        )
        .first()
    )


[docs] class UnpackTask(make_celery_task()): """Base class for unpacking tasks. Keyed on the UnpackOperation id (ADR-0003): retry counting falls to the base uniform ``Operation.retry_count`` reader (the legacy ``unpack_retry_count`` override is gone), and failure/IN_PROGRESS hooks act on the operation row. """
[docs] def __init__(self): super().__init__() self.operation_type = models.OperationKind.UNPACK.value self.max_retries = 3
[docs] def mark_in_progress(self, session, unpack_operation_id): """Flip the UnpackOperation to IN_PROGRESS at task start (the single seam from the base task, ADR-0003).""" op = session.query(models.UnpackOperation).get(unpack_operation_id) if op: op.status = models.Status.IN_PROGRESS
[docs] def reset_state_on_failure(self, session, unpack_operation_id, exc): """Reset the UnpackOperation for retry (own status + own retry_count).""" op = session.query(models.UnpackOperation).get(unpack_operation_id) if op: op.status = models.Status.PENDING data_transfer = _data_transfer_for_unpack(session, op) if data_transfer: for ( raw_data_package ) in data_transfer.data_transfer_package.raw_data_packages: raw_data_package.state = models.PackageState.TRANSFERRING op.failure_error_message = None # error_context mirrors failure_error_message: NULLed on reset; the # durable trail is kept in OperationFailureEvent (#117). op.error_context = None op.retry_count += 1 logger.info( "reset unpack for retry", unpack_operation_id=unpack_operation_id, retry_count=op.retry_count, )
[docs] def mark_permanent_failure(self, session, unpack_operation_id, exc): """Mark the UnpackOperation permanently FAILED.""" op = session.query(models.UnpackOperation).get(unpack_operation_id) if op: op.status = models.Status.FAILED data_transfer = _data_transfer_for_unpack(session, op) if data_transfer: for ( raw_data_package ) in data_transfer.data_transfer_package.raw_data_packages: raw_data_package.state = models.PackageState.FAILED op.failure_error_message = str(exc) # Cache the latest unpack breadcrumb on the row for the UI (#117). op.error_context = self._current_error_context logger.info( "marked unpack as permanently failed", unpack_operation_id=unpack_operation_id, error=str(exc), )
[docs] def get_operation_info(self, args, kwargs): """Get additional context for unpack tasks.""" if not args or len(args) == 0: return {} with self.session_scope() as session: try: op = session.query(models.UnpackOperation).get(args[0]) if op: data_transfer = _data_transfer_for_unpack(session, op) if data_transfer: return { "destination_location": data_transfer.destination_location.name, "package_id": str(op.data_transfer_package_id), } except Exception as e: logger.error("failed to get unpack info", error=str(e)) return {}
@app.task( base=UnpackTask, name="ccat:data_transfer:unpack_data_transfer_package", bind=True, ) def unpack_data_transfer_package( self, unpack_operation_id: int, session: Session = None, ) -> bool: """ Unpack a data transfer package and verify its contents using dynamic queue routing. Parameters ---------- self : celery.Task The bound Celery task instance. unpack_operation_id : int The ID of the UnpackOperation to process (the lifecycle source of truth). session : Session, optional An existing database session to use. If None, a new session will be created. Returns ------- bool True if the unpacking and verification process was successful, False otherwise. """ if session is None: with self.session_scope() as session: return _unpack_data_transfer_package_internal(session, unpack_operation_id) else: return _unpack_data_transfer_package_internal(session, unpack_operation_id) def _cleanup_corrupted_transfer( session: Session, data_transfer: models.DataTransfer ) -> None: """ Clean up a corrupted transfer by removing files and resetting database state. This function: 1. Removes the corrupted archive file from destination 2. Schedules deletion of the source archive file on the primary archive 3. Deregisters the raw data packages from the data transfer package 4. Deletes the data transfer package from the database 5. Raises an ArchiveCorruptionError to trigger proper error handling Raises ------ ArchiveCorruptionError After cleanup is complete, to trigger proper error handling in the task system """ logger.info( "cleaning up corrupted transfer", transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, ) # Get the destination path using the new location system destination_path = _get_location_path( data_transfer.destination_location, data_transfer.data_transfer_package ) # Find the physical copy for the destination location dest_physical_copy = ( session.query(models.DataTransferPackagePhysicalCopy) .filter_by( data_transfer_package_id=data_transfer.data_transfer_package_id, data_location_id=data_transfer.destination_location_id, status=models.PhysicalCopyStatus.PRESENT, ) .first() ) if dest_physical_copy: try: # Mark for deletion and schedule task dest_physical_copy.deletion_status = models.Status.SCHEDULED session.add(dest_physical_copy) session.commit() redis_.publish( "transfer:overview", json.dumps( {"type": "corrupted_transfer_cleanup", "data": data_transfer.id} ), ) # Schedule deletion task using dynamic queue routing from .deletion_manager import delete_physical_copy queue_name = route_task_by_location( OperationType.DELETION, data_transfer.destination_location ) result = delete_physical_copy.apply_async( args=[dest_physical_copy.id], queue=queue_name, ) # Wait for task completion with a timeout result.get(timeout=300) # 5 minute timeout logger.info( "deleted corrupted archive from destination", physical_copy_id=dest_physical_copy.id, path=destination_path, ) destination_deletion_successful = True except Exception as e: logger.error( "failed to delete corrupted archive from destination", physical_copy_id=dest_physical_copy.id, path=destination_path, error=str(e), ) destination_deletion_successful = False else: logger.warning( "could not find physical copy for destination location", path=destination_path, transfer_package_id=data_transfer.data_transfer_package_id, location_id=data_transfer.destination_location_id, ) destination_deletion_successful = False # Schedule deletion of source archive file if it's a secondary transfer In the new # architecture, we determine if it's a secondary transfer by checking if it is a # transfer between two LTA sites lta_sites = ( session.query(models.Site) .join(models.DataLocation) .filter(models.DataLocation.location_type == models.LocationType.LONG_TERM_ARCHIVE) .all() ) lta_sites_ids = [site.id for site in lta_sites] is_secondary_transfer = ( data_transfer.origin_location.site_id in lta_sites_ids and data_transfer.destination_location.site_id in lta_sites_ids ) if is_secondary_transfer: source_path = _get_location_path( data_transfer.origin_location, data_transfer.data_transfer_package ) # Find or create the physical copy for the source location source_physical_copy = ( session.query(models.DataTransferPackagePhysicalCopy) .filter_by( data_transfer_package_id=data_transfer.data_transfer_package_id, data_location_id=data_transfer.origin_location_id, status=models.PhysicalCopyStatus.PRESENT, ) .first() ) if source_physical_copy: # Mark for deletion and schedule task source_physical_copy.deletion_status = models.Status.SCHEDULED session.add(source_physical_copy) session.commit() redis_.publish( "transfer:overview", json.dumps( {"type": "corrupted_transfer_cleanup", "data": data_transfer.id} ), ) # Schedule deletion task using dynamic queue routing from .deletion_manager import delete_physical_copy try: queue_name = route_task_by_location( OperationType.DELETION, data_transfer.origin_location ) result = delete_physical_copy.apply_async( args=[source_physical_copy.id], queue=queue_name, ) # Wait for task completion with a timeout result.get(timeout=300) # 5 minute timeout logger.info( "deleted corrupted archive from source", physical_copy_id=source_physical_copy.id, path=source_path, ) primary_deletion_successful = True except Exception as e: logger.error( "failed to delete corrupted archive from source", physical_copy_id=source_physical_copy.id, path=source_path, error=str(e), ) primary_deletion_successful = False # Continue with cleanup even if deletion fails # The deletion manager will retry the deletion later else: logger.warning( "could not find physical copy for source location", path=source_path, transfer_package_id=data_transfer.data_transfer_package_id, location_id=data_transfer.origin_location_id, ) primary_deletion_successful = False else: # This is a primary transfer, no source cleanup needed primary_deletion_successful = True # Deregister raw data packages from the data transfer package package = data_transfer.data_transfer_package raw_data_packages = package.raw_data_packages.copy() for raw_package in raw_data_packages: raw_package.data_transfer_package_id = None session.add(raw_package) # Delete the data transfer package # we can only delete the package if both deletions were successful # otherwise the automatic cleanup will retry the deletion later and needs the # information from the package if primary_deletion_successful and destination_deletion_successful: session.delete(package) else: logger.info( "keeping data transfer package for retry", transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, primary_deletion_successful=primary_deletion_successful, destination_deletion_successful=destination_deletion_successful, ) # Commit the cleanup changes session.commit() redis_.publish( "transfer:overview", json.dumps({"type": "corrupted_transfer_cleanup", "data": data_transfer.id}), ) logger.info( "completed cleanup of corrupted transfer", transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, ) # Raise the error to trigger proper error handling in the task system raise ArchiveCorruptionError( "Archive corruption detected - transfer package will be recreated", archive_path=destination_path, transfer_id=data_transfer.id, ) def _unpack_data_transfer_package_internal( session: Session, unpack_operation_id: int ) -> None: """ Internal function to unpack a data transfer package and verify its contents. Parameters ---------- session : sqlalchemy.orm.Session The database session. unpack_operation_id : int The ID of the UnpackOperation to process (lifecycle source of truth). Raises ------ FileNotFoundError If the archive file or destination directory is not found. UnpackError If the unpacking process fails. ChecksumVerificationError If the checksum verification fails. ArchiveCorruptionError If the archive is corrupted or incomplete. """ unpack_operation = session.query(models.UnpackOperation).get(unpack_operation_id) if unpack_operation is None: # Non-retryable: a missing/dangling id never reappears, so retrying only # loops forever. A bare ValueError defaulted to retryable in should_retry, # which is what produced the infinite unpack retry. raise OperationNotFoundError( f"Unpack operation not found: {unpack_operation_id}", operation_id=unpack_operation_id, ) data_transfer = _data_transfer_for_unpack(session, unpack_operation) if data_transfer is None: raise ValueError( "No completed DataTransfer to unpack for unpack operation " f"{unpack_operation_id}" ) archive_path, destination = _get_paths(data_transfer) logger.info( "unpack started", unpack_operation_id=unpack_operation_id, transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, path=archive_path, ) try: create_local_folder(destination) except OSError as e: raise FileNotFoundError(f"Failed to create destination directory: {str(e)}") try: success, error = unpack_local(archive_path, destination) if not success: raise UnpackError(f"Failed to unpack archive: {error}") except ArchiveCorruptionError as e: logger.error( "archive corruption detected", unpack_operation_id=unpack_operation_id, transfer_id=data_transfer.id, error=str(e), ) _cleanup_corrupted_transfer(session, data_transfer) raise except Exception as e: raise UnpackError(f"Error during unpacking: {str(e)}") try: if not _verify_checksums(data_transfer, destination): raise ChecksumVerificationError("Checksum verification failed") except Exception as e: raise ChecksumVerificationError(f"Error during checksum verification: {str(e)}") # If we get here, everything succeeded _update_data_transfer_status(session, unpack_operation, data_transfer, True) def _get_data_transfer( session: Session, data_transfer_id: int ) -> Optional[models.DataTransfer]: """Retrieve the DataTransfer object from the database.""" try: return session.query(models.DataTransfer).filter_by(id=data_transfer_id).one() except sqlalchemy.orm.exc.NoResultFound: logger.error("data transfer not found", transfer_id=data_transfer_id) except sqlalchemy.orm.exc.MultipleResultsFound: logger.error("multiple data transfers found", transfer_id=data_transfer_id) return None def _get_paths(data_transfer: models.DataTransfer) -> Tuple[str, str]: """Get the archive path and destination for unpacking using the new location system.""" archive_path = _get_location_path( data_transfer.destination_location, data_transfer.data_transfer_package ) # Get the raw data path from the destination location if isinstance(data_transfer.destination_location, models.DiskDataLocation): # For disk locations, use a subdirectory for raw data # This follows the pattern from the configuration files destination = os.path.join( data_transfer.destination_location.path, ) else: # For non-disk locations, we need to determine the appropriate path # This might need to be configurable per location type raise ValueError( f"Unpacking not yet supported for storage type: {data_transfer.destination_location.storage_type}" ) return archive_path, destination def _get_location_path( data_location: models.DataLocation, data_transfer_package: models.DataTransferPackage, ) -> str: """ Get the full path for a data transfer package at a specific location. Parameters ---------- data_location : models.DataLocation The data location. data_transfer_package : models.DataTransferPackage The data transfer package. Returns ------- str The full path to the package at this location. """ if isinstance(data_location, models.DiskDataLocation): return safe_join(data_location.path, data_transfer_package.relative_path) elif isinstance(data_location, models.S3DataLocation): return f"{data_location.prefix}{data_transfer_package.relative_path}" elif isinstance(data_location, models.TapeDataLocation): return safe_join( data_location.mount_path, data_transfer_package.relative_path ) else: raise ValueError(f"Unsupported storage type: {data_location.storage_type}") def _verify_checksums(data_transfer: models.DataTransfer, destination: str) -> bool: """Verify the checksums of the unpacked files.""" for raw_data_package in data_transfer.data_transfer_package.raw_data_packages: extracted_file_path = safe_join(destination, raw_data_package.relative_path) local_checksum = calculate_checksum(extracted_file_path) if local_checksum is None: logger.error( "checksum calculation failed", transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, path=extracted_file_path, ) return False if raw_data_package.checksum != local_checksum: logger.error( "checksum mismatch", transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, path=extracted_file_path, expected_checksum=raw_data_package.checksum, actual_checksum=local_checksum, ) return False logger.debug( "checksum verified", transfer_id=data_transfer.id, path=extracted_file_path, ) return True def _update_data_transfer_status( session: Session, unpack_operation: models.UnpackOperation, data_transfer: models.DataTransfer, success: bool, ) -> None: """Update the unpack operation in the DB (lifecycle source of truth, ADR-0003).""" try: if success: unpack_operation.status = models.Status.COMPLETED # Consume the transferred destination copy: unpack's input in the # copy-anchored lineage (ADR-0003), produced by the sibling transfer. consumed_copy = ( session.query(models.DataTransferPackagePhysicalCopy) .filter_by( data_transfer_package_id=data_transfer.data_transfer_package_id, data_location_id=data_transfer.destination_location_id, status=models.PhysicalCopyStatus.PRESENT, ) .first() ) if consumed_copy is not None: unpack_operation.consumed_copies.append(consumed_copy) # Add physical copies for each raw data package for ( raw_data_package ) in data_transfer.data_transfer_package.raw_data_packages: # Get the raw data path from the destination location if isinstance( data_transfer.destination_location, models.DiskDataLocation ): # For disk locations, use a subdirectory for raw data # This follows the pattern from the configuration files _ = os.path.join( data_transfer.destination_location.path, "raw_data_packages", raw_data_package.relative_path, ) else: # For non-disk locations, we need to determine the appropriate path raise ValueError( f"Unpacking not yet supported for storage type: {data_transfer.destination_location.storage_type}" ) # Create the physical copy physical_copy = models.RawDataPackagePhysicalCopy( raw_data_package=raw_data_package, data_location=data_transfer.destination_location, checksum=raw_data_package.checksum, ) # Add to the session explicitly session.add(physical_copy) # Also add to the relationship for consistency raw_data_package.physical_copies.append(physical_copy) # Produce the unpacked raw-data-package copy (unpack's output). unpack_operation.produced_copies.append(physical_copy) logger.debug( "raw data package registered at location", transfer_id=data_transfer.id, raw_package_id=raw_data_package.id, ) logger.info( "unpack succeeded", unpack_operation_id=unpack_operation.id, transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, package_count=len( data_transfer.data_transfer_package.raw_data_packages ), ) # Commit the transaction session.commit() # Publish to Redis after successful commit redis_.publish( "transfer:overview", json.dumps({"type": "unpack_completed", "data": data_transfer.id}), ) else: # Handle failure case - mark as failed logger.error( "unpack failed", unpack_operation_id=unpack_operation.id, transfer_id=data_transfer.id, transfer_package_id=data_transfer.data_transfer_package_id, ) unpack_operation.status = models.Status.FAILED session.commit() # Publish to Redis after commit redis_.publish( "transfer:overview", json.dumps({"type": "unpack_failed", "data": data_transfer.id}), ) except Exception as e: # Log the error and rollback the transaction logger.error( "failed to update data transfer status", unpack_operation_id=unpack_operation.id, transfer_id=data_transfer.id, error=str(e), ) session.rollback() raise def _add_raw_data_packages_to_location(data_transfer: models.DataTransfer) -> None: """Add raw data packages to the destination location.""" for raw_data_package in data_transfer.data_transfer_package.raw_data_packages: if raw_data_package not in data_transfer.destination_location.raw_data_packages: data_transfer.destination_location.raw_data_packages.append( raw_data_package ) def _ensure_unpack_operation( session: Session, data_transfer: models.DataTransfer ) -> models.UnpackOperation: """Find or create the sibling UnpackOperation for a transferred copy. Anchored PER DESTINATION on (data_transfer_package_id, destination_location_id): a package transferred to several destinations (primary + secondary) is verified/unpacked independently at each, so Transfer ↔ Unpack is 1:1 per destination (ADR-0003). Idempotent per destination: the find-work loop may re-run before the row is dispatched.""" op = ( session.query(models.UnpackOperation) .filter_by( data_transfer_package_id=data_transfer.data_transfer_package_id, destination_location_id=data_transfer.destination_location_id, ) .first() ) if op is None: op = models.UnpackOperation( data_transfer_package_id=data_transfer.data_transfer_package_id, destination_location_id=data_transfer.destination_location_id, status=models.Status.PENDING, ) session.add(op) # Flush so the row has an id before it is dispatched as the Celery task arg. session.flush() return op def _get_pending_unpackpings(session: Session) -> list[models.DataTransfer]: """Find DataTransfers eligible for unpack, deduped PER (package, destination). Eligibility keys off the SIBLING TransferOperation reaching COMPLETED (ADR-0003) — not the legacy ``unpack_status`` phase column. Each completed TransferOperation owns one destination, so it maps to exactly one (package, destination) unpack slot; a slot whose per-destination UnpackOperation has already moved past PENDING (scheduled/running/done) is skipped. The (package, destination) dedupe guarantees a given destination yields at most one entry per poll even if two completed transfer rows ever pointed at the same destination — closing the in-poll double-dispatch where a package-only anchor let primary + secondary select the same unpack op. """ completed_transfer_ops = ( session.query(models.TransferOperation) .filter(models.TransferOperation.status == models.Status.COMPLETED) .all() ) eligible: list[models.DataTransfer] = [] seen_slots: set[tuple[int, int]] = set() for transfer_op in completed_transfer_ops: slot = ( transfer_op.data_transfer_package_id, transfer_op.destination_location_id, ) if slot in seen_slots: continue unpack_op = ( session.query(models.UnpackOperation) .filter_by( data_transfer_package_id=transfer_op.data_transfer_package_id, destination_location_id=transfer_op.destination_location_id, ) .first() ) if unpack_op is not None and unpack_op.status != models.Status.PENDING: # Already scheduled / running / done — nothing to do. continue data_transfer = ( session.query(models.DataTransfer) .filter_by( data_transfer_package_id=transfer_op.data_transfer_package_id, origin_location_id=transfer_op.origin_location_id, destination_location_id=transfer_op.destination_location_id, ) .first() ) if data_transfer is not None: seen_slots.add(slot) eligible.append(data_transfer) return eligible def _schedule_unpack_task(data_transfer: models.DataTransfer, session: Session) -> None: try: unpack_operation = _ensure_unpack_operation(session, data_transfer) # Re-check just before dispatch: another iteration / poll / crash-recovery # may already have moved this (package, destination) op past PENDING. The # PENDING→SCHEDULED flip is committed BEFORE apply_async so a re-poll # cannot double-queue the same operation. if unpack_operation.status != models.Status.PENDING: logger.debug( "unpack already scheduled; skipping", unpack_operation_id=unpack_operation.id, status=str(unpack_operation.status), ) return unpack_operation.status = models.Status.SCHEDULED session.commit() # Use dynamic queue routing based on destination location queue_name = route_task_by_location( OperationType.DATA_TRANSFER_UNPACKING, data_transfer.destination_location ) logger.info( "unpack scheduled", unpack_operation_id=unpack_operation.id, transfer_id=data_transfer.id, location=data_transfer.destination_location.name, queue=queue_name, ) unpack_data_transfer_package.apply_async( args=[unpack_operation.id], queue=queue_name ) redis_.publish( "transfer:overview", json.dumps({"type": "unpack_scheduled", "data": data_transfer.id}), ) except Exception as e: logger.error( "failed to schedule unpack task", transfer_id=data_transfer.id, error=str(e), ) raise
[docs] def unpack_and_verify_files(verbose: bool = False, session: Session = None) -> None: """ Unpack transferred files and verify their xxHash checksums. This function retrieves all completed data transfers that are pending unpacking, and schedules Celery tasks to unpack and verify each package. Parameters ---------- verbose : bool, optional If True, sets logging level to DEBUG. Default is False. session : Session, optional An existing database session to use. If None, a new session will be created. Returns ------- None Raises ------ SQLAlchemyError If there's an issue with database operations. """ should_close_session = False if session is None: db = DatabaseConnection() session, _ = db.get_connection() should_close_session = True try: pending_data_unpackings = _get_pending_unpackpings(session) if len(pending_data_unpackings) > 0: logger.info( "pending unpacks found", package_count=len(pending_data_unpackings), ) else: logger.debug("no pending unpacks") for data_transfer in pending_data_unpackings: try: _schedule_unpack_task(data_transfer, session) except Exception as e: logger.error( "failed to schedule unpack task", transfer_id=data_transfer.id, error=str(e), ) continue finally: if should_close_session: session.close()