Skip to main content

tomviz-pipeline

A pure-Python node-graph pipeline engine for processing volumetric / tomographic data. It originated as the external-execution runtime of tomviz and is being developed as a standalone library so that other applications can build, run, and persist data processing pipelines without depending on the tomviz application, Qt, or VTK.

Layout

The package is layered:

  • tomviz_pipeline.core — a generic, dependency-free (stdlib only) pipeline engine: Pipeline, Node (SourceNode / TransformNode / SinkNode), typed InputPort / OutputPort, Link, PortData, dirty tracking (New / Stale / Current), execution planning, blocking (DefaultExecutor) and threaded (ThreadedExecutor) pipeline executors, ExecutionFuture, thread-safe Signal / EventQueue observers, and a per-node NodeExecutor abstraction (in-process vs. out-of-process node execution).
  • The rest of tomviz_pipeline — the tomviz-flavored layer: the numpy Dataset model (plus pure-Python Table and Molecule payloads for the corresponding port types), EMD/HDF5 I/O, schema-v2 state files (.tvsm JSON and .tvh5 HDF5 containers), the built-in source/transform/sink nodes, Python operator authoring APIs, progress reporting channels, the ExternalNodeExecutor (run a node under a different Python environment), the batch runner, and the tomviz-pipeline CLI.

Applications that just need a pipeline engine can depend on tomviz_pipeline.core alone; tomviz uses the whole package.

Installation

pip install tomviz-pipeline

Requires Python 3.9+. The runtime dependencies are numpy, h5py, click and tqdm. The optional tomviz-pipeline[vtk] extra lets state files carry tables built directly with VTK by operator scripts; operator scripts bring their own scientific dependencies (scipy, tomopy, ...).

Quick start

