# 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`. 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)