Metadata-Version: 2.4
Name: pasarguard-node-bridge
Version: 0.9.0
Summary: python package to connect your project with PasarGuard node go
Project-URL: Homepage, https://github.com/PasarGuard/node_bridge_py
Project-URL: Repository, https://github.com/PasarGuard/node_bridge_py.git
License-File: LICENSE
Keywords: PasarGuard,PasarGuard API,PasarGuard python,PasarGuard-Node,PasarGuard-panel
Classifier: Programming Language :: Python :: 3.12
Requires-Python: >=3.12
Requires-Dist: aiohttp-socks>=0.11.0
Requires-Dist: aiohttp>=3.13.5
Requires-Dist: grpclib>=0.4.9
Requires-Dist: packaging>=26.2
Requires-Dist: protobuf>=6.33.6
Requires-Dist: python-socks[asyncio]>=2.8.1
Description-Content-Type: text/markdown

# PasarGuard Node Bridge (Python)

Async Python client for connecting to a [PasarGuard node](https://github.com/PasarGuard/node) over `gRPC` or `REST`.

This package provides:
- Strongly typed protobuf models (`service_pb2`)
- Unified node API for both transport types
- User sync helpers (single, batch, and chunked streaming)
- Health/version helpers
- On-demand log streaming
- Node maintenance endpoints (update core/node/geofiles)

## Installation

```bash
pip install pasarguard-node-bridge
```

## Requirements

- Python `>=3.12`
- A reachable PasarGuard node
- Node service port (`port`) for gRPC or protobuf-REST
- Node JSON API port (`api_port`) for maintenance endpoints
- Server CA certificate content (PEM string)
- API key (UUID string)

## Import

```python
import PasarGuardNodeBridge as Bridge
from PasarGuardNodeBridge.common import service_pb2 as service
```

## Create A Node Client

```python
node = Bridge.create_node(
    connection=Bridge.NodeType.grpc,  # Bridge.NodeType.grpc or Bridge.NodeType.rest
    address="127.0.0.1",
    port=2096,                         # gRPC or protobuf-REST port (based on connection)
    api_port=2097,                     # REST JSON API port (used internally for maintenance)
    server_ca=server_ca_pem_string,
    api_key="xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
    name="node-1",                     # optional
    extra={"region": "eu-1"},          # optional
    default_timeout=10,                # optional
    internal_timeout=15,               # optional
    proxy="socks5://user:pass@127.0.0.1:1080",  # optional
)
```

### `create_node(...)` Parameters

- `connection`: `Bridge.NodeType.grpc` or `Bridge.NodeType.rest`
- `address`: node host/IP
- `port`: node service port
- `api_port`: node REST JSON API port
- `server_ca`: PEM certificate content as string
- `api_key`: UUID string
- `name`: optional logger name
- `extra`: optional metadata dictionary
- `logger`: optional custom logger
- `default_timeout`: default timeout for public API methods
- `internal_timeout`: timeout used for internal sync/log operations
- `proxy`: optional upstream proxy URL for node traffic
- `max_message_size`: gRPC only, HTTP/2 window/message sizing

### Proxy Formats

- `socks5://127.0.0.1:1080`
- `socks5://user:pass@127.0.0.1:1080`
- `socks4://127.0.0.1:1080`
- `http://127.0.0.1:3128`
- `http://user:pass@127.0.0.1:3128`
- `https://user:pass@proxy.example.com:443`

### Connection Types

- `Bridge.NodeType.grpc`: gRPC transport via `grpclib`
- `Bridge.NodeType.rest`: protobuf-over-HTTP transport

## User/Proxy Builders

Use helpers for creating protobuf user/proxy payloads.

```python
user = Bridge.create_user(
    email="alice@example.com",
    proxies=Bridge.create_proxy(
        vmess_id="0d59268a-9847-4218-ae09-65308eb52e08",
        vless_id="0d59268a-9847-4218-ae09-65308eb52e08",
        vless_flow="",
        trojan_password="",
        shadowsocks_password="",
        shadowsocks_method="",
        wireguard_public_key="",
        wireguard_peer_ips=["10.10.0.2/32"],
    ),
    inbounds=["inbound-tag-1"],
)
```

## Start/Stop Lifecycle

You should `start()` before calling stats/sync/log methods.

```python
await node.start(
    config=config_json_string,
    backend_type=service.BackendType.XRAY,   # or service.BackendType.WIREGUARD
    users=[user],                             # optional initial user set
    keep_alive=30,                            # optional
    exclude_inbounds=[],                      # optional
    timeout=20,
)

info = await node.info()
print(info.node_version, info.core_version)

await node.stop()
```

## Method Examples

### 1. Queue-Based User Updates (recommended for frequent updates)

`update_user` and `update_users` enqueue users and a background worker handles retries and batching.

```python
await node.update_user(user)

more_users = [user1, user2, user3]
await node.update_users(more_users)
```

#### Shared Storage For Multiple Workers

By default, queued user updates are kept in a process-local in-memory store shared by node instances. This coordinates controllers in a single worker process when they use the same `node_id` (or the same service URL when `node_id` is omitted). For multi-process or multi-host deployments, pass a shared `user_sync_store` implementation so all workers claim from the same pending-user queue. The package only defines the async protocol; Redis, NATS KV, SQL, or any other backend can be implemented by your application.

```python
store = MyRedisUserSyncStore(redis_client)  # implements Bridge.UserSyncStoreProtocol

node = Bridge.create_node(
    connection=Bridge.NodeType.grpc,
    address="127.0.0.1",
    port=2096,
    api_port=2097,
    server_ca=server_ca_pem_string,
    api_key="xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
    node_id="node-1",
    worker_id="worker-a",
    user_sync_store=store,
)
```

A `UserSyncStoreProtocol` implementation must provide these async methods:

- `enqueue_users(node_id, users)` stores latest user payloads by email.
- `claim_users(node_id, worker_id, limit, lease_seconds)` atomically leases work and returns `ClaimedUser` items.
- `ack_users(node_id, tokens)` removes successfully synced claims.
- `requeue_users(node_id, claimed_users)` makes failed claims available again.
- `clear(node_id)` clears pending and claimed updates for a node.

Delivery is at-least-once. A crashed worker may cause the same latest user payload to be synced again after its lease expires, so external adapters should use atomic claim/lease operations such as Redis Lua/transactions or NATS KV revision compare-and-set.

Lifecycle operations are coordinated through the same model. The default process-local coordinator prevents concurrent `start()`, `stop()`, `update_node()`, `update_core()`, and `update_geofiles()` calls from controllers for the same node in one process. Pass a shared `lifecycle_coordinator` in multi-process or multi-host deployments so only one worker can perform a lifecycle operation at a time. Read-only status cron jobs can call stats/info normally; if they write shared observed status, use the current lifecycle epoch so stale cron results cannot overwrite a newer reconnect result.

```python
lifecycle = MyRedisLifecycleCoordinator(redis_client)

node = Bridge.create_node(
    connection=Bridge.NodeType.grpc,
    address="127.0.0.1",
    port=2096,
    api_port=2097,
    server_ca=server_ca_pem_string,
    api_key="xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
    node_id="node-1",
    worker_id="worker-a",
    user_sync_store=store,
    lifecycle_coordinator=lifecycle,
)

state = await node.get_lifecycle_state()
health = await node.get_health()
if state is not None:
    await node.update_observed_lifecycle(
        Bridge.LifecycleStatus.HEALTHY if health is Bridge.Health.HEALTHY else Bridge.LifecycleStatus.BROKEN,
        expected_epoch=state.epoch,
    )
```

A lifecycle adapter must atomically acquire/release leases and fence writes with the returned epoch. This prevents a cron status job or another worker from overwriting the result of a newer `start()`, `stop()`, or reconnect flow.

Node connection configs can also be stored through a registry protocol:

```python
registry = MyNodeRegistry(...)
config = Bridge.NodeConfig(
    connection="grpc",
    address="127.0.0.1",
    port=2096,
    api_port=2097,
    server_ca=server_ca_pem_string,
    api_key="xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
)

await Bridge.save_node_config(registry, "node-1", config)
node = await Bridge.create_node_from_registry(
    registry,
    "node-1",
    user_sync_store=store,
    worker_id="worker-a",
)
```

### 2. Direct User Sync

Use direct sync when you want explicit control in your flow.

```python
await node.sync_users([user1, user2], timeout=15)
```

### 3. Chunked Sync For Large Batches

```python
failed_users = await node.sync_users_chunked(
    users=large_user_list,
    chunk_size=500,
    timeout=30,
)

if failed_users:
    print(f"Failed users: {len(failed_users)}")
```

### 4. Stats APIs

```python
system_stats = await node.get_system_stats()
backend_stats = await node.get_backend_stats()
latencies = await node.get_outbounds_latency()

all_outbounds = await node.get_stats(
    stat_type=service.StatType.Outbounds,
    reset=False,
)

single_user_online = await node.get_user_online_stats("alice@example.com")
single_user_ips = await node.get_user_online_ip_list("alice@example.com")
```

### 5. Health And Version Helpers

```python
health = await node.get_health()            # Bridge.Health enum
node_ver = await node.node_version()
core_ver = await node.core_version()
node_ver2, core_ver2 = await node.get_versions()
meta = await node.get_extra()
```

### 6. On-Demand Log Streaming

`stream_logs()` yields an `asyncio.Queue` that contains log lines (`str`) or `Bridge.NodeAPIError`.

```python
import asyncio

async with node.stream_logs(max_queue_size=200) as log_queue:
    for _ in range(20):
        item = await asyncio.wait_for(log_queue.get(), timeout=2)
        if isinstance(item, Bridge.NodeAPIError):
            raise item
        print(item)
```

### 7. Maintenance Endpoints

These methods use the node REST JSON API (`api_port`).

```python
await node.update_node()
await node.update_core({"version": "latest"})
await node.update_geofiles({"remove_temp": True})
```

### 8. Routing APIs

Routing operations work over both gRPC and REST. They are xray-only: on a non-xray
(e.g. WireGuard) node the call fails with `Bridge.NodeAPIError` code `501`.

```python
rules = await node.list_routing_rules()
balancer = await node.get_balancer_info("balancer-tag")

route = await node.test_route(
    inbound_tag="inbound-1",
    network="tcp",
    target_domain="example.com",
    target_port=443,
)

# `rule` is one xray routing rule as JSON (same shape as a routing.rules[] entry).
# Appended by default (keeps existing rules); pass should_reset=True to clear all
# rules + balancers before adding.
await node.add_routing_rule(
    '{"type":"field","outboundTag":"direct","domain":["example.com"],"ruleTag":"r1"}'
)
await node.remove_routing_rule("r1")
await node.override_balancer_target("balancer-tag", "outbound-tag")
```

## API Reference

### Lifecycle

- `start(config, backend_type, users, keep_alive=0, exclude_inbounds=[], timeout=None)`
- `stop(timeout=None)`
- `info(timeout=None)`

### Health/Version

- `get_health()`
- `node_version()`
- `core_version()`
- `get_versions()`
- `get_extra()`

### Stats

- `get_system_stats(timeout=None)`
- `get_backend_stats(timeout=None)`
- `get_stats(stat_type, reset=True, name="", timeout=None)`
- `get_outbounds_latency(name="", timeout=None)`
- `get_user_online_stats(email, timeout=None)`
- `get_user_online_ip_list(email, timeout=None)`

### User Sync

- `update_user(user)` (queued/background)
- `update_users(users)` (queued/background)
- `sync_users(users, flush_pending=False, timeout=None)` (direct)
- `sync_users_chunked(users, chunk_size=100, flush_pending=False, timeout=None)` (direct streaming)

### Routing

Xray-only (gRPC and REST); on a non-xray backend these raise `NodeAPIError(501)`.

- `list_routing_rules(timeout=None)`
- `get_balancer_info(tag, timeout=None)`
- `test_route(inbound_tag="", network="", target_ip="", target_domain="", target_port=0, protocol="", user="", attributes=None, field_selectors=None, publish_result=False, timeout=None)`
- `add_routing_rule(rule, should_reset=False, timeout=None)`
- `remove_routing_rule(rule_tag, timeout=None)`
- `override_balancer_target(balancer_tag, target, timeout=None)`

### Logging

- `stream_logs(max_queue_size=1000)` async context manager returning an `asyncio.Queue`

### Maintenance

- `update_node()`
- `update_core(json)`
- `update_geofiles(json)`

## Error Handling

All transport and API errors are surfaced as `Bridge.NodeAPIError`:

```python
try:
    await node.get_backend_stats(timeout=5)
except Bridge.NodeAPIError as e:
    print(e.code, e.detail)
```

## Protobuf Access

For direct protobuf usage:

```python
from PasarGuardNodeBridge.common import service_pb2 as service
```

## Complete Minimal Example

```python
import asyncio
import PasarGuardNodeBridge as Bridge
from PasarGuardNodeBridge.common import service_pb2 as service


async def main():
    with open("certs/ssl_cert.pem", "r", encoding="utf-8") as f:
        server_ca = f.read()
    with open("config/xray.json", "r", encoding="utf-8") as f:
        config = f.read()

    node = Bridge.create_node(
        connection=Bridge.NodeType.grpc,
        address="127.0.0.1",
        port=2096,
        api_port=2097,
        server_ca=server_ca,
        api_key="d04d8680-942d-4365-992f-9f482275691d",
        name="example-node",
    )

    await node.start(config=config, backend_type=service.BackendType.XRAY, users=[])
    print(await node.get_system_stats())
    await node.stop()


asyncio.run(main())
```
