Skip to main content

scietex.service

Async worker framework for building background daemon services in Python.

Provides a hierarchy of workers — from basic signal-handling daemons to concurrent task processors with Valkey-backed distributed queues.

Python ≥ 3.10 · License: MIT

Documentation

  • Overview — Core components and architecture
  • BasicWorker — Signal handling, logging, heartbeat & watchdog managers
  • TaskProcessor — Concurrent task processing, handler dispatch, timeout monitoring
  • ValkeyWorker — Valkey stream-based task distribution
  • Task Handler — Pluggable handler architecture, typed schemas

Installation

# Core package (no Valkey)
pip install scietex.service

# With Valkey (Redis-compatible) support
pip install "scietex.service[valkey]"

Dependencies: msgspec>=0.20.0, pyyaml>=6.0, scietex.logging>=2.0.0

Quick Start

Basic Async Worker

A minimal daemon with signal handling, heartbeat, and watchdog. See the full BasicWorker docs for lifecycle, manager system, and configuration details.

import asyncio
import logging
from scietex.service import BasicWorker, WorkerConfig


class MyWorker(BasicWorker):
    async def heartbeat(self) -> None:
        self.logger.info("Worker is alive")

    async def watchdog(self) -> None:
        self.logger.debug("Running watchdog checks")

    async def cleanup(self) -> None:
        self.logger.info("Shutting down gracefully")


async def main() -> None:
    worker = MyWorker(
        WorkerConfig(
            service_name="my_service",
            version="1.0.0",
            logging_level=logging.DEBUG,
            heartbeat_interval=10,
            watchdog_interval=1,
        )
    )
    await worker.start()
    await worker.events["exit"].wait()


if __name__ == "__main__":
    asyncio.run(main())

Send SIGINT (Ctrl+C) or SIGTERM to trigger graceful shutdown.

Task Processor

Register handlers for different task types and process them concurrently. See the full TaskProcessor docs for architecture, task processing flow, and best practices.

import asyncio
import logging
from scietex.service import TaskProcessor, TaskProcessorConfig
from scietex.service.task_handler import TaskData, TaskHandler, TaskResult


class EmailHandler(TaskHandler):
    @property
    def supported_tasks(self) -> list[str]:
        return ["send_email"]

    async def initialize(self) -> bool:
        # Connect to email service, etc.
        self.logger.info("Email handler initialized")
        return True

    async def handle(self, task_data: TaskData) -> TaskResult:
        try:
            # Process task_data.payload
            self.logger.info("Sending email…")
            return TaskResult(status="success", error="")
        except Exception as exc:
            return TaskResult(status="error", error=str(exc))


class MyProcessor(TaskProcessor):
    async def fetch_tasks(self) -> bool:
        # Pull tasks from your source (DB, API, queue, etc.)
        # and enqueue them for processing:
        #     self.enqueue_task(task_id, task_data)
        # Return True when at least one task was enqueued so the
        # task_queue_manager drains a backlog back-to-back.
        return False


async def main() -> None:
    processor = MyProcessor(
        TaskProcessorConfig(
            service_name="email_worker",
            version="1.0.0",
            queue_size=100,
            max_concurrent_tasks=5,
        )
    )
    processor.add_task_handler(EmailHandler)
    await processor.start()
    await processor.events["exit"].wait()


if __name__ == "__main__":
    asyncio.run(main())

Valkey Worker

Distributed task processing backed by a Valkey (Redis-compatible) stream. See the full ValkeyWorker docs for architecture, key naming, and configuration reference.

import asyncio
import logging
from scietex.service import (
    ValkeyAdvancedConfig,
    ValkeyBaseConfig,
    ValkeyConfig,
    ValkeyNode,
    ValkeyWorker,
    ValkeyWorkerConfig,
)


async def main() -> None:
    config = ValkeyConfig(
        base_config=ValkeyBaseConfig(
            nodes=[ValkeyNode(host="localhost", port=6379)],
            request_timeout=10_000,
        ),
        advanced_config=ValkeyAdvancedConfig(
            connection_timeout=10_000,
            tcp_nodelay=True,
        ),
    )
    worker = ValkeyWorker(
        ValkeyWorkerConfig(
            service_name="distributed_worker",
            version="1.0.0",
            logging_level=logging.DEBUG,
            heartbeat_interval=10,
            valkey_config=config,
            queue_size=100,
            max_concurrent_tasks=10,
        )
    )
    await worker.start()
    await worker.events["exit"].wait()


