← back to Exo
fix: Normalize Naming
038cc4cdfa0ab23997464338503114b2c531589a · 2025-07-16 16:11:51 +0100 · Arbion Halili
Files touched
M shared/types/events/common.pyM shared/types/events/events.pyM shared/types/events/registry.pyM shared/types/graphs/common.pyM shared/types/states/master.py
Diff
commit 038cc4cdfa0ab23997464338503114b2c531589a
Author: Arbion Halili <99731180+ToxicPine@users.noreply.github.com>
Date: Wed Jul 16 16:11:51 2025 +0100
fix: Normalize Naming
---
shared/types/events/common.py | 20 ++++++++--------
shared/types/events/events.py | 52 ++++++++++++++++++++---------------------
shared/types/events/registry.py | 10 ++++----
shared/types/graphs/common.py | 32 ++++++++++++-------------
shared/types/states/master.py | 4 ++--
5 files changed, 59 insertions(+), 59 deletions(-)
diff --git a/shared/types/events/common.py b/shared/types/events/common.py
index a451efda..b4c3ae40 100644
--- a/shared/types/events/common.py
+++ b/shared/types/events/common.py
@@ -133,7 +133,7 @@ EventCategories = FrozenSet[EventCategory]
assert_literal_union_covers_enum(EventCategory, EventCategoryEnum)
-class Event[SetMembersT: EventCategories | EventCategory](BaseModel):
+class BaseEvent[SetMembersT: EventCategories | EventCategory](BaseModel):
event_type: EventTypes
event_category: SetMembersT
event_id: EventId
@@ -142,7 +142,7 @@ class Event[SetMembersT: EventCategories | EventCategory](BaseModel):
class EventFromEventLog[SetMembersT: EventCategories | EventCategory](BaseModel):
- event: Event[SetMembersT]
+ event: BaseEvent[SetMembersT]
origin: NodeId
idx_in_log: int = Field(gt=0)
@@ -156,14 +156,14 @@ class EventFromEventLog[SetMembersT: EventCategories | EventCategory](BaseModel)
def narrow_event_type[T: EventCategory, Q: EventCategories | EventCategory](
- event: Event[Q],
+ event: BaseEvent[Q],
target_category: T,
-) -> Event[T]:
+) -> BaseEvent[T]:
if target_category not in event.event_category:
raise ValueError(f"Event Does Not Contain Target Category {target_category}")
narrowed_event = event.model_copy(update={"event_category": {target_category}})
- return cast(Event[T], narrowed_event)
+ return cast(BaseEvent[T], narrowed_event)
def narrow_event_from_event_log_type[
@@ -190,7 +190,7 @@ class State[EventCategoryT: EventCategory](BaseModel):
# Definitions for Type Variables
type Saga[EventCategoryT: EventCategory] = Callable[
[State[EventCategoryT], EventFromEventLog[EventCategoryT]],
- Sequence[Event[EventCategories]],
+ Sequence[BaseEvent[EventCategories]],
]
type Apply[EventCategoryT: EventCategory] = Callable[
[State[EventCategoryT], EventFromEventLog[EventCategoryT]],
@@ -206,19 +206,19 @@ class StateAndEvent[EventCategoryT: EventCategory](NamedTuple):
type EffectHandler[EventCategoryT: EventCategory] = Callable[
[StateAndEvent[EventCategoryT], State[EventCategoryT]], None
]
-type EventPublisher = Callable[[Event[Any]], None]
+type EventPublisher = Callable[[BaseEvent[Any]], None]
# A component that can publish events
class EventPublisherProtocol(Protocol):
- def send(self, events: Sequence[Event[EventCategories]]) -> None: ...
+ def send(self, events: Sequence[BaseEvent[EventCategories]]) -> None: ...
# A component that can fetch events to apply
class EventFetcherProtocol[EventCategoryT: EventCategory](Protocol):
def get_events_to_apply(
self, state: State[EventCategoryT]
- ) -> Sequence[Event[EventCategoryT]]: ...
+ ) -> Sequence[BaseEvent[EventCategoryT]]: ...
# A component that can get the effect handler for a saga
@@ -265,5 +265,5 @@ class Command[
type Decide[EventCategoryT: EventCategory, CommandTypeT: CommandTypes] = Callable[
[State[EventCategoryT], Command[EventCategoryT, CommandTypeT]],
- Sequence[Event[EventCategoryT]],
+ Sequence[BaseEvent[EventCategoryT]],
]
diff --git a/shared/types/events/events.py b/shared/types/events/events.py
index 331ad87f..8644b1d7 100644
--- a/shared/types/events/events.py
+++ b/shared/types/events/events.py
@@ -5,9 +5,9 @@ from typing import Literal, Tuple
from shared.types.common import NodeId
from shared.types.events.chunks import GenerationChunk
from shared.types.events.common import (
+ BaseEvent,
ControlPlaneEventTypes,
DataPlaneEventTypes,
- Event,
EventCategoryEnum,
EventTypes,
InstanceEventTypes,
@@ -39,14 +39,14 @@ from shared.types.worker.common import InstanceId, NodeStatus
from shared.types.worker.instances import InstanceParams, TypeOfInstance
from shared.types.worker.runners import RunnerId, RunnerStatus, RunnerStatusType
-TaskEvent = Event[EventCategoryEnum.MutatesTaskState]
-InstanceEvent = Event[EventCategoryEnum.MutatesInstanceState]
-ControlPlaneEvent = Event[EventCategoryEnum.MutatesControlPlaneState]
-DataPlaneEvent = Event[EventCategoryEnum.MutatesDataPlaneState]
-NodePerformanceEvent = Event[EventCategoryEnum.MutatesNodePerformanceState]
+TaskEvent = BaseEvent[EventCategoryEnum.MutatesTaskState]
+InstanceEvent = BaseEvent[EventCategoryEnum.MutatesInstanceState]
+ControlPlaneEvent = BaseEvent[EventCategoryEnum.MutatesControlPlaneState]
+DataPlaneEvent = BaseEvent[EventCategoryEnum.MutatesDataPlaneState]
+NodePerformanceEvent = BaseEvent[EventCategoryEnum.MutatesNodePerformanceState]
-class TaskCreated(Event[EventCategoryEnum.MutatesTaskState]):
+class TaskCreated(BaseEvent[EventCategoryEnum.MutatesTaskState]):
event_type: EventTypes = TaskEventTypes.TaskCreated
task_id: TaskId
task_params: TaskParams[TaskType]
@@ -55,104 +55,104 @@ class TaskCreated(Event[EventCategoryEnum.MutatesTaskState]):
# Covers Cancellation Of Task, Non-Cancelled Tasks Perist
-class TaskDeleted(Event[EventCategoryEnum.MutatesTaskState]):
+class TaskDeleted(BaseEvent[EventCategoryEnum.MutatesTaskState]):
event_type: EventTypes = TaskEventTypes.TaskDeleted
task_id: TaskId
-class TaskStateUpdated(Event[EventCategoryEnum.MutatesTaskState]):
+class TaskStateUpdated(BaseEvent[EventCategoryEnum.MutatesTaskState]):
event_type: EventTypes = TaskEventTypes.TaskStateUpdated
task_state: TaskState[TaskStatusType, TaskType]
-class InstanceCreated(Event[EventCategoryEnum.MutatesInstanceState]):
+class InstanceCreated(BaseEvent[EventCategoryEnum.MutatesInstanceState]):
event_type: EventTypes = InstanceEventTypes.InstanceCreated
instance_id: InstanceId
instance_params: InstanceParams
instance_type: TypeOfInstance
-class InstanceActivated(Event[EventCategoryEnum.MutatesInstanceState]):
+class InstanceActivated(BaseEvent[EventCategoryEnum.MutatesInstanceState]):
event_type: EventTypes = InstanceEventTypes.InstanceActivated
instance_id: InstanceId
-class InstanceDeactivated(Event[EventCategoryEnum.MutatesInstanceState]):
+class InstanceDeactivated(BaseEvent[EventCategoryEnum.MutatesInstanceState]):
event_type: EventTypes = InstanceEventTypes.InstanceDeactivated
instance_id: InstanceId
-class InstanceDeleted(Event[EventCategoryEnum.MutatesInstanceState]):
+class InstanceDeleted(BaseEvent[EventCategoryEnum.MutatesInstanceState]):
event_type: EventTypes = InstanceEventTypes.InstanceDeleted
instance_id: InstanceId
transition: Tuple[InstanceId, InstanceId]
-class InstanceReplacedAtomically(Event[EventCategoryEnum.MutatesInstanceState]):
+class InstanceReplacedAtomically(BaseEvent[EventCategoryEnum.MutatesInstanceState]):
event_type: EventTypes = InstanceEventTypes.InstanceReplacedAtomically
instance_to_replace: InstanceId
new_instance_id: InstanceId
-class RunnerStatusUpdated(Event[EventCategoryEnum.MutatesRunnerStatus]):
+class RunnerStatusUpdated(BaseEvent[EventCategoryEnum.MutatesRunnerStatus]):
event_type: EventTypes = RunnerStatusEventTypes.RunnerStatusUpdated
instance_id: InstanceId
state_update: Tuple[RunnerId, RunnerStatus[RunnerStatusType]]
-class MLXInferenceSagaPrepare(Event[EventCategoryEnum.MutatesTaskSagaState]):
+class MLXInferenceSagaPrepare(BaseEvent[EventCategoryEnum.MutatesTaskSagaState]):
event_type: EventTypes = TaskSagaEventTypes.MLXInferenceSagaPrepare
task_id: TaskId
instance_id: InstanceId
-class MLXInferenceSagaStartPrepare(Event[EventCategoryEnum.MutatesTaskSagaState]):
+class MLXInferenceSagaStartPrepare(BaseEvent[EventCategoryEnum.MutatesTaskSagaState]):
event_type: EventTypes = TaskSagaEventTypes.MLXInferenceSagaStartPrepare
task_id: TaskId
instance_id: InstanceId
-class NodePerformanceMeasured(Event[EventCategoryEnum.MutatesNodePerformanceState]):
+class NodePerformanceMeasured(BaseEvent[EventCategoryEnum.MutatesNodePerformanceState]):
event_type: EventTypes = NodePerformanceEventTypes.NodePerformanceMeasured
node_id: NodeId
node_profile: NodePerformanceProfile
-class WorkerConnected(Event[EventCategoryEnum.MutatesControlPlaneState]):
+class WorkerConnected(BaseEvent[EventCategoryEnum.MutatesControlPlaneState]):
event_type: EventTypes = ControlPlaneEventTypes.WorkerConnected
edge: DataPlaneEdge
-class WorkerStatusUpdated(Event[EventCategoryEnum.MutatesControlPlaneState]):
+class WorkerStatusUpdated(BaseEvent[EventCategoryEnum.MutatesControlPlaneState]):
event_type: EventTypes = ControlPlaneEventTypes.WorkerStatusUpdated
node_id: NodeId
node_state: NodeStatus
-class WorkerDisconnected(Event[EventCategoryEnum.MutatesControlPlaneState]):
+class WorkerDisconnected(BaseEvent[EventCategoryEnum.MutatesControlPlaneState]):
event_type: EventTypes = ControlPlaneEventTypes.WorkerConnected
vertex_id: ControlPlaneEdgeId
-class ChunkGenerated(Event[EventCategoryEnum.MutatesTaskState]):
+class ChunkGenerated(BaseEvent[EventCategoryEnum.MutatesTaskState]):
event_type: EventTypes = StreamingEventTypes.ChunkGenerated
task_id: TaskId
chunk: GenerationChunk
-class DataPlaneEdgeCreated(Event[EventCategoryEnum.MutatesDataPlaneState]):
+class DataPlaneEdgeCreated(BaseEvent[EventCategoryEnum.MutatesDataPlaneState]):
event_type: EventTypes = DataPlaneEventTypes.DataPlaneEdgeCreated
vertex: ControlPlaneEdgeType
-class DataPlaneEdgeReplacedAtomically(Event[EventCategoryEnum.MutatesDataPlaneState]):
+class DataPlaneEdgeReplacedAtomically(BaseEvent[EventCategoryEnum.MutatesDataPlaneState]):
event_type: EventTypes = DataPlaneEventTypes.DataPlaneEdgeReplacedAtomically
edge_id: DataPlaneEdgeId
edge_profile: DataPlaneEdgeProfile
-class DataPlaneEdgeDeleted(Event[EventCategoryEnum.MutatesDataPlaneState]):
+class DataPlaneEdgeDeleted(BaseEvent[EventCategoryEnum.MutatesDataPlaneState]):
event_type: EventTypes = DataPlaneEventTypes.DataPlaneEdgeDeleted
edge_id: DataPlaneEdgeId
@@ -171,7 +171,7 @@ TEST_EVENT_CATEGORIES = frozenset(
)
-class TestEvent(Event[TEST_EVENT_CATEGORIES_TYPE]):
+class TestEvent(BaseEvent[TEST_EVENT_CATEGORIES_TYPE]):
event_category: TEST_EVENT_CATEGORIES_TYPE = TEST_EVENT_CATEGORIES
test_id: int
"""
\ No newline at end of file
diff --git a/shared/types/events/registry.py b/shared/types/events/registry.py
index b1e8b690..6a9beffd 100644
--- a/shared/types/events/registry.py
+++ b/shared/types/events/registry.py
@@ -5,9 +5,9 @@ from pydantic import Field, TypeAdapter
from shared.constants import get_error_reporting_message
from shared.types.events.common import (
+ BaseEvent,
ControlPlaneEventTypes,
DataPlaneEventTypes,
- Event,
EventCategories,
EventTypes,
InstanceEventTypes,
@@ -102,7 +102,7 @@ def check_union_of_all_events_is_consistent_with_registry(
)
-AllEvents = (
+Event = (
TaskCreated
| TaskStateUpdated
| TaskDeleted
@@ -123,8 +123,8 @@ AllEvents = (
)
# Run the sanity check
-check_union_of_all_events_is_consistent_with_registry(EventRegistry, AllEvents)
+check_union_of_all_events_is_consistent_with_registry(EventRegistry, Event)
-_EventType = Annotated[AllEvents, Field(discriminator="event_type")]
-EventParser: TypeAdapter[Event[EventCategories]] = TypeAdapter(_EventType)
+_EventType = Annotated[Event, Field(discriminator="event_type")]
+EventParser: TypeAdapter[BaseEvent[EventCategories]] = TypeAdapter(_EventType)
diff --git a/shared/types/graphs/common.py b/shared/types/graphs/common.py
index d87fcace..301315af 100644
--- a/shared/types/graphs/common.py
+++ b/shared/types/graphs/common.py
@@ -110,7 +110,7 @@ class MutableGraphProtocol(GraphProtocol[EdgeTypeT, VertexTypeT, EdgeIdT, Vertex
self._add_edge(edge.edge_id, edge.edge_data)
-class Graph(
+class BaseGraph(
Generic[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
MutableGraphProtocol[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
):
@@ -122,51 +122,51 @@ class Graph(
# the first element in the return value is the filtered graph; the second is the
# (possibly empty) set of sub-graphs that were detached during filtering.
def filter_by_edge_data(
- graph: Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ graph: BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
keep: VertexIdT,
predicate: Callable[[EdgeData[EdgeTypeT]], bool],
) -> Tuple[
- Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
- Set[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
+ BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ Set[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
]: ...
# the first element in the return value is the filtered graph; the second is the
# (possibly empty) set of sub-graphs that were detached during filtering.
def filter_by_vertex_data(
- graph: Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ graph: BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
keep: VertexIdT,
predicate: Callable[[VertexData[VertexTypeT]], bool],
) -> Tuple[
- Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
- Set[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
+ BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ Set[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
]: ...
def map_vertices_onto_graph(
vertices: Mapping[VertexIdT, VertexData[VertexTypeT]],
- graph: Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
-) -> Tuple[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[VertexIdT]]: ...
+ graph: BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+) -> Tuple[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[VertexIdT]]: ...
def map_edges_onto_graph(
edges: Mapping[EdgeIdT, EdgeData[EdgeTypeT]],
- graph: Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
-) -> Tuple[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[EdgeIdT]]: ...
+ graph: BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+) -> Tuple[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[EdgeIdT]]: ...
def split_graph_by_edge(
- graph: Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ graph: BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
edge: EdgeIdT,
keep: VertexIdT,
) -> Tuple[
- Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
- Set[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
+ BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT],
+ Set[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
]: ...
def merge_graphs_by_edge(
- graphs: Set[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
+ graphs: Set[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT]],
edge: EdgeIdT,
keep: VertexIdT,
-) -> Tuple[Graph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[EdgeIdT]]: ...
+) -> Tuple[BaseGraph[EdgeTypeT, VertexTypeT, EdgeIdT, VertexIdT], Set[EdgeIdT]]: ...
diff --git a/shared/types/states/master.py b/shared/types/states/master.py
index c9036c5d..8a078d09 100644
--- a/shared/types/states/master.py
+++ b/shared/types/states/master.py
@@ -7,7 +7,7 @@ from pydantic import BaseModel, TypeAdapter
from shared.types.common import NodeId
from shared.types.events.common import (
- Event,
+ BaseEvent,
EventCategory,
EventCategoryEnum,
State,
@@ -97,4 +97,4 @@ def get_shard_assignments(
def get_transition_events(
current_instances: Mapping[InstanceId, InstanceParams],
target_instances: Mapping[InstanceId, InstanceParams],
-) -> Sequence[Event[EventCategory]]: ...
+) -> Sequence[BaseEvent[EventCategory]]: ...
← e2a79350 fix: Fix incorrect logic
·
back to Exo
·
Fixed events issue. cc45c7e9 →