Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions state-manager/app/config/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,14 @@

class Settings(BaseModel):
"""Application settings loaded from environment variables."""

# MongoDB Configuration
mongo_uri: str = Field(..., description="MongoDB connection URI" )
mongo_database_name: str = Field(default="exosphere-state-manager", description="MongoDB database name")
state_manager_secret: str = Field(..., description="Secret key for API authentication")
secrets_encryption_key: str = Field(..., description="Key for encrypting secrets")
trigger_workers: int = Field(default=1, description="Number of workers to run the trigger cron")
trigger_retention_hours: int = Field(default=24, description="Number of hours to retain completed/failed triggers before cleanup")

@classmethod
def from_env(cls) -> "Settings":
Expand All @@ -21,7 +22,8 @@ def from_env(cls) -> "Settings":
mongo_database_name=os.getenv("MONGO_DATABASE_NAME", "exosphere-state-manager"), # type: ignore
state_manager_secret=os.getenv("STATE_MANAGER_SECRET"), # type: ignore
secrets_encryption_key=os.getenv("SECRETS_ENCRYPTION_KEY"), # type: ignore
trigger_workers=int(os.getenv("TRIGGER_WORKERS", 1)) # type: ignore
trigger_workers=int(os.getenv("TRIGGER_WORKERS", 1)), # type: ignore
trigger_retention_hours=int(os.getenv("TRIGGER_RETENTION_HOURS", 24)) # type: ignore
)
Comment thread
NiveditJain marked this conversation as resolved.


Expand Down
19 changes: 18 additions & 1 deletion state-manager/app/models/db/trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,9 @@ class DatabaseTriggers(Document):
namespace: str = Field(..., description="Namespace of the graph")
trigger_time: datetime = Field(..., description="Trigger time of the trigger")
trigger_status: TriggerStatusEnum = Field(..., description="Status of the trigger")
expires_at: Optional[datetime] = Field(default=None, description="Expiration time for automatic cleanup of completed triggers")
Comment thread
NiveditJain marked this conversation as resolved.

class Settings:
class Settings:
indexes = [
IndexModel(
[
Expand All @@ -32,5 +33,21 @@ class Settings:
],
name="uniq_graph_type_expr_time",
unique=True
),
IndexModel(
[
("expires_at", 1),
],
name="ttl_expires_at",
expireAfterSeconds=0, # Delete immediately when expires_at is reached
partialFilterExpression={
"trigger_status": {
"$in": [
TriggerStatusEnum.TRIGGERED,
TriggerStatusEnum.FAILED,
TriggerStatusEnum.CANCELLED
]
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
)
Comment thread
NiveditJain marked this conversation as resolved.
]
29 changes: 20 additions & 9 deletions state-manager/app/tasks/trigger_cron.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from datetime import datetime
from datetime import datetime, timedelta, timezone
from uuid import uuid4
from app.models.db.trigger import DatabaseTriggers
from app.models.trigger_models import TriggerStatusEnum, TriggerTypeEnum
Expand Down Expand Up @@ -34,10 +34,15 @@ async def call_trigger_graph(trigger: DatabaseTriggers):
x_exosphere_request_id=str(uuid4())
)

async def mark_as_failed(trigger: DatabaseTriggers):
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}}
{"$set": {
"trigger_status": TriggerStatusEnum.FAILED,
"expires_at": expires_at
}}
)

async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime):
Expand Down Expand Up @@ -65,24 +70,30 @@ async def create_next_triggers(trigger: DatabaseTriggers, cron_time: datetime):
if next_trigger_time > cron_time:
break

async def mark_as_triggered(trigger: DatabaseTriggers):
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}}
{"$set": {
"trigger_status": TriggerStatusEnum.TRIGGERED,
"expires_at": expires_at
}}
)

async def handle_trigger(cron_time: datetime):
async def handle_trigger(cron_time: datetime, retention_hours: int):
while(trigger:= await get_due_triggers(cron_time)):
try:
await call_trigger_graph(trigger)
await mark_as_triggered(trigger)
await mark_as_triggered(trigger, retention_hours)
except Exception as e:
await mark_as_failed(trigger)
await mark_as_failed(trigger, retention_hours)
logger.error(f"Error calling trigger graph: {e}")
finally:
await create_next_triggers(trigger, cron_time)

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) for _ in range(get_settings().trigger_workers)])
await asyncio.gather(*[handle_trigger(cron_time, settings.trigger_retention_hours) for _ in range(settings.trigger_workers)])
Loading