Skip to content

Structured Events

structured_events.py

Typed event dataclasses for the autoslo structured logging system.

Every structured log emission constructs a :class:BaseStructuredEvent (or :class:QueryRelatedEvent for query-scoped events), passing in an :class:EventType enum member. BaseStructuredEvent.to_dict() serialises the event to a flat dict that the :class:~autoslo.utils.structured_log.StructuredLogHandler can persist to Parquet.

REQUIRED_DETAILS = {EventType.RUN_START: ['workload_name', 'num_queries', 'routing_policy', 'closed_loop'], EventType.RUN_FINISH: ['workload_name'], EventType.COMPLETION: ['success'], EventType.LATENCY_UPDATE: ['old_latency_s', 'latency_s'], EventType.ROUTING: ['slo_violation', 'cost'], EventType.SPIN_UP_DECISION: ['rpu', 'reason', 'autoscaling_policy'], EventType.SPIN_UP_REQUESTED: ['reason'], EventType.SPIN_UP_BLOCKED: ['reason', 'max', 'used', 'reserved', 'available'], EventType.TEAR_DOWN_DECISION: ['reason'], EventType.TEAR_DOWN_REQUESTED: ['reason', 'force'], EventType.SCHEDULED_SPINUP_EXECUTED: ['scheduled_rel_time_s', 'rpu'], EventType.RPU_COUNTERFACTUAL: ['slo_violation', 'cost'], EventType.RPU_SELECTION: ['slo_violation', 'cost'], EventType.QUERY_EXECUTION_FINISH: ['latency_s'], EventType.FORCED_DECISION_POINT: ['force_one_decision_after_query_count'], EventType.SIM_QUERY_ARRIVAL: ['copy_idx', 'phase'], EventType.SIM_QUERY_COMPLETION: ['latency_s', 'started_after_ready']} module-attribute

EventType

Bases: str, Enum

Enumeration of all structured event types.

Source code in src/autoslo/filesystem/structured_events.py
class EventType(str, Enum):
    """Enumeration of all structured event types."""

    # Run lifecycle
    RUN_START = "run_start"
    RUN_FINISH = "run_finish"

    # Query lifecycle
    ARRIVAL = "arrival"
    QUERY_EXECUTION_START = "query_execution_start"
    QUERY_EXECUTION_FINISH = "query_execution_finish"
    COMPLETION = "completion"

    # Routing
    QUERY_ROUTED = "query_routed"
    LATENCY_UPDATE = "latency_update"
    ROUTING_SCORE = "routing_score"
    ROUTING = "routing"

    # Cluster lifecycle
    SPIN_UP_DECISION = "spin_up_decision"
    SPIN_UP_REQUESTED = "spin_up_requested"
    SPIN_UP_STARTED = "spin_up_started"
    SPIN_UP_BLOCKED = "spin_up_blocked"
    CLUSTER_READY = "cluster_ready"

    TEAR_DOWN_DECISION = "tear_down_decision"
    TEAR_DOWN_REQUESTED = "tear_down_requested"
    TEAR_DOWN_BLOCKED = "tear_down_blocked"
    TEAR_DOWN_STARTED = "tear_down_started"
    STATS_COLLECTED = "stats_collected"
    CLUSTER_REMOVED = "cluster_removed"

    # Scheduled spin-up
    SCHEDULED_SPINUP_EXECUTED = "scheduled_spinup_executed"

    # Autoscaler
    RPU_COUNTERFACTUAL = "rpu_counterfactual"
    RPU_SELECTION = "rpu_selection"
    FORCED_DECISION_POINT = "forced_decision_point"

    # Autoscaler counterfactual simulation
    SIM_QUERY_ARRIVAL = "sim_arrival"
    SIM_QUERY_COMPLETION = "sim_completion"

    # ------------------------------------------------------------------
    # Grouped subsets
    # ------------------------------------------------------------------

    @classmethod
    def query_lifecycle_types(cls) -> set[EventType]:
        """Events that track a query from arrival to completion."""
        return {
            cls.ARRIVAL,
            cls.QUERY_EXECUTION_START,
            cls.QUERY_EXECUTION_FINISH,
            cls.COMPLETION,
        }

    @classmethod
    def routing_types(cls) -> set[EventType]:
        """Events emitted during or about query routing."""
        return {
            cls.QUERY_ROUTED,
            cls.LATENCY_UPDATE,
            cls.ROUTING_SCORE,
            cls.ROUTING,
        }

    @classmethod
    def cluster_lifecycle_types(cls) -> set[EventType]:
        """Events that track cluster spin-up, readiness, and tear-down."""
        return {
            cls.SPIN_UP_DECISION,
            cls.SPIN_UP_REQUESTED,
            cls.SPIN_UP_STARTED,
            cls.SPIN_UP_BLOCKED,
            cls.CLUSTER_READY,
            cls.TEAR_DOWN_DECISION,
            cls.TEAR_DOWN_REQUESTED,
            cls.TEAR_DOWN_BLOCKED,
            cls.TEAR_DOWN_STARTED,
            cls.STATS_COLLECTED,
            cls.CLUSTER_REMOVED,
        }

    @classmethod
    def autoscaler_types(cls) -> set[EventType]:
        """Events related to autoscaler RPU decisions."""
        return {
            cls.RPU_COUNTERFACTUAL,
            cls.RPU_SELECTION,
            cls.SIM_QUERY_ARRIVAL,
            cls.SIM_QUERY_COMPLETION,
        }

