Skip to content

Commit 155f7cc

Browse files
committed
Add error handling for trigger processing in trigger_cron.py
- Introduced a new function `mark_as_failed` to update the trigger status to FAILED in case of errors during processing. - Enhanced the `handle_trigger` function to include a try-except block, ensuring that exceptions are logged and the trigger status is updated appropriately when errors occur. - This change improves the robustness of the trigger processing logic by handling failures gracefully.
1 parent 6b2bf85 commit 155f7cc

1 file changed

Lines changed: 13 additions & 3 deletions

File tree

state-manager/app/tasks/trigger_cron.py

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,12 @@ async def call_trigger_graph(trigger: DatabaseTriggers):
3333
x_exosphere_request_id=str(uuid4())
3434
)
3535

36+
async def mark_as_failed(trigger: DatabaseTriggers):
37+
await DatabaseTriggers.get_pymongo_collection().update_one(
38+
{"_id": trigger.id},
39+
{"$set": {"trigger_status": TriggerStatusEnum.FAILED}}
40+
)
41+
3642
async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime):
3743
assert trigger.expression is not None
3844
iter = croniter.croniter(trigger.expression, cron_time)
@@ -55,9 +61,13 @@ async def mark_as_triggered(trigger: DatabaseTriggers):
5561

5662
async def handle_trigger(cron_time: datetime):
5763
while(trigger:= await get_due_triggers(cron_time)):
58-
await call_trigger_graph(trigger)
59-
await create_next_triggers(trigger, cron_time)
60-
await mark_as_triggered(trigger)
64+
try:
65+
await call_trigger_graph(trigger)
66+
await create_next_triggers(trigger, cron_time)
67+
await mark_as_triggered(trigger)
68+
except Exception as e:
69+
await mark_as_failed(trigger)
70+
logger.error(f"Error calling trigger graph: {e}")
6171

6272
async def trigger_cron():
6373
cron_time = datetime.now()

0 commit comments

Comments
 (0)