Metadata-Version: 2.4
Name: crowddrop-pubsub-sdk
Version: 0.2.0
Summary: Framework-agnostic Google Cloud Pub/Sub step-event publisher/subscriber for embodied AI agents.
Project-URL: Repository, https://github.com/crowddrop-ai/crowddrop_ai_agents
Requires-Python: >=3.9
Description-Content-Type: text/markdown
Requires-Dist: google-cloud-pubsub<3.0.0,>=2.21.0
Requires-Dist: google-api-core<3.0.0,>=2.17.1
Provides-Extra: langchain
Requires-Dist: langchain-core>=0.3.0; extra == "langchain"

# pubsub_sdk

A small, framework-agnostic wrapper around Google Cloud Pub/Sub for publishing and
consuming structured step/progress events, and for direct agent-to-agent
instruction/reply exchanges (chat). It has no dependency on LangChain, on any
specific agent framework, or on any particular storage backend — any Python project
(an agent, a backend service, a future embodied-AI-agent codebase) can import it
directly.

## Install

```bash
pip install ./pubsub_sdk          # core publisher/subscriber only
pip install ./pubsub_sdk[langchain]  # + the optional LangChain callback adapter
```

Published to public PyPI as `crowddrop-pubsub-sdk` (the import name stays
`pubsub_sdk` either way):

```bash
pip install crowddrop-pubsub-sdk
pip install crowddrop-pubsub-sdk[langchain]
```

## Configuration (environment variables)

| Variable | Default | Purpose |
|---|---|---|
| `GCP_PROJECT_ID` | `test` | GCP project the client talks to. |
| `PUBSUB_SDK_TOPIC_ID` | `agent-step-events` | Topic events are published to. |
| `PUBSUB_SDK_SUBSCRIPTION_ID` | `agent-step-events-firestore-persister` | Subscription a consumer pulls from. |
| `PUBSUB_SDK_CHAT_TOPIC_ID` | `agent-chat-messages` | Fixed topic for agent-to-agent chat (see below). |
| `PUBSUB_SDK_ENABLED` | `true` | Kill switch — set to `false` to disable publishing/subscribing entirely. |
| `PUBSUB_EMULATOR_HOST` | unset | Standard Google client env var; set to point at a local Pub/Sub emulator instead of real GCP. |

Credentials are picked up automatically via `GOOGLE_APPLICATION_CREDENTIALS`, exactly
like every other `google-cloud-*` client — no SDK-specific credential handling.

## Publishing (any consumer)

```python
from pubsub_sdk import StepEventPublisher

publisher = StepEventPublisher()
publisher.publish_step(
    activity_id="task-123",
    agent_id="MissionCommander",
    event_type="tool_start",
    message="Invoking tool: delegate_to_agent",
    status="in_progress",
    metadata={"tool_name": "delegate_to_agent"},
)
```

`publish_step()` never raises — if Pub/Sub is unavailable it logs the failure and
returns `False`. Check `publisher.is_available()` if you need to know the state
up front.

## Subscribing (any consumer)

```python
from pubsub_sdk import StepEventSubscriber, StepEvent

def handle(event: StepEvent) -> None:
    ...  # e.g. persist it somewhere

subscriber = StepEventSubscriber()
subscriber.pull_forever(handle)  # blocks; acks on success, nacks (redelivers) on exception
```

## Agent-to-agent chat (ChatMessage)

A direct instruction/reply exchange between two agent instances (e.g. a
coordinator persona delegating to a field agent), as an alternative transport to
a live MCP/SSE connection. Both directions of one exchange ride the same fixed
topic (`agent-chat-messages`) and are distinguished by `role`
(`ChatMessageRole.INSTRUCTION` / `.REPLY`), correlated by `message_id`.

