Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

Avtomatika Worker SDK

Official SDK for building workers compatible with the Avtomatika orchestrator. It automates low-level tasks: polling, heartbeats, S3 payload management, and graceful shutdown.

🚀 Key Features

  • Language: Python 3.11+
  • Protocol: Based on RXON (Reverse Axon Protocol) for Hierarchical Logic Networks (Holarchy).
  • Communication Model:
    • PULL: Workers poll tasks from orchestrators (works behind NAT/Firewall).
    • WebSocket: Real-time command channel (cancellation, custom commands).
  • Zero Trust Security:
    • Mandatory HMAC SHA256 signing for all messages using WORKER_TOKEN.
    • Identity Chain and Origin Worker ID support for provenance tracking.
    • Replay protection with timestamp validation.
  • Traffic & Performance Optimization:
    • Telemetry Throttling (Heartbeat Deadband): Telemetry (CPU/RAM/GPU) is only sent when value changes by $>5%$ or after 60s force interval, drastically saving bandwidth.
    • ETag-Based Blob Caching: Heavy assets (e.g. AI model weights) are downloaded only once, cached locally, and symlinked to tasks' workspace.
    • Async Results Uploader: Task results are sent using a non-blocking asyncio.Queue with retry policy, rate limiting wait times, and backoff, instantly freeing the worker for the next task.
    • 3-Tier Skills: Supported (catalog), Available (dynamic limits), and Hot (cached).
    • Stable Hashing: Sends full skill catalog only when changed, using skills_hash for light heartbeats.
  • S3 Streaming: High-performance data transfer using obstore. No OOM on large files.
  • AI-Agent Support: Supports Chain of Thought and Tool Use via OrchestratorClient dependency injection for subtask delegation.
  • Hardware Awareness: Built-in monitoring for CPU, RAM, and NVIDIA GPUs (via psutil and GPUtil).
  • Observability:
    • Built-in support for OpenTelemetry (traces and metrics).
    • Automatic Trace Context Propagation: Workers extract trace_id from tasks and inject it into events (including progress), ensuring end-to-end visibility in Jaeger/Honeycomb.
    • Automatic metrics export via OTLP (Push model) when OTEL_EXPORTER_OTLP_ENDPOINT is configured.

🛡 Resilience & Connectivity

  • Independent Managers: Connection to each orchestrator is managed by a separate background task. One server failure or rate limit doesn't affect others.
  • Smart Backoff: Unified exponential backoff for registration, polling, and heartbeats.
  • Rate Limit Protection: Full support for Retry-After (seconds or HTTP-date). Implements a mandatory 30s safety floor for 429 errors without Retry-After to prevent Retry Storms.
  • Heartbeat Debouncing: Throttles heartbeats to once every 2 seconds. Events are not lost but consolidated and sent after the cooldown period.
  • Infinite Retries: Workers never stop trying to register with an exponential delay.
  • Graceful Shutdown: Handles SIGTERM and SIGINT properly, waiting for active tasks to finish.

🛠 Installation

pip install avtomatika-worker[s3,pydantic]

For development:

pip install -e .[test,dev]

💻 Quick Start

from avtomatika_worker import Worker, TaskFiles, OrchestratorClient

worker = Worker()

@worker.skill("hello_world")
async def my_skill(params: dict, files: TaskFiles):
    """Simple skill that says hello."""
    return {"message": f"Hello, {params.get('name', 'World')}!"}

@worker.skill("ai_agent_reasoning")
async def agent_skill(params: dict, orchestrator_client: OrchestratorClient):
    """AI agent skill delegating a subtask (tool use) via OrchestratorClient."""
    search_result = await orchestrator_client.call_skill("web_search", {"query": params["search_query"]})
    return {"result": f"Based on web search: {search_result['data']}"}

@worker.on_command("reboot")
async def handle_reboot(command):
    print("Rebooting worker...")

if __name__ == "__main__":
    worker.run()

⚙️ Configuration

Controlled via environment variables:

  • ORCHESTRATORS_CONFIG: JSON list of orchestrator configs (URLs, priorities, weights).
  • ORCHESTRATOR_URL: Simple fallback if only one orchestrator is used (default: http://localhost:8080).
  • WORKER_TOKEN: Secret for HMAC signing (Zero Trust).
  • S3_ENDPOINT_URL, S3_ACCESS_KEY, S3_SECRET_KEY, S3_DEFAULT_BUCKET: Storage settings for large payloads.
  • WORKER_BLOB_CACHE_DIR: Directory for caching S3 blobs (default: /tmp/avtomatika_cache).
  • WORKER_TELEMETRY_DEADBAND: Threshold (percent) for throttling telemetry updates (default: 5.0).
  • WORKER_TELEMETRY_FORCE_INTERVAL: Maximum time (seconds) to wait before sending telemetry even if unchanged (default: 60.0).
  • STRICT_EVENT_VALIDATION: (Default: True) Validates events against schemas before emitting.
  • LOG_LEVEL: Logging verbosity (DEBUG, INFO, WARNING, ERROR).
  • POLL_BACKOFF_INITIAL: Initial delay (seconds) after a 429 error or network failure (default: 1.0). Honors Retry-After header.
  • POLL_BACKOFF_MAX: Maximum backoff delay (seconds) (default: 60.0).
  • POLL_BACKOFF_FACTOR: Multiplier for exponential backoff (default: 2.0).
  • MAX_CONCURRENT_TASKS: Global limit for concurrent task execution.

📜 License

Mozilla Public License v. 2.0.

Download files

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

Source Distribution

avtomatika_worker-1.0b18.tar.gz (75.1 kB view details)

Uploaded Source

Built Distribution

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

avtomatika_worker-1.0b18-py3-none-any.whl (42.1 kB view details)

Uploaded Python 3

File details

Details for the file avtomatika_worker-1.0b18.tar.gz.

File metadata

  • Download URL: avtomatika_worker-1.0b18.tar.gz
  • Upload date:
  • Size: 75.1 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.14

File hashes

Hashes for avtomatika_worker-1.0b18.tar.gz
Algorithm Hash digest
SHA256 4f6e3e5eb296d377a8271079cf9cc44f319efd0dd16a1ce094a2103095123238
MD5 21bbfebcef1596cd7c4f382ff276d875
BLAKE2b-256 d3a0854c4184ac922ed13b638dd10ac060d4f17a1679c5bcfee744f4d682dbaf

See more details on using hashes here.

File details

Details for the file avtomatika_worker-1.0b18-py3-none-any.whl.

File metadata

File hashes

Hashes for avtomatika_worker-1.0b18-py3-none-any.whl
Algorithm Hash digest
SHA256 603bc7762e10a4a9ae341cfc4d43e902ead15823312bc6d39786dc716cc4d50f
MD5 32a0b1e1547bd8eff6e3f647bc7b321c
BLAKE2b-256 491d9a9a0d6cad2eb83f7b68c85e9209c1153d73ab93b306317a1d00d0bf713b

See more details on using hashes here.

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