Skip to content

Commit b1a9279

Browse files
committed
refactor: update trigger status filtering and add expiration logic
- Changed the partial filter expression in DatabaseTriggers to exclude PENDING and TRIGGERING statuses. - Introduced an expires_at field in create_crons to set expiration time for triggers based on retention settings, ensuring proper cleanup of triggers after their designated retention period. This enhances the management of trigger states and ensures that only relevant triggers are retained in the database.
1 parent 63665b8 commit b1a9279

2 files changed

Lines changed: 10 additions & 4 deletions

File tree

state-manager/app/models/db/trigger.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -42,9 +42,9 @@ class Settings:
4242
expireAfterSeconds=0, # Delete immediately when expires_at is reached
4343
partialFilterExpression={
4444
"trigger_status": {
45-
"$in": [
46-
TriggerStatusEnum.TRIGGERED,
47-
TriggerStatusEnum.FAILED
45+
"$nin": [
46+
TriggerStatusEnum.PENDING,
47+
TriggerStatusEnum.TRIGGERING
4848
]
4949
}
5050
}

state-manager/app/tasks/verify_graph.py

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,13 @@
1010
from app.singletons.logs_manager import LogsManager
1111
from app.models.trigger_models import TriggerStatusEnum, TriggerTypeEnum
1212
from app.models.db.trigger import DatabaseTriggers
13+
from app.config.settings import get_settings
14+
from datetime import timedelta
1315

1416
logger = LogsManager().get_logger()
1517

18+
settings = get_settings()
19+
1620
async def verify_node_exists(graph_template: GraphTemplate, registered_nodes: list[RegisteredNode]) -> list[str]:
1721
errors = []
1822
template_nodes_set = set([(node.node_name, node.namespace) for node in graph_template.nodes])
@@ -110,6 +114,7 @@ async def create_crons(graph_template: GraphTemplate):
110114
iter = croniter.croniter(expression, current_time)
111115

112116
next_trigger_time = iter.get_next(datetime)
117+
expires_at = next_trigger_time + timedelta(hours=settings.trigger_retention_hours)
113118

114119
new_db_triggers.append(
115120
DatabaseTriggers(
@@ -118,7 +123,8 @@ async def create_crons(graph_template: GraphTemplate):
118123
graph_name=graph_template.name,
119124
namespace=graph_template.namespace,
120125
trigger_status=TriggerStatusEnum.PENDING,
121-
trigger_time=next_trigger_time
126+
trigger_time=next_trigger_time,
127+
expires_at=expires_at
122128
)
123129
)
124130

0 commit comments

Comments
 (0)