```python
from pubsub_sdk import ChatMessage, ChatMessageRole, get_default_chat_publisher, ChatMessageSubscriber

# Sending side: publish an instruction addressed to another agent instance's
# routing identity (its session_id - never a shared persona name; two ephemeral
# instances of the same persona would otherwise collide on the same subscription).
publisher = get_default_chat_publisher()
publisher.publish(ChatMessage(
    message_id="m1", from_agent="MissionCommander__task-42", to_agent="PreciseHumanoid",
    role=ChatMessageRole.INSTRUCTION, content="what is your current battery level?",
    activity_id="task-42",
))

# Receiving side: each participating instance self-provisions its own filtered
# subscription (against real GCP, not just the emulator - a per-instance
# subscription's identity only exists at runtime and can't be pre-provisioned).
subscriber = ChatMessageSubscriber(routing_identity="PreciseHumanoid")
subscriber.pull_forever(lambda msg: ...)  # blocks; ack on success, nack (redeliver) on exception
```

`GenericSubscriber.pull_forever` also accepts a `stop()` call from another thread
to cancel an in-progress streaming pull cleanly.

**Security note:** `agent-chat-messages` is one shared topic with no
per-agent access control beyond GCP IAM on the topic/subscription resources
themselves - any credential with publish/subscribe rights can spoof
`from_agent` or read another agent's traffic via an unfiltered subscription.
Fine within one trust domain (this app's own fleet today); see
`docs/pubsub_chat_access_control.md` before granting a genuinely independent/
external party a credential onto this topic.

## LangChain integration (optional extra)

```python
from pubsub_sdk.langchain_callback import StepEventCallbackHandler

callback = StepEventCallbackHandler(activity_id="task-123", agent_id="MissionCommander")
callback.log_mission_start()
# pass `callback` into a LangChain AgentExecutor's `config={"callbacks": [callback]}` —
# every LLM/tool/chain/agent lifecycle event then publishes automatically.
```

`langchain_callback` is the only module in this package that imports `langchain` —
a consumer who only needs `StepEventPublisher`/`StepEventSubscriber` never pulls
that dependency in.

## Releasing (publishing a new version to PyPI)

Releases are tag-triggered via `.github/workflows/publish-pubsub-sdk.yml`,
using PyPI's **Trusted Publishing** (OIDC) — no API token is stored as a
GitHub secret.

1. Bump `version` in `pubsub_sdk/pyproject.toml`.
2. Commit that change (on a branch, via the normal PR flow).
3. Once merged, tag the merge commit and push the tag:
   ```bash
   git tag pubsub-sdk-v<version>   # e.g. pubsub-sdk-v0.2.1
   git push origin pubsub-sdk-v<version>
   ```
   The tag push is what fires the workflow — it builds `pubsub_sdk/` and
   uploads it to [pypi.org/project/crowddrop-pubsub-sdk](https://pypi.org/project/crowddrop-pubsub-sdk/).
   No other trigger publishes this package.

**One-time setup, not yet done as of this writing — needed before the first
tag push, and again only if this ever moves to a different PyPI
account/org:**
- Register a **pending publisher** for `crowddrop-pubsub-sdk` at
  https://pypi.org/manage/account/publishing/ — this can be done before the
  PyPI project exists, so it covers the *first-ever* release too, not just
  subsequent ones. Fill in: PyPI project name `crowddrop-pubsub-sdk`, repo
  owner `crowddrop-ai`, repo name `crowddrop_ai_agents`, workflow filename
  `publish-pubsub-sdk.yml`, environment name `pypi`. Requires a PyPI account
  with 2FA enabled — no API token to generate or store.
- The `pypi` GitHub Environment referenced by the workflow is created
  automatically the first time the workflow runs against it; create it
  manually in this repo's Settings → Environments beforehand only if you
  want a required-reviewer protection rule (so a tag push pauses for human
  approval before it actually publishes).
- `../scripts/publish_python_packages.sh` (manual `build` + `twine upload`)
  is kept as a fallback/local-dry-run tool only — with a pending publisher
  registered, it's no longer needed even for the first release.

Versioning is manual — nothing cross-checks the tag against
`pyproject.toml`'s `version`. Bump the file first, commit, *then* tag that
exact commit; tagging a commit whose `pyproject.toml` still has an
already-published version will fail the upload (PyPI rejects re-uploading an
existing version).

Note: `crowddrop-sdk`'s `cloud-brain` extra depends on `crowddrop-pubsub-sdk`
(see `../crowddrop_sdk/pyproject.toml`) — release this package first if both
need a coordinated update.
