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.
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
Manager Lifecycle
Managers are async methods decorated with @Manager. The worker
discovers them via the class MRO and runs each as an asyncio.Task:
- Start — Manager loop runs the decorated method in a
while Trueloop until cancelled. - Error — On any exception (except
CancelledError), the error is recorded and the manager is automatically restarted, up tomanager_max_retriesconsecutive failures (default 5), after which the manager gives up and ends in the terminalFAILEDstate (observable viaworker.failed_managers; the watchdog logs CRITICAL but does not auto-shutdown). - Stop — On shutdown, managers are cancelled and their optional
cleanupcallbacks are invoked.
Task Handler System
See the Task Handler docs for the full handler lifecycle, schema details, and best practices.
- 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'ssupported_tasksdeclaration, not by a registration key. An optional keyword-onlyname(processor.add_task_handler(HandlerClass, name="...")) lets multiple instances of one class coexist under distinct keys. Arbitrary**handler_kwargsare also forwarded to the handler constructor on every instantiation, enabling stateful handlers — seeexamples/stateful_handler.py. - Declare support:
Handler.supported_tasksproperty must return a list of task type strings this handler can process. - Dispatch: When a task arrives, the processor calls
handler.supports(task_type)on each registered handler. The first handler returningTruereceives the task. - Initialize:
handler.start()callshandler.initialize()and setshandler.is_readyto the returned value, so it isTrueonly ifinitialize()returnedTrue. - Handle:
await handler.handle(task_data)returns aTaskResultwithstatus("success"/"error"), optionalerrormessage, and optionalpayload. - Timeout: Tasks exceeding their
timeout(default 3s) are canceled and either re-queued or discarded perTaskTimeout.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, partial |
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 |
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:
conf_dirargument (if provided and is a directory)SCIETEX_CONFIG_DIRenvironment variable$XDG_CONFIG_HOME/scietex/~/.config/scietex//etc/scietex//usr/local/etc/scietex/./config/(current working directory)~/.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 |
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, 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) |
TaskData |
Task payload schema |
TaskResult |
Task result schema |
TaskTimeout |
Timeout configuration schema |
TaskTracker |
Internal task tracker 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 |
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
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
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file scietex_service-4.0.0.tar.gz.
File metadata
- Download URL: scietex_service-4.0.0.tar.gz
- Upload date:
- Size: 78.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
079f87724d453e3b7c1cb645be0fed1cd635dc4b4100a6d713b98cc48c757fc9
|
|
| MD5 |
abd7dadbf250d888cf619bbcef7bcca8
|
|
| BLAKE2b-256 |
6b96c70098d33a7ad7a3832e27920d0596c6d682f2e10f1be562c2a6226cfe8f
|
Provenance
The following attestation bundles were made for scietex_service-4.0.0.tar.gz:
Publisher:
python-publish.yml on bond-anton/scietex.service
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
scietex_service-4.0.0.tar.gz -
Subject digest:
079f87724d453e3b7c1cb645be0fed1cd635dc4b4100a6d713b98cc48c757fc9 - Sigstore transparency entry: 2771993263
- Sigstore integration time:
-
Permalink:
bond-anton/scietex.service@2d046e49b560d89f7b5d4d0669a4b934325f49f4 -
Branch / Tag:
refs/tags/v4.0.0 - Owner: https://github.com/bond-anton
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
python-publish.yml@2d046e49b560d89f7b5d4d0669a4b934325f49f4 -
Trigger Event:
release
-
Statement type:
File details
Details for the file scietex_service-4.0.0-py3-none-any.whl.
File metadata
- Download URL: scietex_service-4.0.0-py3-none-any.whl
- Upload date:
- Size: 59.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
234bce776e372c1e11e8e017cd5ce71e6f3fc05ae652ee9da75caa44ea66c9c5
|
|
| MD5 |
810e17750f24d224d92a86a35ac52172
|
|
| BLAKE2b-256 |
12b277ee1365301c02847f13c4aa6ff4a7ad935ebd6484a3530e4cd5a4161b42
|
Provenance
The following attestation bundles were made for scietex_service-4.0.0-py3-none-any.whl:
Publisher:
python-publish.yml on bond-anton/scietex.service
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
scietex_service-4.0.0-py3-none-any.whl -
Subject digest:
234bce776e372c1e11e8e017cd5ce71e6f3fc05ae652ee9da75caa44ea66c9c5 - Sigstore transparency entry: 2771993283
- Sigstore integration time:
-
Permalink:
bond-anton/scietex.service@2d046e49b560d89f7b5d4d0669a4b934325f49f4 -
Branch / Tag:
refs/tags/v4.0.0 - Owner: https://github.com/bond-anton
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
python-publish.yml@2d046e49b560d89f7b5d4d0669a4b934325f49f4 -
Trigger Event:
release
-
Statement type: