4a7f24348c
- Add instance_events and health_checks tables with Alembic migration - InstanceEventBus: typed pub/sub singleton with wildcard support - HealthMonitor: async background loop polling containers every 15s - SSE endpoint GET /events/stream with auth and connection limits - Lifecycle hooks in tool_instances.py (create/start/stop/restart/delete) - Structured JSON logging with correlation IDs - 15 new unit tests (EventBus, HealthMonitor, MonitoringModels) Quality gates: pytest 15 new passed, ruff clean
99 lines
2.8 KiB
Python
99 lines
2.8 KiB
Python
"""Lifecycle hook helpers for instrumenting tool instance transitions."""
|
|
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from src.models.instance_event import InstanceEvent
|
|
from src.models.tool_instance import ToolInstance
|
|
from src.services.correlation import get_correlation_id
|
|
from src.services.event_bus import InstanceEventBus, InstanceEventPayload
|
|
|
|
|
|
def _build_payload(
|
|
event_type: str,
|
|
instance: ToolInstance,
|
|
status: str | None = None,
|
|
message: str | None = None,
|
|
metadata: dict | None = None,
|
|
) -> InstanceEventPayload:
|
|
"""Construct a standard event payload."""
|
|
return {
|
|
"event": event_type,
|
|
"instance_id": str(instance.id),
|
|
"status": status or instance.status,
|
|
"message": message,
|
|
"metadata": metadata or {},
|
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
|
"correlation_id": get_correlation_id(),
|
|
}
|
|
|
|
|
|
async def _write_audit_row(
|
|
session: AsyncSession,
|
|
instance: ToolInstance,
|
|
event_type: str,
|
|
created_by: uuid.UUID | None = None,
|
|
status: str | None = None,
|
|
message: str | None = None,
|
|
metadata: dict | None = None,
|
|
) -> InstanceEvent:
|
|
"""Persist an instance_events audit row."""
|
|
row = InstanceEvent(
|
|
instance_id=instance.id,
|
|
event_type=event_type.replace("instance.", ""),
|
|
status=status or instance.status,
|
|
message=message,
|
|
created_by=created_by,
|
|
event_metadata=metadata or {},
|
|
)
|
|
session.add(row)
|
|
await session.commit()
|
|
return row
|
|
|
|
|
|
async def publish_lifecycle_event(
|
|
event_bus: InstanceEventBus,
|
|
session: AsyncSession,
|
|
instance: ToolInstance,
|
|
event_type: str,
|
|
created_by: uuid.UUID | None = None,
|
|
status: str | None = None,
|
|
message: str | None = None,
|
|
metadata: dict | None = None,
|
|
) -> None:
|
|
"""Publish a lifecycle event and write an audit row after DB commit.
|
|
|
|
Args:
|
|
event_bus: The global event bus.
|
|
session: Active async DB session.
|
|
instance: The affected tool instance.
|
|
event_type: One of instance.created, instance.started, etc.
|
|
created_by: User ID for user-initiated actions; None for system.
|
|
status: Optional status override.
|
|
message: Optional human-readable message.
|
|
metadata: Optional extra metadata.
|
|
"""
|
|
payload = _build_payload(
|
|
event_type=event_type,
|
|
instance=instance,
|
|
status=status,
|
|
message=message,
|
|
metadata=metadata,
|
|
)
|
|
|
|
# Write audit row
|
|
await _write_audit_row(
|
|
session=session,
|
|
instance=instance,
|
|
event_type=event_type,
|
|
created_by=created_by,
|
|
status=status or instance.status,
|
|
message=message,
|
|
metadata=metadata,
|
|
)
|
|
|
|
# Publish to bus
|
|
await event_bus.publish(event_type, payload)
|