This commit is contained in:
dgtlmoon
2026-01-05 19:27:42 +01:00
parent ca5bb8fe93
commit cc6e078a79
15 changed files with 2455 additions and 328 deletions
@@ -12,8 +12,7 @@ Environment Variables:
"""
import os
import struct
import time
from loguru import logger
# Get queue storage type from environment
@@ -208,9 +207,12 @@ def init_huey(datastore_path):
)
# Configure Huey's logger to only show INFO and above (reduce scheduler DEBUG spam)
# Don't do this when running under pytest - tests may want to see DEBUG logs
import logging
huey_logger = logging.getLogger('huey')
huey_logger.setLevel(logging.INFO)
import sys
if 'pytest' not in sys.modules:
huey_logger = logging.getLogger('huey')
huey_logger.setLevel(logging.INFO)
return huey
@@ -279,111 +281,152 @@ def get_pending_notifications(limit=50):
return []
pending = []
import os
import pickle
import time
try:
storage_type = type(huey.storage).__name__
# Use Huey's built-in methods to get queued and scheduled items
# These methods return pickled bytes that need to be unpickled
if storage_type == 'FileStorage' and hasattr(huey.storage, 'path'):
# FileStorage: Read pickled task files
storage_path = huey.storage.path
# Get queued tasks (immediate execution)
if hasattr(huey.storage, 'enqueued_items'):
try:
queued_items = list(huey.storage.enqueued_items(limit=limit))
for queued_bytes in queued_items:
if len(pending) >= limit:
break
try:
message = pickle.loads(queued_bytes)
if hasattr(message, 'args') and message.args:
notification_data = message.args[0]
# Get task ID and metadata for timestamp
task_id = message.id if hasattr(message, 'id') else None
metadata = _get_task_metadata(task_id) if task_id else None
queued_timestamp = metadata.get('timestamp') if metadata else None
# Get queued tasks (immediate)
queue_dir = os.path.join(storage_path, 'queue')
if os.path.exists(queue_dir):
for root, dirs, files in os.walk(queue_dir):
for filename in files:
if filename.startswith('.') or len(pending) >= limit:
continue
filepath = os.path.join(root, filename)
try:
with open(filepath, 'rb') as f:
task_data = pickle.load(f)
notification_data = task_data.get('args', [{}])[0] if task_data.get('args') else {}
pending.append({
'status': 'queued',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
'queued_at': task_data.get('execute_time'),
})
except Exception:
pass
# Format timestamp for display
from changedetectionio.notification_service import timestamp_to_localtime
queued_at_formatted = timestamp_to_localtime(queued_timestamp) if queued_timestamp else 'Unknown'
# Get scheduled tasks (retrying)
schedule_dir = os.path.join(storage_path, 'schedule')
if os.path.exists(schedule_dir):
for root, dirs, files in os.walk(schedule_dir):
for filename in files:
if filename.startswith('.') or len(pending) >= limit:
continue
filepath = os.path.join(root, filename)
try:
with open(filepath, 'rb') as f:
task_data = pickle.load(f)
notification_data = task_data.get('args', [{}])[0] if task_data.get('args') else {}
eta = task_data.get('eta')
pending.append({
'status': 'retrying',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
'retry_at': eta,
'retry_in_seconds': int(eta - time.time()) if eta else 0,
})
except Exception:
pass
pending.append({
'status': 'queued',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
'task_id': task_id,
'queued_at': queued_timestamp,
'queued_at_formatted': queued_at_formatted,
})
except Exception as e:
logger.debug(f"Error processing queued item: {e}")
continue
except Exception as e:
logger.debug(f"Error getting queued items: {e}")
elif storage_type in ['SqliteStorage', 'SqliteHuey'] and hasattr(huey.storage, 'filename'):
# SqliteStorage: Query database
import sqlite3
conn = sqlite3.connect(huey.storage.filename)
cursor = conn.cursor()
# Get scheduled tasks (retrying)
if hasattr(huey.storage, 'scheduled_items'):
try:
scheduled_items = list(huey.storage.scheduled_items(limit=limit))
for scheduled_bytes in scheduled_items:
if len(pending) >= limit:
break
try:
message = pickle.loads(scheduled_bytes)
if hasattr(message, 'args') and message.args:
notification_data = message.args[0]
eta = message.eta if hasattr(message, 'eta') else None
# Calculate seconds until retry (eta is a datetime object)
import datetime
if eta:
now = datetime.datetime.now() if eta.tzinfo is None else datetime.datetime.now(datetime.timezone.utc)
retry_in_seconds = int((eta - now).total_seconds())
# Get queued tasks
cursor.execute("SELECT data FROM queue LIMIT ?", (limit,))
for row in cursor.fetchall():
try:
task_data = pickle.loads(row[0])
notification_data = task_data.get('args', [{}])[0] if task_data.get('args') else {}
pending.append({
'status': 'queued',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
})
except Exception:
pass
# Convert eta to local timezone for display
if eta.tzinfo is not None:
local_tz = datetime.datetime.now().astimezone().tzinfo
eta_local = eta.astimezone(local_tz)
eta_formatted = eta_local.strftime('%Y-%m-%d %H:%M:%S %Z')
else:
eta_formatted = eta.strftime('%Y-%m-%d %H:%M:%S')
else:
retry_in_seconds = 0
eta_formatted = 'Unknown'
# Get scheduled tasks
cursor.execute("SELECT data, eta FROM schedule LIMIT ?", (limit - len(pending),))
for row in cursor.fetchall():
try:
task_data = pickle.loads(row[0])
notification_data = task_data.get('args', [{}])[0] if task_data.get('args') else {}
eta = row[1]
pending.append({
'status': 'retrying',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
'retry_at': eta,
'retry_in_seconds': int(eta - time.time()) if eta else 0,
})
except Exception:
pass
# Get task ID for manual retry button
task_id = message.id if hasattr(message, 'id') else None
conn.close()
# Convert eta to Unix timestamp for JavaScript (with safety check)
retry_at_timestamp = None
if eta and hasattr(eta, 'timestamp'):
try:
# Huey stores ETA as naive datetime in UTC - need to add timezone info
if eta.tzinfo is None:
# Naive datetime - assume it's UTC (Huey's default)
import datetime
eta = eta.replace(tzinfo=datetime.timezone.utc)
retry_at_timestamp = int(eta.timestamp())
logger.debug(f"ETA after timezone fix: {eta}, Timestamp: {retry_at_timestamp}")
except Exception as e:
logger.debug(f"Error converting eta to timestamp: {e}")
# Format timestamps for display
from changedetectionio.notification_service import timestamp_to_localtime
for item in pending:
if item.get('queued_at'):
item['queued_at_formatted'] = timestamp_to_localtime(item['queued_at'])
if item.get('retry_at'):
item['retry_at_formatted'] = timestamp_to_localtime(item['retry_at'])
# Get original queued timestamp from metadata
metadata = _get_task_metadata(task_id) if task_id else None
queued_timestamp = metadata.get('timestamp') if metadata else None
# Format timestamp for display
from changedetectionio.notification_service import timestamp_to_localtime
queued_at_formatted = timestamp_to_localtime(queued_timestamp) if queued_timestamp else 'Unknown'
# Get retry count from retry_attempts directory
# If there are N attempt files, the next attempt will be N+1
retry_number = 1 # Default to 1 (first retry after initial failure)
total_retries = NOTIFICATION_RETRY_COUNT
watch_uuid = notification_data.get('uuid')
if watch_uuid and huey and hasattr(huey.storage, 'path'):
try:
import os
attempts_dir = os.path.join(huey.storage.path, 'retry_attempts')
if os.path.exists(attempts_dir):
attempt_files = [f for f in os.listdir(attempts_dir) if f.startswith(f"{watch_uuid}.")]
if len(attempt_files) > 0:
# Next attempt number = number of previous attempts + 1
retry_number = len(attempt_files) + 1
logger.debug(f"Watch {watch_uuid[:8]}: Found {len(attempt_files)} retry files, next attempt will be #{retry_number}/{total_retries}")
else:
# Directory exists but no files yet - first retry
retry_number = 1
logger.debug(f"Watch {watch_uuid[:8]}: Retry attempts dir exists but empty, first retry (attempt #1/{total_retries})")
else:
# No attempts dir yet - this is first retry (after initial failure)
retry_number = 1
logger.debug(f"Watch {watch_uuid[:8]}: No retry attempts dir, this is first retry (attempt #1/{total_retries})")
except Exception as e:
logger.warning(f"Error reading retry attempts for {watch_uuid}: {e}, defaulting to attempt #1")
retry_number = 1 # Fallback to 1 on error
pending.append({
'status': 'retrying',
'watch_url': notification_data.get('watch_url', 'Unknown'),
'watch_uuid': notification_data.get('uuid'),
'retry_at': eta,
'retry_at_formatted': eta_formatted,
'retry_at_timestamp': retry_at_timestamp,
'retry_in_seconds': retry_in_seconds,
'task_id': task_id,
'queued_at': queued_timestamp,
'queued_at_formatted': queued_at_formatted,
'retry_number': retry_number,
'total_retries': total_retries,
})
except Exception as e:
logger.debug(f"Error processing scheduled item: {e}")
continue
except Exception as e:
logger.debug(f"Error getting scheduled items: {e}")
except Exception as e:
logger.error(f"Error getting pending notifications: {e}", exc_info=True)
logger.debug(f"get_pending_notifications returning {len(pending)} items")
return pending
@@ -471,6 +514,44 @@ def get_failed_notifications(limit=100, max_age_days=30):
for task_id, result in results.items():
if isinstance(result, (Exception, HueyError)):
# This is a failed task (either Exception or Huey Error object)
# Check if task is still scheduled for retry
# If it is, don't include it in failed list (still retrying)
if huey.storage:
try:
# Check if this task is in the schedule queue (still being retried)
task_still_scheduled = False
# Use Huey's built-in scheduled_items() method to get scheduled tasks
try:
if hasattr(huey.storage, 'scheduled_items'):
import pickle
scheduled_items = list(huey.storage.scheduled_items())
for scheduled_bytes in scheduled_items:
try:
# scheduled_items() returns pickled bytes, need to unpickle
scheduled_message = pickle.loads(scheduled_bytes)
# Each item is a Message object with an 'id' attribute
if hasattr(scheduled_message, 'id'):
scheduled_task_id = scheduled_message.id
if scheduled_task_id == task_id:
task_still_scheduled = True
logger.debug(f"Task {task_id[:20]}... IS scheduled")
break
except Exception as e:
logger.debug(f"Error checking scheduled message: {e}")
continue
except Exception as se:
logger.debug(f"Error checking schedule: {se}")
# Skip this task if it's still scheduled for retry
if task_still_scheduled:
logger.debug(f"Task {task_id[:20]}... still scheduled for retry, not counting as failed yet")
continue
else:
logger.debug(f"Task {task_id[:20]}... NOT in schedule, counting as failed")
except Exception as e:
logger.debug(f"Error checking schedule for task {task_id}: {e}")
# Try to extract notification data from task metadata storage
try:
# Get task metadata from our metadata storage
@@ -552,6 +633,179 @@ def _delete_result(task_id):
return task_manager.delete_result(task_id)
def get_task_apprise_log(task_id):
"""
Get the Apprise log for a specific task.
Returns dict with:
- apprise_log: str (the log text)
- task_id: str
- watch_url: str (if available)
- notification_urls: list (if available)
- error: str (if failed)
"""
if huey is None:
return None
try:
# First check task metadata for notification data and logs
metadata = _get_task_metadata(task_id)
# Also check Huey result for error info (failed tasks)
from huey.utils import Error as HueyError
error_info = None
try:
result = huey.result(task_id, preserve=True)
if result and isinstance(result, (Exception, HueyError)):
error_info = str(result)
except Exception as e:
# If huey.result() raises an exception, that IS the error we want
# (Huey raises the stored exception when calling result() on failed tasks)
error_info = str(e)
logger.debug(f"Got error from result for task {task_id}: {type(e).__name__}")
if metadata:
# Get apprise logs from metadata (could be 'apprise_logs' list or 'apprise_log' string)
apprise_logs = metadata.get('apprise_logs', [])
apprise_log_text = '\n'.join(apprise_logs) if isinstance(apprise_logs, list) else metadata.get('apprise_log', '')
# If no logs in metadata but we have error_info, try to extract from error
if not apprise_log_text and error_info and 'Apprise logs:' in error_info:
parts = error_info.split('Apprise logs:', 1)
if len(parts) > 1:
apprise_log_text = parts[1].strip()
# The exception string has escaped newlines (\n), convert to actual newlines
apprise_log_text = apprise_log_text.replace('\\n', '\n')
# Also remove trailing quotes and closing parens from exception repr
apprise_log_text = apprise_log_text.rstrip("')")
logger.debug(f"Extracted Apprise logs from error for task {task_id}: {len(apprise_log_text)} chars")
# Clean up error to not duplicate the Apprise logs
# Only show the main error message, not the logs again
error_parts = error_info.split('\nApprise logs:', 1)
if len(error_parts) > 1:
error_info = error_parts[0] # Keep only the main error message
# Use metadata for apprise_log and notification data, but also include error from result
result = {
'task_id': task_id,
'apprise_log': apprise_log_text if apprise_log_text else 'No log available',
'watch_url': metadata.get('notification_data', {}).get('watch_url'),
'notification_urls': metadata.get('notification_data', {}).get('notification_urls', []),
'error': error_info if error_info else metadata.get('error'),
'timestamp': metadata.get('timestamp')
}
logger.debug(f"Returning log data for task {task_id}: apprise_log length={len(result['apprise_log'])}, has_error={bool(result['error'])}")
return result
# If not in metadata, try to extract from result only
if error_info:
# Try to extract Apprise log from error message
apprise_log = 'No detailed log available'
if 'Apprise logs:' in error_info:
parts = error_info.split('Apprise logs:', 1)
if len(parts) > 1:
apprise_log = parts[1].strip()
return {
'task_id': task_id,
'apprise_log': apprise_log,
'error': error_info
}
return None
except Exception as e:
logger.error(f"Error getting Apprise log for task {task_id}: {e}")
return None
def execute_scheduled_notification(task_id):
"""
Execute a scheduled/retrying notification immediately by canceling the schedule and queuing it now.
Uses Huey's native revoke_by_id() to cancel the scheduled task, then immediately queues it.
Args:
task_id: Huey task ID to execute immediately
Returns:
True if successfully executed, False otherwise
"""
if huey is None:
logger.error("Huey not initialized")
return False
try:
# First, check if task is actually scheduled
scheduled_items = list(huey.storage.scheduled_items())
task_found = False
notification_data = None
import pickle
for scheduled_bytes in scheduled_items:
try:
message = pickle.loads(scheduled_bytes)
if hasattr(message, 'id') and message.id == task_id:
task_found = True
# Extract notification data from scheduled task
if hasattr(message, 'args') and message.args:
notification_data = message.args[0]
break
except Exception as e:
logger.debug(f"Error checking scheduled task: {e}")
continue
if not task_found:
logger.error(f"Task {task_id} not found in schedule")
return False
if not notification_data:
logger.error(f"No notification data found for task {task_id}")
return False
# Execute the notification NOW (synchronously, not queued) by calling the task function directly
logger.info(f"Executing notification for task {task_id} immediately (not queued)...")
try:
# Import here to avoid circular dependency
from changedetectionio.flask_app import datastore
from changedetectionio.notification.handler import process_notification
from changedetectionio.notification_service import NotificationContextData
# Wrap dict in NotificationContextData if needed
if not isinstance(notification_data, NotificationContextData):
notification_data = NotificationContextData(notification_data)
# Call the notification processing function directly (not via Huey queue)
# This executes synchronously in the current thread
sent_obj = process_notification(notification_data, datastore)
# If we get here, notification was sent successfully!
logger.info(f"✓ Notification sent successfully for task {task_id}")
# NOW revoke the scheduled task since we successfully sent it
huey.revoke_by_id(task_id, revoke_once=True)
logger.info(f"Revoked scheduled task {task_id} (no longer needed)")
# Clean up old metadata and result
_delete_result(task_id)
_delete_task_metadata(task_id)
logger.info(f"Cleaned up old result/metadata for task {task_id}")
return True
except Exception as e:
# Notification failed - keep the scheduled task so it retries later
logger.warning(f"Failed to send notification for task {task_id}: {e}")
logger.info(f"Keeping scheduled task {task_id} in queue for automatic retry")
return False
except Exception as e:
logger.error(f"Error executing scheduled notification {task_id}: {e}")
return False
def retry_failed_notification(task_id):
"""
Retry a failed notification by task ID.