if __name__ == "__main__":
    asyncio.run(main())

Tasks are stored in a Valkey stream named scietex:{service_name}:tasks and consumed via a consumer group scietex:{service_name}:task_group.

ValkeyWorker also exposes:

  • client_factory= (keyword-only) — an async callable (GlideClientConfiguration) -> Awaitable[GlideClient] used by connect(); defaults to GlideClient.create. Inject a fake to test without a server.
  • transport_health — a TransportHealth supervisor aggregating connection failures, owning the single reconnect path, and logging one CRITICAL per sustained outage.
  • task_lease_ttl (config field) — lease lifetime in seconds; None derives max(1, int(max(2*heartbeat_interval, 3*watchdog_interval))).

Architecture

Worker Hierarchy

See BasicWorker, TaskProcessor, and ValkeyWorker for detailed architecture diagrams.

BasicWorker          — Signal handling, async logging, heartbeat &
                            watchdog managers, graceful shutdown
    └── TaskProcessor — Task queue, concurrent processing, handler
                            dispatch, timeout watchdog
        └── ValkeyWorker  — Valkey stream integration, connection
                            management, stream-based task fetching

Transport Layer

Task delivery is abstracted behind the TaskTransport protocol (fetch/requeue/release/on_started/ack/on_progress/on_drain). TaskProcessor composes a transport rather than inheriting delivery hooks:

  • InMemoryTransport is the default: a deque-backed in-process transport. Feed it with transport.submit(task_id, task_data); fetch drains it into the processor's queue. A bare TaskProcessor therefore works with no external backend.
  • ValkeyTransport (in scietex.service.valkey) implements the same protocol over a Valkey stream; ValkeyWorker injects it automatically.

Pass a custom transport with the keyword-only transport= argument:

from scietex.service import InMemoryTransport, TaskProcessor, TaskProcessorConfig

transport = InMemoryTransport(logger=logging.getLogger("transport"))
processor = TaskProcessor(TaskProcessorConfig(service_name="svc"), transport=transport)
transport.submit(task_id, task_data)

The legacy template-method hooks (fetch_tasks, return_task_to_queue, on_task_started, on_task_completed, _write_task_progress, _on_queue_drain_task_processing) are retained on TaskProcessor as thin delegators to the transport, so existing subclasses keep working.

Manager Lifecycle

Managers are async methods decorated with @Manager. The worker discovers them via the class MRO and runs each as an asyncio.Task:

  1. Start — Manager loop runs the decorated method in a while True loop until cancelled.
  2. Error — On any exception (except CancelledError), the error is recorded and the manager is automatically restarted, up to manager_max_retries consecutive failures (default 5), after which the manager gives up and ends in the terminal FAILED state (observable via worker.failed_managers; the watchdog logs CRITICAL but does not auto-shutdown).
  3. Stop — On shutdown, managers are cancelled and their optional cleanup callbacks are invoked.

Task Handler System

See the Task Handler docs for the full handler lifecycle, schema details, and best practices.

  1. Register: processor.add_task_handler(HandlerClass) — Registers a handler class under its class name. The processor creates a single handler instance on start. Dispatch is driven by the handler's supported_tasks declaration, not by a registration key. An optional keyword-only name (processor.add_task_handler(HandlerClass, name="...")) lets multiple instances of one class coexist under distinct keys. Arbitrary **handler_kwargs are also forwarded to the handler constructor on every instantiation, enabling stateful handlers — see examples/stateful_handler.py.
  2. Declare support: Handler.supported_tasks property must return a list of task type strings this handler can process.
  3. Dispatch: When a task arrives, the processor calls handler.supports(task_type) on each registered handler. The first handler returning True receives the task.
  4. Initialize: handler.start() calls handler.initialize() and sets handler.is_ready to the returned value, so it is True only if initialize() returned True.
  5. Handle: await handler.handle(task_data) returns a TaskResult with status ("success"/"error"), optional error message, and optional payload.
  6. Timeout: Tasks exceeding their timeout (default 3s) are canceled and either re-queued or discarded per TaskTimeout.timeout_action.
  7. Cancel: A built-in cancel_task handler cancels a running or queued task by id. A deliberate cancel writes status="cancelled" with the original TaskData embedded, so the caller can modify and resubmit it under a new task id.