A python node is built from exactly two artifacts: a definition (the interface — ports, parameters, label; the same JSON vocabulary as tomviz's operator .json sidecar files) and a kernel (the compute — the same class a tomviz operator script contains, so anyone who has written a tomviz Python operator already knows how):

import numpy as np

from tomviz_pipeline import Pipeline, PythonNode
from tomviz_pipeline.dataset import Dataset
from tomviz_pipeline.kernels import SourceKernel, TransformKernel

class ConstantVolume(SourceKernel):
    def produce(self, value=1.0, side=32):
        arr = np.full((side, side, side), value, dtype=np.float32)
        return {'volume': Dataset({'Scalars': arr}, active='Scalars')}

class Multiply(TransformKernel):
    def transform(self, inputs, factor=2.0):
        ds = inputs['volume']
        return {'volume': ds.apply_to_each_scalar_array(
            lambda a: a * factor)}

CONSTANT_VOLUME = {
    'name': 'ConstantVolume',
    'outputs': [{'name': 'volume', 'type': 'ImageData'}],
    'parameters': [{'name': 'value', 'type': 'double', 'default': 1.0},
                   {'name': 'side', 'type': 'int', 'default': 32}],
}

MULTIPLY = {
    'name': 'Multiply',
    'inputs':  [{'name': 'volume', 'type': 'ImageData'}],
    'outputs': [{'name': 'volume', 'type': 'ImageData'}],
    'parameters': [{'name': 'factor', 'type': 'double', 'default': 2.0}],
}

pipeline = Pipeline()
source = pipeline.add_node(PythonNode(CONSTANT_VOLUME,
                                      kernel=ConstantVolume))
scale = pipeline.add_node(PythonNode(MULTIPLY, kernel=Multiply))
pipeline.create_link(source.output_port('volume'),
                     scale.input_port('volume'))
source.set_parameters(value=21.0)

pipeline.execute()
print(scale.output_port('volume').data().payload.active_scalars.max())
# 42.0

Both constructor arguments are polymorphic, with one rule: strings are content, paths are files.

  • definition: a dict, a JSON string, or a pathlib.Path to a .json file. The definition decides the node shape — no inputs means a source.python node, otherwise transform.python.
  • kernel: a kernel class you hold, the source text of an operator script, or a pathlib.Path to a .py file.

That makes the classic sidecar pair — and kernels whose imports only resolve in another Python environment — work without ever importing the kernel here:

from pathlib import Path
from tomviz_pipeline import ExternalNodeExecutor

recon = pipeline.add_node(PythonNode(Path('operators/Reconstruct.json'),
                                     kernel=Path('operators/Reconstruct.py')))
recon.set_parameters(iterations=100)
recon.node_executor = ExternalNodeExecutor('/envs/tomopy')

Notes on the kernel side:

  • Kernels get self.progress, self.canceled, and self.completed for long-running work; a kernel sees deep-copied datasets and keyword parameters — no graph machinery — which is what lets the same class run unchanged inside the tomviz application, in this standalone runtime, and in an external environment.
  • A class-bound kernel executes in-process directly. Whenever the node must be serialized (state files, external execution) the class is re-expressed as a script by capturing its source; for that the class body must be self-contained apart from numpy / tomviz_pipeline imports. Kernels given as script text or a .py path serialize as-is.
  • Kernels also get self.state: a dict the engine preserves across executions of the node (each run gets a fresh instance, so plain attributes don't survive). It round-trips through external execution but is never written to state files. Keep its values JSON-friendly.
  • A kernel may implement should_auto_execute(self, **parameters) -> bool for periodic execution: when node.auto_execute_enabled is set an application polls executor.should_auto_execute(node) every node.auto_execute_interval_seconds — through the node's executor, so an external node is asked inside its own environment — and re-executes the pipeline on True. The hook receives the node's parameters and shares self.state with produce / transform.
  • A kernel can write back to its own parameters with self.set_parameter(name, value) (the name must be declared in the description; the value is coerced to the declared type). Updates are installed on the node when the method returns — quietly: nothing is marked stale and no re-execution is triggered, because the run that made the change is deemed to have consumed it — and surface as the node's parameters_updated(node, changed) signal so an application can refresh its parameter UI. The next run receives the new values; inside should_auto_execute, return True to run with them right away. Updates round-trip through external execution like self.state does.
  • Older operator scripts that import tomviz.nodes (the historical spelling, tomviz.nodes.TransformNode) keep working: those names resolve to the kernel classes through a compatibility alias.

Parameters

Node configuration lives in a private parameters store (so names can never collide with node internals like label or state) and is changed through set_parameters(), which marks the node and everything downstream stale and emits parameters_applied:

scale.set_parameters(factor=3)      # scale (and downstream) now Stale
pipeline.execute()                  # re-runs just what's needed

pipeline.auto_execute = True        # optional: C++-style behavior where
scale.set_parameters(factor=4)      # applying parameters re-executes

Values a kernel writes back with self.set_parameter (see above) take the other door, apply_parameter_updates(): the store changes, nothing is marked stale, and only parameters_updated is emitted.

Using it from an application

A pipeline embedded in an application must not block the UI thread. Install a ThreadedExecutor once and every execute() variant returns immediately with an ExecutionFuture; the plan runs on a worker thread.

Notifications are Signals. A plain connect(handler) is a direct connection — the handler runs on whichever thread emits, which during execution is the worker, so it must be thread-safe and must not touch UI. Connecting through an EventQueue defers delivery: emission only enqueues, and the handler runs when your UI thread drains the queue — the moral equivalent of Qt's queued connections, without Qt. (The optional second argument to connect() is a dispatcherEventQueue here, AsyncioDispatcher for event-loop apps below.)

from tomviz_pipeline import EventQueue, ThreadedExecutor

# Fresh pipeline with the ConstantVolume / Multiply kernels from the
# quick start.
pipeline = Pipeline()
source = pipeline.add_node(PythonNode(CONSTANT_VOLUME,
                                      kernel=ConstantVolume))
scale = pipeline.add_node(PythonNode(MULTIPLY, kernel=Multiply))
pipeline.create_link(source.output_port('volume'),
                     scale.input_port('volume'))

pipeline.set_executor(ThreadedExecutor())

# Everything below runs on the thread that drains `events` — never on
# the worker.
events = EventQueue()
pipeline.executor.node_execution_started.connect(
    lambda node: print(f'started:  {node.label}'), events)
pipeline.executor.node_execution_finished.connect(
    lambda node, ok: print(f'finished: {node.label} ok={ok}'), events)
scale.progress_step_changed.connect(
    lambda node, step: print(f'progress: {step}'), events)
pipeline.execution_finished.connect(
    lambda future: print(f'done, succeeded={future.succeeded()}'), events)

future = pipeline.execute()          # returns immediately

# A real application drains the queue from its main loop (see the Qt
# sketch below). A minimal stand-in loop:
while not future.is_finished():
    events.process(block=True, timeout=0.05)
events.process()                     # deliver anything queued after finish

Interactive control, all cooperative and non-blocking:

future = pipeline.execute()          # user starts a run...
pipeline.cancel_execution()          # ...and changes their mind: stops at
future.wait(10)                      # the next node boundary; running
                                     # kernels observe self.canceled

pipeline.set_paused(True)            # batch edits without re-running
scale.set_parameters(factor=3.0)
pipeline.set_paused(False)           # un-pausing re-executes what's stale

pipeline.auto_execute = True         # or: every set_parameters() re-runs,
scale.set_parameters(factor=4.0)     # as in the tomviz application

pipeline.executor.cancel_and_wait()  # teardown: join the worker.
                                     # pipeline.clear() does this for you.

Calling execute() while a run is in flight never queues behind it: the new plan becomes pending, the in-flight run is canceled at its next node boundary, and the pending plan starts when the worker exits — the right semantics for "the user dragged the slider again".

In a Qt application the drain is a timer on the GUI thread (illustrative sketch):

class PipelinePanel(QWidget):
    def __init__(self, pipeline):
        ...
        self.events = EventQueue()
        pipeline.set_executor(ThreadedExecutor(sync_queue=self.events))
        pipeline.execution_finished.connect(self._on_finished,
                                            self.events)
        self._timer = QTimer(self, interval=16)   # ~60 Hz
        self._timer.timeout.connect(self.events.process)
        self._timer.start()

The optional sync_queue makes the worker wait, after each node, until the application thread has drained the queue — use it when per-node UI updates (e.g. re-rendering) must complete before the next node starts, mirroring the C++ tomviz executor's barrier.

In an asyncio application there is no drain timer to write: an AsyncioDispatcher plays the EventQueue's role, delivering handlers onto the event loop. Create it once, on the loop, and pass it to any number of connect() calls; plain callbacks run as loop callbacks, and coroutine-function handlers are scheduled as tasks. Await a run's completion by parking future.wait() on a thread-pool thread:

import asyncio

from tomviz_pipeline import AsyncioDispatcher

async def run_pipeline(pipeline):
    dispatcher = AsyncioDispatcher()        # bound to this event loop

    connection = pipeline.executor.node_execution_finished.connect(
        lambda node, ok: print(f'finished: {node.label} ok={ok}'),
        dispatcher)

    future = pipeline.execute()             # returns immediately
    await asyncio.to_thread(future.wait)    # yields until the run ends
    connection.disconnect()
    return future.succeeded()

scale.set_parameters(factor=5.0)
assert asyncio.run(run_pipeline(pipeline)) is True

Cancelling an asyncio task that is awaiting a run does not cancel the pipeline — awaiting is observation, not ownership. Couple them explicitly when you want to:

async def run_and_own(pipeline):
    future = pipeline.execute()
    try:
        await asyncio.to_thread(future.wait)
    except asyncio.CancelledError:
        pipeline.cancel_execution()   # propagate deliberately
        raise

Output persistence

Each output port is either transient (port.persistent = False — its payload lives only while some consumer holds a handle; the planner re-runs the producer when the data is needed again) or persistent, in which case port.persistence_mode picks the medium: InMemory pins the payload on the port, OnDisk spills it to a temp cache file ($TOMVIZ_PORT_CACHE_DIR or the system temp dir) once the last handle drops and reloads it lazily via port.materialize(). port.data() is a non-loading peek. Sources default to persistent-InMemory; transform outputs follow the pipeline-wide default:

from tomviz_pipeline import PipelineSettings, TransformPersistenceDefault

PipelineSettings.instance().transform_persistence_default = \
    TransformPersistenceDefault.OnDisk   # keep intermediates out of RAM

Modes can be switched at runtime and round-trip through schema-v2 state files (persistent / persistenceMode keys, identical to the C++ implementation).

Graph nodes (application integration)

Kernels cover custom computation. Applications embedding the engine that need custom graph machinery — their own port layouts, payload types, serialization, or state handling — subclass the core graph classes directly and register them under their own type strings:

from tomviz_pipeline import (
    DefaultExecutor, NodeFactory, Pipeline, PortData, SourceNode,
    TransformNode,
)

class Constant(SourceNode):
    type_name = 'example.constant'

    def __init__(self, value=0):
        super().__init__()
        self._parameters['value'] = value
        self.add_output('output', 'ImageData')

    def execute(self):
        self.output_port('output').set_data(
            PortData(self.parameter('value'), 'ImageData'))
        return True

class Scale(TransformNode):
    type_name = 'example.scale'

    def __init__(self, factor=1):
        super().__init__()
        self._parameters['factor'] = factor
        self.add_input('input', 'ImageData')
        self.add_output('output', 'ImageData')

    def transform(self, inputs):
        value = inputs['input'].payload * self.parameter('factor')
        return {'output': PortData(value, 'ImageData')}

NodeFactory.register('example.constant', Constant)
NodeFactory.register('example.scale', Scale)

pipeline = Pipeline()
source = pipeline.add_node(Constant(21))
scale = pipeline.add_node(Scale(2))
pipeline.create_link(source.output_port('output'),
                     scale.input_port('input'))
DefaultExecutor(pipeline).execute()
print(scale.output_port('output').data().payload)  # 42

This is the layer the built-in tomviz node types (readers, crops, reconstructions, ...) are made of. Note the trade-off versus kernels: graph nodes speak PortData and manage their own ports, but only deserialize in processes where the application's classes are importable — a kernel-hosted node round-trips anywhere because it travels with its source.

Running a state file

tomviz-pipeline -s pipeline.tvsm -o output/
tomviz-pipeline -s pipeline.tvsm -o output/ --input 'data/*.emd'

Development

git clone https://github.com/OpenChemistry/tomviz-pipeline
cd tomviz-pipeline
pip install -e .[dev]
flake8 --config setup.cfg src tests
pytest

Releases are built and uploaded to PyPI by the Publish to PyPI workflow when a GitHub release is published; bump __version__ in src/tomviz_pipeline/__init__.py (the single source of the package version) and tag v<version>.

License

BSD 3-Clause. See LICENSE.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

tomviz_pipeline-3.1.3.tar.gz (155.7 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

tomviz_pipeline-3.1.3-py3-none-any.whl (133.0 kB view details)

Uploaded Python 3

File details

Details for the file tomviz_pipeline-3.1.3.tar.gz.

File metadata

  • Download URL: tomviz_pipeline-3.1.3.tar.gz
  • Upload date:
  • Size: 155.7 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for tomviz_pipeline-3.1.3.tar.gz
Algorithm Hash digest
SHA256 278b8fff86d679f88155bc11a59d136f219710e2a4e9839acaee51be6e201acb
MD5 60a5a01f6cc05aa148041e9e77179a63
BLAKE2b-256 371ee10e106b336585c431c8a1d422bfd7d9d301faf7372e0f70835236d6e55c

See more details on using hashes here.

Provenance

The following attestation bundles were made for tomviz_pipeline-3.1.3.tar.gz:

Publisher: publish.yml on OpenChemistry/tomviz-pipeline

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file tomviz_pipeline-3.1.3-py3-none-any.whl.

File metadata

  • Download URL: tomviz_pipeline-3.1.3-py3-none-any.whl
  • Upload date:
  • Size: 133.0 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for tomviz_pipeline-3.1.3-py3-none-any.whl
Algorithm Hash digest
SHA256 73fe6191bac4967a69bfaaf73abf0f5bef1569fede0c5fcf97f645f5aa67351a
MD5 c604b3d4e81f0f679c8f046fc7090b1f
BLAKE2b-256 04e62fb5433533f7584a04bba495c26c497d94f9e67d6f9ed2bf80a6c2d025b7

See more details on using hashes here.

Provenance

The following attestation bundles were made for tomviz_pipeline-3.1.3-py3-none-any.whl:

Publisher: publish.yml on OpenChemistry/tomviz-pipeline

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Release history Release notifications | RSS feed

3.1.4

2 files

This release

3.1.3 This release

2 files

3.1.2

2 files

3.1.1

2 files

3.1.0

2 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