distributed_eval
Optional, fault-tolerant distributed evaluation for Puzzletron.
This package is intentionally isolated from the existing local evaluator and
rpc_eval scaffold. Import it explicitly when distributed evaluation is
desired; importing modelopt.torch.puzzletron does not start services or
change scoring behavior.
Classes
Blocking facade backed by one persistent asyncio loop thread. |
|
Functions
Stable identity reserved for a future topology-specific prefix cache. |
- class AsyncEvaluationClient
Bases:
object- __init__(coordinator)
- Parameters:
coordinator (AsyncEvaluationCoordinator)
- async as_completed(handles)
- Parameters:
handles (Iterable[EvaluationHandle])
- Return type:
AsyncIterator[EvaluationResult]
- async cancel(handle)
- Parameters:
handle (EvaluationHandle | str)
- Return type:
bool
- async close()
- Return type:
None
- classmethod from_campaign(campaign_dir, **kwargs)
- Parameters:
campaign_dir (str)
- Return type:
- lookup(handle)
- Parameters:
handle (EvaluationHandle | str)
- Return type:
EvaluationResult | None
- status()
- Return type:
dict[str, Any]
- async submit(request)
- Parameters:
request (EvaluationRequest)
- Return type:
- async submit_many(requests)
- Parameters:
requests (Iterable[EvaluationRequest])
- Return type:
list[EvaluationHandle]
- async wait(handle, *, timeout=None)
- Parameters:
handle (EvaluationHandle | str)
timeout (float | None)
- Return type:
- ModeloptConfig AttemptRecord
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "schema_version": 1, "request_id": null, "attempt_id": null, "worker_id": null, "worker_boot_id": null, "status": null, "leased_at": null, "finished_at": null, "error": null, "metadata": null }
- field attempt_id: str [Required]
- field error: EvaluationError | None
- field finished_at: datetime | None
- field leased_at: datetime [Optional]
- field metadata: dict[str, Any] [Optional]
- field request_id: str [Required]
- field schema_version: Literal[1]
- field status: AttemptStatus [Required]
- field worker_boot_id: str | None
- field worker_id: str | None
- class AttemptStatus
Bases:
str,Enum- CANCELLED = 'cancelled'
- CONFLICT = 'conflict'
- DUPLICATE = 'duplicate'
- FAILED = 'failed'
- LEASED = 'leased'
- RETRY = 'retry'
- SUCCEEDED = 'succeeded'
- __new__(value)
- class Campaign
Bases:
object- __init__(root, manifest)
- Parameters:
root (str | Path)
manifest (CampaignManifest)
- property campaign_id: str
- coordinator_lease()
- Return type:
Iterator[None]
- classmethod create(root, manifest)
- Parameters:
root (str | Path)
manifest (CampaignManifest)
- Return type:
- property root: Path
- validate_request(request)
- Parameters:
request (EvaluationRequest)
- Return type:
None
- ModeloptConfig CampaignManifest
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "schema_version": 1, "name": null, "model": null, "descriptor": null, "force_hf": null, "parallelism": null, "precision": null, "automodel_recipe": null, "data": null, "metrics": null, "evaluator_revision": null, "result_atol": 0.0001, "result_rtol": 0.0001, "metadata": null }
- field automodel_recipe: dict[str, Any] [Required]
- field data: dict[str, Any] [Required]
- field descriptor: str [Required]
- field evaluator_revision: str [Required]
- field force_hf: bool [Required]
- field metadata: dict[str, Any] [Optional]
- field metrics: dict[str, Any] [Required]
- field model: dict[str, Any] [Required]
- field name: str [Required]
- field parallelism: ParallelismSpec [Required]
- field precision: dict[str, Any] [Optional]
- field result_atol: float
- Constraints:
ge = 0.0
- field result_rtol: float
- Constraints:
ge = 0.0
- field schema_version: Literal[1]
- property campaign_id: str
- class ErrorKind
Bases:
str,Enum- CANCELLED = 'cancelled'
- CANDIDATE = 'candidate'
- INTERNAL = 'internal'
- INVALID_REQUEST = 'invalid_request'
- PROCESS_GROUP = 'process_group'
- RESOURCE_EXHAUSTED = 'resource_exhausted'
- TIMEOUT = 'timeout'
- TRANSPORT = 'transport'
- UNSUPPORTED = 'unsupported'
- WORKER_LOST = 'worker_lost'
- __new__(value)
- class EvaluationClient
Bases:
objectBlocking facade backed by one persistent asyncio loop thread.
- __init__(campaign_dir, **kwargs)
- Parameters:
campaign_dir (str)
- batch(requests, *, timeout=None)
- Parameters:
requests (Iterable[EvaluationRequest])
timeout (float | None)
- Return type:
list[EvaluationResult]
- cancel(handle)
- Parameters:
handle (EvaluationHandle | str)
- Return type:
bool
- close()
- Return type:
None
- evaluate(request, *, timeout=None)
- Parameters:
request (EvaluationRequest)
timeout (float | None)
- Return type:
- classmethod from_campaign(campaign_dir, **kwargs)
- Parameters:
campaign_dir (str)
- Return type:
- lookup(handle)
- Parameters:
handle (EvaluationHandle | str)
- Return type:
EvaluationResult | None
- status()
- Return type:
dict[str, Any]
- submit(request)
- Parameters:
request (EvaluationRequest)
- Return type:
- ModeloptConfig EvaluationError
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "kind": null, "message": null, "traceback": null, "retryable": null }
- field message: str [Required]
- field retryable: bool | None
- field traceback: str | None
- ModeloptConfig EvaluationHandle
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "campaign_id": null, "request_id": null }
- field campaign_id: str [Required]
- field request_id: str [Required]
- ModeloptConfig EvaluationRequest
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "schema_version": 1, "campaign_id": null, "handler": null, "payload": null, "model": null, "data": null, "metrics": null, "precision": null, "evaluator_revision": null, "metadata": null }
- field campaign_id: str [Required]
- field data: dict[str, Any] [Required]
- field evaluator_revision: str [Required]
- field handler: str [Required]
- field metadata: dict[str, Any] [Optional]
- field metrics: dict[str, Any] [Required]
- field model: dict[str, Any] [Required]
- field payload: dict[str, Any] [Required]
- field precision: dict[str, Any] [Optional]
- field schema_version: Literal[1]
- classmethod from_wire(payload)
- Parameters:
payload (dict[str, Any])
- Return type:
- to_wire()
- Return type:
dict[str, Any]
- property request_id: str
- ModeloptConfig EvaluationResult
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "schema_version": 1, "request_id": null, "campaign_id": null, "metrics": null, "counts": null, "reduction_state": null, "artifacts": null, "timing": null, "provenance": null, "completed_at": null }
- field artifacts: dict[str, Any] [Optional]
- field campaign_id: str [Required]
- field completed_at: datetime [Optional]
- field counts: dict[str, int | float] [Optional]
- field metrics: dict[str, Any] [Required]
- field provenance: dict[str, Any] [Optional]
- field reduction_state: dict[str, Any] [Optional]
- field request_id: str [Required]
- field schema_version: Literal[1]
- field timing: dict[str, float] [Optional]
- ModeloptConfig ParallelismSpec
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "tp_size": 1, "ep_size": 1, "cp_size": 1, "pp_size": 1, "dp_size": null, "sequence_parallel": false, "fsdp": true, "distributed_backend": "nccl", "world_size": 1, "gpus_per_task": null }
- field cp_size: int
- Constraints:
ge = 1
- field distributed_backend: str
- field dp_size: int | None
- Constraints:
ge = 1
- field ep_size: int
- Constraints:
ge = 1
- field fsdp: bool
- field gpus_per_task: int | None
- Constraints:
ge = 1
- field pp_size: int
- Constraints:
ge = 1
- field sequence_parallel: bool
- field tp_size: int
- Constraints:
ge = 1
- field world_size: int
- Constraints:
ge = 1
- ModeloptConfig WorkerRecord
Bases:
StrictModelShow default config as JSON
- Default config (JSON):
{ "schema_version": 1, "worker_id": null, "boot_id": null, "campaign_id": null, "host": null, "port": null, "parallelism": null, "capabilities": null, "state": "starting", "current_request_id": null, "started_at": null, "heartbeat_at": null }
- field boot_id: str [Required]
- field campaign_id: str [Required]
- field capabilities: dict[str, Any] [Optional]
- field current_request_id: str | None
- field heartbeat_at: datetime [Optional]
- field host: str [Required]
- field parallelism: ParallelismSpec [Required]
- field port: int [Required]
- Constraints:
ge = 1
le = 65535
- field schema_version: Literal[1]
- field started_at: datetime [Optional]
- field state: WorkerState
- field worker_id: str [Required]
- property endpoint: str
- class WorkerState
Bases:
str,Enum- BUSY = 'busy'
- DRAINING = 'draining'
- FAILED = 'failed'
- IDLE = 'idle'
- STARTING = 'starting'
- __new__(value)
- prefix_cache_id(*, campaign_id, model, data, prefix_signature, prefix_length, parallelism)
Stable identity reserved for a future topology-specific prefix cache.
- Parameters:
campaign_id (str)
model (Any)
data (Any)
prefix_signature (Any)
prefix_length (int)
parallelism (Any)
- Return type:
str