RUN_START = 'run_start' class-attribute instance-attribute

RUN_FINISH = 'run_finish' class-attribute instance-attribute

ARRIVAL = 'arrival' class-attribute instance-attribute

QUERY_EXECUTION_START = 'query_execution_start' class-attribute instance-attribute

QUERY_EXECUTION_FINISH = 'query_execution_finish' class-attribute instance-attribute

COMPLETION = 'completion' class-attribute instance-attribute

QUERY_ROUTED = 'query_routed' class-attribute instance-attribute

LATENCY_UPDATE = 'latency_update' class-attribute instance-attribute

ROUTING_SCORE = 'routing_score' class-attribute instance-attribute

ROUTING = 'routing' class-attribute instance-attribute

SPIN_UP_DECISION = 'spin_up_decision' class-attribute instance-attribute

SPIN_UP_REQUESTED = 'spin_up_requested' class-attribute instance-attribute

SPIN_UP_STARTED = 'spin_up_started' class-attribute instance-attribute

SPIN_UP_BLOCKED = 'spin_up_blocked' class-attribute instance-attribute

CLUSTER_READY = 'cluster_ready' class-attribute instance-attribute

TEAR_DOWN_DECISION = 'tear_down_decision' class-attribute instance-attribute

TEAR_DOWN_REQUESTED = 'tear_down_requested' class-attribute instance-attribute

TEAR_DOWN_BLOCKED = 'tear_down_blocked' class-attribute instance-attribute

TEAR_DOWN_STARTED = 'tear_down_started' class-attribute instance-attribute

STATS_COLLECTED = 'stats_collected' class-attribute instance-attribute

CLUSTER_REMOVED = 'cluster_removed' class-attribute instance-attribute

SCHEDULED_SPINUP_EXECUTED = 'scheduled_spinup_executed' class-attribute instance-attribute

RPU_COUNTERFACTUAL = 'rpu_counterfactual' class-attribute instance-attribute

RPU_SELECTION = 'rpu_selection' class-attribute instance-attribute

FORCED_DECISION_POINT = 'forced_decision_point' class-attribute instance-attribute

SIM_QUERY_ARRIVAL = 'sim_arrival' class-attribute instance-attribute

SIM_QUERY_COMPLETION = 'sim_completion' class-attribute instance-attribute

query_lifecycle_types() classmethod

Events that track a query from arrival to completion.

Source code in src/autoslo/filesystem/structured_events.py
@classmethod
def query_lifecycle_types(cls) -> set[EventType]:
    """Events that track a query from arrival to completion."""
    return {
        cls.ARRIVAL,
        cls.QUERY_EXECUTION_START,
        cls.QUERY_EXECUTION_FINISH,
        cls.COMPLETION,
    }

routing_types() classmethod

Events emitted during or about query routing.

Source code in src/autoslo/filesystem/structured_events.py
@classmethod
def routing_types(cls) -> set[EventType]:
    """Events emitted during or about query routing."""
    return {
        cls.QUERY_ROUTED,
        cls.LATENCY_UPDATE,
        cls.ROUTING_SCORE,
        cls.ROUTING,
    }

cluster_lifecycle_types() classmethod

Events that track cluster spin-up, readiness, and tear-down.

