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
  • BasicAsyncWorker — Signal handling, logging, heartbeat & watchdog managers
  • AsyncTaskProcessor — 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>=1.1.0

Quick Start

Basic Async Worker

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

import asyncio
import logging
from scietex.service import BasicAsyncWorker


class MyWorker(BasicAsyncWorker):
    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(
        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 AsyncTaskProcessor docs for architecture, task processing flow, and best practices.

import asyncio
import logging
from scietex.service import AsyncTaskProcessor
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(AsyncTaskProcessor):
    async def fetch_tasks(self) -> None:
        # Pull tasks from your source (DB, API, queue, etc.)
        # and enqueue them for processing:
        #     self.enqueue_task(task_id, task_data)
        pass


async def main() -> None:
    processor = MyProcessor(
        service_name="email_worker",
        version="1.0.0",
        queue_size=100,
        max_concurrent_tasks=5,
    )
    processor.add_task_handler("send_email", 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,
)


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(
        service_name="distributed_worker",
        version="1.0.0",
        worker_id=1,
        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}:{worker_id}:tasks and consumed via a consumer group scietex:{service_name}:{worker_id}:task_group.

Architecture

Worker Hierarchy

See BasicAsyncWorker, AsyncTaskProcessor, and ValkeyWorker for detailed architecture diagrams.

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

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.
  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("type", HandlerClass) — Registers a handler class under a name. The processor creates handler instances on start. The name is validated against the handler's supported_tasks (a warning is logged when it isn't among them, since such a name can never be dispatched to).
  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 = 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.

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, retry_count, partial, requeue
TaskTimeout Timeout config: timeout (seconds, None for default 3s), timeout_action ("requeue"/"discard")
TaskTracker Internal: tracks running asyncio.Task, associated TaskData, and monotonic start time

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.

API Reference

Exported from scietex.service

Symbol Description
BasicAsyncWorker Base async daemon worker
AsyncTaskProcessor Concurrent task processor
Manager Decorator for creating managed async loop methods
ValkeyWorker Valkey-backed distributed worker
__version__ Package version string

Exported from scietex.service.task_handler

Symbol Description
TaskHandler Abstract base class for task handlers
TaskData Task payload schema
TaskResult Task result schema
TaskTimeout Timeout configuration schema
TaskTracker Internal task tracker schema

Exported from scietex.service.valkey

Symbol Description
ValkeyConfig Top-level Valkey configuration
ValkeyBaseConfig Basic connection settings
ValkeyAdvancedConfig Advanced connection settings
ValkeyNode Server node address
ValkeyUserCredentials Authentication credentials
ValkeyBackoffStrategy Reconnection backoff config
ValkeyTlsAdvancedConfiguration TLS settings

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

python -m examples.async_service
python -m examples.async_task_processor
python -m examples.valkey_async_service   # requires valkey-glide

License

MIT

Download files

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

Source Distribution

scietex_service-3.2.0.tar.gz (55.4 kB view details)

Uploaded Source

Built Distribution

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

scietex_service-3.2.0-py3-none-any.whl (48.5 kB view details)

Uploaded Python 3

File details

Details for the file scietex_service-3.2.0.tar.gz.

File metadata

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

File hashes

Hashes for scietex_service-3.2.0.tar.gz
Algorithm Hash digest
SHA256 2b9a3123553f37223c2557f59e2087c16e6f994a84b4a99662f0dbbf2b0e758c
MD5 b541be91e0c1e37d1b9bbbf8fa7a0474
BLAKE2b-256 55440e2f000751443ea08460560bab692402d77987ff737ef0c94e06daf4ab5d

See more details on using hashes here.

Provenance

The following attestation bundles were made for scietex_service-3.2.0.tar.gz:

Publisher: python-publish.yml on bond-anton/scietex.service

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

File details

Details for the file scietex_service-3.2.0-py3-none-any.whl.

File metadata

File hashes

Hashes for scietex_service-3.2.0-py3-none-any.whl
Algorithm Hash digest
SHA256 2ffacfa36a109161d69a9203585841b3077969bc8ca2d62a75e96ce25267a575
MD5 3aa5deb9f00b42be3ea77b58c5c7cd74
BLAKE2b-256 09ef5ded937785c0f885f33dfefa7f3188ed99a003fafc4115c01c050be821fb

See more details on using hashes here.

Provenance

The following attestation bundles were made for scietex_service-3.2.0-py3-none-any.whl:

Publisher: python-publish.yml on bond-anton/scietex.service

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

Release history Release notifications | RSS feed

4.0.0

2 files

This release

3.2.0 This release

2 files

3.1.0

2 files

3.0.0

2 files

2.0.0

2 files

1.0.7

2 files

1.0.6

2 files

1.0.5

2 files

1.0.4

2 files

1.0.3

2 files

1.0.2

2 files

1.0.1

2 files

1.0.0

2 files

0.2.1

2 files

0.2.0

2 files

0.1.6

2 files

0.1.5

2 files

0.1.4

2 files

0.1.3

2 files

0.1.2

2 files

0.1.1

2 files

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