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

AsyncEvaluationClient

AttemptStatus

Campaign

ErrorKind

EvaluationClient

Blocking facade backed by one persistent asyncio loop thread.

WorkerState

Functions

prefix_cache_id

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:

AsyncEvaluationClient

lookup(handle)
Parameters:

handle (EvaluationHandle | str)

Return type:

EvaluationResult | None

status()
Return type:

dict[str, Any]

async submit(request)
Parameters:

request (EvaluationRequest)

Return type:

EvaluationHandle

async submit_many(requests)
Parameters:

requests (Iterable[EvaluationRequest])

Return type:

list[EvaluationHandle]

async wait(handle, *, timeout=None)
Parameters:
Return type:

EvaluationResult

ModeloptConfig AttemptRecord

Bases: StrictModel

Show 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:
property campaign_id: str
coordinator_lease()
Return type:

Iterator[None]

classmethod create(root, manifest)
Parameters:
Return type:

Campaign

classmethod open(root)
Parameters:

root (str | Path)

Return type:

Campaign

property root: Path
validate_request(request)
Parameters:

request (EvaluationRequest)

Return type:

None

ModeloptConfig CampaignManifest

Bases: StrictModel

Show 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: object

Blocking facade backed by one persistent asyncio loop thread.

__init__(campaign_dir, **kwargs)
Parameters:

campaign_dir (str)

batch(requests, *, timeout=None)
Parameters:
Return type:

list[EvaluationResult]

cancel(handle)
Parameters:

handle (EvaluationHandle | str)

Return type:

bool

close()
Return type:

None

evaluate(request, *, timeout=None)
Parameters:
Return type:

EvaluationResult

classmethod from_campaign(campaign_dir, **kwargs)
Parameters:

campaign_dir (str)

Return type:

EvaluationClient

lookup(handle)
Parameters:

handle (EvaluationHandle | str)

Return type:

EvaluationResult | None

status()
Return type:

dict[str, Any]

submit(request)
Parameters:

request (EvaluationRequest)

Return type:

EvaluationHandle

ModeloptConfig EvaluationError

Bases: StrictModel

Show default config as JSON
Default config (JSON):

{
   "kind": null,
   "message": null,
   "traceback": null,
   "retryable": null
}

field kind: ErrorKind [Required]
field message: str [Required]
field retryable: bool | None
field traceback: str | None
ModeloptConfig EvaluationHandle

Bases: StrictModel

Show 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: StrictModel

Show 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:

EvaluationRequest

to_wire()
Return type:

dict[str, Any]

property request_id: str
ModeloptConfig EvaluationResult

Bases: StrictModel

Show 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: StrictModel

Show 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: StrictModel

Show 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