Task Schemas

All schemas are frozen msgspec.Struct instances (immutable).

Type Description
TaskData Immutable task payload: task (type string), payload (bytes), timeout (TaskTimeout), canceled_action ("requeue"/"discard")
TaskResult Handler result: status ("success"/"error"), error (message), processed_at (UTC datetime), payload (bytes), plus error-taxonomy fields error_code, retryable, partial
TaskTimeout Timeout config: timeout (seconds, None for default 3s), timeout_action ("requeue"/"discard")
TaskStatus Per-task tracking record: task_id, service, task, status ("queued"/"running"/"completed"/"failed"/"cancelled"), progress, result, data (original TaskData embedded on a deliberate cancel), error, error_code, timestamps
TaskTracker Internal runtime handle (in task_handler/runtime.py): tracks running asyncio.Task, associated TaskData, and monotonic start time
TaskEnvelope Versioned transport envelope: version (int, 1) wrapping data (serialized TaskData bytes) — the durable on-the-wire format

Configuration

Config Directory Precedence

The worker searches for a config directory in this order:

  1. conf_dir argument (if provided and is a directory)
  2. SCIETEX_CONFIG_DIR environment variable
  3. $XDG_CONFIG_HOME/scietex/
  4. ~/.config/scietex/
  5. /etc/scietex/
  6. /usr/local/etc/scietex/
  7. ./config/ (current working directory)
  8. ~/.config/scietex/ — created if none of the above exist

The first existing directory is used. If none exist, ~/.config/scietex/ is created.

Valkey Configuration

ValkeyWorker reads valkey.yml from the config directory:

base_config:
  nodes:
    - host: localhost
      port: 6379
  user_credentials: null
  use_tls: false
  request_timeout: 5000
  database_id: null
  client_name: null
  inflight_requests_limit: null
  client_az: null
  lazy_connect: null
  read_from: PRIMARY
  backoff_strategy: null
  protocol: RESP3

advanced_config:
  connection_timeout: 10000
  tcp_nodelay: null
  tls_config:
    use_insecure_tls: false
    root_pem_cacerts: null

If the file is missing, it is created with default values. If the file is present but invalid, a RuntimeError is raised and the file is left untouched. The read (and the default-file write) is deferred to the first connect()/initialize() call — constructing ValkeyWorker() with no explicit valkey_config does not touch the filesystem (AR-066).

API Reference

Exported from scietex.service

Symbol Description
BasicWorker Base async daemon worker
TaskProcessor Concurrent task processor
Manager Decorator for creating managed async loop methods
TaskTransport Protocol for the task-delivery backend (fetch/requeue/release/on_started/ack/on_progress/on_drain)
TaskSink Protocol for the enqueue surface a transport delivers into (task_queue_full/enqueue_task)
InMemoryTransport Default in-process transport (deque-backed; submit() feeds it)
ValkeyWorker Valkey-backed distributed worker
WorkerConfig Immutable msgspec.Struct configuration for BasicWorker
TaskProcessorConfig Immutable configuration for TaskProcessor (extends WorkerConfig)
ValkeyWorkerConfig Immutable configuration for ValkeyWorker (extends TaskProcessorConfig)
__version__ Package version string

The Valkey configuration classes (ValkeyConfig, ValkeyNode, ValkeyUserCredentials, ValkeyBackoffStrategy, ValkeyBaseConfig, ValkeyAdvancedConfig, ValkeyPubSubConfig, ValkeyTlsAdvancedConfiguration, ValkeyWorkerConfig) are top-level re-exports: they are importable directly from scietex.service (as in the Valkey quick-start above), not only from scietex.service.valkey.

Exported from scietex.service.task_handler

