Skip to main content

Pipeline Toolkit

A small functional pipeline toolkit for Python.

pipeline-toolkit provides a simple way to build sequential pipelines from ordinary Python callables. Each step receives the result of the previous step, while the pipeline runs asynchronously in a worker thread.

Features

  • Sequential functional pipeline execution
  • Asynchronous execution using a worker thread
  • Positional and keyword arguments for pipeline steps
  • Stop, skip, wait, and rerun execution
  • Manual synchronous step execution
  • Result and error history using a stack
  • Pipeline modification with add(), insert(), pop(), and clear()
  • Small utility modules for functional workflows

Installation

Install from PyPI:

pip install pipeline-toolkit

Or install directly from GitHub:

pip install git+https://github.com/Hoang-Long2012/pipeline-toolkit.git

Quick Start

A pipeline is created from an iterable of steps. Each step is a tuple whose first item is a callable.

from pipeline import Pipeline

def add(value, amount):
	return value + amount

def multiply(value, factor):
	return value * factor

pipeline = Pipeline([
	(add, (5,)),
	(multiply, (2,)),
])

pipeline.run(10).wait()

print(pipeline.results.get())

The execution flow is:

10
 ↓
add(10, 5)
 ↓
15
 ↓
multiply(15, 2)
 ↓
30

The final result is 30.

Pipeline Steps

Each step can use one of four supported forms.

Callable only

(function,)

The callable receives the previous result:

pipeline = Pipeline([
	(str.upper,),
])

pipeline.run("hello").wait()

Positional arguments

(function, args)

where args is a tuple:

pipeline = Pipeline([
	(add, (5,)),
	(multiply, (2,)),
])

A step such as:

(add, (5,))

is executed as:

add(previous_result, 5)

Keyword arguments

(function, kwargs)

where kwargs is a mapping:

pipeline = Pipeline([
	(pow, {"exp": 2}),
])

The step is executed as:

pow(previous_result, exp=2)

Positional and keyword arguments

(function, args, kwargs)

For example:

pipeline = Pipeline([
	(my_function, (1, 2), {"option": True}),
])

The callable receives the previous result followed by the supplied positional and keyword arguments.

Execution

run()

Start the pipeline asynchronously.

pipeline.run(default=None, delay=0, daemon=False, stop_on_error=True)

The default value becomes the initial result and is passed to the first step.

delay specifies the delay between steps.

daemon controls whether the worker thread is a daemon thread.

stop_on_error controls whether execution stops after the first exception.

run() returns the pipeline instance, allowing calls such as:

pipeline.run(10).wait()

The pipeline is snapshotted when execution starts. Changes made to pipeline.pipeline after run() begins do not affect the current execution.

wait()

Wait for the current execution to finish.

pipeline.wait()

It returns the pipeline instance.

stop()

Request the running pipeline to stop and wait for its worker thread to terminate.

step = pipeline.stop()

The return value is the current one-based step index when execution is stopped, or 0 if the pipeline was not running.

skip()

Request the worker to skip the next step that reaches its skip check.

pipeline.skip()

The method returns the pipeline instance.

rerun()

Stop the current execution and start the pipeline again.

pipeline.rerun(10)

Arguments are passed directly to run().

Manual Step Execution

run_step() executes one configured step synchronously.

result = pipeline.run_step(2, 10)

Unlike run(), this method:

  • does not create a worker thread
  • does not modify the worker thread or pipeline execution state
  • does not store the result in results
  • does not store exceptions in errors
  • allows exceptions to propagate to the caller

This makes it useful when a single pipeline step needs to be executed manually.

Results and Errors

The pipeline provides two Stack instances:

pipeline.results
pipeline.errors

results contains the initial value and the results produced by executed steps.

For example:

pipeline.run(10).wait()

print(pipeline.results.get())

errors contains exceptions raised by pipeline steps.

When stop_on_error=True, execution stops after the first exception.

When stop_on_error=False, the exception is stored in errors and execution continues with the previous result.

Managing Pipeline Steps

Pipeline steps can be modified before or between executions.

add()

Append a step:

pipeline.add((str.upper,))

insert()

Insert a step at a one-based position:

pipeline.insert(2, (str.strip,))

pop()

