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()