How to Use Conditional Task Execution with the Condition Parameter in Sieves

You can skip documents on a per-task basis by passing a callable to the condition parameter of any Task in Sieves, which evaluates each document and stores None for skipped items while preserving the original document order.

Sieves is a modular document processing framework that supports conditional task execution through a flexible condition parameter available on all tasks. This feature allows you to filter documents dynamically based on custom logic, ensuring that expensive operations like model inference only run on relevant documents. By leveraging the condition parameter in the mantisai/sieves repository, you can build efficient pipelines that skip unnecessary processing without breaking the document stream.

How Conditional Task Execution Works in Sieves

The Task Base Class and Condition Storage

The conditional logic is implemented in the base Task class defined in sieves/tasks/core.py. When initializing any task, the condition parameter is stored as self._condition:


# sieves/tasks/core.py

class Task(abc.ABC):
    def __init__(self, task_id, include_meta, batch_size, condition=None):
        # ...

        self._condition = condition                         # ← stored here

This design ensures that all subclasses—including built-in tasks like Classification and Chunking—inherit conditional execution capabilities without additional implementation.

The Execution Flow and Document Filtering

When a pipeline invokes a task, the __call__ method (lines 68-88 in sieves/tasks/core.py) handles conditional filtering automatically. The method evaluates the condition for each document in the batch to determine which items to process:


# sieves/tasks/core.py (lines 68-88)

while docs_batch := [doc for doc in itertools.islice(docs, batch_size)]:
    passing_indices = {
        idx for idx, doc in enumerate(docs_batch)
        if self._condition is None or self._condition(doc)
    }
    processed = self._call(d for i, d in enumerate(docs_batch) if i in passing_indices)
    for idx, doc in enumerate(docs_batch):
        if idx in passing_indices:
            yield next(processed)          # run the task

        else:
            doc.results[self.id] = None    # skipped → store None

            yield doc

Documents that fail the condition check receive None in their results dictionary under the task ID, while passing documents are processed normally. The original document order is preserved throughout the pipeline.

The pipeline itself simply forwards documents to each task. Because each task handles its own conditional logic in sieves/pipeline/core.py, no extra pipeline configuration is required:


# sieves/pipeline/core.py

for i, task in enumerate(self._tasks):
    processed_docs = task(processed_docs)   # ← each task decides what to do

When to Use Conditional Task Execution

Conditional task execution is particularly valuable when processing heterogeneous document collections. Common scenarios include:

Situation Example Condition
Skip short texts lambda d: len(d.text or "") > 20
Process only PDFs lambda d: d.uri and d.uri.endswith('.pdf')
Rate limiting lambda d, i=itertools.count(): next(i) < N
External data validation Any callable that queries external APIs or databases to determine eligibility

Implementation Examples

Custom Task with a Condition

The following example demonstrates a custom task that only processes documents longer than 20 characters. This pattern is tested in sieves/tests/tasks/test_conditional_execution.py.

from sieves.tasks.core import Task
from sieves.data import Doc

class DummyTask(Task):
    def __init__(self, condition=None):
        super().__init__(task_id="DummyTask", include_meta=False,
                         batch_size=-1, condition=condition)

    def _call(self, docs):
        for doc in docs:
            doc.results[self.id] = {"processed": True}
            yield doc

To use this task with conditional execution:

from sieves.pipeline import Pipeline
from sieves import Doc

docs = [Doc(text="short"), Doc(text="this is a much longer document")]
task = DummyTask(condition=lambda d: len(d.text or "") > 20)

pipe = Pipeline([task])
for d in pipe(docs):
    print(d.results["DummyTask"])

# Output: None for the short doc, {"processed": True} for the long one

Built-in Classification with Conditional Logic

Built-in predictive tasks like Classification inherit the condition parameter from the base Task class. This example from sieves/tests/docs/test_pipeline_docs.py shows how to conditionally classify only long documents.

