main
py 61 lines 2.29 KB
Raw
1 from datetime import datetime
2
3 from loguru import logger
4 from sqlalchemy.future import select
5
6 from app.connectors.wazuh_indexer.services.snapshot_and_restore import (
7 execute_all_enabled_schedules,
8 )
9 from app.db.db_session import get_db_session
10 from app.schedulers.models.scheduler import JobMetadata
11
12
13 async def invoke_snapshot_schedules():
14 """
15 Scheduled job to execute all enabled snapshot schedules.
16 This function is called by the scheduler at configured intervals (e.g., hourly).
17 """
18 logger.info("Starting scheduled snapshot execution")
19
20 try:
21 results = await execute_all_enabled_schedules()
22
23 successful = sum(1 for r in results if r.success)
24 failed = len(results) - successful
25
26 if results:
27 logger.info(
28 f"Scheduled snapshot execution completed: " f"{successful} successful, {failed} failed out of {len(results)} schedules",
29 )
30
31 # Log details for each execution
32 for result in results:
33 if result.success:
34 logger.info(
35 f" - {result.schedule_name}: SUCCESS "
36 f"(snapshot: {result.snapshot_name}, "
37 f"indices: {len(result.indices_snapshotted)}, "
38 f"skipped: {len(result.skipped_write_indices)}, ",
39 )
40 else:
41 logger.error(f" - {result.schedule_name}: FAILED - {result.message}")
42 else:
43 logger.info("No enabled snapshot schedules to execute")
44
45 # Update job metadata with last success timestamp
46 async with get_db_session() as session:
47 stmt = select(JobMetadata).where(JobMetadata.job_id == "invoke_snapshot_schedules")
48 result = await session.execute(stmt)
49 job_metadata = result.scalars().first()
50
51 if job_metadata:
52 job_metadata.last_success = datetime.utcnow()
53 session.add(job_metadata)
54 await session.commit()
55 logger.info("Updated job metadata with the last success timestamp.")
56 else:
57 logger.warning("JobMetadata for 'invoke_snapshot_schedules' not found.")
58
59 except Exception as e:
60 logger.error(f"Failed to execute scheduled snapshots: {e}")
61 raise