Metadata-Version: 2.4
Name: ironflock
Version: 1.9.0
Summary: IronFlock Python SDK for connecting to the IronFlock Platform
License-Expression: MIT
License-File: LICENSE
Author: Marko Petzold, IronFlock GmbH
Author-email: info@ironflock.com
Requires-Python: >=3.8
Classifier: Development Status :: 5 - Production/Stable
Classifier: Intended Audience :: Developers
Provides-Extra: dev
Provides-Extra: docs
Requires-Dist: autobahn[asyncio,serialization] (==25.12.2)
Requires-Dist: mypy ; extra == "dev"
Requires-Dist: pydantic (>=2.0.0)
Requires-Dist: pydantic ; extra == "dev"
Requires-Dist: pytest ; extra == "dev"
Requires-Dist: pytest-asyncio ; extra == "dev"
Requires-Dist: pytest-cov ; extra == "dev"
Requires-Dist: ruff ; extra == "dev"
Requires-Dist: sphinx (>=8.2.3) ; extra == "docs"
Requires-Dist: sphinx-autodoc-typehints (>=2.0) ; extra == "docs"
Requires-Dist: sphinx-book-theme (>=1.1.4) ; extra == "docs"
Project-URL: Homepage, https://github.com/RecordEvolution/ironflock-py
Project-URL: Repository, https://github.com/RecordEvolution/ironflock-py
Description-Content-Type: text/markdown

# ironflock

![Coverage](https://img.shields.io/badge/coverage-73%25-yellow)

## About

With this library you can publish data from your apps on your IoT edge hardware to the fleet data storage of the [IronFlock](https://studio.ironflock.com) devops platform.
When this library is used on a certain device the library automatically uses the private messaging realm (Unified Name Space)
of the device's fleet and the data is collected in the respective fleet database.

So if you use the library in your app, the data collection will always be private to the app user's fleet.

For more information on the IronFlock IoT Devops Platform for engineers and developers visit our [IronFlock](https://www.ironflock.com) home page.

## Requirements

- Python 3.8 or higher

## Installation

Install from PyPI:

```shell
pip install ironflock
```
## Usage

```python
import asyncio
from ironflock import IronFlock

# create an IronFlock instance to connect to the IronFlock platform data infrastructure.
# The IronFlock instance handles authentication when run on a device registered in IronFlock.
ironflock = IronFlock()

async def main():
    while True:
        # publish an event (if connection is not established the publish is skipped)
        publication = await ironflock.publish("test.publish.example", {"temperature": 20})
        print(publication)
        await asyncio.sleep(3)


if __name__ == "__main__":
    ironflock = IronFlock(mainFunc=main)
    ironflock.run()
```

## Options

The `IronFlock` `__init__` function can be configured with the following options:

```ts
{
    serial_number: string;
}
```

**serial_number**: Used to set the serial_number of the device if the `DEVICE_SERIAL_NUMBER` environment variable does not exist. It can also be used if the user wishes to authenticate as another device.

**reconnect_window** (seconds, default `60`): how long a table operation (`append_to_table`, `append_rows_to_table`/`publish_rows_to_table`, `publish_to_table`, `getHistory`, `get_series_history`, `reveal_secrets`, `verify_secret`, and a consumed app's history reads) rides out a platform restart before it raises. It waits that long for the connection to come back and for the platform to serve the app's tables again, so an app keeps running through a router or platform update. `0` turns this off (10 s connection wait, no retries).

## API Reference

### `publish(topic, *args, **kwargs)`

Publishes an event to a topic on the IronFlock message router.

```python
publication = await ironflock.publish("com.myapp.mytopic", {"temperature": 20})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `topic` | `str` | The URI of the topic to publish to |
| `*args` | positional | Payload arguments |
| `**kwargs` | keyword | Payload keyword arguments |

**Returns:** `Publication` — The publication object (an acknowledged publish receipt).

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the publish fails (e.g. not connected).

---

### `publish_to_table(tablename, *args, **kwargs)`

Convenience function to publish data to a fleet table in the IronFlock platform. Automatically constructs the table's write topic `<SWARM_KEY>.<APP_KEY>.<tablename>` from the injected environment variables. No app or dashboard may subscribe to that topic; read the rows back with `getHistory` or `subscribe_to_table` (see [URIs an app may use](#uris-an-app-may-use)).

```python
await ironflock.publish_to_table("sensordata", {"temperature": 22.5, "humidity": 60})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table, e.g. `"sensordata"` |
| `*args` | positional | Row data to publish |
| `**kwargs` | keyword | Row data as keyword arguments |

**Returns:** `Publication` — The publication object (an acknowledged publish receipt).

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the publish fails (e.g. not connected).

---

### `append_to_table(tablename, *args, **kwargs)`

Appends data to a fleet table by calling the registered append procedure at `append.<SWARM_KEY>.<APP_KEY>.<tablename>`. Unlike `publish_to_table`, this uses a remote procedure call rather than a pub/sub event.

```python
await ironflock.append_to_table("sensordata", {"temperature": 22.5, "humidity": 60})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table, e.g. `"sensordata"` |
| `*args` | positional | Row data to append |
| `**kwargs` | keyword | Row data as keyword arguments |

**Returns:** `Any` — The result of the remote procedure call.

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the call fails (e.g. not connected, or the procedure is not registered).

---

### `publish_rows_to_table(tablename, rows, **kwargs)`

Publishes **many rows in a single message** (bulk insert) to the dedicated topic `bulk.<SWARM_KEY>.<APP_KEY>.<tablename>`. The platform inserts the whole batch atomically (all-or-nothing) in one operation. Use this for high-frequency data where one round-trip per row is too costly. Like `publish_to_table`, this is fire-and-forget — the ack confirms delivery to the router, not the DB insert.

```python
await ironflock.publish_rows_to_table("sensordata", [
    {"tsp": "2024-01-15T10:30:00.000Z", "temperature": 22.5},
    {"tsp": "2024-01-15T10:30:01.000Z", "temperature": 22.7},
])
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table, e.g. `"sensordata"` |
| `rows` | `List[dict]` | Non-empty list of row dicts to insert |
| `**kwargs` | keyword | Extra arguments shared by the whole batch |

**Returns:** `Publication` — The publication object (an acknowledged publish receipt).

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the publish fails (e.g. not connected).

---

### `append_rows_to_table(tablename, rows, **kwargs)`

Appends **many rows in a single RPC** (bulk insert) by calling the dedicated procedure `appendBulk.<SWARM_KEY>.<APP_KEY>.<tablename>`. The platform inserts the whole batch atomically (all-or-nothing): if any row is invalid the entire batch is rejected and nothing is persisted. Prefer this over `publish_rows_to_table` when you need the insert outcome.

```python
result = await ironflock.append_rows_to_table("sensordata", [
    {"tsp": "2024-01-15T10:30:00.000Z", "temperature": 22.5},
    {"tsp": "2024-01-15T10:30:01.000Z", "temperature": 22.7},
])
# result -> {"success": True, "count": 2}
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table, e.g. `"sensordata"` |
| `rows` | `List[dict]` | Non-empty list of row dicts to insert |
| `**kwargs` | keyword | Extra arguments shared by the whole batch |

**Returns:** `Any` — The result of the remote procedure call (e.g. `{"success": True, "count": N}`).

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the bulk append fails (all-or-nothing — no rows were persisted).

---

### `report_error(error, level="error", append=False, tsp=None)`

Reports an application error into the fleet's `error-logs` table. This is a convenience wrapper over `publish_to_table` / `append_to_table`: it stamps the row with `source="app"`, a severity `level`, and a timestamp, then writes it like any normal table row. The error lands in the same per-databackend `error-logs` table that fleetdb system errors use (tagged `source="system"`), so it is queryable with `getHistory`, streamable with `subscribe_to_table`, usable in board-templates, and delivered in realtime on `transformed.error-logs` — without firing the platform's system-error toast.

```python
# Fire-and-forget (default): publishes to the error-logs table
await ironflock.report_error("Sensor read timed out", level="warn")

# Pass an Exception to capture its traceback (falls back to the message)
try:
    risky_operation()
except Exception as err:
    await ironflock.report_error(err)

# Use the append RPC when you want to await the insert outcome
await ironflock.report_error("Calibration failed", level="error", append=True)
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `error` | `str \| BaseException` | The error message, or an exception whose traceback (or message) is recorded |
| `level` | `str`, optional | Severity: `"error"`, `"warn"`, `"info"` or `"debug"`. Defaults to `"error"` |
| `append` | `bool`, optional | When `True`, use the append RPC (returns the insert outcome). Defaults to `False` (fire-and-forget publish) |
| `tsp` | `str`, optional | ISO-8601 timestamp override. Defaults to the current time |

**Returns:** `Publication | Any` — The publication object (or, with `append=True`, the RPC result).

**Raises:** `RuntimeError` with a descriptive message if the underlying publish/append fails (e.g. not connected).

---

### `subscribe(topic, handler, options=None)`

Subscribes to a topic on the IronFlock message router.

```python
def on_message(*args, **kwargs):
    print("Received:", args, kwargs)

subscription = await ironflock.subscribe("com.myapp.mytopic", on_message)
```

Use your own topic names here; [URIs an app may use](#uris-an-app-may-use) lists the names the router reserves. For table rows use `subscribe_to_table`.

| Parameter | Type | Description |
|-----------|------|-------------|
| `topic` | `str` | The URI of the topic to subscribe to |
| `handler` | callable | Function called when a message is received |
| `options` | `SubscribeOptions`, optional | Subscription options |

**Returns:** `Subscription` — The subscription object.

**Raises:** `RuntimeError` with a descriptive message if the subscription fails (e.g. not connected).

---

### `subscribe_to_table(tablename, handler, options=None)`

Convenience function to subscribe to a fleet table. Subscribes to the data backend's realtime feed `transformed.<tablename>` and its bulk counterpart `transformed.bulk.<tablename>`. Each event is the row as stored: values typed to the data-template columns, secret columns masked. Receives rows written via both the single-row and the bulk insert paths — rows from a bulk insert are delivered to your handler one at a time, so handler code stays the same.

```python
def on_table_data(*args, **kwargs):
    print("New row:", args, kwargs)

await ironflock.subscribe_to_table("sensordata", on_table_data)
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table to subscribe to |
| `handler` | callable | Function called when new data arrives |
| `options` | `SubscribeOptions`, optional | Subscription options |

**Returns:** `Subscription` — The subscription object.

**Raises:** `RuntimeError` with a descriptive message if the subscription fails (e.g. not connected).

---

### `getHistory(tablename, queryParams)`

Retrieves historical data from a fleet table.

```python
# Simple query with limit
data = await ironflock.getHistory("sensordata", {"limit": 100})

# Query with time range and filters
data = await ironflock.getHistory("sensordata", {
    "limit": 500,
    "offset": 0,
    "timeRange": {
        "start": "2026-01-01T00:00:00Z",
        "end": "2026-03-01T00:00:00Z"
    },
    "filterAnd": [
        {"column": "temperature", "operator": ">", "value": 20},
        {"column": "humidity", "operator": "<=", "value": 80}
    ]
})

# Current value(s) only: add the `latest` marker. The data backend derives
# the latest row per entity in SQL (entity = the table's maintainLatestFlagFor
# columns from the data-template; without one, the single most recent row).
current = await ironflock.getHistory("sensordata", {
    "limit": 100,
    "filterAnd": [{"latest": True}]
})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table to query |
| `queryParams` | `dict` or `TableQueryParams` | Query parameters (see below) |

**queryParams fields:**

| Field | Type | Description |
|-------|------|-------------|
| `limit` | `int` | Maximum number of rows to return (1–10000, required) |
| `offset` | `int`, optional | Offset for pagination |
| `timeRange` | `dict`, optional | `{"start": "<ISO datetime>", "end": "<ISO datetime>"}` |
| `filterAnd` | `list`, optional | List of AND filter conditions `{"column": str, "operator": str, "value": ...}`, and/or the `{"latest": True}` mode marker (see below) |
| `columns` | `List[str]`, optional | Columns to return (`tsp`, `device_key` and `authid` are always included). Omit for all columns |

Supported filter operators: `=`, `!=`, `<>`, `>`, `<`, `>=`, `<=`, `LIKE`, `ILIKE`, `NOT LIKE`, `NOT ILIKE`, `IN`, `NOT IN`, `IS NULL`, `IS NOT NULL`.

`IS NULL` and `IS NOT NULL` take no `value`, and they are the only way to ask about NULL: every other operator yields unknown on a NULL column and excludes the row, exactly as in SQL. `LIKE` is case-sensitive; use `ILIKE` to match case-insensitively.

**Combining with OR:** entries of `filterAnd` are AND-ed. A group entry combines its own entries with one operator, so the soft-delete pattern becomes:

```python
filterAnd=[
    {
        "combinator": "OR",
        "filters": [
            {"column": "deleted", "operator": "IS NULL"},
            {"column": "deleted", "operator": "=", "value": False},
        ],
    }
]
```

A data backend that predates groups rejects this shape rather than applying part of it, so it fails loudly instead of quietly returning rows the filter excludes.

**Latest values:** a `{"latest": True}` entry in `filterAnd` is not a WHERE predicate but a mode switch: the data backend returns only the latest row per entity, derived on the fly in SQL (`DISTINCT ON` over the entity key declared as `maintainLatestFlagFor` in the table's data-template; a table without an entity key yields the single most recent row). The former physical `latest_flag` column no longer exists — a legacy `{"column": "latest_flag", "operator": "=", "value": True}` filter is still accepted and treated as the marker, but new code should use `{"latest": True}`. Other predicates combine with the marker as expected: entity-key predicates narrow which entities are returned, all other predicates and `timeRange` filter the resulting latest rows.

**Returns:** `Any` — The query result data (typically a list of row dicts). Columns your data template declares `secret: true` come back as `SECRET_PLACEHOLDER` — see [Secret Columns](#secret-columns).

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message when the history procedure is not registered (table not declared / data backend not running) or the call fails.

---

### `get_series_history(tablename, params)`

Retrieves **down-sampled time-series** data from a fleet table — numeric columns aggregated into time namespaces (e.g. hourly averages). Ideal for charts over long time ranges. Available for tables (not transforms).

```python
series = await ironflock.get_series_history("sensordata", {
    "metrics": ["temperature", "humidity"],
    "method": "AVG",
    "limit": 500,
    "timeRange": ["2026-01-01T00:00:00Z", "2026-03-01T00:00:00Z"],
    "groupBy": ["device_id"]
})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table to query |
| `params` | `SeriesQueryParams` or `dict` | Series query parameters (see below) |

**params fields:**

| Field | Type | Description |
|-------|------|-------------|
| `metrics` | `List[str]` | Numeric columns to down-sample |
| `method` | `str` | Aggregation per namespace: `"AVG"`, `"SUM"`, `"COUNT"`, `"MIN"`, `"MAX"`, `"FIRST"` or `"LAST"` |
| `limit` | `int` | Maximum number of namespaces (1–10000) |
| `timeRange` | `list` | `[start, end]` — ISO strings or epoch-ms numbers; `None` = open end (required) |
| `groupBy` | `List[str]`, optional | Columns to group the series by |
| `filterAnd` | `list`, optional | AND filter conditions (WHERE predicates only — the `{"latest": True}` marker is not supported in series queries; use `getHistory` for latest values) |

**Returns:** `Any` — The down-sampled series rows.

**Raises:** `ValueError` on invalid parameters (including a latest marker in `filterAnd`); `RuntimeError` with a descriptive message when the series procedure is not registered or the call fails.

---

### `call(topic, args=None, kwargs=None, options=None)`

Calls a remote procedure on the IronFlock message router using a full WAMP topic URI. An app cannot register a name like `some.full.wamp.topic` (see [URIs an app may use](#uris-an-app-may-use)), so such a name only reaches procedures that platform services register on the app's data realm, such as the data backend's and the file storage's. To call a function another device of your app registered, use `call_device_function`.

```python
result = await ironflock.call("some.full.wamp.topic", args=[42])
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `topic` | `str` | The full WAMP URI of the procedure to call |
| `args` | `list`, optional | Positional arguments |
| `kwargs` | `dict`, optional | Keyword arguments |
| `options` | `CallOptions`, optional | Call options |

**Returns:** `Any` — The result of the remote procedure call.

**Raises:** `ValueError` on invalid parameters; `RuntimeError` with a descriptive message if the call fails (e.g. not connected, or the procedure is not registered).

---

### `call_device_function(device_key, topic, args=None, kwargs=None, options=None)`

Calls a remote procedure registered by another IronFlock device. Automatically assembles the full WAMP topic as `{SWARM_KEY}.{device_key}.{APP_KEY}.{STAGE}.{topic}`, the same URI `register_device_function` registers on the target device. `STAGE` is `DEV` or `PROD`, the stage of the realm the app joined (`ENV=PROD` in any case selects `PROD`; anything else, unset included, selects `DEV`).

```python
result = await ironflock.call_device_function(42, "com.myapp.myprocedure", args=[42])
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `device_key` | `int` | The device key of the target device |
| `topic` | `str` | The URI of the procedure to call |
| `args` | `list`, optional | Positional arguments |
| `kwargs` | `dict`, optional | Keyword arguments |
| `options` | `CallOptions`, optional | Call options |

**Returns:** `Any` — The result of the remote procedure call.

**Raises:** `ValueError` when `topic` is empty or blank, or when `device_key`, `SWARM_KEY` or `APP_KEY` is not set (missing, empty, blank or `0`); `RuntimeError` with a descriptive message if the call fails (e.g. not connected, or the procedure is not registered).

> `call_function()` is a deprecated alias for `call_device_function()`.

---

### `register_device_function(topic, endpoint, options=None)`

Registers a procedure that can be called by other devices in the fleet, and by dashboard widget actions. Automatically constructs the full WAMP topic as `{SWARM_KEY}.{DEVICE_KEY}.{APP_KEY}.{STAGE}.{topic}` (`STAGE` as for `call_device_function`). Register procedures this way: the router refuses to register names such as `com.myapp.proc` (see [URIs an app may use](#uris-an-app-may-use) for what it accepts), and accepts single registrations only (`RegisterOptions` with an `invoke` policy other than `"single"` is refused).

```python
def add(a, b):
    return a + b

await ironflock.register_device_function("com.myapp.add", add)
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `topic` | `str` | The URI of the procedure to register |
| `endpoint` | callable | The function to register |
| `options` | `RegisterOptions`, optional | Registration options |

**Returns:** `Registration` — The registration object.

**Raises:** `ValueError` when `topic` is empty or blank, or when `SWARM_KEY`, `DEVICE_KEY` or `APP_KEY` is not set (missing, empty, blank or `0`); `RuntimeError` with a descriptive message if the registration fails (e.g. not connected, or the router refuses it).

> **Apps that override `ENV`:** the stage segment is always `DEV` or `PROD`. ironflock-py 1.8.6 and earlier put the raw `ENV` value there. That differs only when `ENV` is missing (older versions sent `None`) or your app sets it in its own environment settings with another spelling, e.g. `prod`. With such an override, devices on 1.8.6 or earlier cannot call functions that devices on a newer version register, and the other way round. Dashboard widget actions call the `DEV`/`PROD` spelling, so they now reach these functions.

> `register()` is an alias for `register_device_function()`. `register_function()` is a deprecated alias.

---

### `set_device_location(long, lat)`

Updates the device's location in the platform master data. The maps in device or group overviews will reflect the new location in realtime.

```python
await ironflock.set_device_location(long=8.6821, lat=50.1109)
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `long` | `float` | Longitude (-180 to 180) |
| `lat` | `float` | Latitude (-90 to 90) |

**Returns:** `Any` — The result of the location update call.

**Raises:** `ValueError` on invalid coordinates; `RuntimeError` with a descriptive message if the location service call fails.

> **Note:** Location history is not stored. If you need location history, create a dedicated table and use `publish_to_table`.

> **Not available at the moment:** no platform service answers this call on the app's data realm, so it raises `RuntimeError` (`wamp.error.no_such_procedure`). Keep locations in a table of your own until it is.

---

### `getRemoteAccessUrlForPort(port)`

Returns the remote access URL for a given port on the device.

```python
url = ironflock.getRemoteAccessUrlForPort(8080)
# e.g. "https://<device_key>-<app_name>-8080.app.ironflock.com"
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `port` | `int` | The port number |

**Returns:** `str | None` — The remote access URL string (e.g. `"https://<device_key>-<app_name>-8080.app.ironflock.com"`), or `None` if the device key or app name is not available.

---

### Properties

| Property | Type | Description |
|----------|------|-------------|
| `is_connected` | `bool` | `True` if the connection to the platform is established |
| `files` | `FileStore` | The app's managed object storage — see [Managed File Storage](#managed-file-storage). Safe to access before `start()`; each call waits for the connection just like the table API |
| `connection` | `CrossbarConnection` | The underlying connection instance (for advanced use) |

### Lifecycle Methods

| Method | Description |
|--------|-------------|
| `run()` | Starts the connection and runs the main function (blocking, synchronous) |
| `await start()` | Starts the connection asynchronously |
| `await stop()` | Stops the connection and cancels the main task |
| `await run_async()` | Starts the connection and keeps it running asynchronously |

## URIs an app may use

Your app container joins its own data realm, `realm-<SWARM_KEY>-<APP_KEY>-<dev|prod>`, with the router role `swarm_device` (the device agent's role, which app containers share until the platform gives them a role of their own). The router checks every publish, subscribe, call and register there against that role's allow-list. The SDK's own methods stay inside that list. It matters when you pass your own URIs to `publish()`, `subscribe()`, `call()` or `ironflock.connection`. A refused action fails with `wamp.error.not_authorized`, which the `IronFlock` methods raise as a `RuntimeError`.

| URI on your own realm | Allowed | SDK method |
|-----------------------|---------|------------|
| Your own names, e.g. `com.myapp.status` | publish, subscribe (a prefix subscription such as `com.myapp.` too), call; never register | `publish()`, `subscribe()`, `call()` |
| `<SWARM_KEY>.<APP_KEY>.<table>`, `bulk.<SWARM_KEY>.<APP_KEY>.<table>` | publish (never subscribe) | `publish_to_table()`, `report_error()`, `publish_rows_to_table()` |
| `transformed.<table>`, `transformed.bulk.<table>` | subscribe only | `subscribe_to_table()` |
| `databackend.errors.…`, `files.events…` | subscribe only | – |
| `<SWARM_KEY>.<DEVICE_KEY>.<APP_KEY>.<STAGE>.<name>` | register (single), call, publish | `register_device_function()`, `call_device_function()` |
| `append.…`, `appendBulk.…`, `history.transformed.…`, `secret.reveal.…`, `secret.verify.…`, `files.read.…`, `files.write.…` | as for your own names (the SDK calls them) | table, history, secret and `files` APIs |
| `sys.appaccess.list`, `sys.appaccess.resolve` | call | `list_consumable_apps()`, `connect_to_app()`, `connect_to_all_apps()` |

Every name that starts with a digit (the table write topics and the device-function URIs) allows publish, call and single registration, never subscribe; `bulk.…` allows publish only.

**Reserved first segments.** These belong to the platform and the device agent. Apart from a few names that the agent and the SDKs use, the router refuses them: `sys.`, `wamp.`, `reswarm.`, `reacct.`, `re.mgmt.`, `re.tunnel.`, `ironflock.`, `instance.`, `cloud.instance.`, `alarms.`, `tunnel.`, `svc.`, `dev.`, `uplink.`, `auth.`, `truncate_tables.`, `table_info.` and `history.errors.`. `user.`, `pub.`, `store.` and `agent.` are platform audiences as well: publish and subscribe there still pass but will be closed, so do not build on them. Do not start your own names with a digit, `bulk.`, `transformed.`, `databackend.errors.` or `files.events` either: those are the table, function and data-backend topics above.

**Nobody subscribes to raw table topics.** `publish_to_table()` and `publish_rows_to_table()` write to `<SWARM_KEY>.<APP_KEY>.<table>` and `bulk.<SWARM_KEY>.<APP_KEY>.<table>`. These topics carry the row exactly as sent, including the plaintext of columns declared `secret: true`, which the data backend encrypts only when it stores the row. So no app or dashboard may subscribe to them. Use `subscribe_to_table()`: the data backend publishes every stored row on `transformed.<table>` (bulk inserts on `transformed.bulk.<table>`), typed, with secret columns masked.

**Register functions only through the device-function API.** Register procedures with `register_device_function()`. It builds the one shape the platform addresses, `<SWARM_KEY>.<DEVICE_KEY>.<APP_KEY>.<STAGE>.<name>`, where `STAGE` is `DEV` or `PROD`, the stage of the realm. Other devices of your app call it with `call_device_function()`; dashboard widget actions call the same URI. Today the role grants register only on names that start with a digit and on the device agent's own `re.mgmt.<serial>.…` names, which are not for apps; any other name, such as `com.myapp.proc`, is refused, and only single registrations are accepted. Let the SDK build the URI: the platform's upcoming identity check (Stage 2) will accept a registration only when the URI carries your own swarm key, device key, app key and stage. Because no app can register names like `com.myapp.proc`, `call()` on such a name only reaches procedures that platform services register on the data realm.

**Cross-app access is read-only.** `connect_to_app()` and `connect_to_all_apps()` open a second connection to the provider app's realm with the `app_reader` role (this needs the per-app credential the agent injects as `APP_AUTH_ID`/`APP_AUTH_SECRET`; with the legacy device credential the provider realm is refused, or joined with the provider app's own rights when that app runs on the same device). There your app may subscribe to `transformed.…` and `files.events…`, and call `history.transformed.…` (history and series) and `files.read.…`. Nothing else is allowed, besides the heartbeat probe `wamp.session.get` (used by the JS SDK): no publish, no register, no writes and no secret procedures. The provider's private tables are filtered out as well.

## Connection reliability

The connection reconnects on its own. When the socket drops, the SDK retries
until the router is back, then restores every subscription and every registered
device function — you do not need to re-subscribe after a reconnect. If one
topic cannot be restored — a permission that changed, a transient router error —
the rest still are, and the failed one is retried on the next reconnect.

A special case is a router that is reachable but has no realm for the app yet.
The concurrent-install race boots the app container before the databackend has
provisioned its realm, so the first joins are refused with
`wamp.error.no_such_realm`; the SDK keeps retrying every couple of seconds and
the app comes up the moment the realm does. A realm still missing after a
minute is most likely never going to appear — the app was deleted, or has no
databackend for this stage — so from then on the SDK slows to one attempt every
two minutes rather than hitting the router every second forever. The fast
cadence returns as soon as the realm appears, or the router itself goes away
and comes back.

The harder case is a socket that dies *silently*: an idle NAT or proxy cuts the
connection without telling either side, or a router restarts behind a load
balancer. Nothing arrives to signal the loss. A connection that only publishes
finds out on its next failing write, but a connection that only subscribes never
writes at all, so without help it sits dead until the process restarts.

**WebSocket keepalive** closes that gap. The SDK pings the router every 30
seconds while the link is idle and expects a pong within 10 seconds; if none
arrives, the transport is dropped and the normal reconnect path runs. Incoming
traffic counts as proof of life, so a busy connection is never probed — the ping
only goes out when the link has actually gone quiet.

This is on by default and needs no configuration. Cross-app connections from
`connect_to_app()` are typically subscribe-only, which is exactly the case
keepalive protects, and they are covered automatically.

> The JavaScript SDK solves the same problem differently: the browser WebSocket
> API exposes no ping, so it uses a WAMP-level heartbeat with tunable
> `heartbeatIntervalMs` / `heartbeatTimeoutMs` properties instead.

## Secret Columns

A data-template column declared `secret: true` is encrypted by the data backend (AES-256-GCM) before it is written, so what Postgres stores is an `ifsec:1:<base64url>` envelope, not the value. The flag is allowed on `dataType: string` columns only, never on `tsp`, and never on a column used as an entity key (`maintainLatestFlagFor`).

No read path hands the plaintext back by accident:

| Read path | What a non-null secret value looks like |
|-----------|------------------------------------------|
| `getHistory`, `subscribe_to_table`, cross-app reads | `SECRET_PLACEHOLDER` (the string `"__secret__"`), exported from `ironflock` |
| SQL — `sys.dataservice.select_query`, the FleetDB Access Postgres login, backups | the raw `ifsec:1:...` ciphertext |
| `reveal_secrets` / `verify_secret` | the decrypted value / a boolean |

`None` stays `None` everywhere, so "set but hidden" remains distinguishable from "never written".

Both procedures below are registered per app realm and callable **only by your app's own containers** — the router denies every browser role and every consuming app, so another app can never decrypt your secrets.

Encryption uses a random IV, so two rows holding the same plaintext hold **different** ciphertext. Equality filters, `GROUP BY`, `ORDER BY`, `DISTINCT` and entity keys on a secret column are impossible, not merely discouraged — select the rows by their other columns.

### `reveal_secrets(tablename, query_params={"limit": 10})`

Reads rows of your own table with the secret columns decrypted.

```python
# The most recent row, secrets in clear
rows = await ironflock.reveal_secrets("credentials", {
    "limit": 1,
    "filterAnd": [{"latest": True}],
})
print(rows[0]["api_key"])  # the plaintext, not "__secret__"

# Narrow by a NON-secret column — filtering by the secret itself can never match
rows = await ironflock.reveal_secrets("credentials", {
    "limit": 10,
    "filterAnd": [{"column": "gateway_id", "operator": "=", "value": 471}],
})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table to read |
| `query_params` | `dict` or `SecretQueryParams` | Same fields as `getHistory` (`limit`, `offset`, `timeRange`, `filterAnd`, `columns`), except that `limit` must be 1–100 — the data backend **rejects** a larger value for this procedure rather than clamping it |

**Returns:** `Any` — The rows, in the same shape `getHistory` returns, with the secret columns decrypted.

**Raises:** `ValueError` on invalid parameters; `RuntimeError` when the procedure is not registered (the table declares no secret column / data backend not running) or the call fails.

### `verify_secret(tablename, column, candidate, query_params={"limit": 1})`

Checks a candidate against a stored secret **without reading it back** — the comparison happens inside the data backend, in constant time, and nothing decrypted leaves it. Prefer this whenever you don't actually need the plaintext.

```python
# Is this the current API key? (default query = the most recent row)
ok = await ironflock.verify_secret("credentials", "api_key", received_key)

# Check against a particular entity's rows
ok = await ironflock.verify_secret("credentials", "api_key", received_key, {
    "limit": 1,
    "filterAnd": [
        {"latest": True},
        {"column": "gateway_id", "operator": "=", "value": 471},
    ],
})
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `tablename` | `str` | The name of the table to check |
| `column` | `str` | The secret column to compare |
| `candidate` | `str` | The plaintext to compare the stored value against |
| `query_params` | `dict` or `SecretQueryParams` | Selects WHICH rows are checked — same fields as `getHistory`; `limit` must be 1–100, the data backend **rejects** a larger value rather than clamping it |

**Returns:** `bool` — `True` when **any** selected row's column decrypts to `candidate`.

**Raises:** `ValueError` on invalid parameters; `RuntimeError` when the procedure is not registered or the call fails.

---

## Cross-App Data Access

Read another app's fleet data from within your app, in the same project and fleet. The provider app must list your app in its data-template `consumes:` section, and the project user must grant access. Access is **read-only**: you can query history and subscribe to realtime rows of the tables and transforms (views) the provider shares — you cannot write to them.

```python
from ironflock import IronFlock, CrossAppAccessError

ironflock = IronFlock()
await ironflock.start()

# Open a read-only handle on another app's data backend
weather = await ironflock.connect_to_app("weather-app")

# Inspect what the provider shares (non-private tables / transforms)
print([t["tablename"] for t in weather.tables])

# Query history, just like your own tables
rows = await weather.get_history("forecasts", {"limit": 100})

# Subscribe to realtime rows
def on_forecast(*args, **kwargs):
    print("New forecast:", args)

await weather.subscribe_to_table("forecasts", on_forecast)

# Access errors carry a machine-readable code
try:
    await ironflock.connect_to_app("unshared-app")
except CrossAppAccessError as err:
    print(err.code)  # e.g. "NO_GRANT"
```

Consumed-app connections are cached per app + stage and are closed automatically by `ironflock.stop()`.

### `connect_to_app(app_name, stage=None, on_error=None)`

Opens a read-only connection to another app's data backend in the same project and returns a `ConsumedApp` handle. Resolves the provider and connects to its realm using this device's credentials.

```python
weather = await ironflock.connect_to_app("weather-app", stage="prod")
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `app_name` | `str` | Provider app name, as declared in your `consumes:` section |
| `stage` | `str`, optional | Provider stage to connect to (`"dev"` or `"prod"`). Defaults to the stage of the realm this app joined (`ENV=PROD` in any case gives `prod`, anything else `dev`) |
| `on_error` | callable, optional | Called with a `CrossAppAccessError` when the connection is fatally denied *after* `connect_to_app` resolved (e.g. the grant is later revoked) |

**Returns:** `ConsumedApp` — A read-only handle on the provider's data backend.

**Raises:** `CrossAppAccessError` (`code`: `NO_GRANT`, `PROVIDER_NOT_INSTALLED`, `UNKNOWN_APP`, or `NOT_AUTHORIZED`).

---

### `ConsumedApp` handle

Returned by `connect_to_app`. A read-only view of a provider app's shared tables and transforms.

**Properties:**

| Property | Type | Description |
|----------|------|-------------|
| `app` | `str` | Provider app name |
| `stage` | `str` | Provider stage this handle is connected to (`"dev"` or `"prod"`) |
| `tables` | `List[dict]` | Non-private tables the provider shares (dicts with `tablename`, optional `description`/`columns`) |
| `transforms` | `List[dict]` | Non-private transforms (views) the provider shares (same shape as tables) |
| `is_connected` | `bool` | `True` while the connection to the provider is open |
| `connection` | `CrossbarConnection` | The underlying connection (advanced use) |

#### `consumed_app.get_history(tablename, query_params=None)`

Queries history rows of a shared table or transform. Takes the same query parameters as `getHistory` (`limit`, `offset`, `timeRange`, `filterAnd`, `columns`) — including the `{"latest": True}` marker in `filterAnd` for reading the provider's current values.

```python
rows = await weather.get_history("forecasts", {"limit": 100})
current = await weather.get_history("forecasts", {
    "limit": 100,
    "filterAnd": [{"latest": True}],
})
```

Columns the provider marks `secret: true` (see [Secret Columns](#secret-columns)) are never readable across apps — a non-null value arrives as `SECRET_PLACEHOLDER`, and only the owning app's own containers can decrypt it. Filtering on or selecting such a column raises `CrossAppAccessError` with code `SECRET_COLUMN` client-side rather than returning an empty or masked result.

**Returns:** `Any` — The query result rows.

#### `consumed_app.subscribe_to_table(tablename, handler, options=None)`

Subscribes to realtime rows of a shared table or transform. Bulk-inserted rows are delivered one at a time, exactly like `subscribe_to_table` on your own app.

```python
def on_forecast(*args, **kwargs):
    print("New forecast:", args)

await weather.subscribe_to_table("forecasts", on_forecast)
```

**Returns:** `Any` — The subscription object.

#### `consumed_app.get_series_history(tablename, params)`

Queries down-sampled time-series history of a shared **table** (not available for transforms).

```python
series = await weather.get_series_history("forecasts", {
    "metrics": ["temperature"],
    "method": "AVG",
    "limit": 500,
    "timeRange": ["2026-01-01T00:00:00Z", "2026-03-01T00:00:00Z"],
})
```

**params fields (`SeriesQueryParams`):**

| Field | Type | Description |
|-------|------|-------------|
| `metrics` | `List[str]` | Numeric columns to down-sample |
| `method` | `str` | Aggregation per namespace: `"AVG"`, `"SUM"`, `"COUNT"`, `"MIN"`, `"MAX"`, `"FIRST"` or `"LAST"` |
| `limit` | `int` | Maximum number of rows (1–10000) |
| `timeRange` | `list` | `[start, end]` — ISO strings or epoch-ms numbers; `None` = open end (required) |
| `groupBy` | `List[str]`, optional | Columns to group the series by |
| `filterAnd` | `list`, optional | AND filter conditions (WHERE predicates only — no `{"latest": True}` marker) |

**Returns:** `Any` — The down-sampled series rows.

#### `consumed_app.close()`

Closes this handle's connection to the provider. (All consumed-app connections are also closed by `ironflock.stop()`.)

**Returns:** `None`

---

### `CrossAppAccessError`

Raised by `connect_to_app` and the `ConsumedApp` methods when cross-app access is denied or misused. Exposes a machine-readable `code`:

| `code` | Meaning |
|--------|---------|
| `NO_GRANT` | The project user has not granted your app access to the provider |
| `PROVIDER_NOT_INSTALLED` | The provider app has no data backend for that stage in this project |
| `UNKNOWN_APP` | No app by that name |
| `PRIVATE_TABLE` | The requested table/transform is not in the provider's shared catalog |
| `SECRET_COLUMN` | The query filters on or projects a column the provider marks `secret` — its values are encrypted per row and never readable by a consumer |
| `NOT_AUTHORIZED` | The router/provider denied access (e.g. the grant was revoked) |

## Managed File Storage

Every app databackend gets private object storage alongside its tables. With no
`files:` section in the data template you still get one namespace named `default`,
so this works against apps released before FleetFiles existed.

```python
# Store an object and get a permanent URL back in the same call
info = await ironflock.files.put("part-1.jpg", jpeg_bytes, content_type="image/jpeg")

# That URL is safe to put in a FleetDB column — a dashboard widget can render
# <img src="{{photo_url}}"> and it just works
await ironflock.publish_to_table("inspections", part_id="1", photo_url=info.url)

data = await ironflock.files.get("part-1.jpg")
async for obj in ironflock.files.iter(prefix="2026/"):
    print(obj.key, obj.size)
```

The URL never expires, yet stays readable **only** to an authenticated requestor
holding READ on this databackend — an auth proxy re-checks on every request, so
it is safe to store but not a public link.

A **namespace** is a key prefix that carries policy — retention, sharing,
allowed content types. It is not a separate S3 bucket; every namespace lives in
the app's one bucket. Declare one only when a set of objects needs *different
rules*; otherwise stay in `default` and organise with key paths.

Declare additional namespaces in `data-template.yml`:

```yaml
files:
  # Budget the app suggests for itself; the project user's setting is what gets
  # enforced. catalog() reports both as quota_bytes and suggested_quota_bytes.
  quotaBytes: 5368709120
  namespaces:
    - name: frames
      description: Raw camera frames, one JPEG per inspected part.
      contentTypes: ["image/jpeg"]
      maxObjectBytes: 20971520
      retention: { deleteAfter: 30 days }
```

The budget is app-wide, not per namespace — a namespace is only a key prefix inside the app's one storage area, so there is nothing for a per-prefix budget to be enforced against. `maxObjectBytes` *is* per namespace: it caps a single object, not a total.

### `files.put(key, data, namespace="default", content_type=None)`

Stores bytes. Returns `ObjectInfo` with `key`, `size`, `etag`, `content_type`
and `url`.

### `files.get(key, namespace="default")`

Returns `bytes`.

### `files.list(prefix="", namespace="default", limit=200, cursor="")`

One page: `.objects`, `.prefixes`, `.is_truncated`, `.cursor`.

### `files.iter(prefix="", namespace="default")`

Async iterator over every object under a prefix, paginating for you.

### `files.stat(key)` / `files.exists(key)`

Metadata without transferring; `exists` returns a bool.

### `files.delete(key)` / `files.copy(key, to)` / `files.move(key, to)`

`move` is copy-then-delete on the client (the service has no move verb), so it
is **not atomic** — a failed delete leaves both copies.

### `files.put_file(key, path)` / `files.get_to_file(key, path)`

Convenience wrappers over a local file.

### `files.url(key, namespace="default", version=None)`

The permanent URL. Returns `None` where the deployment has no HTTP edge (a
plain-HTTP appliance) — the signal to fall back to `files.get()`. Pass an ETag as
`version` to let browsers cache it immutably.

### `files.usage(detail=False)`

What the object store reports, in one call: `size_bytes`, `object_count`,
`quota_bytes` (the **enforced** budget) and `free_bytes`. `free_bytes` is `-1`
when there is no quota — a quota of `0` means unlimited, and `0` free would read
as full.

Pass `detail=True` for a `per_namespace` breakdown. That costs one listing per
namespace (the store accounts per bucket; a namespace is a prefix), so it is off
by default.

### `files.namespaces()` / `files.catalog()`

What this app may use. `catalog()` reports `quota_bytes` (**enforced**, read from
the object store) alongside `suggested_quota_bytes` (what the app's data template
asked for) — the quota is a user setting, so those two can differ. `catalog()` also reports
`inline_max_bytes` and `public_base_url`; it is cached after the first call.

### `files.share_url(key, namespace="default", ttl=900)`

An expiring link anyone holding it can fetch. Unlike `files.url()` this is a
**bearer capability** — nothing re-checks authorization when it is used. Hand it
to a person; do not store it in a column. The server clamps `ttl`.

### `files.upload_url(key, namespace="default", ttl=3600, content_type=None, size=None)`

An expiring URL that accepts a direct upload. Returns `url`, `method`, `headers`
and `expires_in`; send exactly those headers or the signature will not verify.

### Large objects

The SDK picks the transport by **size**, automatically:

| Size | Path |
| --- | --- |
| ≤ `inline_max_bytes` (6 MiB) | one WAMP call |
| larger | direct to object storage over HTTPS, bypassing the router |

`put_file()` and `get_to_file()` stream from and to disk on the direct path, so a
multi-gigabyte object never has to fit in memory. `put()` and `get()` work on
`bytes` and therefore do hold the object in memory — prefer the file variants
above a few hundred MB.

Two ceilings remain, and both raise `TOO_LARGE` with a reason that says which
one you hit:

- **5 GiB** — S3's single-upload limit. Multipart is not implemented yet.
- **`inline_max_bytes` where there is no direct endpoint** — an air-gapped
  appliance cannot transfer a large object at all. The message says so, because
  no retry or smaller chunk can help.

The direct path needs the device to reach the object store host, not just the
router. Two field failures have their own codes rather than looking like auth
problems: `PRESIGN_UNREACHABLE` (a proxy allowing only the router) and
`CLOCK_SKEW` (S3 rejects requests more than 15 minutes out of step — check NTP).

### `FileStoreError`

Every failure raises `FileStoreError` with a stable `.code` and a human-readable
`.reason`. **Branch on `.code`, never on `.reason`.**

```python
from ironflock.filestore import FileStoreError

try:
    await ironflock.files.put("huge.bin", payload)
except FileStoreError as e:
    if e.code == "QUOTA_EXCEEDED":
        ...
```

| Code | Meaning |
| --- | --- |
| `NOT_AUTHORIZED` | Caller may not perform this operation |
| `NO_SUCH_NAMESPACE` | Namespace is not declared in the data template |
| `NO_SUCH_OBJECT` | Key does not exist |
| `TOO_LARGE` | Exceeds the single-call transfer limit |
| `OBJECT_TOO_LARGE` | Exceeds the namespace's own `maxObjectBytes` |
| `QUOTA_EXCEEDED` | Filestore is full |
| `CONTENT_TYPE_NOT_ALLOWED` | Namespace restricts `contentTypes` |
| `NOT_SUPPORTED` | Backend cannot do this |
| `NOT_AVAILABLE` | No file service on this deployment |
| `PRESIGN_UNREACHABLE` | Object store not reachable directly (proxy?) |
| `CLOCK_SKEW` | Device clock too far out of step for S3 |
| `INTERNAL` | Anything else |

A newer server may add codes; unknown ones pass through as `.code` rather than
crashing, so treat anything unrecognised as a generic failure.

## Advanced Usage

If you need more control, `ironflock.connection` is the underlying `CrossbarConnection`;
[crossbar_connection_example.py](https://github.com/RecordEvolution/ironflock-py/tree/main/examples/crossbar_connection_example.py)
shows it used directly. The rules in [URIs an app may use](#uris-an-app-may-use) apply to it as well.


## Development

This project uses [uv](https://docs.astral.sh/uv/) for dependency management and building.

Install uv if you don't have it:

```shell
curl -LsSf https://astral.sh/uv/install.sh | sh
```

Install dependencies (including dev dependencies):

```shell
uv sync --extra dev
```

Run tests:

```shell
just test-unit    # Run unit tests only
just test         # Run all tests
just test-docker  # Run integration tests with Docker
```

Build and publish a new pypi package:

```shell
just publish
```

Or manually:

```shell
# Clean previous builds
rm -rf dist

# Build the package
uv build

# Upload to PyPI
uv publish
```

Check the package at https://pypi.org/project/ironflock/.

## Test Deployment

To test the package before deploying to PyPI you can use test.pypi.

```shell
just publish-test
```

Or manually:

```shell
uv build
uv publish --publish-url https://test.pypi.org/legacy/
```

Once the package is published you can install it from TestPyPI:

```shell
pip install --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ ironflock
```

Once the package is published you can use it in other code by putting
these lines at the top of the requirements.txt

```
--index-url https://test.pypi.org/simple/
--extra-index-url https://pypi.org/simple/
```