Symbol Description
TaskHandler Abstract base class for task handlers
TaskHandlerContext Narrow read-only context passed to handlers (service name, instance id, logger)
CancelTaskHandler Built-in handler for the cancel_task task type
CancelTaskRequest Payload schema for a cancel_task task (target_task_id, reason)
CancelTaskResponse Success payload schema for a cancel_task task (target_task_id, outcome)
CancelOutcome Cancellation outcome literal (cancelled/not_running/ignored/not_found)
CancelReason Why a task was cancelled (deliberate/timeout/shutdown)
CANCEL_TASK_TYPE Task type string that selects the built-in cancel handler ("cancel_task")
TaskData Task payload schema
TaskResult Task result schema
TaskTimeout Timeout configuration schema
TaskStatus Per-task tracking record schema
TaskProgress Granular progress payload embedded in TaskStatus.progress (progress, value)
TaskTracker Internal runtime handle for running tasks (not a wire schema)
TaskEnvelope Versioned transport envelope (version + serialized payload bytes)
encode_task_envelope Wrap a TaskData in a versioned envelope and msgpack-encode it
decode_task_envelope Decode an envelope back to a TaskData (returns None on invalid/unknown version)

Exported from scietex.service.valkey

Symbol Description
ValkeyConfig Top-level Valkey configuration
ValkeyBaseConfig Basic connection settings
ValkeyAdvancedConfig Advanced connection settings
ValkeyPubSubConfig PubSub control-channel settings (listening, parse_control_message)
ValkeyNode Server node address
ValkeyUserCredentials Authentication credentials
ValkeyBackoffStrategy Reconnection backoff config
ValkeyTlsAdvancedConfiguration TLS settings
ValkeyWorkerConfig Immutable configuration for ValkeyWorker (extends TaskProcessorConfig)
purge_task_stream Standalone operational utility to purge a task stream

Development

Setup

# Clone the repository and install all dependencies
uv sync --all-extras

# Or install specific extras
uv sync --extra dev --extra test --extra lint

Commands

Command Description
uv run ruff check src/ Lint (auto-fix: ruff check --fix)
uv run ty check src/ Type check
uv run ruff format src/ Format code
uv run pytest tests/ Run tests
tox Run tests with coverage

Running Examples

The examples/ directory contains runnable blueprints; see examples/README.md for what each demonstrates.

python -m examples.basic_worker
python -m examples.manager_cleanup
python -m examples.manager_collision
python -m examples.task_processor
python -m examples.named_task_handlers
python -m examples.stateful_handler
python -m examples.valkey_async_service      # requires valkey-glide
python -m examples.valkey_pubsub_worker      # requires valkey-glide
python -m examples.valkey_perf               # requires valkey-glide
python -m examples.progress_and_cancel       # requires valkey-glide

License

MIT

Release files for scietex.service 4.2.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 scietex.service 4.2.0
File Size Uploaded
scietex_service-4.2.0.tar.gz 82.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for scietex.service 4.2.0
File Interpreter ABI Platform
scietex_service-4.2.0-py3-none-any.whl Python 3 none any Details

Total release size: 161.8 kB

Release files / scietex_service-4.2.0.tar.gz

Download URL scietex_service-4.2.0.tar.gz
Size 82.2 kB
Tags Source
SHA-256 checksum
How to use checksums
472f73f8a6d52e554c73bada05acaf40d0817bbcd34812a409f56b04618eab08
BLAKE2b-256 checksum
How to use checksums
2c82018a402996db18dad3a88ece24a3cf3a23904ce6516551c6bb220686d708
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 16, 2026.

Transparency log

Release files / scietex_service-4.2.0-py3-none-any.whl

Download URL scietex_service-4.2.0-py3-none-any.whl
Size 79.6 kB
Tags Python 3
SHA-256 checksum
How to use checksums
b391867091f3bc6f031f487a7e3afe08d4618a1c8d146be56f223b28fe171f97
BLAKE2b-256 checksum
How to use checksums
d74daee09207cfcfaf79d00fb4f8579af32d9c6df000008fdb5f9a3befc82689
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 16, 2026.

Transparency log

Release history Release notifications | RSS feed

5.0.0

2 release files

4.6.0

2 release files

4.5.0

2 release files

4.4.0

2 release files

4.3.0

2 release files

This release

4.2.0 This release

2 release files

4.1.0

2 release files

4.0.0

2 release files

3.2.0

2 release files

3.1.0

2 release files

3.0.0

2 release files

2.0.0

2 release files

1.0.7

2 release files

1.0.6

2 release files

1.0.5

2 release files

1.0.4

2 release files

1.0.3

2 release files

1.0.2

2 release files

1.0.1

2 release files

1.0.0

2 release files

0.2.1

2 release files

0.2.0

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

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