Remove and return a step:

step = pipeline.pop(1)

Pipeline indexes are one-based.

clear()

Remove all configured steps:

pipeline.clear()

Pipeline State

The running property indicates whether the worker thread is currently running:

if pipeline.running:
	print("Pipeline is running")

The step attribute contains the one-based index of the currently executing step. It is 0 when the pipeline is not running.

A Pipeline instance can also be used as a boolean:

if pipeline:
	print("Pipeline is running")

Calling a pipeline instance is equivalent to calling run():

pipeline(10)

is equivalent to:

pipeline.run(10)

The length of a pipeline is the number of configured steps:

len(pipeline)

A callable can be checked with the in operator:

if add in pipeline:
	print("add is part of the pipeline")

Callable membership uses identity comparison.

Utilities

Stack

Stack is a simple LIFO stack container with optional capacity limits.

It supports common stack operations such as pushing, retrieving, peeking, and removing items, with dedicated exceptions for overflow and underflow conditions.

Import it directly from its submodule:

from pipeline.stack import Stack

Stack is also used internally by Pipeline for storing results and errors.

For example:

from pipeline.stack import Stack

stack = Stack()

stack.push("first")
stack.push("second")

print(stack.get())

For detailed stack operations and behavior, see the pipeline.stack module.

tap

tap is a small functional utility for performing a side effect while keeping the pipeline value available for subsequent processing.

tap performs a side effect on a deep copy of the current value and returns the original value unchanged.

Import it directly from its submodule:

from pipeline.tap import tap

For example:

from pipeline import Pipeline
from pipeline.tap import tap

def add(value, amount):
	return value + amount

pipeline = Pipeline([
	(add, (5,)),
	(tap, (print,)),
	(add, (10,)),
])

pipeline.run(10).wait()

Utilities are provided as separate submodules rather than being exported from the top-level pipeline package.

API Overview

Pipeline

Member Description
run() Start asynchronous pipeline execution
run_step() Execute one step synchronously
stop() Stop the current execution
skip() Request the next step to be skipped
wait() Wait for the current execution
rerun() Restart the pipeline
add() Append a step
insert() Insert a step
pop() Remove and return a step
clear() Remove all steps
running Whether the worker is running
step Current one-based step index
results Stack of initial value and results
errors Stack of raised exceptions

Requirements

  • Python 3.8 or newer

Changelog

See changelog from: CHANGELOG.md

License

This project is licensed under the MIT License. See LICENSE for details.

Contribution

If you'd like to contribute, feel free to submit a pull request.
If you'd like to report a bug or request a feature, please open an issue.

Copyright (C) 2026 Hoàng Long

Release files for pipeline-toolkit 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for pipeline-toolkit 0.1.0
File Size Uploaded
pipeline_toolkit-0.1.0.tar.gz 12.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for pipeline-toolkit 0.1.0
File Interpreter ABI Platform
pipeline_toolkit-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 23.5 kB

Release files / pipeline_toolkit-0.1.0.tar.gz

Download URL pipeline_toolkit-0.1.0.tar.gz
Size 12.4 kB
Tags Source
SHA-256 checksum
How to use checksums
a73fc2ecb268cfcfa64826eb04974fe5a7f3991fafd4d616b31e61607356e43e
BLAKE2b-256 checksum
How to use checksums
0458847e9bfca88979e911c936b1ad1f9af59ef75c30079bf71f49e5f21f5608
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.14.7

Release files / pipeline_toolkit-0.1.0-py3-none-any.whl

Download URL pipeline_toolkit-0.1.0-py3-none-any.whl
Size 11.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
c3d03dc7e708f771f86f209d5d517dea6a38ec104d60c07d877d129fb56a6498
BLAKE2b-256 checksum
How to use checksums
7fbe0591587772f3bb130482813cb76b9a621a9ce0a66e6d582685db936d14a9
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.14.7

Release history Release notifications | RSS feed

0.5.0

2 release files

0.4.1

2 release files

0.4.0

2 release files

0.3.2

2 release files

0.3.1

2 release files

0.3.0

2 release files

0.2.4

2 release files

0.2.3

2 release files

0.2.2

2 release files

0.2.1

2 release files

0.2.0

2 release files

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page