Skip to content

Commit 9fdda88

Browse files
committed
Implement trigger cron logic for processing due triggers
- Added functions to retrieve and update due triggers in the database, changing their status from PENDING to TRIGGERED. - Implemented the `call_trigger_graph` function to execute the trigger graph with a unique request ID. - Enhanced the `trigger_cron` function to log the start time and process due triggers in a loop, calling the necessary functions for each trigger. - Introduced a placeholder for `create_next_triggers` to facilitate future trigger creation logic.
1 parent 7f52fe3 commit 9fdda88

1 file changed

Lines changed: 39 additions & 1 deletion

File tree

Lines changed: 39 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,42 @@
11
from datetime import datetime
2+
from uuid import uuid4
3+
from app.models.db.trigger import DatabaseTriggers
4+
from app.models.trigger_models import TriggerStatusEnum
5+
from app.singletons.logs_manager import LogsManager
6+
from app.controller.trigger_graph import trigger_graph
7+
from app.models.trigger_graph_model import TriggerGraphRequestModel
8+
from pymongo import ReturnDocument
9+
10+
logger = LogsManager().get_logger()
11+
12+
async def get_due_triggers(cron_time: datetime) -> DatabaseTriggers | None:
13+
data = await DatabaseTriggers.get_pymongo_collection().find_one_and_update(
14+
{
15+
"trigger_time": {"$lte": cron_time},
16+
"trigger_status": TriggerStatusEnum.PENDING
17+
},
18+
{
19+
"$set": {"trigger_status": TriggerStatusEnum.TRIGGERED}
20+
},
21+
return_document=ReturnDocument.AFTER
22+
)
23+
return DatabaseTriggers(**data) if data else None
24+
25+
async def call_trigger_graph(trigger: DatabaseTriggers):
26+
await trigger_graph(
27+
namespace_name=trigger.namespace,
28+
graph_name=trigger.graph_name,
29+
body=TriggerGraphRequestModel(),
30+
x_exosphere_request_id=str(uuid4())
31+
)
32+
33+
async def create_next_triggers(trigger: DatabaseTriggers):
34+
pass
235

336
async def trigger_cron():
4-
print(f"From trigger_cron: {datetime.now()}")
37+
cron_time = datetime.now()
38+
logger.info(f"starting trigger_cron: {cron_time}")
39+
40+
while(trigger:= await get_due_triggers(cron_time)):
41+
await call_trigger_graph(trigger)
42+
await create_next_triggers(trigger)

0 commit comments

Comments
 (0)