← back to Exo
feat: single entrypoint and logging rework
5efe5562d7d5cecec888ab79c42a8c579db4fc59 · 2025-08-26 11:08:09 +0100 · Evan Quiney
Files touched
D mlx-lm-checkM pyproject.tomlA src/exo/__main__.pyM src/exo/main.pyM src/exo/master/api.pyM src/exo/master/election_callback.pyM src/exo/master/forwarder_supervisor.pyM src/exo/master/main.pyM src/exo/master/tests/test_forwarder_supervisor.pyM src/exo/master/tests/test_master.pyM src/exo/shared/constants.pyM src/exo/shared/db/sqlite/connector.pyM src/exo/shared/db/sqlite/event_log_manager.pyM src/exo/shared/ipc/pipe_duplex.pyA src/exo/shared/logging.pyD src/exo/shared/logging/common.pyM src/exo/shared/models/model_meta.pyM src/exo/worker/download/huggingface_utils.pyM src/exo/worker/main.pyM src/exo/worker/runner/communication.pyM src/exo/worker/runner/runner_supervisor.pyM src/exo/worker/runner/utils.pyM src/exo/worker/tests/test_handlers/conftest.pyM src/exo/worker/tests/test_integration/conftest.pyM src/exo/worker/tests/test_integration/test_inference.pyM src/exo/worker/tests/test_multimodel/test_inference_llama70B.pyM src/exo/worker/tests/test_plan/test_worker_plan.pyM src/exo/worker/tests/test_runner_connection.pyM src/exo/worker/tests/test_supervisor/test_memory.pyM src/exo/worker/tests/test_supervisor/test_oom.pyM src/exo/worker/tests/test_supervisor/test_supervisor.pyM src/exo/worker/tests/test_supervisor/test_supervisor_sad.pyM src/exo/worker/utils/profile.pyM src/exo/worker/utils/system_info.pyM src/exo/worker/worker.pyM uv.lock
Diff
commit 5efe5562d7d5cecec888ab79c42a8c579db4fc59
Author: Evan Quiney <evanev7@gmail.com>
Date: Tue Aug 26 11:08:09 2025 +0100
feat: single entrypoint and logging rework
---
mlx-lm-check | 1 -
pyproject.toml | 5 +-
src/exo/__main__.py | 4 +
src/exo/main.py | 41 ++++++++++-
src/exo/master/api.py | 7 +-
src/exo/master/election_callback.py | 9 +--
src/exo/master/forwarder_supervisor.py | 23 +++---
src/exo/master/main.py | 65 ++++++++---------
src/exo/master/tests/test_forwarder_supervisor.py | 23 ++++--
src/exo/master/tests/test_master.py | 5 +-
src/exo/shared/constants.py | 2 +
src/exo/shared/db/sqlite/connector.py | 37 ++++------
src/exo/shared/db/sqlite/event_log_manager.py | 27 ++++---
src/exo/shared/ipc/pipe_duplex.py | 14 ++--
src/exo/shared/logging.py | 61 ++++++++++++++++
src/exo/shared/logging/common.py | 18 -----
src/exo/shared/models/model_meta.py | 5 +-
src/exo/worker/download/huggingface_utils.py | 3 +-
src/exo/worker/main.py | 56 +++++++-------
src/exo/worker/runner/communication.py | 3 +
src/exo/worker/runner/runner_supervisor.py | 45 ++++++------
src/exo/worker/runner/utils.py | 6 +-
src/exo/worker/tests/test_handlers/conftest.py | 5 +-
src/exo/worker/tests/test_integration/conftest.py | 13 +---
.../tests/test_integration/test_inference.py | 20 +++--
.../test_multimodel/test_inference_llama70B.py | 20 +++--
src/exo/worker/tests/test_plan/test_worker_plan.py | 3 +-
src/exo/worker/tests/test_runner_connection.py | 13 ++--
.../worker/tests/test_supervisor/test_memory.py | 4 +-
src/exo/worker/tests/test_supervisor/test_oom.py | 3 +-
.../tests/test_supervisor/test_supervisor.py | 11 +--
.../tests/test_supervisor/test_supervisor_sad.py | 9 ++-
src/exo/worker/utils/profile.py | 11 +--
src/exo/worker/utils/system_info.py | 11 ++-
src/exo/worker/worker.py | 14 ++--
uv.lock | 85 +++++++++++++---------
36 files changed, 390 insertions(+), 292 deletions(-)
diff --git a/mlx-lm-check b/mlx-lm-check
deleted file mode 160000
index d5bdab1a..00000000
--- a/mlx-lm-check
+++ /dev/null
@@ -1 +0,0 @@
-Subproject commit d5bdab1a22b053d75194ce4d225df9fc1635a400
diff --git a/pyproject.toml b/pyproject.toml
index d43868ef..788405ff 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -27,16 +27,17 @@ dependencies = [
"greenlet>=3.2.4",
"huggingface-hub>=0.33.4",
"mlx==0.26.3",
- "mlx-lm @ https://github.com/ml-explore/mlx-lm.git",
+ "mlx-lm==0.26.4",
"psutil>=7.0.0",
"transformers>=4.55.2",
"cobs>=1.2.2",
+ "loguru>=0.7.3",
]
[project.scripts]
exo-master = "exo.master.main:main"
exo-worker = "exo.worker.main:main"
-#exo = "exo.main:main"
+exo = "exo.main:main"
# dependencies only required for development
[dependency-groups]
diff --git a/src/exo/__main__.py b/src/exo/__main__.py
new file mode 100644
index 00000000..6cfe06a5
--- /dev/null
+++ b/src/exo/__main__.py
@@ -0,0 +1,4 @@
+from exo.main import main
+
+if __name__ == "__main__":
+ main()
diff --git a/src/exo/main.py b/src/exo/main.py
index 46b4ca54..bbcc08c9 100644
--- a/src/exo/main.py
+++ b/src/exo/main.py
@@ -1,2 +1,41 @@
+import argparse
+import multiprocessing as mp
+
+from loguru import logger
+
+from exo.master.main import main as master_main
+from exo.shared.constants import EXO_LOG
+from exo.shared.logging import logger_cleanup, logger_setup
+from exo.worker.main import main as worker_main
+
+
def main():
- print("Hello world!")
+ parser = argparse.ArgumentParser(prog="exo")
+ parser.add_argument(
+ "-v", "--verbose", action="store_const", const=1, dest="verbosity", default=0
+ )
+ parser.add_argument(
+ "-vv",
+ "--very-verbose",
+ action="store_const",
+ const=2,
+ dest="verbosity",
+ default=0,
+ )
+ args = parser.parse_args()
+ if type(args.verbosity) is not int: # type: ignore
+ raise TypeError("Verbosity was parsed incorrectly")
+ logger_setup(EXO_LOG, args.verbosity)
+ logger.info("starting exo")
+
+ # This is for future PyInstaller compatibility
+ mp.set_start_method("spawn", force=True)
+
+ worker = mp.Process(target=worker_main, args=(EXO_LOG, args.verbosity))
+ master = mp.Process(target=master_main, args=(EXO_LOG, args.verbosity))
+ worker.start()
+ master.start()
+ worker.join()
+ master.join()
+
+ logger_cleanup()
diff --git a/src/exo/master/api.py b/src/exo/master/api.py
index a347f7d4..f37418a4 100644
--- a/src/exo/master/api.py
+++ b/src/exo/master/api.py
@@ -9,6 +9,7 @@ from fastapi import FastAPI, HTTPException
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse
from fastapi.staticfiles import StaticFiles
+from loguru import logger
from exo.shared.db.sqlite.connector import AsyncSQLiteEventStorage
from exo.shared.models.model_cards import MODEL_CARDS
@@ -184,7 +185,7 @@ class API:
chunk_response: ChatCompletionResponse = chunk_to_response(
event.chunk
)
- print(chunk_response)
+ logger.debug(chunk_response)
yield f"data: {chunk_response.model_dump_json()}\n\n"
if event.chunk.finish_reason is not None:
@@ -197,7 +198,9 @@ class API:
return
async def _trigger_notify_user_to_download_model(self, model_id: str) -> None:
- print("TODO: we should send a notification to the user to download the model")
+ logger.warning(
+ "TODO: we should send a notification to the user to download the model"
+ )
async def chat_completions(
self, payload: ChatCompletionTaskParams
diff --git a/src/exo/master/election_callback.py b/src/exo/master/election_callback.py
index 92569f3b..0d2ad65c 100644
--- a/src/exo/master/election_callback.py
+++ b/src/exo/master/election_callback.py
@@ -1,4 +1,4 @@
-from logging import Logger
+from loguru import logger
from exo.master.forwarder_supervisor import ForwarderRole, ForwarderSupervisor
@@ -9,16 +9,15 @@ class ElectionCallbacks:
No event system involvement - just direct forwarder control.
"""
- def __init__(self, forwarder_supervisor: ForwarderSupervisor, logger: Logger):
+ def __init__(self, forwarder_supervisor: ForwarderSupervisor):
self._forwarder_supervisor = forwarder_supervisor
- self._logger = logger
async def on_became_master(self) -> None:
"""Called when this node is elected as master"""
- self._logger.info("Node elected as master")
+ logger.info("Node elected as master")
await self._forwarder_supervisor.notify_role_change(ForwarderRole.MASTER)
async def on_became_replica(self) -> None:
"""Called when this node becomes a replica"""
- self._logger.info("Node demoted to replica")
+ logger.info("Node demoted to replica")
await self._forwarder_supervisor.notify_role_change(ForwarderRole.REPLICA)
diff --git a/src/exo/master/forwarder_supervisor.py b/src/exo/master/forwarder_supervisor.py
index a1fb6120..1ff87d5d 100644
--- a/src/exo/master/forwarder_supervisor.py
+++ b/src/exo/master/forwarder_supervisor.py
@@ -2,9 +2,10 @@ import asyncio
import contextlib
import os
from enum import Enum
-from logging import Logger
from pathlib import Path
+from loguru import logger
+
from exo.shared.constants import (
EXO_GLOBAL_EVENT_DB,
EXO_WORKER_EVENT_DB,
@@ -40,12 +41,10 @@ class ForwarderSupervisor:
self,
node_id: NodeId,
forwarder_binary_path: Path,
- logger: Logger,
health_check_interval: float = 5.0,
):
self.node_id = node_id
self._binary_path = forwarder_binary_path
- self._logger = logger
self._health_check_interval = health_check_interval
self._current_role: ForwarderRole | None = None
self._process: asyncio.subprocess.Process | None = None
@@ -57,10 +56,11 @@ class ForwarderSupervisor:
This is the main public interface.
"""
if self._current_role == new_role:
- self._logger.debug(f"Role unchanged: {new_role}")
+ logger.debug(f"Role unchanged: {new_role}")
return
-
- self._logger.info(f"Role changing from {self._current_role} to {new_role}")
+ logger.bind(user_facing=True).info(
+ f"Node changing from {self._current_role} to {new_role}"
+ )
self._current_role = new_role
await self._restart_with_role(new_role)
@@ -119,8 +119,7 @@ class ForwarderSupervisor:
stderr=None,
env=env_vars,
)
-
- self._logger.info(f"Starting forwarder with forwarding pairs: {pairs}")
+ logger.info(f"Starting forwarder with forwarding pairs: {pairs}")
# Start health monitoring
self._health_check_task = asyncio.create_task(self._monitor_health())
@@ -141,7 +140,9 @@ class ForwarderSupervisor:
self._process.terminate()
await asyncio.wait_for(self._process.wait(), timeout=5.0)
except asyncio.TimeoutError:
- self._logger.warning("Forwarder didn't terminate, killing")
+ logger.bind(user_facing=True).warning(
+ "Forwarder didn't terminate, killing"
+ )
self._process.kill()
await self._process.wait()
except ProcessLookupError:
@@ -158,7 +159,9 @@ class ForwarderSupervisor:
self._process.wait(), timeout=self._health_check_interval
)
# Process exited
- self._logger.error(f"Forwarder exited with code {retcode}")
+ logger.bind(user_facing=True).error(
+ f"Forwarder died with code {retcode}"
+ )
# Auto-restart
await asyncio.sleep(0.2) # Brief delay before restart
diff --git a/src/exo/master/main.py b/src/exo/master/main.py
index e7f982cb..18d77c4a 100644
--- a/src/exo/master/main.py
+++ b/src/exo/master/main.py
@@ -1,20 +1,21 @@
import asyncio
-import logging
import os
import threading
-import traceback
from pathlib import Path
-from typing import List
+
+from loguru import logger
from exo.master.api import start_fastapi_server
from exo.master.election_callback import ElectionCallbacks
from exo.master.forwarder_supervisor import ForwarderRole, ForwarderSupervisor
from exo.master.placement import get_instance_placements, get_transition_events
from exo.shared.apply import apply
+from exo.shared.constants import EXO_MASTER_LOG
from exo.shared.db.sqlite.config import EventLogConfig
from exo.shared.db.sqlite.connector import AsyncSQLiteEventStorage
from exo.shared.db.sqlite.event_log_manager import EventLogManager
from exo.shared.keypair import Keypair, get_node_id_keypair
+from exo.shared.logging import logger_cleanup, logger_setup
from exo.shared.types.common import CommandId, NodeId
from exo.shared.types.events import (
Event,
@@ -46,7 +47,6 @@ class Master:
global_events: AsyncSQLiteEventStorage,
worker_events: AsyncSQLiteEventStorage,
forwarder_binary_path: Path,
- logger: logging.Logger,
):
self.state = State()
self.node_id_keypair = node_id_keypair
@@ -56,10 +56,10 @@ class Master:
self.worker_events = worker_events
self.command_task_mapping: dict[CommandId, TaskId] = {}
self.forwarder_supervisor = ForwarderSupervisor(
- self.node_id, forwarder_binary_path=forwarder_binary_path, logger=logger
+ self.node_id,
+ forwarder_binary_path=forwarder_binary_path,
)
- self.election_callbacks = ElectionCallbacks(self.forwarder_supervisor, logger)
- self.logger = logger
+ self.election_callbacks = ElectionCallbacks(self.forwarder_supervisor)
@property
def event_log_for_reads(self) -> AsyncSQLiteEventStorage:
@@ -85,7 +85,10 @@ class Master:
):
# for now we do one command at a time
next_command = self.command_buffer.pop(0)
- self.logger.info(f"got command: {next_command}")
+
+ logger.bind(user_facing=True).info(f"Executing command: {next_command}")
+ logger.info(f"Got command: {next_command}")
+
# TODO: validate the command
match next_command:
case ChatCompletionCommand():
@@ -152,13 +155,17 @@ class Master:
if len(events) == 0:
await asyncio.sleep(0.01)
return
- self.logger.debug(f"got events: {events}")
+
+ if len(events) == 1:
+ logger.debug(f"Master received event: {events[0]}")
+ else:
+ logger.debug(f"Master received events: {events}")
# 3. for each event, apply it to the state
for event_from_log in events:
- self.logger.debug(f"applying event: {event_from_log}")
+ logger.trace(f"Applying event: {event_from_log}")
self.state = apply(self.state, event_from_log)
- self.logger.debug(f"state: {self.state.model_dump_json()}")
+ logger.trace(f"State: {self.state.model_dump_json()}")
# TODO: This can be done in a better place. But for now, we use this to check if any running instances have been broken.
write_events: list[Event] = []
@@ -216,46 +223,33 @@ class Master:
try:
await self._run_event_loop_body()
except Exception as e:
- self.logger.error(f"Error in _run_event_loop_body: {e}")
- traceback.print_exc()
+ logger.opt(exception=e).error(f"Error in _run_event_loop_body: {e}")
await asyncio.sleep(0.1)
async def async_main():
- logger = logging.getLogger("master_logger")
- logger.setLevel(logging.INFO)
- if not logger.handlers:
- handler = logging.StreamHandler()
- handler.setFormatter(
- logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
- )
- logger.addHandler(handler)
-
node_id_keypair = get_node_id_keypair()
node_id = NodeId(node_id_keypair.to_peer_id().to_base58())
- event_log_manager = EventLogManager(EventLogConfig(), logger=logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
global_events: AsyncSQLiteEventStorage = event_log_manager.global_events
worker_events: AsyncSQLiteEventStorage = event_log_manager.worker_events
- command_buffer: List[Command] = []
+ command_buffer: list[Command] = []
+ logger.info("Starting EXO Master")
logger.info(f"Starting Master with node_id: {node_id}")
+ api_port = int(os.environ.get("API_PORT", 8000))
+
api_thread = threading.Thread(
target=start_fastapi_server,
- args=(
- command_buffer,
- global_events,
- lambda: master.state,
- "0.0.0.0",
- int(os.environ.get("API_PORT", 8000)),
- ),
+ args=(command_buffer, global_events, lambda: master.state, "0.0.0.0", api_port),
daemon=True,
)
api_thread.start()
- logger.info("Running FastAPI server in a separate thread. Listening on port 8000.")
+ logger.bind(user_facing=True).info(f"Dashboard started on port {api_port}.")
master = Master(
node_id_keypair,
@@ -263,13 +257,14 @@ async def async_main():
command_buffer,
global_events,
worker_events,
- forwarder_binary_path=Path(os.environ["GO_BUILD_DIR"]) / "forwarder",
- logger=logger,
+ Path(os.environ["GO_BUILD_DIR"]) / "forwarder",
)
await master.run()
+ logger_cleanup() # pyright: ignore[reportUnreachable]
-def main():
+def main(logfile: Path = EXO_MASTER_LOG, verbosity: int = 1):
+ logger_setup(logfile, verbosity)
asyncio.run(async_main())
diff --git a/src/exo/master/tests/test_forwarder_supervisor.py b/src/exo/master/tests/test_forwarder_supervisor.py
index 1ac45bbd..dabdf5cb 100644
--- a/src/exo/master/tests/test_forwarder_supervisor.py
+++ b/src/exo/master/tests/test_forwarder_supervisor.py
@@ -1,5 +1,5 @@
"""
-Comprehensive unit tests for Forwardersupervisor.
+Comprehensive unit tests for ForwarderSupervisor.
Tests basic functionality, process management, and edge cases.
"""
@@ -25,6 +25,7 @@ from exo.shared.constants import (
LIBP2P_GLOBAL_EVENTS_TOPIC,
LIBP2P_WORKER_EVENTS_TOPIC,
)
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import NodeId
# Mock forwarder script content
@@ -191,10 +192,11 @@ class TestForwardersupervisorBasic:
],
) -> None:
"""Test starting forwarder in replica mode."""
+ logger_test_install(test_logger)
# Set environment
os.environ.update(mock_env_vars)
- supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script, test_logger)
+ supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script)
await supervisor.start_as_replica()
# Track the process for cleanup
@@ -236,9 +238,10 @@ class TestForwardersupervisorBasic:
],
) -> None:
"""Test changing role from replica to master."""
+ logger_test_install(test_logger)
os.environ.update(mock_env_vars)
- supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script, test_logger)
+ supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script)
await supervisor.start_as_replica()
if supervisor.process:
@@ -282,9 +285,10 @@ class TestForwardersupervisorBasic:
],
) -> None:
"""Test that setting the same role twice doesn't restart the process."""
+ logger_test_install(test_logger)
os.environ.update(mock_env_vars)
- supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script, test_logger)
+ supervisor = ForwarderSupervisor(NodeId(), mock_forwarder_script)
await supervisor.start_as_replica()
original_pid = supervisor.process_pid
@@ -312,6 +316,7 @@ class TestForwardersupervisorBasic:
],
) -> None:
"""Test that Forwardersupervisor restarts the process if it crashes."""
+ logger_test_install(test_logger)
# Configure mock to exit after 1 second
mock_env_vars["MOCK_EXIT_AFTER"] = "1"
mock_env_vars["MOCK_EXIT_CODE"] = "1"
@@ -320,7 +325,6 @@ class TestForwardersupervisorBasic:
supervisor = ForwarderSupervisor(
NodeId(),
mock_forwarder_script,
- test_logger,
health_check_interval=0.5, # Faster health checks for testing
)
await supervisor.start_as_replica()
@@ -361,9 +365,10 @@ class TestForwardersupervisorBasic:
self, test_logger: logging.Logger, temp_dir: Path
) -> None:
"""Test behavior when forwarder binary doesn't exist."""
+ logger_test_install(test_logger)
nonexistent_path = temp_dir / "nonexistent_forwarder"
- supervisor = ForwarderSupervisor(NodeId(), nonexistent_path, test_logger)
+ supervisor = ForwarderSupervisor(NodeId(), nonexistent_path)
# Should raise FileNotFoundError
with pytest.raises(FileNotFoundError):
@@ -376,10 +381,11 @@ class TestElectionCallbacks:
@pytest.mark.asyncio
async def test_on_became_master(self, test_logger: logging.Logger) -> None:
"""Test callback when becoming master."""
+ logger_test_install(test_logger)
mock_supervisor = MagicMock(spec=ForwarderSupervisor)
mock_supervisor.notify_role_change = AsyncMock()
- callbacks = ElectionCallbacks(mock_supervisor, test_logger)
+ callbacks = ElectionCallbacks(mock_supervisor)
await callbacks.on_became_master()
mock_supervisor.notify_role_change.assert_called_once_with(ForwarderRole.MASTER) # type: ignore
@@ -387,10 +393,11 @@ class TestElectionCallbacks:
@pytest.mark.asyncio
async def test_on_became_replica(self, test_logger: logging.Logger) -> None:
"""Test callback when becoming replica."""
+ logger_test_install(test_logger)
mock_supervisor = MagicMock(spec=ForwarderSupervisor)
mock_supervisor.notify_role_change = AsyncMock()
- callbacks = ElectionCallbacks(mock_supervisor, test_logger)
+ callbacks = ElectionCallbacks(mock_supervisor)
await callbacks.on_became_replica()
mock_supervisor.notify_role_change.assert_called_once_with( # type: ignore
diff --git a/src/exo/master/tests/test_master.py b/src/exo/master/tests/test_master.py
index 5e63ce52..cc0c02ad 100644
--- a/src/exo/master/tests/test_master.py
+++ b/src/exo/master/tests/test_master.py
@@ -11,6 +11,7 @@ from exo.shared.db.sqlite.config import EventLogConfig
from exo.shared.db.sqlite.connector import AsyncSQLiteEventStorage
from exo.shared.db.sqlite.event_log_manager import EventLogManager
from exo.shared.keypair import Keypair
+from exo.shared.logging import logger_test_install
from exo.shared.types.api import ChatCompletionMessage, ChatCompletionTaskParams
from exo.shared.types.common import NodeId
from exo.shared.types.events import Event, EventFromEventLog, Heartbeat, TaskCreated
@@ -53,7 +54,8 @@ def _create_forwarder_dummy_binary() -> Path:
@pytest.mark.asyncio
async def test_master():
logger = Logger(name="test_master_logger")
- event_log_manager = EventLogManager(EventLogConfig(), logger=logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
global_events: AsyncSQLiteEventStorage = event_log_manager.global_events
await global_events.delete_all_events()
@@ -85,7 +87,6 @@ async def test_master():
command_buffer=command_buffer,
global_events=global_events,
forwarder_binary_path=forwarder_binary_path,
- logger=logger,
worker_events=global_events,
)
asyncio.create_task(master.run())
diff --git a/src/exo/shared/constants.py b/src/exo/shared/constants.py
index eb7b7ba9..2be7d1f2 100644
--- a/src/exo/shared/constants.py
+++ b/src/exo/shared/constants.py
@@ -10,6 +10,8 @@ EXO_MASTER_STATE = EXO_HOME / "master_state.json"
EXO_WORKER_STATE = EXO_HOME / "worker_state.json"
EXO_MASTER_LOG = EXO_HOME / "master.log"
EXO_WORKER_LOG = EXO_HOME / "worker.log"
+EXO_LOG = EXO_HOME / "exo.log"
+EXO_TEST_LOG = EXO_HOME / "exo_test.log"
EXO_NODE_ID_KEYPAIR = EXO_HOME / "node_id.keypair"
diff --git a/src/exo/shared/db/sqlite/connector.py b/src/exo/shared/db/sqlite/connector.py
index 7a6d0767..5cb514b8 100644
--- a/src/exo/shared/db/sqlite/connector.py
+++ b/src/exo/shared/db/sqlite/connector.py
@@ -4,10 +4,10 @@ import json
import random
from asyncio import Queue, Task
from collections.abc import Sequence
-from logging import Logger, getLogger
from pathlib import Path
from typing import Any, cast
+from loguru import logger
from sqlalchemy import text
from sqlalchemy.exc import OperationalError
from sqlalchemy.ext.asyncio import AsyncConnection, AsyncSession, create_async_engine
@@ -41,15 +41,12 @@ class AsyncSQLiteEventStorage:
batch_timeout_ms: int,
debounce_ms: int,
max_age_ms: int,
- logger: Logger | None = None,
):
self._db_path = Path(db_path)
self._batch_size = batch_size
self._batch_timeout_s = batch_timeout_ms / 1000.0
self._debounce_s = debounce_ms / 1000.0
self._max_age_s = max_age_ms / 1000.0
- self._logger = logger or getLogger(__name__)
-
self._write_queue: Queue[tuple[Event, NodeId]] = Queue()
self._batch_writer_task: Task[None] | None = None
self._engine = None
@@ -65,7 +62,7 @@ class AsyncSQLiteEventStorage:
# Start batch writer
self._batch_writer_task = asyncio.create_task(self._batch_writer())
- self._logger.info(f"Started SQLite event storage: {self._db_path}")
+ logger.info(f"Started SQLite event storage: {self._db_path}")
async def append_events(self, events: Sequence[Event], origin: NodeId) -> None:
"""Append events to the log (fire-and-forget). The writes are batched and committed
@@ -162,7 +159,7 @@ class AsyncSQLiteEventStorage:
if self._engine is not None:
await self._engine.dispose()
- self._logger.info("Closed SQLite event storage")
+ logger.info("Closed SQLite event storage")
async def delete_all_events(self) -> None:
"""Delete all events from the database."""
@@ -239,10 +236,10 @@ class AsyncSQLiteEventStorage:
)
)
- self._logger.info("Events table and indexes created successfully")
+ logger.info("Events table and indexes created successfully")
except OperationalError as e:
# Even with IF NOT EXISTS, log any unexpected errors
- self._logger.error(f"Error creating table: {e}")
+ logger.error(f"Error creating table: {e}")
# Re-check if table exists now
result = await conn.execute(
text(
@@ -252,11 +249,11 @@ class AsyncSQLiteEventStorage:
if result.fetchone() is None:
raise RuntimeError(f"Failed to create events table: {e}") from e
else:
- self._logger.info(
+ logger.info(
"Events table exists (likely created by another process)"
)
else:
- self._logger.debug("Events table already exists")
+ logger.debug("Events table already exists")
# Enable WAL mode and other optimizations with retry logic
await self._execute_pragma_with_retry(
@@ -338,20 +335,18 @@ class AsyncSQLiteEventStorage:
await session.commit()
if len([ev for ev in batch if not isinstance(ev[0], Heartbeat)]) > 0:
- self._logger.debug(f"Committed batch of {len(batch)} events")
+ logger.debug(f"Committed batch of {len(batch)} events")
except OperationalError as e:
if "database is locked" in str(e):
- self._logger.warning(
- f"Database locked during batch commit, will retry: {e}"
- )
+ logger.warning(f"Database locked during batch commit, will retry: {e}")
# Retry with exponential backoff
await self._commit_batch_with_retry(batch)
else:
- self._logger.error(f"Failed to commit batch: {e}")
+ logger.error(f"Failed to commit batch: {e}")
raise
except Exception as e:
- self._logger.error(f"Failed to commit batch: {e}")
+ logger.error(f"Failed to commit batch: {e}")
raise
async def _execute_pragma_with_retry(
@@ -372,13 +367,13 @@ class AsyncSQLiteEventStorage:
float,
base_delay * (2**retry_count) + random.uniform(0, 0.1),
)
- self._logger.warning(
+ logger.warning(
f"Database locked on '{pragma}', retry {retry_count + 1}/{max_retries} after {delay:.2f}s"
)
await asyncio.sleep(delay)
retry_count += 1
else:
- self._logger.error(
+ logger.error(
f"Failed to execute '{pragma}' after {retry_count + 1} attempts: {e}"
)
raise
@@ -407,7 +402,7 @@ class AsyncSQLiteEventStorage:
await session.commit()
if len([ev for ev in batch if not isinstance(ev[0], Heartbeat)]) > 0:
- self._logger.debug(
+ logger.debug(
f"Committed batch of {len(batch)} events after {retry_count} retries"
)
return
@@ -417,13 +412,13 @@ class AsyncSQLiteEventStorage:
delay = cast(
float, base_delay * (2**retry_count) + random.uniform(0, 0.1)
)
- self._logger.warning(
+ logger.warning(
f"Database locked on batch commit, retry {retry_count + 1}/{max_retries} after {delay:.2f}s"
)
await asyncio.sleep(delay)
retry_count += 1
else:
- self._logger.error(
+ logger.error(
f"Failed to commit batch after {retry_count + 1} attempts: {e}"
)
raise
diff --git a/src/exo/shared/db/sqlite/event_log_manager.py b/src/exo/shared/db/sqlite/event_log_manager.py
index 571d6c8c..00144ffc 100644
--- a/src/exo/shared/db/sqlite/event_log_manager.py
+++ b/src/exo/shared/db/sqlite/event_log_manager.py
@@ -1,7 +1,7 @@
import asyncio
-from logging import Logger
from typing import Dict, Optional, cast
+from loguru import logger
from sqlalchemy.exc import OperationalError
from exo.shared.constants import EXO_HOME
@@ -20,9 +20,8 @@ class EventLogManager:
- Master (replica): writes to worker_events, tails global_events
"""
- def __init__(self, config: EventLogConfig, logger: Logger):
+ def __init__(self, config: EventLogConfig):
self._config = config
- self._logger = logger
self._connectors: Dict[EventLogType, AsyncSQLiteEventStorage] = {}
# Ensure base directory exists
@@ -45,20 +44,20 @@ class EventLogManager:
if "database is locked" in str(e) and retry_count < max_retries - 1:
retry_count += 1
delay = cast(float, 0.5 * (2**retry_count))
- self._logger.warning(
+ logger.warning(
f"Database locked while initializing {log_type.value}, retry {retry_count}/{max_retries} after {delay}s"
)
await asyncio.sleep(delay)
else:
- self._logger.error(
- f"Failed to initialize {log_type.value} after {retry_count + 1} attempts: {e}"
+ logger.opt(exception=e).error(
+ f"Failed to initialize {log_type.value} after {retry_count + 1} attempts"
)
raise RuntimeError(
f"Could not initialize {log_type.value} database after {retry_count + 1} attempts"
) from e
except Exception as e:
- self._logger.error(
- f"Unexpected error initializing {log_type.value}: {e}"
+ logger.opt(exception=e).error(
+ f"Unexpected error initializing {log_type.value}"
)
raise
@@ -66,8 +65,7 @@ class EventLogManager:
raise RuntimeError(
f"Could not initialize {log_type.value} database after {max_retries} attempts"
) from last_error
-
- self._logger.info("Initialized all event log connectors")
+ logger.bind(user_facing=True).info("Initialized all event log connectors")
async def get_connector(self, log_type: EventLogType) -> AsyncSQLiteEventStorage:
"""Get or create a connector for the specified log type"""
@@ -81,18 +79,19 @@ class EventLogManager:
batch_timeout_ms=self._config.batch_timeout_ms,
debounce_ms=self._config.debounce_ms,
max_age_ms=self._config.max_age_ms,
- logger=self._logger,
)
# Start the connector (creates tables if needed)
await connector.start()
self._connectors[log_type] = connector
- self._logger.info(
+ logger.bind(user_facing=True).info(
f"Initialized {log_type.value} connector at {db_path}"
)
except Exception as e:
- self._logger.error(f"Failed to create {log_type.value} connector: {e}")
+ logger.bind(user_facing=True).opt(exception=e).error(
+ f"Failed to create {log_type.value} connector"
+ )
raise
return self._connectors[log_type]
@@ -119,5 +118,5 @@ class EventLogManager:
"""Close all open connectors"""
for log_type, connector in self._connectors.items():
await connector.close()
- self._logger.info(f"Closed {log_type.value} connector")
+ logger.bind(user_facing=True).info(f"Closed {log_type.value} connector")
self._connectors.clear()
diff --git a/src/exo/shared/ipc/pipe_duplex.py b/src/exo/shared/ipc/pipe_duplex.py
index 3ba5a98e..0f1f3178 100644
--- a/src/exo/shared/ipc/pipe_duplex.py
+++ b/src/exo/shared/ipc/pipe_duplex.py
@@ -69,10 +69,10 @@ class PipeDuplex:
"""
def __init__(
- self,
- in_pipe: StrPath,
- out_pipe: StrPath,
- in_callback: Callable[[bytes], None],
+ self,
+ in_pipe: StrPath,
+ out_pipe: StrPath,
+ in_callback: Callable[[bytes], None],
):
assert in_pipe != out_pipe # they must be different files
@@ -156,7 +156,7 @@ def _ensure_fifo_exists(path: StrPath):
def _pipe_buffer_reader(
- path: StrPath, mq: MQueueT[bytes], started: MEventT, kill: MEventT
+ path: StrPath, mq: MQueueT[bytes], started: MEventT, kill: MEventT
):
# TODO: right now the `kill` control flow is somewhat haphazard -> ensure every loop-y or blocking part always
# checks for kill.is_set() and returns/cleans up early if so
@@ -241,7 +241,7 @@ def _pipe_buffer_reader(
def _binary_object_dispatcher(
- mq: MQueueT[bytes], callback: Callable[[bytes], None], kill: TEventT
+ mq: MQueueT[bytes], callback: Callable[[bytes], None], kill: TEventT
):
while not kill.is_set():
# try to get with timeout (to allow to read the kill-flag)
@@ -255,7 +255,7 @@ def _binary_object_dispatcher(
def _pipe_buffer_writer(
- path: StrPath, mq: MQueueT[bytes], started: MEventT, kill: MEventT
+ path: StrPath, mq: MQueueT[bytes], started: MEventT, kill: MEventT
):
# TODO: right now the `kill` control flow is somewhat haphazard -> ensure every loop-y or blocking part always
# checks for kill.is_set() and returns/cleans up early if so
diff --git a/src/exo/shared/logging.py b/src/exo/shared/logging.py
new file mode 100644
index 00000000..4946f1ad
--- /dev/null
+++ b/src/exo/shared/logging.py
@@ -0,0 +1,61 @@
+from __future__ import annotations
+
+import sys
+from logging import Logger
+from pathlib import Path
+
+import loguru
+from loguru import logger
+
+from exo.shared.constants import EXO_TEST_LOG
+
+
+def is_user_facing(record: loguru.Record) -> bool:
+ return ("user_facing" in record["extra"]) and record["extra"]["user_facing"]
+
+
+def logger_setup(log_file: Path, verbosity: int = 0):
+ """Set up logging for this process - formatting, file handles, verbosity and output"""
+ logger.remove()
+ if verbosity == 0:
+ _ = logger.add( # type: ignore
+ sys.__stderr__, # type: ignore
+ format="[ {time:hh:mmA} | <level>{level: <8}</level>] <level>{message}</level>",
+ level="INFO",
+ colorize=True,
+ enqueue=True,
+ filter=is_user_facing,
+ )
+ elif verbosity == 1:
+ _ = logger.add( # type: ignore
+ sys.__stderr__, # type: ignore
+ format="[ {time:hh:mmA} | <level>{level: <8}</level>] <level>{message}</level>",
+ level="INFO",
+ colorize=True,
+ enqueue=True,
+ )
+ else:
+ _ = logger.add( # type: ignore
+ sys.__stderr__, # type: ignore
+ format="[ {time:HH:mm:ss.SSS} | <level>{level: <8}</level> | {name}:{function}:{line} ] <level>{message}</level>",
+ level="DEBUG",
+ colorize=True,
+ )
+ _ = logger.add(
+ log_file,
+ format="[ {time:YYYY-MM-DD HH:mm:ss.SSS} | {level: <8} | {name}:{function}:{line} ] {message}",
+ level="DEBUG",
+ enqueue=True,
+ )
+
+
+def logger_cleanup():
+ """Flush all queues before shutting down so any in-flight logs are written to disk"""
+ logger.complete()
+
+
+def logger_test_install(py_logger: Logger):
+ """Installs a default python logger into the Loguru environment by capturing all its handlers - intended to be used for pytest compatibility, not within the main codebase"""
+ logger_setup(EXO_TEST_LOG, 3)
+ for handler in py_logger.handlers:
+ logger.add(handler)
diff --git a/src/exo/shared/logging/common.py b/src/exo/shared/logging/common.py
deleted file mode 100644
index 52e01f49..00000000
--- a/src/exo/shared/logging/common.py
+++ /dev/null
@@ -1,18 +0,0 @@
-from collections.abc import Set
-from enum import Enum
-from typing import Generic, TypeVar
-
-from pydantic import BaseModel
-
-LogEntryTypeT = TypeVar("LogEntryTypeT", bound=str)
-
-
-class LogEntryType(str, Enum):
- telemetry = "telemetry"
- metrics = "metrics"
- cluster = "cluster"
-
-
-class LogEntry(BaseModel, Generic[LogEntryTypeT]):
- entry_destination: Set[LogEntryType]
- entry_type: LogEntryTypeT
diff --git a/src/exo/shared/models/model_meta.py b/src/exo/shared/models/model_meta.py
index 31260eae..de54536f 100644
--- a/src/exo/shared/models/model_meta.py
+++ b/src/exo/shared/models/model_meta.py
@@ -3,6 +3,7 @@ from typing import Annotated, Dict, Optional
import aiofiles
import aiofiles.os as aios
from huggingface_hub import model_info
+from loguru import logger
from pydantic import BaseModel, Field
from exo.shared.types.models import ModelMetadata
@@ -56,7 +57,7 @@ async def get_config_data(model_id: str) -> ConfigData:
"main",
"config.json",
target_dir,
- lambda curr_bytes, total_bytes: print(
+ lambda curr_bytes, total_bytes: logger.info(
f"Downloading config.json for {model_id}: {curr_bytes}/{total_bytes}"
),
)
@@ -73,7 +74,7 @@ async def get_safetensors_size(model_id: str) -> int:
"main",
"model.safetensors.index.json",
target_dir,
- lambda curr_bytes, total_bytes: print(
+ lambda curr_bytes, total_bytes: logger.info(
f"Downloading model.safetensors.index.json for {model_id}: {curr_bytes}/{total_bytes}"
),
)
diff --git a/src/exo/worker/download/huggingface_utils.py b/src/exo/worker/download/huggingface_utils.py
index 837d5bc3..2e3df1b8 100644
--- a/src/exo/worker/download/huggingface_utils.py
+++ b/src/exo/worker/download/huggingface_utils.py
@@ -5,6 +5,7 @@ from typing import Callable, Dict, Generator, Iterable, List, Optional, TypeVar,
import aiofiles
import aiofiles.os as aios
+from loguru import logger
from exo.shared.types.worker.shards import ShardMetadata
@@ -112,5 +113,5 @@ def get_allow_patterns(weight_map: Dict[str, str], shard: ShardMetadata) -> List
shard_specific_patterns.add(sorted_file_names[-1])
else:
shard_specific_patterns = set(["*.safetensors"])
- print(f"get_allow_patterns {shard=} {shard_specific_patterns=}")
+ logger.info(f"get_allow_patterns {shard=} {shard_specific_patterns=}")
return list(default_patterns | shard_specific_patterns)
diff --git a/src/exo/worker/main.py b/src/exo/worker/main.py
index abd9af78..2db7eedb 100644
--- a/src/exo/worker/main.py
+++ b/src/exo/worker/main.py
@@ -1,9 +1,13 @@
import asyncio
-import logging
+from pathlib import Path
+
+from loguru import logger
from exo.shared.apply import apply
+from exo.shared.constants import EXO_WORKER_LOG
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
from exo.shared.keypair import Keypair, get_node_id_keypair
+from exo.shared.logging import logger_setup, logger_cleanup
from exo.shared.types.common import NodeId
from exo.shared.types.events import (
NodePerformanceMeasured,
@@ -19,46 +23,45 @@ from exo.worker.utils.profile import start_polling_node_metrics
from exo.worker.worker import Worker
-async def run(worker_state: Worker, logger: logging.Logger):
- assert worker_state.global_events is not None
+async def run(worker: Worker):
+ assert worker.global_events is not None
while True:
# 1. get latest events
- events = await worker_state.global_events.get_events_since(
- worker_state.state.last_event_applied_idx
+ events = await worker.global_events.get_events_since(
+ worker.state.last_event_applied_idx
)
# 2. for each event, apply it to the state and run sagas
for event_from_log in events:
- worker_state.state = apply(worker_state.state, event_from_log)
+ worker.state = apply(worker.state, event_from_log)
# 3. based on the updated state, we plan & execute an operation.
op: RunnerOp | None = plan(
- worker_state.assigned_runners,
- worker_state.node_id,
- worker_state.state.instances,
- worker_state.state.runners,
- worker_state.state.tasks,
+ worker.assigned_runners,
+ worker.node_id,
+ worker.state.instances,
+ worker.state.runners,
+ worker.state.tasks,
)
- if op is not None:
- worker_state.logger.info(f"!!! plan result: {op}")
# run the op, synchronously blocking for now
if op is not None:
logger.info(f"Executing op {op}")
+ logger.bind(user_facing=True).debug(f"Worker executing op: {op}")
try:
- async for event in worker_state.execute_op(op):
- await worker_state.event_publisher(event)
+ async for event in worker.execute_op(op):
+ await worker.event_publisher(event)
except Exception as e:
if isinstance(op, ExecuteTaskOp):
- generator = worker_state.fail_task(
+ generator = worker.fail_task(
e, runner_id=op.runner_id, task_id=op.task.task_id
)
else:
- generator = worker_state.fail_runner(e, runner_id=op.runner_id)
+ generator = worker.fail_runner(e, runner_id=op.runner_id)
async for event in generator:
- await worker_state.event_publisher(event)
+ await worker.event_publisher(event)
await asyncio.sleep(0.01)
@@ -66,16 +69,8 @@ async def run(worker_state: Worker, logger: logging.Logger):
async def async_main():
node_id_keypair: Keypair = get_node_id_keypair()
node_id = NodeId(node_id_keypair.to_peer_id().to_base58())
- logger: logging.Logger = logging.getLogger("worker_logger")
- logger.setLevel(logging.DEBUG)
- if not logger.handlers:
- handler = logging.StreamHandler()
- handler.setFormatter(
- logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
- )
- logger.addHandler(handler)
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = exo_shard_downloader()
@@ -96,16 +91,17 @@ async def async_main():
worker = Worker(
node_id,
- logger,
shard_downloader,
event_log_manager.worker_events,
event_log_manager.global_events,
)
- await run(worker, logger)
+ await run(worker)
+ logger_cleanup()
-def main():
+def main(logfile: Path = EXO_WORKER_LOG, verbosity: int = 1):
+ logger_setup(logfile, verbosity)
asyncio.run(async_main())
diff --git a/src/exo/worker/runner/communication.py b/src/exo/worker/runner/communication.py
index 544bf4e8..0b889aa4 100644
--- a/src/exo/worker/runner/communication.py
+++ b/src/exo/worker/runner/communication.py
@@ -2,6 +2,8 @@ import asyncio
import sys
import traceback
+from loguru import logger
+
from exo.shared.types.worker.commands_runner import (
ErrorResponse,
PrintResponse,
@@ -96,3 +98,4 @@ def runner_write_error(error: Exception) -> None:
traceback=traceback.format_exc(),
)
runner_write_response(error_response)
+ logger.opt(exception=error).exception("Critical Runner error")
diff --git a/src/exo/worker/runner/runner_supervisor.py b/src/exo/worker/runner/runner_supervisor.py
index bb9106d9..87cfd7d7 100644
--- a/src/exo/worker/runner/runner_supervisor.py
+++ b/src/exo/worker/runner/runner_supervisor.py
@@ -2,11 +2,11 @@ import asyncio
import contextlib
import traceback
from collections.abc import AsyncGenerator
-from logging import Logger
from types import CoroutineType
from typing import Any, Callable, Optional
import psutil
+from loguru import logger
from exo.shared.types.common import CommandId, Host
from exo.shared.types.events.chunks import GenerationChunk, TokenChunk
@@ -44,13 +44,10 @@ class RunnerSupervisor:
model_shard_meta: ShardMetadata,
hosts: list[Host],
runner_process: asyncio.subprocess.Process,
- logger: Logger,
read_queue: asyncio.Queue[RunnerResponse],
write_queue: asyncio.Queue[RunnerMessage],
stderr_queue: asyncio.Queue[str],
):
- self.logger = logger
-
self.model_shard_meta = model_shard_meta
self.hosts = hosts
self.runner_process = runner_process
@@ -68,7 +65,6 @@ class RunnerSupervisor:
cls,
model_shard_meta: ShardMetadata,
hosts: list[Host],
- logger: Logger,
initialize_timeout: Optional[float] = None,
) -> "RunnerSupervisor":
"""
@@ -91,13 +87,12 @@ class RunnerSupervisor:
model_shard_meta=model_shard_meta,
hosts=hosts,
runner_process=runner_process,
- logger=logger,
read_queue=read_queue,
write_queue=write_queue,
stderr_queue=stderr_queue,
)
- self.logger.info(f"initializing mlx instance with {model_shard_meta=}")
+ logger.info(f"Initializing mlx instance with {model_shard_meta=}")
await self.write_queue.put(
SetupMessage(
model_shard_meta=model_shard_meta,
@@ -111,7 +106,7 @@ class RunnerSupervisor:
response = await self._read_with_error_check(initialize_timeout)
assert isinstance(response, InitializedResponse)
- self.logger.info(f"Runner initialized in {response.time_taken} seconds")
+ logger.info(f"Runner initialized in {response.time_taken} seconds")
return self
@@ -143,7 +138,7 @@ class RunnerSupervisor:
if self.read_task in done:
await self.read_task # Re-raises any exception from read_task
- self.logger.error(
+ logger.error(
"Unreachable code run. We should have raised an error on the read_task being done."
)
@@ -183,14 +178,16 @@ class RunnerSupervisor:
prefil_timeout = get_prefil_timeout(self.model_shard_meta)
token_timeout = get_token_generate_timeout(self.model_shard_meta)
timeout = prefil_timeout
- self.logger.info(f"starting chat completion with timeout {timeout}")
+ logger.bind(user_facing=True).info(
+ f"Starting chat completion with timeout {timeout}"
+ )
while True:
try:
response = await self._read_with_error_check(timeout)
except asyncio.TimeoutError as e:
- self.logger.info(
- f"timed out from timeout duration {timeout} - {'prefil' if timeout == prefil_timeout else 'decoding stage'}"
+ logger.bind(user_facing=True).info(
+ f"Generation timed out during {'prefil' if timeout == prefil_timeout else 'decoding stage'}"
)
raise e
@@ -235,7 +232,8 @@ class RunnerSupervisor:
match response:
case PrintResponse():
- self.logger.info(f"runner printed: {response.text}")
+ # TODO: THIS IS A REALLY IMPORTANT LOG MESSAGE, AND SHOULD BE MADE PRETTIER
+ logger.bind(user_facing=True).info(f"{response.text}")
case ErrorResponse():
## Failure case #1: a crash happens Python, so it's neatly handled by passing an ErrorResponse with the details
await self.read_queue.put(response)
@@ -255,7 +253,7 @@ class RunnerSupervisor:
await await_task(self.write_task)
# Kill the process and all its children
- await kill_process_tree(self.runner_process, self.logger)
+ await kill_process_tree(self.runner_process)
# Wait to make sure that the model has been unloaded from memory
async def wait_for_memory_release() -> None:
@@ -266,8 +264,8 @@ class RunnerSupervisor:
if available_memory_bytes >= required_memory_bytes:
break
if asyncio.get_event_loop().time() - start_time > 30.0:
- self.logger.warning(
- "Timeout waiting for memory release after 30 seconds"
+ logger.warning(
+ "Runner memory not released after 30 seconds - exiting"
)
break
await asyncio.sleep(0.1)
@@ -276,8 +274,8 @@ class RunnerSupervisor:
def __del__(self) -> None:
if self.runner_process.returncode is None:
- print(
- "Warning: RunnerSupervisor was not stopped cleanly before garbage collection. Force killing process tree."
+ logger.warning(
+ "RunnerSupervisor was not stopped cleanly before garbage collection. Force killing process tree."
)
# Can't use async in __del__, so use psutil directly
try:
@@ -321,10 +319,9 @@ class RunnerSupervisor:
except asyncio.QueueEmpty:
break
- # print('STDERR OUTPUT IS')
- # print(stderr_output)
-
- self.logger.error(f"Error {self.runner_process.returncode}: {stderr_output}")
+ logger.bind(user_facing=True).error(
+ f"Runner Error {self.runner_process.returncode}: {stderr_output}"
+ )
return RunnerError(
error_type="MLXCrash",
error_message=stderr_output,
@@ -341,7 +338,7 @@ class RunnerSupervisor:
line = line_bytes.decode("utf-8").strip()
await self.stderr_queue.put(line)
- self.logger.warning(f"Runner stderr read: {line}")
+ logger.warning(f"Runner stderr read: {line}")
except Exception as e:
- self.logger.warning(f"Error reading runner stderr: {e}")
+ logger.warning(f"Error reading runner stderr: {e}")
break
diff --git a/src/exo/worker/runner/utils.py b/src/exo/worker/runner/utils.py
index c5c480ca..328d1a07 100644
--- a/src/exo/worker/runner/utils.py
+++ b/src/exo/worker/runner/utils.py
@@ -1,17 +1,15 @@
import asyncio
import contextlib
import sys
-from logging import Logger
import psutil
+from loguru import logger
from exo.shared.constants import LB_DISK_GBPS, LB_MEMBW_GBPS, LB_TFLOPS
from exo.shared.types.worker.shards import ShardMetadata
-async def kill_process_tree(
- runner_process: asyncio.subprocess.Process, logger: Logger
-) -> None:
+async def kill_process_tree(runner_process: asyncio.subprocess.Process) -> None:
"""Kill the process and all its children forcefully."""
if runner_process.returncode is not None:
return # Process already dead
diff --git a/src/exo/worker/tests/test_handlers/conftest.py b/src/exo/worker/tests/test_handlers/conftest.py
index 7707754a..ccd1b75b 100644
--- a/src/exo/worker/tests/test_handlers/conftest.py
+++ b/src/exo/worker/tests/test_handlers/conftest.py
@@ -4,6 +4,7 @@ from typing import Callable
import pytest
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import NodeId
from exo.shared.types.worker.common import InstanceId
from exo.shared.types.worker.instances import Instance
@@ -24,13 +25,13 @@ def user_message():
@pytest.fixture
async def worker(logger: Logger):
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
shard_downloader = NoopShardDownloader()
await event_log_manager.initialize()
return Worker(
NODE_A,
- logger,
shard_downloader,
worker_events=event_log_manager.global_events,
global_events=event_log_manager.global_events,
diff --git a/src/exo/worker/tests/test_integration/conftest.py b/src/exo/worker/tests/test_integration/conftest.py
index 2f1888ec..b4e0ee7f 100644
--- a/src/exo/worker/tests/test_integration/conftest.py
+++ b/src/exo/worker/tests/test_integration/conftest.py
@@ -6,18 +6,13 @@ import pytest
from exo.shared.db.sqlite.connector import AsyncSQLiteEventStorage
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import NodeId
from exo.worker.download.shard_downloader import NoopShardDownloader
from exo.worker.main import run
from exo.worker.worker import Worker
-@pytest.fixture
-def user_message():
- """Override this fixture in tests to customize the message"""
- return "What is the capital of Japan?"
-
-
@pytest.fixture
def worker_running(
logger: Logger,
@@ -25,7 +20,8 @@ def worker_running(
async def _worker_running(
node_id: NodeId,
) -> tuple[Worker, AsyncSQLiteEventStorage]:
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
global_events = event_log_manager.global_events
@@ -34,12 +30,11 @@ def worker_running(
shard_downloader = NoopShardDownloader()
worker = Worker(
node_id,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker, logger))
+ asyncio.create_task(run(worker))
return worker, global_events
diff --git a/src/exo/worker/tests/test_integration/test_inference.py b/src/exo/worker/tests/test_integration/test_inference.py
index 8262af4a..53e40abe 100644
--- a/src/exo/worker/tests/test_integration/test_inference.py
+++ b/src/exo/worker/tests/test_integration/test_inference.py
@@ -2,9 +2,9 @@ import asyncio
from logging import Logger
from typing import Awaitable, Callable
-# TaskStateUpdated and ChunkGenerated are used in test_worker_integration_utils.py
from exo.shared.db.sqlite.connector import AsyncSQLiteEventStorage
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
+from exo.shared.logging import logger_test_install
from exo.shared.types.api import ChatCompletionMessage, ChatCompletionTaskParams
from exo.shared.types.common import CommandId, Host, NodeId
from exo.shared.types.events import (
@@ -96,7 +96,8 @@ async def test_2_runner_inference(
hosts: Callable[[int], list[Host]],
chat_completion_task: Callable[[InstanceId, TaskId], Task],
):
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = NoopShardDownloader()
@@ -105,21 +106,19 @@ async def test_2_runner_inference(
worker1 = Worker(
NODE_A,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker1, logger))
+ asyncio.create_task(run(worker1))
worker2 = Worker(
NODE_B,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker2, logger))
+ asyncio.create_task(run(worker2))
## Instance
model_id = ModelId("mlx-community/Llama-3.2-1B-Instruct-4bit")
@@ -182,7 +181,8 @@ async def test_2_runner_multi_message(
pipeline_shard_meta: Callable[[int, int], PipelineShardMetadata],
hosts: Callable[[int], list[Host]],
):
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = NoopShardDownloader()
@@ -191,21 +191,19 @@ async def test_2_runner_multi_message(
worker1 = Worker(
NODE_A,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker1, logger))
+ asyncio.create_task(run(worker1))
worker2 = Worker(
NODE_B,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker2, logger))
+ asyncio.create_task(run(worker2))
## Instance
model_id = ModelId("mlx-community/Llama-3.2-1B-Instruct-4bit")
diff --git a/src/exo/worker/tests/test_multimodel/test_inference_llama70B.py b/src/exo/worker/tests/test_multimodel/test_inference_llama70B.py
index b38ab16d..b98d9c5f 100644
--- a/src/exo/worker/tests/test_multimodel/test_inference_llama70B.py
+++ b/src/exo/worker/tests/test_multimodel/test_inference_llama70B.py
@@ -5,8 +5,8 @@ from typing import Callable
import pytest
-# TaskStateUpdated and ChunkGenerated are used in test_worker_integration_utils.py
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
+from exo.shared.logging import logger_test_install
from exo.shared.models.model_meta import get_model_meta
from exo.shared.types.api import ChatCompletionMessage, ChatCompletionTaskParams
from exo.shared.types.common import Host
@@ -90,7 +90,8 @@ async def test_2_runner_inference(
hosts: Callable[[int], list[Host]],
chat_completion_task: Callable[[InstanceId, TaskId], Task],
):
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = NoopShardDownloader()
@@ -99,21 +100,19 @@ async def test_2_runner_inference(
worker1 = Worker(
NODE_A,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker1, logger))
+ asyncio.create_task(run(worker1))
worker2 = Worker(
NODE_B,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker2, logger))
+ asyncio.create_task(run(worker2))
## Instance
model_id = ModelId(MODEL_ID)
@@ -199,7 +198,8 @@ async def test_parallel_inference(
hosts: Callable[[int], list[Host]],
chat_completion_task: Callable[[InstanceId, TaskId], Task],
):
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ logger_test_install(logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = NoopShardDownloader()
@@ -208,21 +208,19 @@ async def test_parallel_inference(
worker1 = Worker(
NODE_A,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker1, logger))
+ asyncio.create_task(run(worker1))
worker2 = Worker(
NODE_B,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
- asyncio.create_task(run(worker2, logger))
+ asyncio.create_task(run(worker2))
## Instance
model_id = ModelId(MODEL_ID)
diff --git a/src/exo/worker/tests/test_plan/test_worker_plan.py b/src/exo/worker/tests/test_plan/test_worker_plan.py
index dd304bd1..bbb59fc1 100644
--- a/src/exo/worker/tests/test_plan/test_worker_plan.py
+++ b/src/exo/worker/tests/test_plan/test_worker_plan.py
@@ -4,6 +4,7 @@ import logging
import pytest
+from exo.shared.logging import logger_test_install
from exo.shared.types.api import ChatCompletionMessage
from exo.shared.types.state import State
from exo.shared.types.tasks import (
@@ -507,13 +508,13 @@ def test_worker_plan(case: PlanTestCase) -> None:
node_id = NODE_A
logger = logging.getLogger("test_worker_plan")
+ logger_test_install(logger)
shard_downloader = NoopShardDownloader()
worker = Worker(
node_id=node_id,
shard_downloader=shard_downloader,
worker_events=None,
global_events=None,
- logger=logger,
)
runner_config: InProcessRunner
diff --git a/src/exo/worker/tests/test_runner_connection.py b/src/exo/worker/tests/test_runner_connection.py
index 196c2401..a561de85 100644
--- a/src/exo/worker/tests/test_runner_connection.py
+++ b/src/exo/worker/tests/test_runner_connection.py
@@ -6,6 +6,7 @@ from typing import Callable
import pytest
from exo.shared.db.sqlite.event_log_manager import EventLogConfig, EventLogManager
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import Host
from exo.shared.types.events import InstanceCreated, InstanceDeleted
from exo.shared.types.models import ModelId
@@ -39,12 +40,13 @@ async def check_runner_connection(
pipeline_shard_meta: Callable[[int, int], PipelineShardMetadata],
hosts: Callable[[int], list[Host]],
) -> bool:
+ logger_test_install(logger)
# Track all tasks and workers for cleanup
tasks: list[asyncio.Task[None]] = []
workers: list[Worker] = []
try:
- event_log_manager = EventLogManager(EventLogConfig(), logger)
+ event_log_manager = EventLogManager(EventLogConfig())
await event_log_manager.initialize()
shard_downloader = NoopShardDownloader()
@@ -53,24 +55,22 @@ async def check_runner_connection(
worker1 = Worker(
NODE_A,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
workers.append(worker1)
- task1 = asyncio.create_task(run(worker1, logger))
+ task1 = asyncio.create_task(run(worker1))
tasks.append(task1)
worker2 = Worker(
NODE_B,
- logger=logger,
shard_downloader=shard_downloader,
worker_events=global_events,
global_events=global_events,
)
workers.append(worker2)
- task2 = asyncio.create_task(run(worker2, logger))
+ task2 = asyncio.create_task(run(worker2))
tasks.append(task2)
model_id = ModelId("mlx-community/Llama-3.2-1B-Instruct-4bit")
@@ -162,6 +162,7 @@ async def check_runner_connection(
# hosts: Callable[[int], list[Host]],
# chat_completion_task: Callable[[InstanceId, str], Task],
# ) -> None:
+# logger_test_install(logger)
# total_runs = 100
# successes = 0
@@ -176,7 +177,6 @@ async def check_runner_connection(
# try:
# result = loop.run_until_complete(check_runner_connection(
-# logger=logger,
# pipeline_shard_meta=pipeline_shard_meta,
# hosts=hosts,
# chat_completion_task=chat_completion_task,
@@ -190,7 +190,6 @@ async def check_runner_connection(
# task.cancel()
# try:
# result = loop.run_until_complete(check_runner_connection(
-# logger=logger,
# pipeline_shard_meta=pipeline_shard_meta,
# hosts=hosts,
# chat_completion_task=chat_completion_task,
diff --git a/src/exo/worker/tests/test_supervisor/test_memory.py b/src/exo/worker/tests/test_supervisor/test_memory.py
index 58b3238a..c7c494ba 100644
--- a/src/exo/worker/tests/test_supervisor/test_memory.py
+++ b/src/exo/worker/tests/test_supervisor/test_memory.py
@@ -5,6 +5,7 @@ from typing import Callable
import psutil
import pytest
+from exo.shared.logging import logger_test_install
from exo.shared.models.model_meta import get_model_meta
from exo.shared.types.common import Host
from exo.shared.types.models import ModelMetadata
@@ -36,13 +37,12 @@ async def test_supervisor_inference_exception(
chat_completion_task: Callable[[InstanceId, TaskId], Task],
logger: Logger,
):
- """Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
process: Process = supervisor.runner_process
diff --git a/src/exo/worker/tests/test_supervisor/test_oom.py b/src/exo/worker/tests/test_supervisor/test_oom.py
index afade315..9b1b4778 100644
--- a/src/exo/worker/tests/test_supervisor/test_oom.py
+++ b/src/exo/worker/tests/test_supervisor/test_oom.py
@@ -3,6 +3,7 @@ from typing import Callable
import pytest
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import Host
from exo.shared.types.tasks import (
Task,
@@ -30,13 +31,13 @@ async def test_supervisor_catches_oom(
chat_completion_task: Callable[[InstanceId, TaskId], Task],
logger: Logger,
):
+ logger_test_install(logger)
"""Test that asking for the capital of France returns 'Paris' in the response"""
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
task = chat_completion_task(INSTANCE_1_ID, TASK_1_ID)
diff --git a/src/exo/worker/tests/test_supervisor/test_supervisor.py b/src/exo/worker/tests/test_supervisor/test_supervisor.py
index 3452044c..17756c18 100644
--- a/src/exo/worker/tests/test_supervisor/test_supervisor.py
+++ b/src/exo/worker/tests/test_supervisor/test_supervisor.py
@@ -4,6 +4,7 @@ from typing import Callable
import pytest
+from exo.shared.logging import logger_test_install
from exo.shared.openai_compat import FinishReason
from exo.shared.types.common import Host
from exo.shared.types.events.chunks import TokenChunk
@@ -32,6 +33,7 @@ async def test_supervisor_single_node_response(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
instance_id = InstanceId()
@@ -40,7 +42,6 @@ async def test_supervisor_single_node_response(
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
try:
@@ -73,13 +74,13 @@ async def test_supervisor_two_node_response(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
instance_id = InstanceId()
async def create_supervisor(shard_idx: int) -> RunnerSupervisor:
supervisor = await RunnerSupervisor.create(
model_shard_meta=pipeline_shard_meta(2, shard_idx),
hosts=hosts(2, offset=15),
- logger=logger,
)
return supervisor
@@ -138,13 +139,13 @@ async def test_supervisor_early_stopping(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
instance_id = InstanceId()
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
task = chat_completion_task(instance_id, TaskId())
@@ -192,12 +193,12 @@ async def test_supervisor_handles_terminated_runner(
logger: Logger,
):
"""Test that the supervisor handles a terminated runner"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
# Terminate the runner
@@ -217,12 +218,12 @@ async def test_supervisor_handles_killed_runner(
logger: Logger,
):
"""Test that the supervisor handles a killed runner"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
assert supervisor.healthy
diff --git a/src/exo/worker/tests/test_supervisor/test_supervisor_sad.py b/src/exo/worker/tests/test_supervisor/test_supervisor_sad.py
index bfaf1580..959e41b2 100644
--- a/src/exo/worker/tests/test_supervisor/test_supervisor_sad.py
+++ b/src/exo/worker/tests/test_supervisor/test_supervisor_sad.py
@@ -4,6 +4,7 @@ from typing import Callable
import pytest
+from exo.shared.logging import logger_test_install
from exo.shared.types.common import Host
from exo.shared.types.tasks import Task, TaskId
from exo.shared.types.worker.common import InstanceId, RunnerError
@@ -19,6 +20,7 @@ async def test_supervisor_instantiation_exception(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
model_shard_meta.immediate_exception = True
@@ -26,7 +28,6 @@ async def test_supervisor_instantiation_exception(
_ = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
@@ -37,6 +38,7 @@ async def test_supervisor_instantiation_timeout(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
model_shard_meta.should_timeout = 10 # timeout after 10s
@@ -44,7 +46,6 @@ async def test_supervisor_instantiation_timeout(
_ = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
@@ -56,12 +57,12 @@ async def test_supervisor_inference_exception(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
task = chat_completion_task(INSTANCE_1_ID, TASK_1_ID)
@@ -79,12 +80,12 @@ async def test_supervisor_inference_timeout(
logger: Logger,
):
"""Test that asking for the capital of France returns 'Paris' in the response"""
+ logger_test_install(logger)
model_shard_meta = pipeline_shard_meta(1, 0)
supervisor = await RunnerSupervisor.create(
model_shard_meta=model_shard_meta,
hosts=hosts(1, offset=10),
- logger=logger,
)
task = chat_completion_task(INSTANCE_1_ID, TASK_1_ID)
diff --git a/src/exo/worker/utils/profile.py b/src/exo/worker/utils/profile.py
index be4a17ea..ab4d3e33 100644
--- a/src/exo/worker/utils/profile.py
+++ b/src/exo/worker/utils/profile.py
@@ -3,6 +3,8 @@ import os
import platform
from typing import Any, Callable, Coroutine
+from loguru import logger
+
from exo.shared.types.profiling import (
MemoryPerformanceProfile,
NodePerformanceProfile,
@@ -20,10 +22,6 @@ from exo.worker.utils.system_info import (
get_network_interface_info_async,
)
-# from exo.infra.event_log import EventLog
-# from exo.app.config import ResourceMonitorConfig
-# from exo.utils.mlx.mlx_utils import profile_flops_fp16
-
async def get_metrics_async() -> Metrics:
"""Return detailed Metrics on macOS or a minimal fallback elsewhere.
@@ -66,7 +64,6 @@ async def start_polling_node_metrics(
)
# Run heavy FLOPs profiling only if enough time has elapsed
-
override_memory_env = os.getenv("OVERRIDE_MEMORY")
override_memory: int | None = (
int(override_memory_env) * 2**30 if override_memory_env else None
@@ -121,11 +118,11 @@ async def start_polling_node_metrics(
except asyncio.TimeoutError:
# One of the operations took too long; skip this iteration but keep the loop alive.
- print(
+ logger.warning(
"[resource_monitor] Operation timed out after 30s, skipping this cycle."
)
except Exception as e:
# Catch-all to ensure the monitor keeps running.
- print(f"[resource_monitor] Encountered error: {e}")
+ logger.opt(exception=e).error("Resource Monitor encountered error")
finally:
await asyncio.sleep(poll_interval_s)
diff --git a/src/exo/worker/utils/system_info.py b/src/exo/worker/utils/system_info.py
index 1aa69325..0c818241 100644
--- a/src/exo/worker/utils/system_info.py
+++ b/src/exo/worker/utils/system_info.py
@@ -3,6 +3,7 @@ import re
import sys
from typing import Dict, List, Optional
+from loguru import logger
from pydantic import BaseModel, Field
from exo.shared.types.profiling import NetworkInterfaceInfo
@@ -22,7 +23,7 @@ async def get_mac_friendly_name_async() -> str | None:
Returns the name as a string, or None if an error occurs or not on macOS.
"""
if sys.platform != "darwin": # 'darwin' is the platform name for macOS
- print("This function is designed for macOS only.")
+ logger.warning("Mac friendly name is designed for macOS only.")
return None
try:
@@ -204,6 +205,14 @@ async def get_mac_system_info_async() -> SystemInfo:
memory_val = 0
network_interfaces_info_list: List[NetworkInterfaceInfo] = []
+ if sys.platform != "darwin":
+ return SystemInfo(
+ model_id=model_id_val,
+ chip_id=chip_id_val,
+ memory=memory_val,
+ network_interfaces=network_interfaces_info_list,
+ )
+
try:
process = await asyncio.create_subprocess_exec(
"system_profiler",
diff --git a/src/exo/worker/worker.py b/src/exo/worker/worker.py
index 0c66dc76..a05b2aae 100644
--- a/src/exo/worker/worker.py
+++ b/src/exo/worker/worker.py
@@ -1,10 +1,11 @@
import asyncio
-import logging
import time
from asyncio import Queue
from functools import partial
from typing import AsyncGenerator, Optional
+from loguru import logger
+
from exo.shared.db.sqlite import AsyncSQLiteEventStorage
from exo.shared.types.common import NodeId
from exo.shared.types.events import (
@@ -52,7 +53,6 @@ class Worker:
def __init__(
self,
node_id: NodeId,
- logger: logging.Logger,
shard_downloader: ShardDownloader,
worker_events: AsyncSQLiteEventStorage | None,
global_events: AsyncSQLiteEventStorage | None,
@@ -64,7 +64,6 @@ class Worker:
worker_events # worker_events is None in some tests.
)
self.global_events: AsyncSQLiteEventStorage | None = global_events
- self.logger: logging.Logger = logger
self.assigned_runners: dict[RunnerId, AssignedRunner] = {}
self._task: asyncio.Task[None] | None = None
@@ -233,7 +232,6 @@ class Worker:
assigned_runner.runner = await RunnerSupervisor.create(
model_shard_meta=assigned_runner.shard_metadata,
hosts=assigned_runner.hosts,
- logger=self.logger,
initialize_timeout=initialize_timeout,
)
@@ -255,9 +253,7 @@ class Worker:
if runner.runner_process.stdout is None:
health_issues.append("runner_process.stdout is None")
- self.logger.warning(
- f"Runner status is not healthy: {', '.join(health_issues)}"
- )
+ logger.warning(f"Runner status is not healthy: {', '.join(health_issues)}")
assigned_runner.status = FailedRunnerStatus()
yield self.assigned_runners[op.runner_id].status_update_event()
@@ -375,7 +371,7 @@ class Worker:
try:
await asyncio.wait_for(task, timeout=5)
except asyncio.TimeoutError:
- self.logger.warning(
+ logger.warning(
"Timed out waiting for task cleanup after inference execution."
)
@@ -436,4 +432,4 @@ class Worker:
async def event_publisher(self, event: Event) -> None:
assert self.worker_events is not None
await self.worker_events.append_events([event], self.node_id)
- self.logger.info(f"published event: {event}")
+ logger.info(f"published event: {event}")
diff --git a/uv.lock b/uv.lock
index bb4af869..9abbbc8c 100644
--- a/uv.lock
+++ b/uv.lock
@@ -256,6 +256,7 @@ dependencies = [
{ name = "filelock", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "greenlet", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "huggingface-hub", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
+ { name = "loguru", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "mlx", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "mlx-lm", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "networkx", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
@@ -298,9 +299,10 @@ requires-dist = [
{ name = "filelock", specifier = ">=3.18.0" },
{ name = "greenlet", specifier = ">=3.2.4" },
{ name = "huggingface-hub", specifier = ">=0.33.4" },
+ { name = "loguru", specifier = ">=0.7.3" },
{ name = "mlx", specifier = "==0.26.3" },
{ name = "mlx", marker = "extra == 'darwin'" },
- { name = "mlx-lm", git = "https://github.com/ml-explore/mlx-lm.git" },
+ { name = "mlx-lm", specifier = "==0.26.4" },
{ name = "networkx", specifier = ">=3.5" },
{ name = "openai", specifier = ">=1.99.9" },
{ name = "pathlib", specifier = ">=1.0.1" },
@@ -565,6 +567,15 @@ wheels = [
{ url = "https://files.pythonhosted.org/packages/b3/4a/4175a563579e884192ba6e81725fc0448b042024419be8d83aa8a80a3f44/jiter-0.10.0-cp314-cp314t-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:3aa96f2abba33dc77f79b4cf791840230375f9534e5fac927ccceb58c5e604a5", size = 354213, upload-time = "2025-05-18T19:04:41.894Z" },
]
+[[package]]
+name = "loguru"
+version = "0.7.3"
+source = { registry = "https://pypi.org/simple" }
+sdist = { url = "https://files.pythonhosted.org/packages/3a/05/a1dae3dffd1116099471c643b8924f5aa6524411dc6c63fdae648c4f1aca/loguru-0.7.3.tar.gz", hash = "sha256:19480589e77d47b8d85b2c827ad95d49bf31b0dcde16593892eb51dd18706eb6", size = 63559, upload-time = "2024-12-06T11:20:56.608Z" }
+wheels = [
+ { url = "https://files.pythonhosted.org/packages/0c/29/0348de65b8cc732daa3e33e67806420b2ae89bdce2b04af740289c5c6c8c/loguru-0.7.3-py3-none-any.whl", hash = "sha256:31a33c10c8e1e10422bfd431aeb5d351c7cf7fa671e3c4df004162264b28220c", size = 61595, upload-time = "2024-12-06T11:20:54.538Z" },
+]
+
[[package]]
name = "markdown-it-py"
version = "4.0.0"
@@ -623,8 +634,8 @@ wheels = [
[[package]]
name = "mlx-lm"
-version = "0.26.3"
-source = { git = "https://github.com/ml-explore/mlx-lm.git#e7f241094c6f95b6b058f270db7fe6d413411a2c" }
+version = "0.26.4"
+source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "jinja2", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "mlx", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
@@ -633,6 +644,10 @@ dependencies = [
{ name = "pyyaml", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "transformers", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
]
+sdist = { url = "https://files.pythonhosted.org/packages/88/20/f3af9d99a5ad6ac42419a3d381290a28bf6d9899ed517a7ccc9fea08546e/mlx_lm-0.26.4.tar.gz", hash = "sha256:1bf21ede1d2d7b660ae312868790df9d73a8553dc50655cf7ae867a36ebcc08c", size = 176384, upload-time = "2025-08-25T15:57:41.723Z" }
+wheels = [
+ { url = "https://files.pythonhosted.org/packages/de/6a/4d20d1b20cd690a3eeaf609c7cb9058f2d52c6d1081394f0d91bd12d08f7/mlx_lm-0.26.4-py3-none-any.whl", hash = "sha256:79bf3afb399ae3bb6073bf0fa6c04f33d70c831ccc6bbbc206c10567d4eef162", size = 242038, upload-time = "2025-08-25T15:57:40.181Z" },
+]
[[package]]
name = "multidict"
@@ -724,7 +739,7 @@ wheels = [
[[package]]
name = "openai"
-version = "1.100.2"
+version = "1.101.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "anyio", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
@@ -736,9 +751,9 @@ dependencies = [
{ name = "tqdm", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "typing-extensions", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
]
-sdist = { url = "https://files.pythonhosted.org/packages/e7/36/e2e24d419438a5e66aa6445ec663194395226293d214bfe615df562b2253/openai-1.100.2.tar.gz", hash = "sha256:787b4c3c8a65895182c58c424f790c25c790cc9a0330e34f73d55b6ee5a00e32", size = 507954, upload-time = "2025-08-19T15:32:47.854Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/00/7c/eaf06b62281f5ca4f774c4cff066e6ddfd6a027e0ac791be16acec3a95e3/openai-1.101.0.tar.gz", hash = "sha256:29f56df2236069686e64aca0e13c24a4ec310545afb25ef7da2ab1a18523f22d", size = 518415, upload-time = "2025-08-21T21:11:01.645Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/db/8d/9ab1599c7942b3d04784ac5473905dc543aeb30a1acce3591d0b425682db/openai-1.100.2-py3-none-any.whl", hash = "sha256:54d3457b2c8d7303a1bc002a058de46bdd8f37a8117751c7cf4ed4438051f151", size = 787755, upload-time = "2025-08-19T15:32:46.252Z" },
+ { url = "https://files.pythonhosted.org/packages/c8/a6/0e39baa335bbd1c66c7e0a41dbbec10c5a15ab95c1344e7f7beb28eee65a/openai-1.101.0-py3-none-any.whl", hash = "sha256:6539a446cce154f8d9fb42757acdfd3ed9357ab0d34fcac11096c461da87133b", size = 810772, upload-time = "2025-08-21T21:10:59.215Z" },
]
[[package]]
@@ -989,25 +1004,25 @@ wheels = [
[[package]]
name = "ruff"
-version = "0.12.9"
+version = "0.12.10"
source = { registry = "https://pypi.org/simple" }
-sdist = { url = "https://files.pythonhosted.org/packages/4a/45/2e403fa7007816b5fbb324cb4f8ed3c7402a927a0a0cb2b6279879a8bfdc/ruff-0.12.9.tar.gz", hash = "sha256:fbd94b2e3c623f659962934e52c2bea6fc6da11f667a427a368adaf3af2c866a", size = 5254702, upload-time = "2025-08-14T16:08:55.2Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/3b/eb/8c073deb376e46ae767f4961390d17545e8535921d2f65101720ed8bd434/ruff-0.12.10.tar.gz", hash = "sha256:189ab65149d11ea69a2d775343adf5f49bb2426fc4780f65ee33b423ad2e47f9", size = 5310076, upload-time = "2025-08-21T18:23:22.595Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/ad/20/53bf098537adb7b6a97d98fcdebf6e916fcd11b2e21d15f8c171507909cc/ruff-0.12.9-py3-none-linux_armv6l.whl", hash = "sha256:fcebc6c79fcae3f220d05585229463621f5dbf24d79fdc4936d9302e177cfa3e", size = 11759705, upload-time = "2025-08-14T16:08:12.968Z" },
- { url = "https://files.pythonhosted.org/packages/20/4d/c764ee423002aac1ec66b9d541285dd29d2c0640a8086c87de59ebbe80d5/ruff-0.12.9-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:aed9d15f8c5755c0e74467731a007fcad41f19bcce41cd75f768bbd687f8535f", size = 12527042, upload-time = "2025-08-14T16:08:16.54Z" },
- { url = "https://files.pythonhosted.org/packages/8b/45/cfcdf6d3eb5fc78a5b419e7e616d6ccba0013dc5b180522920af2897e1be/ruff-0.12.9-py3-none-macosx_11_0_arm64.whl", hash = "sha256:5b15ea354c6ff0d7423814ba6d44be2807644d0c05e9ed60caca87e963e93f70", size = 11724457, upload-time = "2025-08-14T16:08:18.686Z" },
- { url = "https://files.pythonhosted.org/packages/72/e6/44615c754b55662200c48bebb02196dbb14111b6e266ab071b7e7297b4ec/ruff-0.12.9-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:d596c2d0393c2502eaabfef723bd74ca35348a8dac4267d18a94910087807c53", size = 11949446, upload-time = "2025-08-14T16:08:21.059Z" },
- { url = "https://files.pythonhosted.org/packages/fd/d1/9b7d46625d617c7df520d40d5ac6cdcdf20cbccb88fad4b5ecd476a6bb8d/ruff-0.12.9-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:1b15599931a1a7a03c388b9c5df1bfa62be7ede6eb7ef753b272381f39c3d0ff", size = 11566350, upload-time = "2025-08-14T16:08:23.433Z" },
- { url = "https://files.pythonhosted.org/packages/59/20/b73132f66f2856bc29d2d263c6ca457f8476b0bbbe064dac3ac3337a270f/ruff-0.12.9-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:3d02faa2977fb6f3f32ddb7828e212b7dd499c59eb896ae6c03ea5c303575756", size = 13270430, upload-time = "2025-08-14T16:08:25.837Z" },
- { url = "https://files.pythonhosted.org/packages/a2/21/eaf3806f0a3d4c6be0a69d435646fba775b65f3f2097d54898b0fd4bb12e/ruff-0.12.9-py3-none-manylinux_2_17_ppc64.manylinux2014_ppc64.whl", hash = "sha256:17d5b6b0b3a25259b69ebcba87908496e6830e03acfb929ef9fd4c58675fa2ea", size = 14264717, upload-time = "2025-08-14T16:08:27.907Z" },
- { url = "https://files.pythonhosted.org/packages/d2/82/1d0c53bd37dcb582b2c521d352fbf4876b1e28bc0d8894344198f6c9950d/ruff-0.12.9-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:72db7521860e246adbb43f6ef464dd2a532ef2ef1f5dd0d470455b8d9f1773e0", size = 13684331, upload-time = "2025-08-14T16:08:30.352Z" },
- { url = "https://files.pythonhosted.org/packages/3b/2f/1c5cf6d8f656306d42a686f1e207f71d7cebdcbe7b2aa18e4e8a0cb74da3/ruff-0.12.9-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:a03242c1522b4e0885af63320ad754d53983c9599157ee33e77d748363c561ce", size = 12739151, upload-time = "2025-08-14T16:08:32.55Z" },
- { url = "https://files.pythonhosted.org/packages/47/09/25033198bff89b24d734e6479e39b1968e4c992e82262d61cdccaf11afb9/ruff-0.12.9-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:9fc83e4e9751e6c13b5046d7162f205d0a7bac5840183c5beebf824b08a27340", size = 12954992, upload-time = "2025-08-14T16:08:34.816Z" },
- { url = "https://files.pythonhosted.org/packages/52/8e/d0dbf2f9dca66c2d7131feefc386523404014968cd6d22f057763935ab32/ruff-0.12.9-py3-none-manylinux_2_31_riscv64.whl", hash = "sha256:881465ed56ba4dd26a691954650de6ad389a2d1fdb130fe51ff18a25639fe4bb", size = 12899569, upload-time = "2025-08-14T16:08:36.852Z" },
- { url = "https://files.pythonhosted.org/packages/a0/bd/b614d7c08515b1428ed4d3f1d4e3d687deffb2479703b90237682586fa66/ruff-0.12.9-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:43f07a3ccfc62cdb4d3a3348bf0588358a66da756aa113e071b8ca8c3b9826af", size = 11751983, upload-time = "2025-08-14T16:08:39.314Z" },
- { url = "https://files.pythonhosted.org/packages/58/d6/383e9f818a2441b1a0ed898d7875f11273f10882f997388b2b51cb2ae8b5/ruff-0.12.9-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:07adb221c54b6bba24387911e5734357f042e5669fa5718920ee728aba3cbadc", size = 11538635, upload-time = "2025-08-14T16:08:41.297Z" },
- { url = "https://files.pythonhosted.org/packages/20/9c/56f869d314edaa9fc1f491706d1d8a47747b9d714130368fbd69ce9024e9/ruff-0.12.9-py3-none-musllinux_1_2_i686.whl", hash = "sha256:f5cd34fabfdea3933ab85d72359f118035882a01bff15bd1d2b15261d85d5f66", size = 12534346, upload-time = "2025-08-14T16:08:43.39Z" },
- { url = "https://files.pythonhosted.org/packages/bd/4b/d8b95c6795a6c93b439bc913ee7a94fda42bb30a79285d47b80074003ee7/ruff-0.12.9-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:f6be1d2ca0686c54564da8e7ee9e25f93bdd6868263805f8c0b8fc6a449db6d7", size = 13017021, upload-time = "2025-08-14T16:08:45.889Z" },
+ { url = "https://files.pythonhosted.org/packages/24/e7/560d049d15585d6c201f9eeacd2fd130def3741323e5ccf123786e0e3c95/ruff-0.12.10-py3-none-linux_armv6l.whl", hash = "sha256:8b593cb0fb55cc8692dac7b06deb29afda78c721c7ccfed22db941201b7b8f7b", size = 11935161, upload-time = "2025-08-21T18:22:26.965Z" },
+ { url = "https://files.pythonhosted.org/packages/d1/b0/ad2464922a1113c365d12b8f80ed70fcfb39764288ac77c995156080488d/ruff-0.12.10-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:ebb7333a45d56efc7c110a46a69a1b32365d5c5161e7244aaf3aa20ce62399c1", size = 12660884, upload-time = "2025-08-21T18:22:30.925Z" },
+ { url = "https://files.pythonhosted.org/packages/d7/f1/97f509b4108d7bae16c48389f54f005b62ce86712120fd8b2d8e88a7cb49/ruff-0.12.10-py3-none-macosx_11_0_arm64.whl", hash = "sha256:d59e58586829f8e4a9920788f6efba97a13d1fa320b047814e8afede381c6839", size = 11872754, upload-time = "2025-08-21T18:22:34.035Z" },
+ { url = "https://files.pythonhosted.org/packages/12/ad/44f606d243f744a75adc432275217296095101f83f966842063d78eee2d3/ruff-0.12.10-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:822d9677b560f1fdeab69b89d1f444bf5459da4aa04e06e766cf0121771ab844", size = 12092276, upload-time = "2025-08-21T18:22:36.764Z" },
+ { url = "https://files.pythonhosted.org/packages/06/1f/ed6c265e199568010197909b25c896d66e4ef2c5e1c3808caf461f6f3579/ruff-0.12.10-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:37b4a64f4062a50c75019c61c7017ff598cb444984b638511f48539d3a1c98db", size = 11734700, upload-time = "2025-08-21T18:22:39.822Z" },
+ { url = "https://files.pythonhosted.org/packages/63/c5/b21cde720f54a1d1db71538c0bc9b73dee4b563a7dd7d2e404914904d7f5/ruff-0.12.10-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:2c6f4064c69d2542029b2a61d39920c85240c39837599d7f2e32e80d36401d6e", size = 13468783, upload-time = "2025-08-21T18:22:42.559Z" },
+ { url = "https://files.pythonhosted.org/packages/02/9e/39369e6ac7f2a1848f22fb0b00b690492f20811a1ac5c1fd1d2798329263/ruff-0.12.10-py3-none-manylinux_2_17_ppc64.manylinux2014_ppc64.whl", hash = "sha256:059e863ea3a9ade41407ad71c1de2badfbe01539117f38f763ba42a1206f7559", size = 14436642, upload-time = "2025-08-21T18:22:45.612Z" },
+ { url = "https://files.pythonhosted.org/packages/e3/03/5da8cad4b0d5242a936eb203b58318016db44f5c5d351b07e3f5e211bb89/ruff-0.12.10-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:1bef6161e297c68908b7218fa6e0e93e99a286e5ed9653d4be71e687dff101cf", size = 13859107, upload-time = "2025-08-21T18:22:48.886Z" },
+ { url = "https://files.pythonhosted.org/packages/19/19/dd7273b69bf7f93a070c9cec9494a94048325ad18fdcf50114f07e6bf417/ruff-0.12.10-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:4f1345fbf8fb0531cd722285b5f15af49b2932742fc96b633e883da8d841896b", size = 12886521, upload-time = "2025-08-21T18:22:51.567Z" },
+ { url = "https://files.pythonhosted.org/packages/c0/1d/b4207ec35e7babaee62c462769e77457e26eb853fbdc877af29417033333/ruff-0.12.10-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:1f68433c4fbc63efbfa3ba5db31727db229fa4e61000f452c540474b03de52a9", size = 13097528, upload-time = "2025-08-21T18:22:54.609Z" },
+ { url = "https://files.pythonhosted.org/packages/ff/00/58f7b873b21114456e880b75176af3490d7a2836033779ca42f50de3b47a/ruff-0.12.10-py3-none-manylinux_2_31_riscv64.whl", hash = "sha256:141ce3d88803c625257b8a6debf4a0473eb6eed9643a6189b68838b43e78165a", size = 13080443, upload-time = "2025-08-21T18:22:57.413Z" },
+ { url = "https://files.pythonhosted.org/packages/12/8c/9e6660007fb10189ccb78a02b41691288038e51e4788bf49b0a60f740604/ruff-0.12.10-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:f3fc21178cd44c98142ae7590f42ddcb587b8e09a3b849cbc84edb62ee95de60", size = 11896759, upload-time = "2025-08-21T18:23:00.473Z" },
+ { url = "https://files.pythonhosted.org/packages/67/4c/6d092bb99ea9ea6ebda817a0e7ad886f42a58b4501a7e27cd97371d0ba54/ruff-0.12.10-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:7d1a4e0bdfafcd2e3e235ecf50bf0176f74dd37902f241588ae1f6c827a36c56", size = 11701463, upload-time = "2025-08-21T18:23:03.211Z" },
+ { url = "https://files.pythonhosted.org/packages/59/80/d982c55e91df981f3ab62559371380616c57ffd0172d96850280c2b04fa8/ruff-0.12.10-py3-none-musllinux_1_2_i686.whl", hash = "sha256:e67d96827854f50b9e3e8327b031647e7bcc090dbe7bb11101a81a3a2cbf1cc9", size = 12691603, upload-time = "2025-08-21T18:23:06.935Z" },
+ { url = "https://files.pythonhosted.org/packages/ad/37/63a9c788bbe0b0850611669ec6b8589838faf2f4f959647f2d3e320383ae/ruff-0.12.10-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:ae479e1a18b439c59138f066ae79cc0f3ee250712a873d00dbafadaad9481e5b", size = 13164356, upload-time = "2025-08-21T18:23:10.225Z" },
]
[[package]]
@@ -1097,14 +1112,14 @@ wheels = [
[[package]]
name = "starlette"
-version = "0.47.2"
+version = "0.47.3"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "anyio", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
]
-sdist = { url = "https://files.pythonhosted.org/packages/04/57/d062573f391d062710d4088fa1369428c38d51460ab6fedff920efef932e/starlette-0.47.2.tar.gz", hash = "sha256:6ae9aa5db235e4846decc1e7b79c4f346adf41e9777aebeb49dfd09bbd7023d8", size = 2583948, upload-time = "2025-07-20T17:31:58.522Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/15/b9/cc3017f9a9c9b6e27c5106cc10cc7904653c3eec0729793aec10479dd669/starlette-0.47.3.tar.gz", hash = "sha256:6bc94f839cc176c4858894f1f8908f0ab79dfec1a6b8402f6da9be26ebea52e9", size = 2584144, upload-time = "2025-08-24T13:36:42.122Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/f7/1f/b876b1f83aef204198a42dc101613fefccb32258e5428b5f9259677864b4/starlette-0.47.2-py3-none-any.whl", hash = "sha256:c5847e96134e5c5371ee9fac6fdf1a67336d5815e09eb2a01fdb57a351ef915b", size = 72984, upload-time = "2025-07-20T17:31:56.738Z" },
+ { url = "https://files.pythonhosted.org/packages/ce/fd/901cfa59aaa5b30a99e16876f11abe38b59a1a2c51ffb3d7142bb6089069/starlette-0.47.3-py3-none-any.whl", hash = "sha256:89c0778ca62a76b826101e7c709e70680a1699ca7da6b44d38eb0a7e61fe4b51", size = 72991, upload-time = "2025-08-24T13:36:40.887Z" },
]
[[package]]
@@ -1141,7 +1156,7 @@ wheels = [
[[package]]
name = "transformers"
-version = "4.55.2"
+version = "4.55.4"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "filelock", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
@@ -1155,9 +1170,9 @@ dependencies = [
{ name = "tokenizers", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
{ name = "tqdm", marker = "sys_platform == 'darwin' or sys_platform == 'linux'" },
]
-sdist = { url = "https://files.pythonhosted.org/packages/70/a5/d8b8a1f3a051daeb5f11253bb69fc241f193d1c0566e299210ed9220ff4e/transformers-4.55.2.tar.gz", hash = "sha256:a45ec60c03474fd67adbce5c434685051b7608b3f4f167c25aa6aeb1cad16d4f", size = 9571466, upload-time = "2025-08-13T18:25:43.767Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/2b/43/3cb831d5f28cc723516e5bb43a8c6042aca3038bb36b6bd6016b40dfd1e8/transformers-4.55.4.tar.gz", hash = "sha256:574a30559bc273c7a4585599ff28ab6b676e96dc56ffd2025ecfce2fd0ab915d", size = 9573015, upload-time = "2025-08-22T15:18:43.192Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/db/5a/022ac010bedfb5119734cf9d743cf1d830cb4c604f53bb1552216f4344dc/transformers-4.55.2-py3-none-any.whl", hash = "sha256:097e3c2e2c0c9681db3da9d748d8f9d6a724c644514673d0030e8c5a1109f1f1", size = 11269748, upload-time = "2025-08-13T18:25:40.394Z" },
+ { url = "https://files.pythonhosted.org/packages/fa/0a/8791a6ee0529c45f669566969e99b75e2ab20eb0bfee8794ce295c18bdad/transformers-4.55.4-py3-none-any.whl", hash = "sha256:df28f3849665faba4af5106f0db4510323277c4bb595055340544f7e59d06458", size = 11269659, upload-time = "2025-08-22T15:18:40.025Z" },
]
[[package]]
@@ -1174,20 +1189,20 @@ wheels = [
[[package]]
name = "types-aiofiles"
-version = "24.1.0.20250809"
+version = "24.1.0.20250822"
source = { registry = "https://pypi.org/simple" }
-sdist = { url = "https://files.pythonhosted.org/packages/03/b8/34a4f9da445a104d240bb26365a10ef68953bebdc812859ea46847c7fdcb/types_aiofiles-24.1.0.20250809.tar.gz", hash = "sha256:4dc9734330b1324d9251f92edfc94fd6827fbb829c593313f034a77ac33ae327", size = 14379, upload-time = "2025-08-09T03:14:41.555Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/19/48/c64471adac9206cc844afb33ed311ac5a65d2f59df3d861e0f2d0cad7414/types_aiofiles-24.1.0.20250822.tar.gz", hash = "sha256:9ab90d8e0c307fe97a7cf09338301e3f01a163e39f3b529ace82466355c84a7b", size = 14484, upload-time = "2025-08-22T03:02:23.039Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/28/78/0d8ffa40e9ec6cbbabe4d93675092fea1cadc4c280495375fc1f2fa42793/types_aiofiles-24.1.0.20250809-py3-none-any.whl", hash = "sha256:657c83f876047ffc242b34bfcd9167f201d1b02e914ee854f16e589aa95c0d45", size = 14300, upload-time = "2025-08-09T03:14:40.438Z" },
+ { url = "https://files.pythonhosted.org/packages/bc/8e/5e6d2215e1d8f7c2a94c6e9d0059ae8109ce0f5681956d11bb0a228cef04/types_aiofiles-24.1.0.20250822-py3-none-any.whl", hash = "sha256:0ec8f8909e1a85a5a79aed0573af7901f53120dd2a29771dd0b3ef48e12328b0", size = 14322, upload-time = "2025-08-22T03:02:21.918Z" },
]
[[package]]
name = "typing-extensions"
-version = "4.14.1"
+version = "4.15.0"
source = { registry = "https://pypi.org/simple" }
-sdist = { url = "https://files.pythonhosted.org/packages/98/5a/da40306b885cc8c09109dc2e1abd358d5684b1425678151cdaed4731c822/typing_extensions-4.14.1.tar.gz", hash = "sha256:38b39f4aeeab64884ce9f74c94263ef78f3c22467c8724005483154c26648d36", size = 107673, upload-time = "2025-07-04T13:28:34.16Z" }
+sdist = { url = "https://files.pythonhosted.org/packages/72/94/1a15dd82efb362ac84269196e94cf00f187f7ed21c242792a923cdb1c61f/typing_extensions-4.15.0.tar.gz", hash = "sha256:0cea48d173cc12fa28ecabc3b837ea3cf6f38c6d1136f85cbaaf598984861466", size = 109391, upload-time = "2025-08-25T13:49:26.313Z" }
wheels = [
- { url = "https://files.pythonhosted.org/packages/b5/00/d631e67a838026495268c2f6884f3711a15a9a2a96cd244fdaea53b823fb/typing_extensions-4.14.1-py3-none-any.whl", hash = "sha256:d1e1e3b58374dc93031d6eda2420a48ea44a36c2b4766a4fdeb3710755731d76", size = 43906, upload-time = "2025-07-04T13:28:32.743Z" },
+ { url = "https://files.pythonhosted.org/packages/18/67/36e9267722cc04a6b9f15c7f3441c2363321a3ea07da7ae0c0707beb2a9c/typing_extensions-4.15.0-py3-none-any.whl", hash = "sha256:f0fa19c6845758ab08074a0cfa8b7aecb71c999ca73d62883bc25cc018c4e548", size = 44614, upload-time = "2025-08-25T13:49:24.86Z" },
]
[[package]]
← ef5c5b96 changes include: ipc, general utilities, flakes stuff w/ jus
·
back to Exo
·
feat: mlx memory cache for faster ttft 84c90a6d →