Source code in src/autoslo/filesystem/structured_events.py
@classmethod
def cluster_lifecycle_types(cls) -> set[EventType]:
    """Events that track cluster spin-up, readiness, and tear-down."""
    return {
        cls.SPIN_UP_DECISION,
        cls.SPIN_UP_REQUESTED,
        cls.SPIN_UP_STARTED,
        cls.SPIN_UP_BLOCKED,
        cls.CLUSTER_READY,
        cls.TEAR_DOWN_DECISION,
        cls.TEAR_DOWN_REQUESTED,
        cls.TEAR_DOWN_BLOCKED,
        cls.TEAR_DOWN_STARTED,
        cls.STATS_COLLECTED,
        cls.CLUSTER_REMOVED,
    }

autoscaler_types() classmethod

Events related to autoscaler RPU decisions.

Source code in src/autoslo/filesystem/structured_events.py
@classmethod
def autoscaler_types(cls) -> set[EventType]:
    """Events related to autoscaler RPU decisions."""
    return {
        cls.RPU_COUNTERFACTUAL,
        cls.RPU_SELECTION,
        cls.SIM_QUERY_ARRIVAL,
        cls.SIM_QUERY_COMPLETION,
    }

BaseStructuredEvent dataclass

Structured log event.

Use this class directly for events that are not query-related. For query-related events, use :class:QueryRelatedEvent.

Source code in src/autoslo/filesystem/structured_events.py
@dataclass
class BaseStructuredEvent:
    """Structured log event.

    Use this class directly for events that are not query-related.
    For query-related events, use :class:`QueryRelatedEvent`.
    """

    rel_time_s: float
    event_type: EventType
    source: str
    cluster_name: str = ""
    details: dict[str, Any] = field(default_factory=dict)
    wall_clock_s: float = field(init=False, default_factory=wall_clock_utc)

    def __post_init__(self) -> None:
        required = REQUIRED_DETAILS.get(self.event_type, [])
        for key in required:
            if key not in self.details:
                raise ValueError(
                    f"Missing required detail '{key}' in {self.event_type} event."
                )

    def to_dict(self) -> dict[str, Any]:
        # vars(self) is a direct __dict__ lookup — faster than iterating
        # dataclass fields.  details is kept as a plain dict; the log
        # handler serialises it to JSON in bulk at flush time.
        d = vars(self).copy()
        d["event_type"] = self.event_type.value
        return d

rel_time_s instance-attribute

event_type instance-attribute

source instance-attribute

cluster_name = '' class-attribute instance-attribute

details = field(default_factory=dict) class-attribute instance-attribute

wall_clock_s = field(init=False, default_factory=wall_clock_utc) class-attribute instance-attribute

__init__(rel_time_s, event_type, source, cluster_name='', details=dict())

__post_init__()

Source code in src/autoslo/filesystem/structured_events.py
def __post_init__(self) -> None:
    required = REQUIRED_DETAILS.get(self.event_type, [])
    for key in required:
        if key not in self.details:
            raise ValueError(
                f"Missing required detail '{key}' in {self.event_type} event."
            )

to_dict()

Source code in src/autoslo/filesystem/structured_events.py
def to_dict(self) -> dict[str, Any]:
    # vars(self) is a direct __dict__ lookup — faster than iterating
    # dataclass fields.  details is kept as a plain dict; the log
    # handler serialises it to JSON in bulk at flush time.
    d = vars(self).copy()
    d["event_type"] = self.event_type.value
    return d

QueryRelatedEvent dataclass

Bases: BaseStructuredEvent

Structured log event related to a specific query.

Source code in src/autoslo/filesystem/structured_events.py
@dataclass
class QueryRelatedEvent(BaseStructuredEvent):
    """Structured log event related to a specific query."""

    query_id: str = ""
    query_text_id: QueryTextId = QueryTextId("")

query_id = '' class-attribute instance-attribute

query_text_id = QueryTextId('') class-attribute instance-attribute

__init__(rel_time_s, event_type, source, cluster_name='', details=dict(), query_id='', query_text_id=QueryTextId(''))

wall_clock_utc()

Return the current UTC wall-clock time as epoch seconds.

Source code in src/autoslo/filesystem/structured_events.py
def wall_clock_utc() -> float:
    """
    Return the current UTC wall-clock time as epoch seconds.
    """
    return datetime.now(tz=timezone.utc).timestamp()