Fix rebase
This commit is contained in:
@@ -1,6 +1,5 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime
|
||||
@@ -102,8 +101,8 @@ class DownloadCoordinator:
|
||||
shard: ShardMetadata,
|
||||
found: Path,
|
||||
total: Memory,
|
||||
) -> DownloadCompleted:
|
||||
return DownloadCompleted(
|
||||
) -> ModelReady:
|
||||
return ModelReady(
|
||||
shard_metadata=shard,
|
||||
node_id=self.node_id,
|
||||
total=total,
|
||||
@@ -124,7 +123,7 @@ class DownloadCoordinator:
|
||||
callback_shard, found, progress.total
|
||||
)
|
||||
else:
|
||||
completed = DownloadCompleted(
|
||||
completed = ModelReady(
|
||||
shard_metadata=callback_shard,
|
||||
node_id=self.node_id,
|
||||
total=progress.total,
|
||||
@@ -291,7 +290,7 @@ class DownloadCoordinator:
|
||||
shard, found, initial_progress.total
|
||||
)
|
||||
else:
|
||||
completed = DownloadCompleted(
|
||||
completed = ModelReady(
|
||||
shard_metadata=shard,
|
||||
node_id=self.node_id,
|
||||
total=initial_progress.total,
|
||||
@@ -353,9 +352,6 @@ class DownloadCoordinator:
|
||||
await self.event_sender.send(
|
||||
NodeDownloadProgress(download_progress=failed)
|
||||
)
|
||||
except anyio.get_cancelled_exc_class():
|
||||
# ignore cancellation - let cleanup do its thing
|
||||
pass
|
||||
finally:
|
||||
self.active_downloads.pop(model_id, None)
|
||||
|
||||
@@ -438,11 +434,11 @@ class DownloadCoordinator:
|
||||
resolve_existing_model, model_id
|
||||
)
|
||||
if found is not None:
|
||||
status: DownloadProgress = self._completed_from_path(
|
||||
status: ModelStatus = self._completed_from_path(
|
||||
progress.shard, found, progress.total
|
||||
)
|
||||
else:
|
||||
status = DownloadCompleted(
|
||||
status = ModelReady(
|
||||
node_id=self.node_id,
|
||||
shard_metadata=progress.shard,
|
||||
total=progress.total,
|
||||
@@ -494,7 +490,7 @@ class DownloadCoordinator:
|
||||
end_layer=card.n_layers,
|
||||
n_layers=card.n_layers,
|
||||
)
|
||||
path_completed: DownloadProgress = (
|
||||
path_completed: ModelStatus = (
|
||||
self._completed_from_path(
|
||||
path_shard, found, card.storage_size
|
||||
)
|
||||
|
||||
@@ -107,7 +107,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_evicts_oldest_model_to_fit_new_download(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -135,7 +135,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_evicts_multiple_in_lru_order(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -168,7 +168,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_rejects_when_cannot_free_enough_space(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -191,7 +191,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_no_eviction_when_space_available(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -216,7 +216,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_manual_policy_rejects_instead_of_evicting(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -237,7 +237,7 @@ class TestStartDownloadAutoEviction:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_eviction_emits_not_downloading_event_for_evicted_model(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -280,7 +280,7 @@ class TestActiveModelProtection:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_active_model_not_evicted(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -311,7 +311,7 @@ class TestActiveModelProtection:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_all_active_models_rejected(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -340,7 +340,7 @@ class TestDiskDeleteFailure:
|
||||
new_callable=AsyncMock,
|
||||
return_value=False,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_eviction_rejected_on_disk_delete_failure(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -418,7 +418,7 @@ class TestEvictionEvents:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_eviction_emits_not_downloading_event(
|
||||
self, _mock_resolve: AsyncMock, mock_delete: AsyncMock
|
||||
) -> None:
|
||||
@@ -453,7 +453,7 @@ class TestEvictionEvents:
|
||||
new_callable=AsyncMock,
|
||||
return_value=True,
|
||||
)
|
||||
@patch("exo.download.coordinator.resolve_model_in_path", return_value=None)
|
||||
@patch("exo.download.coordinator.resolve_existing_model", return_value=None)
|
||||
async def test_multi_eviction_emits_event_per_model(
|
||||
self, _mock_resolve: AsyncMock, _mock_delete: AsyncMock
|
||||
) -> None:
|
||||
|
||||
@@ -184,7 +184,7 @@ async def load_storage_config(
|
||||
if raw.strip():
|
||||
data = tomllib.loads(raw)
|
||||
base = StorageConfig.from_disk(data)
|
||||
except Exception:
|
||||
except (OSError, tomllib.TOMLDecodeError, ValueError, KeyError):
|
||||
logger.warning("Failed to read storage config from config file, using defaults")
|
||||
|
||||
resolved_max_storage = (
|
||||
|
||||
@@ -17,7 +17,10 @@ class StorageConfig(FrozenModel):
|
||||
"""Parse from a TOML config dict (e.g. from tomllib)."""
|
||||
max_storage: Memory | None = None
|
||||
if "max_storage_gb" in data:
|
||||
max_storage = Memory.from_gb(float(data["max_storage_gb"])) # pyright: ignore[reportAny]
|
||||
gb = float(data["max_storage_gb"]) # pyright: ignore[reportAny]
|
||||
if gb < 0:
|
||||
raise ValueError(f"max_storage_gb must be non-negative, got {gb}")
|
||||
max_storage = Memory.from_gb(gb)
|
||||
policy: StoragePolicy = data.get("storage_policy", "manual") # pyright: ignore[reportAny]
|
||||
return cls(max_storage=max_storage, storage_policy=policy)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user