Metadata-Version: 2.4
Name: pgfsm-async-worker-sdk
Version: 0.2.0
Summary: Python worker SDK for the pgfsm Activity Gateway: registers actors and serves invocations over the sidecar gRPC stream.
Author: Niraj Kashyap
Author-email: Niraj Kashyap <niraj.38.re@gmail.com>
License-Expression: Apache-2.0
License-File: LICENSE
Requires-Dist: pgfsm-proto-codegen>=0.1.1,<0.2
Requires-Python: >=3.10
Project-URL: Repository, https://github.com/pgfsm/fsm
Project-URL: Source, https://github.com/pgfsm/fsm/tree/main/packages/fsm-async-worker-sdk-python
Description-Content-Type: text/markdown

# pgfsm-async-worker-sdk

Python worker SDK for the pgfsm Activity Gateway. A worker process built on it
connects to the gateway's sidecar Unix socket, registers a set of actors, and
serves the invocations the gateway routes to them over the
`pgfsm.sidecargateway.v1.SidecarGatewayService` gRPC stream (stubs from
[`pgfsm-proto-codegen`](https://pypi.org/project/pgfsm-proto-codegen/)).

It never opens a database connection — that stays in the gateway.

Python counterpart of the TypeScript
[`@pgfsm/async-worker-sdk`](https://www.npmjs.com/package/@pgfsm/async-worker-sdk).

## Usage

You normally don't write against this package directly. `@pgfsm/compiler`'s
`generate-async-logic` writes a small `run_async_worker.py` plus a
`pyproject.toml` that pins this package:

```python
# async-worker/python/run_async_worker.py (generated)
import logging
import sys

from pgfsm.async_worker_sdk import run_actor_worker_cli
from python_actors_registry_generated import ACTOR_REGISTRATIONS

if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    raise SystemExit(run_actor_worker_cli(ACTOR_REGISTRATIONS, sys.argv[1:]))
```

Run it from that directory with [uv](https://docs.astral.sh/uv/):

```bash
uv run run_async_worker.py list    # print the actors in the registry, no gateway needed
uv run run_async_worker.py start   # connect to the gateway and serve invocations
```

or with pip: `python3 -m pip install pgfsm-async-worker-sdk`, then
`python3 run_async_worker.py start`.

### Options

```
-g, --gateway-socket <path>   Sidecar socket to connect to (default: /tmp/pgfsm-activity-gateway-workers.sock)
-i, --worker-id <id>          Stable worker identity (default: python-<random>)
    --heartbeat-ms <ms>       Heartbeat interval (default: 5000)
    --reconnect-initial-delay-ms <ms>
                              First reconnect backoff step (default: 250)
    --reconnect-max-delay-ms <ms>
                              Reconnect backoff cap (default: 30000)
    --reconnect-max-attempts <n>
                              Exit after n consecutive failed attempts (default: 0 = retry forever)
-h, --help                    Show this help message
```

SIGINT/SIGTERM stop the worker gracefully (it unregisters from the gateway).

`start` doesn't need the gateway to be up first: it retries the connection with
exponential backoff (full jitter, 250 ms doubling up to 30 s), and if a session
drops (e.g. the gateway restarts) it reconnects and re-registers on its own. A
session only resets the backoff once it has stayed up for 10 s, so a gateway
that accepts and immediately drops still gets backed off from. `start` ends with
exit code 1 only on what reconnecting can't fix: an explicit registration
rejection; a gRPC `UNAUTHENTICATED`, `PERMISSION_DENIED`, `UNIMPLEMENTED` or
`INVALID_ARGUMENT` (a misconfiguration, so it fails fast instead of retrying);
or `--reconnect-max-attempts` consecutive failed attempts. An invoke result that
can't be sent because its session ended is logged and dropped; the gateway has
already failed that invoke.

## API

- `run_actor_worker_cli(registrations, args, invocation=None) -> int` — the
  `list`/`start` CLI. Returns the process exit code instead of exiting.
- `ActorWorker(worker_id, gateway_socket_path, registrations, heartbeat_ms=5000,
  reconnect_initial_delay_ms=250, reconnect_max_delay_ms=30000,
  reconnect_max_attempts=0)`
  — `run()` registers every actor and serves invocations until `stop()` is
  called, reconnecting and re-registering whenever a session drops. It raises
  `RegistrationRejectedError` on an explicit rejection, or `ConnectionError`
  after `reconnect_max_attempts` consecutive failed attempts.

A registration is a dict:

```python
{
    "parent_fsm_name": "creditCheck",
    "parent_fsm_version": "v01",
    "async_operation_type": "internalAsyncOperation",
    "async_operation_name": "checkBureau",
    "async_operation_version": "v01",
    "async_operation_language": "python",
    "handler": check_bureau,  # (input) -> JSON-serializable output; may be async
}
```

A handler that raises is reported to the gateway as an `INTERNAL` invoke error;
an invoke for an unregistered actor is reported as `NOT_FOUND`.

Logging goes through the standard `logging` module (`pgfsm.async_worker_sdk`
loggers); the library never configures logging itself.

## Releasing (maintainers)

Released from the [pgfsm/fsm](https://github.com/pgfsm/fsm) monorepo by
`.github/workflows/pypi-publish.yml`. Pushing an
`async-worker-sdk-py-v<version>` tag publishes to PyPI with trusted publishing
(no API token). This package releases independently of
[`pgfsm-proto-codegen`](https://pypi.org/project/pgfsm-proto-codegen/).

1. **Pick the version.** While below 1.0: breaking API change → minor (`0.1.0` →
   `0.2.0`); new backward-compatible features → minor; fixes only → patch
   (`0.1.0` → `0.1.1`).
2. **Bump it in a PR.** From `packages/fsm-async-worker-sdk-python`:

   ```bash
   uv version --bump minor     # or: --bump patch, or an exact version: uv version 0.2.0
   ```

   This updates `pyproject.toml` and `uv.lock` together; commit both. If the new
   version needs a newer `pgfsm-proto-codegen`, release that first, then raise
   the `pgfsm-proto-codegen>=` pin here.
3. **Tag the merge commit and push the tag.** The tag must be exactly
   `async-worker-sdk-py-v` + `uv version --short`:

   ```bash
   git fetch origin
   git tag async-worker-sdk-py-v0.2.0 origin/main
   git push origin async-worker-sdk-py-v0.2.0
   ```

4. **Check the release.** `gh run list --workflow pypi-publish.yml` shows the
   run, which checks the tag against `pyproject.toml`, runs the tests, builds,
   and uploads. Then confirm https://pypi.org/project/pgfsm-async-worker-sdk/
   lists the version and it installs:
   `pip install pgfsm-async-worker-sdk==0.2.0`.

For prereleases, uv writes the version in PEP 440 form
(`uv version 0.2.0-alpha.0` stores `0.2.0a0`), so tag
`async-worker-sdk-py-v0.2.0a0`. pip only installs a prerelease if asked
explicitly (`--pre` or an exact `==` version).

A published version can never be re-uploaded. If a bad version ships, release
the next patch and yank the bad one on pypi.org. Yanked versions stay
installable when pinned exactly, but resolvers skip them otherwise.

## License

Apache-2.0
