A library for creating custom beVault or AWS Step Functions workers and data stores
Project description
README.md
beVault Workers Library
A Python library for creating beVault's States workers with pluggable data stores. This library provides everything for building scalable, distributed workers that can process Step Functions activities with support for multiple data backends.
Table of Contents
Features
- 🚀 Easy Worker Creation: Simple base class for implementing custom workers
- 🔌 Pluggable Data Stores: Support for PostgreSQL, S3, SFTP, and custom stores. Plus, the possibility to create your own data stores
- ⚙️ Configurable: JSON-based configuration for stores
- 🔄 Auto-Discovery: Automatic worker discovery and registration
- 📊 Built-in Logging: Comprehensive logging with configurable outputs
- 🛡️ Robust Error Handling: Graceful shutdown and error recovery
- 📈 Scalable: Multi-process architecture with configurable concurrency
Installation
From PyPI
Install the latest published package from PyPI:
pip install bevault_workers
For development
Clone this project and install the library in editable mode with development extras:
git clone https://github.com/depfac/bevault-workers-python.git
cd python-worker-framework
python -m venv .venv
.venv\Scripts\Activate.ps1
python -m pip install --upgrade pip
pip install -e ".[dev]"
Run the local entrypoint to test workers from dev_workers and validate store behavior:
python main.py
After that, you can import the library in your project.
Warning The library is named bevault_workers (PyPI package bevault_workers)
import bevault_workers
Project Structure
Here's a recommended structure for a basic project using the library:
python-worker-framework/
├── main.py # Dev entrypoint (loads dev_workers)
├── config.json # Local store configuration
├── logging_config.json # Logging configuration (optional)
├── dev_workers/ # Workers used for local testing
├── src/
│ └── bevault_workers/ # Published library package
│ ├── workers/
│ ├── stores/
│ └── utils/
├── tests/ # Test suite
├── pyproject.toml # Package metadata and dependencies
└── README.md
main.py example
from bevault_workers import WorkerManager
def main():
"""Main entry point for the worker application"""
try:
# Initialize the worker manager with configuration
manager = WorkerManager(
config_path="config.json",
workers_module="workers" # Discover workers from the 'workers' module
)
# Start the worker manager (this will block until stopped)
manager.start()
except KeyboardInterrupt:
print("Application interrupted by user")
except Exception as e:
print(f"Application error: {e}")
raise
if __name__ == "__main__":
main()
Add your custom workers
CFR here to see how to create a custom worker
.env file
Create a .env file in the root directory with the following variables:
stepFunctions__authenticationKey=YOUR_AUTH_KEY
stepFunctions__authenticationSecret=YOUR_AUTH_SECRET
stepFunctions__awsRegion=us-east-1
stepFunctions__DefaultHeartbeatDelay=5
stepFunctions__DefaultMaxConcurrency=3
stepFunctions__EnvironmentName=python
stepFunctions__roleArn=YOUR_ROLE_ARN
stepFunctions__serviceUrl=df2-states
stepFunctions__enableStatesStoreSync=false
stepFunctions__statesStoreBaseUrl=
stepFunctions__statesPollTimeoutSeconds=70
stepFunctions__statesStatusHeartbeatSeconds=60
stepFunctions__statesRequestTimeoutSeconds=15
Environment Variables Details
| Variable | Description |
|---|---|
stepFunctions__authenticationKey |
Authentication key for AWS Step Functions |
stepFunctions__authenticationSecret |
Authentication secret for AWS Step Functions |
stepFunctions__awsRegion |
AWS region (default: us-east-1) |
stepFunctions__DefaultHeartbeatDelay |
Heartbeat delay in seconds (default: 5) |
stepFunctions__DefaultMaxConcurrency |
Maximum concurrency (default: 3) |
stepFunctions__EnvironmentName |
Environment name (default: python) |
stepFunctions__roleArn |
AWS IAM role ARN |
stepFunctions__serviceUrl |
Service URL for Step Functions (default: df2-states) |
stepFunctions__enableStatesStoreSync |
Enable dFakto States store synchronization (default: false) |
stepFunctions__statesStoreBaseUrl |
Optional override for the store sync API base URL |
stepFunctions__statesPollTimeoutSeconds |
Long-poll timeout in seconds (default: 70) |
stepFunctions__statesStatusHeartbeatSeconds |
Store status heartbeat period in seconds (default: 60) |
stepFunctions__statesRequestTimeoutSeconds |
Non-long-poll request timeout in seconds (default: 15) |
States store synchronization
When stepFunctions__enableStatesStoreSync=true, the worker manager starts a background synchronization service that:
- Uses AWS-signed Step Functions extension calls (
X-Amz-Target) for dFakto store APIs. - Sends local stores to
DfaktoStatesSyncStores. - Long-polls for states-defined stores and refreshes
StoreRegistryin memory. - Long-polls force-check requests with
DfaktoStatesGetStoreForceCheckRequests. - Posts periodic and on-demand store status updates to
DfaktoStatesPostStoreStatus.
If a local and a states-defined store share the same name, the states store is registered with the internal prefix states:: (for example states::myStore) so both instances remain addressable.
config.json file
You can configure stores in two ways:
- Locally in
config.json(recommended for development and local runs) - In beVault States by enabling store synchronization (
stepFunctions__enableStatesStoreSync=true)
For full store configuration details and supported store types, refer to the official beVault documentation: Stores reference.
If you use local configuration, create a config.json file in the root directory. Example:
[
{
"Name": "postgresqlStore",
"Type": "postgresql",
"Config": {
"host": "",
"port": "5432",
"user": "",
"password": "",
"dbname": ""
}
},
{
"Name": "sftpStore",
"Type": "sftp",
"Config": {
"Host": "sftp.example.com",
"Port": 22,
"Username": "myuser",
"Password": "secret",
"Prefix": "/uploads/"
}
}
]
For key-based authentication with SFTP, use KeyFilename instead of Password:
"KeyFilename": "/path/to/private_key"
States/.NET compatibility for existing stores
The stores keep backward compatibility with the legacy Python config keys and also accept States/.NET-style keys for:
postgresqls3sftp
For database stores, when both parameter fields and a connection string are present:
- if
connectionString(orConnectionString) is non-empty, it is used; - if it is empty or null, parameter fields are used.
Example States-style PostgreSQL config:
{
"host": "localhost",
"port": 5432,
"database": "northwind__mig__pgsql_migration",
"username": "metavault",
"password": "xxx",
"connectionString": ""
}
logging_config.json file
Create a logging_config.json file to configure the logging behavior:
{
"logging": {
"minimumLevel": {
"default": "Information",
"override": {
"urllib3": "Warning",
"boto3": "Warning",
"botocore": "Warning",
"paramiko": "Warning",
"requests": "Warning"
}
},
"writeTo": [
{
"name": "Console",
"args": {
"outputTemplate": "[{Timestamp:yyyy-MM-dd HH:mm:ss.fff} {Level}] {Message:lj} <s:{SourceContext}>{NewLine}{Exception}"
}
},
{
"name": "File",
"args": {
"path": "/var/logs/bevault_workers/bevault_workers.log",
"rollOnFileSizeLimit": true,
"fileSizeLimitBytes": 10000000,
"retainedFileCountLimit": 10
}
}
]
}
}
Workers
Your custom workers will be automatically loaded from the package you specified while creating a new WorkerManager in your main.py. If not specified, the package "workers" is scanned by default.
All your custom workers should extend the BaseWorker class from this bevault_workers library. Here is a simple example:
from bevault_workers import BaseWorker
class DataProcessorWorker(BaseWorker):
name = "my_custom_worker"
def handle(self, input_data):
"""Process the input data"""
try:
# Your business logic here
processed_data = self.process_data(input_data)
return {
"status": "success",
"result": processed_data
}
except Exception as e:
return {
"status": "error",
"error_message": str(e)
}
def process_data(self, data):
# Your custom processing logic
return {"processed": True, "data": data}
Stores
Stores are data endpoints that you can use to extract from or send data to with your workers. The make a store available in your workers, you need to configure them in your config.json file.
Here is an example of worker that uses a postgresql store:
from bevault_workers import BaseWorker
class DataProcessorWorker(BaseWorker):
name = "my_custom_worker"
def handle(self, input_data):
"""
extract version number of PostgreSQL
Args:
dbStore: The PostgreSQL store
Returns:
Version number of postgresql
"""
try:
postgreStore = StoreRegistry.get(input_data["dbStore"])
postgreStore.connect()
result = postgreStore.execute(query = f"SELECT VERSION()")
return {
"status": "success",
"result": result
}
except Exception as e:
return {
"status": "error",
"error_message": str(e)
}
Create your own Stores
You can create your own Store by creating a file with a class Store that extends the DbStore of FileStore class. Here is an example of custom implementation of the postgresql store: custom_postgresql.py
import psycopg
from bevault_workers import DbStore
class Store(DbStore):
def __init__(self, config):
self.config = config
self.connection = None
def connect(self):
self.connection = psycopg.connect(**self.config)
def execute(self, query, params=None):
with self.connection as conn:
with conn.cursor() as cur:
cur.execute(query, params)
# Check if the query might return results (has a description)
if cur.description:
results = cur.fetchall()
# Return results only if there are any
return results if results else None
# For non-SELECT queries, return affected row count
return cur.rowcount
Here is an example of a worker that uses an SFTP store to read and write files:
from bevault_workers import BaseWorker
from bevault_workers.stores.store_registry import StoreRegistry
class FileHandlerWorker(BaseWorker):
name = "file_handler"
def handle(self, input_data):
sftp_store = StoreRegistry.get(input_data["fileStore"])
sftp_store.connect()
token = sftp_store.createFileToken("example.txt")
sftp_store.openWrite(token, b"Hello from worker")
content = sftp_store.openRead(token)
return {"status": "success", "content": content.decode()}
Note that you can override the implementation of a built-in store by placing a module in your project whose basename matches the bare Type (for example, Type: "postgresql" resolves to stores.postgresql, then bevault_workers.stores.postgresql). The store Name in JSON is only the instance identifier passed to StoreRegistry.get(...) in workers, not the implementation selector.
To keep neutral filenames (for example stores/custom_db_store.py), set Type to a fully qualified path such as stores.custom_db_store or stores.custom_db_store:Store instead of a bare name like postgresql.
Project details
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 bevault_workers-0.1.2.tar.gz.
File metadata
- Download URL: bevault_workers-0.1.2.tar.gz
- Upload date:
- Size: 26.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
519cedd8de03c0d92c66425f3a7a400056b7daf888855635bc6f9b6f286a60ac
|
|
| MD5 |
cdc28dc070517eef7e6db85f6dd0cd39
|
|
| BLAKE2b-256 |
4cd4eb030bd255f92f9d3a8d291d68f6ab465025630675c463248ec4277cb898
|
Provenance
The following attestation bundles were made for bevault_workers-0.1.2.tar.gz:
Publisher:
publish-pypi.yml on depfac/bevault-workers-python
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
bevault_workers-0.1.2.tar.gz -
Subject digest:
519cedd8de03c0d92c66425f3a7a400056b7daf888855635bc6f9b6f286a60ac - Sigstore transparency entry: 1205758080
- Sigstore integration time:
-
Permalink:
depfac/bevault-workers-python@632b64c7f2c2b0effc5db83103807425a4f429f2 -
Branch / Tag:
refs/tags/v0.1.2 - Owner: https://github.com/depfac
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish-pypi.yml@632b64c7f2c2b0effc5db83103807425a4f429f2 -
Trigger Event:
push
-
Statement type:
File details
Details for the file bevault_workers-0.1.2-py3-none-any.whl.
File metadata
- Download URL: bevault_workers-0.1.2-py3-none-any.whl
- Upload date:
- Size: 34.3 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
e7b8a98b0587f2f2eb2181f5bb0304907aa54466e64aff7030de3b6e422d661b
|
|
| MD5 |
e06ce6d2b4648710a70b85903544f47d
|
|
| BLAKE2b-256 |
b2d243fdba119e51e08ea83bf937a38a97caede72ea32d3d05970b9a2a048516
|
Provenance
The following attestation bundles were made for bevault_workers-0.1.2-py3-none-any.whl:
Publisher:
publish-pypi.yml on depfac/bevault-workers-python
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
bevault_workers-0.1.2-py3-none-any.whl -
Subject digest:
e7b8a98b0587f2f2eb2181f5bb0304907aa54466e64aff7030de3b6e422d661b - Sigstore transparency entry: 1205758081
- Sigstore integration time:
-
Permalink:
depfac/bevault-workers-python@632b64c7f2c2b0effc5db83103807425a4f429f2 -
Branch / Tag:
refs/tags/v0.1.2 - Owner: https://github.com/depfac
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish-pypi.yml@632b64c7f2c2b0effc5db83103807425a4f429f2 -
Trigger Event:
push
-
Statement type: