221 lines
7.9 KiB
Python
221 lines
7.9 KiB
Python
"""
|
|
Cleanup activities for removing stale files from local filesystem.
|
|
|
|
This module provides activities for cleaning up temporary files and directories
|
|
that are older than the configured retention period. It operates independently
|
|
of the database, using timestamps embedded in filenames.
|
|
"""
|
|
|
|
from temporalio import activity, workflow
|
|
|
|
with workflow.unsafe.imports_passed_through():
|
|
import os
|
|
import re
|
|
import shutil
|
|
import traceback
|
|
from datetime import datetime, timedelta
|
|
from typing import Any
|
|
|
|
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
|
|
from sientia_do.notifications.models import NotificationLevel
|
|
from sientia_do.observability.logger import Logger
|
|
from sientia_do.observability.metrics_controller import MetricsController
|
|
from sientia_do.observability.sientia_monitoring import SientiaMonitoring
|
|
|
|
from model_manager.runtime_paths import REPORTS_TEMP_DIR
|
|
|
|
RETENTION_HOURS = int(os.getenv('CLEANUP_RETENTION_HOURS', '24'))
|
|
DRY_RUN = os.getenv('CLEANUP_DRY_RUN', 'false').lower() == 'true'
|
|
|
|
|
|
class Cleanup(SientiaMonitoring):
|
|
"""
|
|
Activity for cleaning up stale files and directories.
|
|
|
|
This activity extends SientiaMonitoring and handles cleanup of:
|
|
- Local temporary directories with timestamp suffixes
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
logger: Logger,
|
|
notification_handler: NotificationHandler,
|
|
metrics_controller: MetricsController,
|
|
):
|
|
"""
|
|
Initialize Cleanup activity.
|
|
|
|
Args:
|
|
logger: Logger instance for observability
|
|
notification_handler: Handler for sending notifications
|
|
metrics_controller: Controller for metrics emission
|
|
"""
|
|
SientiaMonitoring.__init__(self, logger, notification_handler, metrics_controller)
|
|
|
|
# Configuration from environment variables
|
|
self.retention_hours = RETENTION_HOURS
|
|
self.dry_run = DRY_RUN
|
|
|
|
# Regex patterns for timestamp extraction
|
|
self.dir_timestamp_pattern = re.compile(
|
|
r'^(.+)_(\d{8}_\d{6}_\d{6})$'
|
|
) # name_YYYYMMDD_HHMMSS_microseconds
|
|
|
|
def _process_directory_item(
|
|
self,
|
|
item_path: str,
|
|
item_name: str,
|
|
cutoff_time: datetime,
|
|
metadata: dict[str, Any],
|
|
) -> tuple[bool, str | None]:
|
|
"""
|
|
Handle a single temp-directory entry: skip if the name doesn't match the
|
|
timestamp pattern, otherwise delete (or dry-run log) it when stale.
|
|
|
|
Args:
|
|
item_path: Full path to the directory being evaluated
|
|
item_name: Directory name (used to extract the embedded timestamp)
|
|
cutoff_time: Directories older than this are considered stale
|
|
metadata: Workflow execution metadata for logging
|
|
|
|
Return:
|
|
tuple[bool, str | None]: (deleted, error_message). `deleted` is True
|
|
if the directory was removed (or would be, in dry-run mode).
|
|
"""
|
|
match = self.dir_timestamp_pattern.match(item_name)
|
|
if not match:
|
|
self.debug(f'Skipping directory without timestamp pattern: {item_name}', metadata)
|
|
return False, None
|
|
|
|
timestamp_str = match.group(2)
|
|
try:
|
|
# Parse YYYYMMDD_HHMMSS_microseconds
|
|
dir_time = datetime.strptime(timestamp_str, '%Y%m%d_%H%M%S_%f')
|
|
except ValueError as e:
|
|
error_msg = f'Invalid timestamp format in directory {item_name}: {str(e)}'
|
|
self.error(error_msg, metadata)
|
|
return False, error_msg
|
|
|
|
return self._delete_if_stale(item_path, item_name, dir_time, cutoff_time, metadata)
|
|
|
|
def _delete_if_stale(
|
|
self,
|
|
item_path: str,
|
|
item_name: str,
|
|
dir_time: datetime,
|
|
cutoff_time: datetime,
|
|
metadata: dict[str, Any],
|
|
) -> tuple[bool, str | None]:
|
|
"""
|
|
Delete (or dry-run log) a directory whose embedded timestamp is older than
|
|
cutoff_time; otherwise leave it alone.
|
|
|
|
Args:
|
|
item_path: Full path to the directory
|
|
item_name: Directory name (for logging)
|
|
dir_time: Timestamp parsed from the directory name
|
|
cutoff_time: Directories older than this are considered stale
|
|
metadata: Workflow execution metadata for logging
|
|
|
|
Return:
|
|
tuple[bool, str | None]: (deleted, error_message)
|
|
"""
|
|
if dir_time >= cutoff_time:
|
|
age_hours = (datetime.now() - dir_time).total_seconds() / 3600
|
|
self.debug(f'Keeping recent directory: {item_name} (age: {age_hours:.1f}h)', metadata)
|
|
return False, None
|
|
|
|
age_hours = (datetime.now() - dir_time).total_seconds() / 3600
|
|
|
|
if self.dry_run:
|
|
self.info(
|
|
f'[DRY RUN] Would delete directory: {item_name} (age: {age_hours:.1f}h)', metadata
|
|
)
|
|
return True, None
|
|
|
|
try:
|
|
shutil.rmtree(item_path)
|
|
self.info(f'Deleted stale directory: {item_name} (age: {age_hours:.1f}h)', metadata)
|
|
return True, None
|
|
except OSError as e:
|
|
error_msg = f'Failed to delete directory {item_name}: {str(e)}'
|
|
self.error(error_msg, metadata)
|
|
return False, error_msg
|
|
|
|
@activity.defn(name='cleanup_temp_directories')
|
|
def cleanup_temp_directories(self, input_data: dict[str, Any]) -> None:
|
|
"""
|
|
Clean up stale temporary directories based on timestamp in directory name.
|
|
|
|
This activity scans the reports/temp directory for subdirectories following
|
|
the pattern '{name}_{timestamp}' where timestamp is in YYYYMMDD_HHMMSS_microseconds format.
|
|
Directories older than the retention period are deleted.
|
|
|
|
Args:
|
|
input_data: Cleanup configuration containing:
|
|
- metadata (dict): Workflow execution metadata
|
|
- temp_path (str): Path to temp directory (optional, defaults to reports/temp)
|
|
|
|
Returns:
|
|
None: Results are logged and tracked via metrics
|
|
|
|
Raises:
|
|
Exception: If cleanup fails (after sending notification)
|
|
"""
|
|
metadata = input_data.get('metadata', {})
|
|
temp_path = input_data.get('temp_path', REPORTS_TEMP_DIR)
|
|
|
|
cutoff_time = datetime.now() - timedelta(hours=self.retention_hours)
|
|
|
|
try:
|
|
self.info(
|
|
f'Starting local directory cleanup - Path: {temp_path}, '
|
|
f'Retention: {self.retention_hours}h, Dry run: {self.dry_run}',
|
|
metadata,
|
|
)
|
|
|
|
if not os.path.exists(temp_path):
|
|
self.warning(f'Temp directory does not exist: {temp_path}', metadata)
|
|
return
|
|
|
|
directories_scanned = 0
|
|
directories_deleted = 0
|
|
errors = []
|
|
|
|
for item_name in os.listdir(temp_path):
|
|
item_path = os.path.join(temp_path, item_name)
|
|
|
|
if not os.path.isdir(item_path):
|
|
continue
|
|
|
|
directories_scanned += 1
|
|
|
|
deleted, error = self._process_directory_item(
|
|
item_path, item_name, cutoff_time, metadata
|
|
)
|
|
if error:
|
|
errors.append(error)
|
|
if deleted:
|
|
directories_deleted += 1
|
|
|
|
self.info(
|
|
f'Directory cleanup completed - Scanned: {directories_scanned}, '
|
|
f'Deleted: {directories_deleted}, Errors: {len(errors)}',
|
|
metadata,
|
|
)
|
|
|
|
except Exception as e:
|
|
error_msg = f'Error in directory cleanup: {str(e)}'
|
|
trace = traceback.format_exc()
|
|
|
|
self.send_notification(
|
|
metadata=metadata,
|
|
notification_id='CLEANUP_DIRECTORIES_ERROR',
|
|
message=error_msg,
|
|
block='cleanup_temp_directories',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=trace,
|
|
)
|
|
|
|
raise
|