Skip to content

About

Production-grade mini workflow orchestration with a Python DAG scheduler, Go workers, ZeroMQ transport, SQLite persistence, and a PySide6 desktop UI.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

12 Commits

Folders and files

Repository files navigation

Mini Workflow Orchestration System

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

Screenshots

Desktop UI Overview

Desktop UI Overview

Orchestrator Log Output

Orchestrator Log Output

Retry View in UI

Retry View in UI

Architecture

1. Python Orchestrator

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.

2. DAG Core

workflow/core/dag.py

  • Directed acyclic graph with cycle detection on dependency insertion.
  • Deterministic topological sort.
  • Computes ready nodes from completed + in-progress sets.

3. Worker Engine (Go)

go_worker/main.go, go_worker/engine.go

  • Pulls tasks from port 5555 and pushes results to port 5556.
  • 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).

4. Persistence Layer

workflow/db/repository.py

  • SQLite schema for tasks and dependencies.
  • Save/load whole workflow and incremental task upserts.
  • Recovery path restores persisted workflow and resumes orchestration.

5. Desktop UI

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.

6. Logging

workflow/logging_config.py

  • Centralized logging with console + rotating file handler.
  • Logs workflow lifecycle and task events (dispatch/success/retry/failure).

Communication Contract

Task Message (Python -> Go)

{
  "id": "transform_data",
  "retry": 1,
  "max_retry": 3,
  "payload": {
    "fail_probability": 0.35
  }
}

Result Message (Go -> Python)

{
  "id": "transform_data",
  "status": "SUCCESS",
  "retry": 1,
  "worker": "worker-2"
}

Project Structure

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

Run Instructions

1. Start Go workers

cd go_worker
go run . -workers 4

2. Start Python app (new terminal)

python -m venv .venv
# Windows PowerShell
.venv\Scripts\Activate.ps1
pip install -r requirements.txt
python -m workflow.main

You can start either side first. The Go worker retries connecting to the Python result endpoint.

Run Tests

Python tests

$env:PYTEST_DISABLE_PLUGIN_AUTOLOAD='1'; pytest -q

Go tests

cd go_worker
go test ./...

CI/CD and Releases

  • CI workflow runs automatically on every push and pull request:
    • Python dependency install + pytest
    • Go worker go test ./...
  • Release workflow runs automatically when you push a version tag like v1.0.0.
  • For each release tag, GitHub Release assets include:
    • Cross-platform worker binaries (linux, windows, darwin, amd64/arm64)
    • ZIP packages containing Python app + matching worker binary
    • A source ZIP package

Create a release

git tag v1.0.0
git push origin v1.0.0

After the workflow completes, check: https://github.com/mongo-driver/mini-workflow/releases

Architecture Decisions

  • 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.

About

Production-grade mini workflow orchestration with a Python DAG scheduler, Go workers, ZeroMQ transport, SQLite persistence, and a PySide6 desktop UI.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages