165 lines
9.1 KiB
Markdown
165 lines
9.1 KiB
Markdown
# OPC UA communication (Laborious)
|
|
|
|
Laborious exports predictions to OPC UA servers through `OpcRepository` ([`laborious/utils/repository/opc_repository.py`](../laborious/utils/repository/opc_repository.py)) and the synchronous Temporal activity layer in [`laborious/activities/opc.py`](../laborious/activities/opc.py). The repository uses `asyncua.sync.Client` (asyncio on a background thread) so activities remain blocking without `async def`.
|
|
|
|
OPC reconnect, write error classification (`opc_error_kind`), and activity confidence/comment behavior are converted from the **async** implementation on `main` at `fcc8920a8be4` (`asyncua.Client` + `asyncio` reconnect task → `threading` reconnect thread). Re-convert with `scripts/convert_opc_async_to_sync.py` when `main` OPC files change.
|
|
|
|
Implementation plan for session/channel recovery on Tier-1 `Bad*` errors: [`.cursor/plans/opc_bad_reconnect_ac4c6045.plan.md`](../.cursor/plans/opc_bad_reconnect_ac4c6045.plan.md).
|
|
|
|
## Architecture
|
|
|
|
```text
|
|
Worker (long-lived)
|
|
└── OpcRepository per OPC server id (from OPC_CONFIG / env)
|
|
├── connect / disconnect / validate_connection (read-only)
|
|
├── _connect_locked / _reconnect_locked (under _connection_lock)
|
|
├── write_data (single attempt per call)
|
|
└── background reconnect on Tier-1 Bad*, closed protocol, or stale session
|
|
|
|
Temporal activity write_opc_data
|
|
└── OPC.manage_output_tags → write_data per tag (sequential per activity)
|
|
```
|
|
|
|
One worker process holds one `OpcRepository` instance per configured server. Multiple Temporal activities can call `write_data` concurrently on the same repository.
|
|
|
|
## Connection lifecycle
|
|
|
|
| Phase | Behavior |
|
|
|-------|----------|
|
|
| Startup | `init_opc()` creates repositories and calls `connect()` → `_connect_locked()` |
|
|
| Steady state | `validate_connection()` is read-only (`protocol.state` only); `_session_ready` is checked in `write_data` |
|
|
| Tier-1 Bad* / protocol closed / session not ready | `_start_reconnect(reason)` → `_run_reconnect` (thread) → `_reconnect_locked()` (respects `reconnection_interval`) |
|
|
| Write | `write_data()` checks in-flight reconnect thread, `_session_ready`, validates, then one `get_node` + `write_value` |
|
|
| Shutdown | `disconnect()` sets `_allow_reconnect = False`, then tears down session |
|
|
|
|
### Session and channel timeouts
|
|
|
|
Requested session and secure-channel lifetime: **10 minutes** (`OPC_UA_SESSION_AND_CHANNEL_TIMEOUT_MS` in `opc_repository.py`). The server may revise these values; negotiated values are logged after connect and exposed as `opc_session_revised_timeout_milliseconds`.
|
|
|
|
### Reconnection interval
|
|
|
|
`OPC_RECONNECTION_INTERVAL` is in **seconds** (default `120`). It gates **background** reconnect after Tier-1 `Bad*` (`last_reconnection_time` is updated only in `_reconnect_locked()`). It limits load on the OPC server when many workflows fail at once.
|
|
|
|
## Concurrency: connection lock and session readiness
|
|
|
|
To allow **multiple concurrent writes** when the session is healthy, but **block all writes** while the connection is being torn down or re-established:
|
|
|
|
| Primitive | Role |
|
|
|-----------|------|
|
|
| `_connection_lock` (`threading.Lock`) | Held for the entire `disconnect` → `connect` path. Only one connection-maintenance task at a time. |
|
|
| `_session_ready` (`threading.Event`) | Set when a session is ready for writes; cleared before reconnect starts and set again after a successful connect. |
|
|
| `_allow_reconnect` | Cleared in `disconnect()` so shutdown does not spawn reconnect threads |
|
|
|
|
**Connection methods (caller holds `_connection_lock` for `_*_locked` helpers):**
|
|
|
|
| Method | Role |
|
|
|--------|------|
|
|
| `_create_client()` | Create asyncua `Client` + optional `set_security`; raises if `client` already exists |
|
|
| `_open_session()` | `client.connect()` + metrics; raises if session already open or client missing |
|
|
| `_connect_locked()` | `_create_client()` (when needed) + `_open_session()`; raises if already connected |
|
|
| `_disconnect_locked()` | Teardown session and clear `client` |
|
|
| `_reconnect_locked()` | `_disconnect_locked()` + `_connect_locked()`; sets `last_reconnection_time` |
|
|
|
|
Public `connect()` / `disconnect()` acquire the lock and call `_connect_locked()` / `_disconnect_locked()`.
|
|
|
|
**Write path (`write_data`):**
|
|
|
|
1. If a reconnect **thread** is alive → `reconnect_in_progress`.
|
|
2. If `_session_ready` is cleared → schedule `SessionNotReady` reconnect; return `reconnect_in_progress` or `connection_lost`.
|
|
3. If `validate_connection()` fails (protocol closed) → schedule `ProtocolClosed` reconnect; return `connection_lost`.
|
|
4. Single `get_node` + `write_value` (no retry in the same call).
|
|
|
|
**Reconnect path (`_run_reconnect`):**
|
|
|
|
1. `_start_reconnect` clears `_session_ready` and starts a daemon thread when the interval allows and `_allow_reconnect` is true.
|
|
2. `with _connection_lock:` → `_reconnect_locked()`.
|
|
3. `_session_ready` is set on successful `_open_session()`.
|
|
|
|
A second `_connect_locked()` while a session is already open raises `OpcSessionAlreadyConnectedError` (disconnect first).
|
|
|
|
**asyncua note:** Concurrent `write_value` on the same session is only safe if the stack tolerates it. If production shows issues, serialize writes while keeping the connection lock semantics above.
|
|
|
|
## Tier-1 `Bad*` errors and reconnect
|
|
|
|
When the server invalidates the session (e.g. `BadSessionIdInvalid`) but the client still sees transport as open, `write_data` fails once, records the OPC status in metrics, and **schedules** reconnect if:
|
|
|
|
- The exception is a `UaStatusCodeError` whose name is in `RECONNECTABLE_OPC_BAD_NAMES` (see plan), and
|
|
- `reconnection_interval` has elapsed since `last_reconnection_time`, and
|
|
- No reconnect task is already running.
|
|
|
|
There is **no write retry**: the failed export is not sent again in the same activity.
|
|
|
|
## Prediction confidence and PostgreSQL comments
|
|
|
|
| `prediction_confidence` | Meaning |
|
|
|-------------------------|---------|
|
|
| (unchanged) | Successful OPC export |
|
|
| **12** | Generic OPC write failure (`OPC_WRITTING_ERROR_CONFIDENCE`) |
|
|
| **14** | Tier-1 session/channel `Bad*` on export (`OPC_SESSION_BAD_CONFIDENCE`) |
|
|
| **14** | Write while reconnect in progress (`OPC_SESSION_BAD_CONFIDENCE`, comment `OPC UA reconnect in progress`) |
|
|
| **13** | PI Web API write failure (separate path) |
|
|
|
|
Session/channel errors use a stable comment for counting:
|
|
|
|
```text
|
|
OPC UA session/channel error: BadSessionIdInvalid
|
|
```
|
|
|
|
Reconnect-in-progress exports use:
|
|
|
|
```text
|
|
OPC UA reconnect in progress
|
|
```
|
|
|
|
Example SQL:
|
|
|
|
```sql
|
|
SELECT count(*) FROM predictions WHERE prediction_confidence = 14;
|
|
SELECT count(*) FROM predictions WHERE comments LIKE 'OPC UA session/channel error:%';
|
|
```
|
|
|
|
## Prometheus metrics (`opc_*`)
|
|
|
|
Defined in [`laborious/metrics.py`](../laborious/metrics.py). Do not rename in production without a dashboard migration.
|
|
|
|
| Metric | Purpose |
|
|
|--------|---------|
|
|
| `opc_connections_initiated_total` | Connection attempts |
|
|
| `opc_connections_failed_total` | Failed connects |
|
|
| `opc_connection_status` | Gauge 1=connected, 0=disconnected |
|
|
| `opc_session_created_total` | Session established after connect |
|
|
| `opc_session_closed_total` | Disconnect initiated |
|
|
| `opc_session_revised_timeout_milliseconds` | Negotiated session timeout (ms) |
|
|
| `opc_write_attempts_total` | Per write; label `result` = `OK` or exception name |
|
|
| `opc_write_inter_arrival_over_session_timeout_total` | Successful writes spaced longer than revised session timeout |
|
|
|
|
Legacy activity metrics: `laborious_prediction_opc_writing_count`, `laborious_prediction_opc_writing_response_time_monitor`.
|
|
|
|
## Environment variables
|
|
|
|
| Variable | Default | Description |
|
|
|----------|---------|-------------|
|
|
| `OPC_CONFIG` | — | JSON map of server configs (overrides single-server env) |
|
|
| `OPC_ID` | `1` | Server id |
|
|
| `OPC_URL` | `opc.tcp://localhost:4840` | Endpoint |
|
|
| `OPC_SERVER_NAME` | `default_server` | Label for metrics/logs |
|
|
| `OPC_SERVER_URI` | same as URL | Application URI / cert SAN |
|
|
| `OPC_CERT_PATH` | — | Client certificate (secure mode) |
|
|
| `OPC_PRIVATE_KEY_PATH` | — | Client private key |
|
|
| `OPC_SERVER_CERT_PATH` | — | Server certificate |
|
|
| `OPC_RECONNECTION_INTERVAL` | `120` | Minimum seconds between reconnects |
|
|
|
|
## Operations checklist
|
|
|
|
- Correlate `BadSessionIdInvalid` in `opc_write_attempts_total` with `opc_session_closed_total` / `opc_session_created_total` (reconnect may finish after the row is stored with confidence 14).
|
|
- Use confidence **14** and comment prefix for session invalidation rates; use **12** for other OPC failures.
|
|
- Respect `OPC_RECONNECTION_INTERVAL` under parallel load; bursts of confidence 14 are expected until the next successful cycle.
|
|
|
|
## Related tests
|
|
|
|
- Unit: [`tests/laborious/utils/repository/test_opc_repository.py`](../tests/laborious/utils/repository/test_opc_repository.py)
|
|
- Unit: [`tests/laborious/activities/test_opc.py`](../tests/laborious/activities/test_opc.py)
|
|
- E2E (mock OPC): [`e2e/test_predictions_batch_format_export.py`](../e2e/test_predictions_batch_format_export.py)
|
|
- E2E (in-process asyncua server + real `OpcRepository`): [`e2e/test_opc_real_server.py`](../e2e/test_opc_real_server.py) — scenarios 3.1.2, 3.2.2, 3.2.4, 3.2.5
|
|
- Scenarios: [`e2e/scenarios.md`](../e2e/scenarios.md)
|