Skip to content

Commit 45cb95d

Browse files
committed
Refactor trigger management and settings configuration
- Updated the Settings model to replace `trigger_ahead_time` with `trigger_workers`, allowing for configurable worker count for trigger processing. - Added a new status `TRIGGERING` to the TriggerStatusEnum to better represent the trigger lifecycle. - Enhanced the `trigger_cron` function to utilize the new worker configuration, enabling concurrent processing of due triggers. - Refined the `create_next_triggers` function to calculate the next trigger time based on the cron expression, improving trigger scheduling logic.
1 parent 9fdda88 commit 45cb95d

4 files changed

Lines changed: 44 additions & 26 deletions

File tree

state-manager/app/config/settings.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ class Settings(BaseModel):
1212
mongo_database_name: str = Field(default="exosphere-state-manager", description="MongoDB database name")
1313
state_manager_secret: str = Field(..., description="Secret key for API authentication")
1414
secrets_encryption_key: str = Field(..., description="Key for encrypting secrets")
15-
trigger_ahead_time: int = Field(default=10, description="Time in minutes to trigger the graph ahead of the current time")
15+
trigger_workers: int = Field(default=1, description="Number of workers to run the trigger cron")
1616

1717
@classmethod
1818
def from_env(cls) -> "Settings":
@@ -21,7 +21,7 @@ def from_env(cls) -> "Settings":
2121
mongo_database_name=os.getenv("MONGO_DATABASE_NAME", "exosphere-state-manager"), # type: ignore
2222
state_manager_secret=os.getenv("STATE_MANAGER_SECRET"), # type: ignore
2323
secrets_encryption_key=os.getenv("SECRETS_ENCRYPTION_KEY"), # type: ignore
24-
trigger_ahead_time=int(os.getenv("TRIGGER_AHEAD_TIME", 10)) # type: ignore
24+
trigger_workers=int(os.getenv("TRIGGER_WORKERS", 1)) # type: ignore
2525
)
2626

2727

state-manager/app/models/trigger_models.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ class TriggerStatusEnum(str, Enum):
1111
FAILED = "FAILED"
1212
CANCELLED = "CANCELLED"
1313
TRIGGERED = "TRIGGERED"
14+
TRIGGERING = "TRIGGERING"
1415

1516
class CronTrigger(BaseModel):
1617
expression: str = Field(..., description="Cron expression for the trigger")

state-manager/app/tasks/trigger_cron.py

Lines changed: 31 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,14 @@
11
from datetime import datetime
22
from uuid import uuid4
33
from app.models.db.trigger import DatabaseTriggers
4-
from app.models.trigger_models import TriggerStatusEnum
4+
from app.models.trigger_models import TriggerStatusEnum, TriggerTypeEnum
55
from app.singletons.logs_manager import LogsManager
66
from app.controller.trigger_graph import trigger_graph
77
from app.models.trigger_graph_model import TriggerGraphRequestModel
88
from pymongo import ReturnDocument
9+
from app.config.settings import get_settings
10+
import croniter
11+
import asyncio
912

1013
logger = LogsManager().get_logger()
1114

@@ -16,7 +19,7 @@ async def get_due_triggers(cron_time: datetime) -> DatabaseTriggers | None:
1619
"trigger_status": TriggerStatusEnum.PENDING
1720
},
1821
{
19-
"$set": {"trigger_status": TriggerStatusEnum.TRIGGERED}
22+
"$set": {"trigger_status": TriggerStatusEnum.TRIGGERING}
2023
},
2124
return_document=ReturnDocument.AFTER
2225
)
@@ -30,13 +33,33 @@ async def call_trigger_graph(trigger: DatabaseTriggers):
3033
x_exosphere_request_id=str(uuid4())
3134
)
3235

33-
async def create_next_triggers(trigger: DatabaseTriggers):
34-
pass
36+
async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime):
37+
assert trigger.expression is not None
38+
iter = croniter.croniter(trigger.expression, cron_time)
39+
next_trigger_time = iter.get_next(datetime)
3540

36-
async def trigger_cron():
37-
cron_time = datetime.now()
38-
logger.info(f"starting trigger_cron: {cron_time}")
41+
await DatabaseTriggers(
42+
type=TriggerTypeEnum.CRON,
43+
expression=trigger.expression,
44+
graph_name=trigger.graph_name,
45+
namespace=trigger.namespace,
46+
trigger_time=next_trigger_time,
47+
trigger_status=TriggerStatusEnum.PENDING
48+
).insert()
49+
50+
async def mark_as_triggered(trigger: DatabaseTriggers):
51+
await DatabaseTriggers.get_pymongo_collection().update_one(
52+
{"_id": trigger.id},
53+
{"$set": {"trigger_status": TriggerStatusEnum.TRIGGERED}}
54+
)
3955

56+
async def handle_trigger(cron_time: datetime):
4057
while(trigger:= await get_due_triggers(cron_time)):
4158
await call_trigger_graph(trigger)
42-
await create_next_triggers(trigger)
59+
await create_next_triggers(trigger, cron_time)
60+
await mark_as_triggered(trigger)
61+
62+
async def trigger_cron():
63+
cron_time = datetime.now()
64+
logger.info(f"starting trigger_cron: {cron_time}")
65+
await asyncio.gather(*[handle_trigger(cron_time) for _ in range(get_settings().trigger_workers)])

state-manager/app/tasks/verify_graph.py

Lines changed: 10 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -127,30 +127,24 @@ async def create_crons(graph_template: GraphTemplate, old_triggers: list[Trigger
127127

128128
crons_to_create = new_crons - old_crons
129129

130-
trigger_ahead_time = get_settings().trigger_ahead_time
131130
current_time = datetime.now()
132-
limit_time = datetime.now() + timedelta(minutes=trigger_ahead_time)
133131

134132
new_db_triggers = []
135133
for cron in crons_to_create:
136134
iter = croniter.croniter(cron.expression, current_time)
137135

138-
while(True):
139-
next_trigger_time = iter.get_next(datetime)
136+
next_trigger_time = iter.get_next(datetime)
140137

141-
# at least one event should be inserted for the cron
142-
new_db_triggers.append(
143-
DatabaseTriggers(
144-
type=TriggerTypeEnum.CRON,
145-
expression=cron.expression,
146-
graph_name=graph_template.name,
147-
namespace=graph_template.namespace,
148-
trigger_status=TriggerStatusEnum.PENDING,
149-
trigger_time=next_trigger_time
150-
)
138+
new_db_triggers.append(
139+
DatabaseTriggers(
140+
type=TriggerTypeEnum.CRON,
141+
expression=cron.expression,
142+
graph_name=graph_template.name,
143+
namespace=graph_template.namespace,
144+
trigger_status=TriggerStatusEnum.PENDING,
145+
trigger_time=next_trigger_time
151146
)
152-
if next_trigger_time > limit_time:
153-
break
147+
)
154148

155149
if len(new_db_triggers) > 0:
156150
await DatabaseTriggers.insert_many(new_db_triggers)

0 commit comments

Comments
 (0)