nemo_curator.tasks.sentinels

View as Markdown

Payload-less marker tasks.

EmptyTask seeds a pipeline (the implicit root id "0"). The resumability layer adds two more markers on the same :class:SentinelTask base:

  • NoneTask — this slot was intentionally filtered. The resumability counter treats it as a consumed branch (decrements). The adapter auto-wraps a returned None as a NoneTask.
  • FailedTask — this slot failed and should be retried on resume. The counter is NOT decremented, so its source stays pending and reruns.

All carry no payload (data is None) and get their task_id assigned by the executor adapter; sentinels are stripped before the next stage. Construct with EmptyTask() / NoneTask() / FailedTask().

Module Contents

Classes

NameDescription
EmptyTaskSeeds a pipeline with task_id="0" — the implicit root every task
FailedTaskMarks a slot as failed → retried on resume (counter does NOT decrement).
NoneTaskMarks a slot as intentionally filtered (resumability counter decrements).
SentinelTaskBase for payload-less marker tasks: no data, framework-assigned task_id.

API

class nemo_curator.tasks.sentinels.EmptyTask(
dataset_name: str = 'empty',
data: None = None,
_stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: SentinelTask

Seeds a pipeline with task_id="0" — the implicit root every task descends from (so all ids share the "0" prefix).

dataset_name
str = 'empty'
task_id
str = field(init=False, default='0')
class nemo_curator.tasks.sentinels.FailedTask(
dataset_name: str = 'failed',
data: None = None,
_stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: SentinelTask

Marks a slot as failed → retried on resume (counter does NOT decrement).

dataset_name
str = 'failed'
class nemo_curator.tasks.sentinels.NoneTask(
dataset_name: str = 'none',
data: None = None,
_stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: SentinelTask

Marks a slot as intentionally filtered (resumability counter decrements).

dataset_name
str = 'none'
class nemo_curator.tasks.sentinels.SentinelTask(
dataset_name: str,
data: None = None,
_stage_perf: list[nemo_curator.utils.performance_utils.StagePerfStats] = list(),
_metadata: dict[str, typing.Any] = dict()
)
Dataclass

Bases: Task[None]

Base for payload-less marker tasks: no data, framework-assigned task_id.

data
None = None
num_items
int
nemo_curator.tasks.sentinels.SentinelTask.__post_init__() -> None
nemo_curator.tasks.sentinels.SentinelTask.validate() -> bool