from sieves import Pipeline, tasks, Doc
from sieves.models import SmallTransformer  # any ModelWrapper implementation

model = SmallTransformer()
docs = [
    Doc(text="short"),
    Doc(text="this is a much longer document that will be processed"),
    Doc(text="med")
]

def is_long(doc: Doc) -> bool:
    return len(doc.text or "") > 20

# Pass condition directly to the task constructor

task = tasks.Classification(labels=["science", "politics"],
                            model=model,
                            condition=is_long)

pipe = Pipeline([task])
for d in pipe(docs):
    print(d.results[task.id])   # None for short/med docs, prediction dict for the long one

Multiple Tasks with Independent Conditions

You can assign different conditions to different tasks within the same pipeline. Each task evaluates its condition independently, allowing complex processing workflows.

from sieves import Pipeline, tasks, Doc
from sieves.models import SmallTransformer
from sieves.chunkers import example_chunker

model = SmallTransformer()
docs = [
    Doc(text="short"),
    Doc(text="this is a much longer document"),
    Doc(text="medium text here")
]

# Task 1: Only chunk documents longer than 10 characters

task1 = tasks.Chunking(example_chunker,
                       condition=lambda d: len(d.text or "") > 10)

# Task 2: Only classify documents longer than 20 characters

task2 = tasks.Classification(labels=["science", "politics"],
                             model=model,
                             condition=lambda d: len(d.text or "") > 20)

pipe = Pipeline([task1, task2])
for d in pipe(docs):
    print(d.chunks, d.results.get(task2.id))

Key Implementation Files

The conditional execution system is implemented across these core files in the mantisai/sieves repository:

File Purpose Link
sieves/tasks/core.py Defines the base Task class with condition parameter and the __call__ method that filters documents per batch. View source
sieves/pipeline/core.py Orchestrates task execution; forwards documents to each task which handles its own conditional logic independently. View source
sieves/tests/tasks/test_conditional_execution.py Unit tests demonstrating custom task behavior with conditional execution. View source
sieves/tests/docs/test_pipeline_docs.py Integration examples showing built-in tasks with conditions in pipeline contexts. View source
sieves/tasks/predictive/classification/core.py Example of a predictive task inheriting conditional capabilities from the base Task class. View source

Summary

  • The condition parameter in Sieves accepts any callable that takes a Doc object and returns a boolean, enabling conditional task execution on a per-document basis.
  • When a document fails the condition check, the task stores None in doc.results[task.id] and yields the document unchanged, preserving the original stream order.
  • All tasks—including custom subclasses and built-in predictive tasks like Classification—inherit this functionality from the base Task class in sieves/tasks/core.py.
  • Conditions are evaluated independently for each task in a pipeline, allowing different filtering logic for chunking, classification, or custom processing steps.
  • Skipped documents remain in the pipeline for subsequent tasks, making the system cache-friendly and composable.

Frequently Asked Questions

What happens to documents that don't meet the condition?

Documents that fail the condition check are not processed by the task's _call method. Instead, the task automatically sets doc.results[self.id] = None and yields the document unchanged. This ensures the document remains in the pipeline for subsequent tasks while maintaining the original order of the document stream.

Can I use different conditions for different tasks in the same pipeline?

Yes, each task in a Sieves pipeline maintains its own independent condition callable. You can configure Task A to process only PDFs while Task B processes only long documents, and both will execute their conditions separately as the pipeline forwards documents through each step.

Does the condition parameter affect batch processing?

The condition is evaluated within the batch processing loop in Task.__call__. For each batch, the task identifies which documents satisfy the condition (passing_indices) and only processes those, while still yielding skipped documents to maintain stream continuity. This means batching efficiency is preserved while allowing per-document filtering.

How do I access the document index in my condition function?

The condition callable receives only the Doc object as its argument by default. If you need access to the document index, you can use a closure with itertools.count() or a class-based callable that maintains state. For example: lambda d, i=itertools.count(): next(i) < 100 processes only the first 100 documents.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →