nemo_curator.utils.resumability_client

View as Markdown

Worker-side helpers to talk to the resumability actor; all no-ops when no actor is registered, so unchecked pipelines pay nothing.

Module Contents

Functions

NameDescription
_resumability_actorThe resumability actor handle, or None if Ray is down / no actor registered.
completed_resumability_sourcesSubset of source_ids already marked complete; the source stage uses it to skip them.
flush_resumability_deltasFire-and-forget per-task deltas (task_id, source_id, delta). No
is_resumability_actor_activeTrue if a resumability actor is registered in this Ray cluster.

Data

ACTOR_NAME

API

nemo_curator.utils.resumability_client._resumability_actor() -> ray.actor.ActorHandle | None

The resumability actor handle, or None if Ray is down / no actor registered.

nemo_curator.utils.resumability_client.completed_resumability_sources(
source_ids: list[str]
) -> set[str]

Subset of source_ids already marked complete; the source stage uses it to skip them.

nemo_curator.utils.resumability_client.flush_resumability_deltas(
per_task: list[tuple[str, str, int]]
) -> None

Fire-and-forget per-task deltas (task_id, source_id, delta). No ray.get — the actor never raises, so there’s no synchronous error path.

nemo_curator.utils.resumability_client.is_resumability_actor_active() -> bool

True if a resumability actor is registered in this Ray cluster.

nemo_curator.utils.resumability_client.ACTOR_NAME = 'nemo_curator_resumability'