Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 45
fix: wire mark_as_cancelled into trigger cancellation flow & replace destructive deletes with TTL-based cleanup#590
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base:main
Are you sure you want to change the base?
Uh oh!
There was an error while loading. Please reload this page.
Changes from all commits
File filter
Filter by extension
Conversations
Uh oh!
There was an error while loading. Please reload this page.
Jump to
Uh oh!
There was an error while loading. Please reload this page.
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,20 +1,54 @@ | ||
| # tasks to run when the server starts | ||
| from datetime import datetime, timedelta, timezone | ||
| import asyncio | ||
| from app.config.settings import get_settings | ||
| from app.models.db.trigger import DatabaseTriggers | ||
| from app.models.trigger_models import TriggerStatusEnum | ||
| import asyncio | ||
| from app.singletons.logs_manager import LogsManager | ||
| logger = LogsManager().get_logger() | ||
| async def mark_old_triggers_cancelled() -> None: | ||
| """ | ||
| Migrate legacy TRIGGERED/FAILED triggers that predate TTL. | ||
| async def delete_old_triggers(): | ||
| await DatabaseTriggers.get_pymongo_collection().delete_many( | ||
| These documents have expires_at = None, so we mark them as CANCELLED and | ||
| set expires_at so the TTL index can eventually clean them up. | ||
| """ | ||
| settings = get_settings() | ||
| retention_hours = settings.trigger_retention_hours | ||
| expires_at = datetime.now(timezone.utc) + timedelta(hours=retention_hours) | ||
| # Use the same filter used before by delete_many() | ||
| filter_query = { | ||
| "trigger_status": { | ||
| "$in": [ | ||
| TriggerStatusEnum.TRIGGERED.value, | ||
| TriggerStatusEnum.FAILED.value, | ||
| ] | ||
| }, | ||
| "expires_at": None, | ||
| } | ||
| logger.info( | ||
| "Init task marking legacy TRIGGERED/FAILED triggers as CANCELLED " | ||
| f"for filter={filter_query}, expires_at={expires_at.isoformat()}" | ||
| ) | ||
| await DatabaseTriggers.get_pymongo_collection().update_many( | ||
| filter_query, | ||
| { | ||
| "trigger_status": { | ||
| "$in": [TriggerStatusEnum.TRIGGERED, TriggerStatusEnum.FAILED] | ||
| }, | ||
| "expires_at": None | ||
| } | ||
| "$set": { | ||
| "trigger_status": TriggerStatusEnum.CANCELLED.value, | ||
| "expires_at": expires_at, | ||
| } | ||
| }, | ||
| ) | ||
Brijesh-Thakkar marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| async def init_tasks(): | ||
| async def init_tasks() -> None: | ||
| await asyncio.gather( | ||
| *[ | ||
| delete_old_triggers() | ||
| ]) | ||
| mark_old_triggers_cancelled(), | ||
| ) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -13,19 +13,21 @@ | ||
| logger = LogsManager().get_logger() | ||
| async def get_due_triggers(cron_time: datetime) -> DatabaseTriggers | None: | ||
| data = await DatabaseTriggers.get_pymongo_collection().find_one_and_update( | ||
| { | ||
| "trigger_time": {"$lte": cron_time}, | ||
| "trigger_status": TriggerStatusEnum.PENDING | ||
| "trigger_status": TriggerStatusEnum.PENDING.value | ||
| }, | ||
| { | ||
| "$set": {"trigger_status": TriggerStatusEnum.TRIGGERING} | ||
| "$set": {"trigger_status": TriggerStatusEnum.TRIGGERING.value} | ||
| }, | ||
| return_document=ReturnDocument.AFTER | ||
| ) | ||
| return DatabaseTriggers(**data) if data else None | ||
| async def call_trigger_graph(trigger: DatabaseTriggers): | ||
| await trigger_graph( | ||
| namespace_name=trigger.namespace, | ||
| @@ -34,17 +36,45 @@ async def call_trigger_graph(trigger: DatabaseTriggers): | ||
| x_exosphere_request_id=str(uuid4()) | ||
| ) | ||
| async def mark_as_failed(trigger: DatabaseTriggers, retention_hours: int): | ||
| expires_at = datetime.now(timezone.utc) + timedelta(hours=retention_hours) | ||
| await DatabaseTriggers.get_pymongo_collection().update_one( | ||
| {"_id": trigger.id}, | ||
| {"$set": { | ||
| "trigger_status": TriggerStatusEnum.FAILED, | ||
| "trigger_status": TriggerStatusEnum.FAILED.value, | ||
| "expires_at": expires_at | ||
| }} | ||
| ) | ||
| async def mark_as_cancelled(trigger: DatabaseTriggers, retention_hours: int): | ||
| """ | ||
| Mark a trigger as CANCELLED and set expires_at so MongoDB TTL will remove it | ||
| after `retention_hours`. | ||
| """ | ||
| expires_at = datetime.now(timezone.utc) + timedelta(hours=retention_hours) | ||
| await DatabaseTriggers.get_pymongo_collection().update_one( | ||
| {"_id": trigger.id}, | ||
| {"$set": { | ||
| "trigger_status": TriggerStatusEnum.CANCELLED.value, # keep .value ✅ | ||
| "expires_at": expires_at | ||
| }} | ||
Brijesh-Thakkar marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| ) | ||
Brijesh-Thakkar marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
Brijesh-Thakkar marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| async def cancel_trigger(trigger: DatabaseTriggers, retention_hours: int): | ||
| """ | ||
| Cancel a trigger using the mark_as_cancelled helper. | ||
| This is intended to be used by other modules instead of duplicating | ||
| the cancellation logic inline. | ||
| """ | ||
| await mark_as_cancelled(trigger, retention_hours) | ||
| async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime, retention_hours: int): | ||
| assert trigger.expression is not None | ||
| iter = croniter.croniter(trigger.expression, trigger.trigger_time) | ||
| @@ -60,7 +90,7 @@ async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime, r | ||
| graph_name=trigger.graph_name, | ||
| namespace=trigger.namespace, | ||
| trigger_time=next_trigger_time, | ||
| trigger_status=TriggerStatusEnum.PENDING, | ||
| trigger_status=TriggerStatusEnum.PENDING, # OK because insert() converts | ||
| expires_at=expires_at | ||
| ).insert() | ||
| except DuplicateKeyError: | ||
| @@ -72,19 +102,21 @@ async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime, r | ||
| if next_trigger_time > cron_time: | ||
| break | ||
| async def mark_as_triggered(trigger: DatabaseTriggers, retention_hours: int): | ||
| expires_at = datetime.now(timezone.utc) + timedelta(hours=retention_hours) | ||
| await DatabaseTriggers.get_pymongo_collection().update_one( | ||
| {"_id": trigger.id}, | ||
| {"$set": { | ||
| "trigger_status": TriggerStatusEnum.TRIGGERED, | ||
| "trigger_status": TriggerStatusEnum.TRIGGERED.value, | ||
| "expires_at": expires_at | ||
| }} | ||
| ) | ||
| async def handle_trigger(cron_time: datetime, retention_hours: int): | ||
| while(trigger:= await get_due_triggers(cron_time)): | ||
| while(trigger:= await get_due_triggers(cron_time)): | ||
| try: | ||
| await call_trigger_graph(trigger) | ||
| await mark_as_triggered(trigger, retention_hours) | ||
| @@ -94,8 +126,12 @@ async def handle_trigger(cron_time: datetime, retention_hours: int): | ||
| finally: | ||
| await create_next_triggers(trigger, cron_time, retention_hours) | ||
| async def trigger_cron(): | ||
| cron_time = datetime.now() | ||
| settings = get_settings() | ||
| logger.info(f"starting trigger_cron: {cron_time}") | ||
| await asyncio.gather(*[handle_trigger(cron_time, settings.trigger_retention_hours) for _ in range(settings.trigger_workers)]) | ||
| await asyncio.gather(*[ | ||
| handle_trigger(cron_time, settings.trigger_retention_hours) | ||
| for _ in range(settings.trigger_workers) | ||
| ]) | ||
Uh oh!
There was an error while loading. Please reload this page.