Production-grade mini workflow orchestration built with:
- Python orchestrator (DAG scheduling, retry policy, persistence, logging)
- Go worker engine (parallel processing with goroutines)
- ZeroMQ transport (PUSH/PULL)
- SQLite task state store
- PySide6 desktop UI for live monitoring
workflow/services/orchestrator.py
- Maintains in-memory DAG + task state.
- Dispatches ready tasks to workers using ZeroMQ (
tcp://127.0.0.1:5555). - Receives worker results from ZeroMQ (
tcp://*:5556). - Implements retry logic (
retry < max_retry) to simulate at-least-once processing. - Propagates permanent failures to blocked downstream tasks so workflows always terminate.
- Synchronizes in-memory state with SQLite on every task change.
workflow/core/dag.py
- Directed acyclic graph with cycle detection on dependency insertion.
- Deterministic topological sort.
- Computes ready nodes from completed + in-progress sets.
go_worker/main.go, go_worker/engine.go
- Pulls tasks from port
5555and pushes results to port5556. - Uses goroutine pool (
-workers N) for concurrent execution. - Retries outbound result-socket connection on startup to avoid strict startup ordering.
- Simulates failures with configurable payload (
fail_probability,force_fail).
workflow/db/repository.py
- SQLite schema for
tasksanddependencies. - Save/load whole workflow and incremental task upserts.
- Recovery path restores persisted workflow and resumes orchestration.
workflow/ui/main_window.py
- Read-only table with
task_id,status,retry/max_retry. - "Create & Run Demo Workflow" button.
- Auto-refresh every second via
QTimer. - UI remains responsive because orchestration runs in a background thread.
workflow/logging_config.py
- Centralized logging with console + rotating file handler.
- Logs workflow lifecycle and task events (dispatch/success/retry/failure).
{
"id": "transform_data",
"retry": 1,
"max_retry": 3,
"payload": {
"fail_probability": 0.35
}
}{
"id": "transform_data",
"status": "SUCCESS",
"retry": 1,
"worker": "worker-2"
}mini workflow/
workflow/
core/
dag.py
models.py
db/
repository.py
services/
orchestrator.py
workflow_builder.py
zmq_gateway.py
ui/
main_window.py
logging_config.py
main.py
go_worker/
engine.go
engine_test.go
go.mod
main.go
tests/
test_dag.py
test_retry.py
requirements.txt
README.md
cd go_worker
go run . -workers 4python -m venv .venv
# Windows PowerShell
.venv\Scripts\Activate.ps1
pip install -r requirements.txt
python -m workflow.mainYou can start either side first. The Go worker retries connecting to the Python result endpoint.
$env:PYTEST_DISABLE_PLUGIN_AUTOLOAD='1'; pytest -qcd go_worker
go test ./...CIworkflow runs automatically on every push and pull request:- Python dependency install +
pytest - Go worker
go test ./...
- Python dependency install +
Releaseworkflow runs automatically when you push a version tag likev1.0.0.- For each release tag, GitHub Release assets include:
- Cross-platform
workerbinaries (linux,windows,darwin,amd64/arm64) - ZIP packages containing Python app + matching worker binary
- A source ZIP package
- Cross-platform
git tag v1.0.0
git push origin v1.0.0After the workflow completes, check:
https://github.com/mongo-driver/mini-workflow/releases
- PUSH/PULL over ZeroMQ: simple load balancing across multiple Go workers with no custom broker.
- Single writer sockets in Go: avoids ZeroMQ socket thread-safety issues while still allowing worker pool concurrency.
- SQLite for local durability: lightweight, zero-admin persistence suitable for desktop workflow orchestration.
- Background orchestrator loop: keeps UI responsive and allows continuous scheduling/result handling.
- Retry as at-least-once simulation: failed tasks are re-dispatched until
max_retry. - Fail-fast DAG completion: downstream tasks blocked by permanently failed dependencies are explicitly marked
FAILED.


