SIENTIAPDE-1811
Enhance OPC UA communication and metrics tracking - Updated README.md to include new OPC UA Communication section and detailed metrics for session and write diagnostics. - Added new metrics in laborious/metrics.py for tracking OPC UA session states and write attempts. - Refactored OPC activity in laborious/activities/opc.py to handle session errors and improve error reporting. - Updated e2e tests to cover new scenarios for OPC session/channel errors and reconnect handling. - Modified .gitignore to include relatorio files and mlruns directory. - Added ipykernel to requirements-dev.txt for Jupyter notebook support.
This commit is contained in:
2
.gitignore
vendored
2
.gitignore
vendored
@@ -52,3 +52,5 @@ catboost_info/
|
|||||||
.ruff_cache/
|
.ruff_cache/
|
||||||
.mypy_cache/
|
.mypy_cache/
|
||||||
mlruns/
|
mlruns/
|
||||||
|
|
||||||
|
relatorio*
|
||||||
24
README.md
24
README.md
@@ -48,6 +48,7 @@ A comprehensive, Temporal-based ML orchestration system for industrial data proc
|
|||||||
- [Prediction Operation Metrics](#prediction-operation-metrics)
|
- [Prediction Operation Metrics](#prediction-operation-metrics)
|
||||||
- [OPC Export Metrics](#opc-export-metrics)
|
- [OPC Export Metrics](#opc-export-metrics)
|
||||||
- [Data Quality Metrics](#data-quality-metrics)
|
- [Data Quality Metrics](#data-quality-metrics)
|
||||||
|
- [OPC UA Communication](#opc-ua-communication)
|
||||||
- [Configuration](#configuration-1)
|
- [Configuration](#configuration-1)
|
||||||
- [Environment Variables](#environment-variables)
|
- [Environment Variables](#environment-variables)
|
||||||
- [OPC Configuration](#opc-configuration)
|
- [OPC Configuration](#opc-configuration)
|
||||||
@@ -164,7 +165,7 @@ Laborious uses a Temporal-based architecture with strong separation of concerns
|
|||||||
- `connectors_config.py`: Env-driven configuration builders
|
- `connectors_config.py`: Env-driven configuration builders
|
||||||
- `models/minio_dataframe_payload.py`: MinIO-offloaded DataFrame payload model
|
- `models/minio_dataframe_payload.py`: MinIO-offloaded DataFrame payload model
|
||||||
- `repository/model_repository.py`: MLFlow operations and retraining
|
- `repository/model_repository.py`: MLFlow operations and retraining
|
||||||
- `repository/opc_repository.py`: OPC communication and writes
|
- `repository/opc_repository.py`: OPC UA client, writes, session recovery (see [OPC UA Communication](#opc-ua-communication))
|
||||||
- `repository/minio_manager.py`: MinIO object storage operations
|
- `repository/minio_manager.py`: MinIO object storage operations
|
||||||
- `filters/conditional_filters.py` and `filters/mlflow_filters.py`
|
- `filters/conditional_filters.py` and `filters/mlflow_filters.py`
|
||||||
|
|
||||||
@@ -779,6 +780,15 @@ The Laborious system exposes comprehensive Prometheus metrics for operational vi
|
|||||||
- Labels: `pod_id`, `model_name`, `workflow_name`, `opc_server_id`
|
- Labels: `pod_id`, `model_name`, `workflow_name`, `opc_server_id`
|
||||||
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
||||||
|
|
||||||
|
OPC UA session and write diagnostics (Prometheus, `laborious/metrics.py`):
|
||||||
|
|
||||||
|
- `opc_connections_initiated_total`, `opc_connections_failed_total`, `opc_connection_status`
|
||||||
|
- `opc_session_created_total`, `opc_session_closed_total`, `opc_session_revised_timeout_milliseconds`
|
||||||
|
- `opc_write_attempts_total` (label `result`: `OK` or exception name, e.g. `BadSessionIdInvalid`)
|
||||||
|
- `opc_write_inter_arrival_over_session_timeout_total`
|
||||||
|
|
||||||
|
See [OPC UA Communication](#opc-ua-communication) for semantics, concurrency, and confidence codes **12** / **14**.
|
||||||
|
|
||||||
### Data Quality Metrics
|
### Data Quality Metrics
|
||||||
- Filter pass/fail rates through notification system
|
- Filter pass/fail rates through notification system
|
||||||
- MLFlow API response validation metrics
|
- MLFlow API response validation metrics
|
||||||
@@ -810,7 +820,7 @@ The Laborious system exposes comprehensive Prometheus metrics for operational vi
|
|||||||
| `OPC_CERT_PATH` | OPC client certificate path | `None` | No |
|
| `OPC_CERT_PATH` | OPC client certificate path | `None` | No |
|
||||||
| `OPC_PRIVATE_KEY_PATH` | OPC private key path | `None` | No |
|
| `OPC_PRIVATE_KEY_PATH` | OPC private key path | `None` | No |
|
||||||
| `OPC_SERVER_CERT_PATH` | OPC server certificate path | `None` | No |
|
| `OPC_SERVER_CERT_PATH` | OPC server certificate path | `None` | No |
|
||||||
| `OPC_RECONNECTION_INTERVAL` | OPC reconnection interval (ms) | `120` | No |
|
| `OPC_RECONNECTION_INTERVAL` | Minimum seconds between OPC reconnects | `120` | No |
|
||||||
| `PI_WEB_API_BASE_URL` | PI Web API server base URL | `None` | No |
|
| `PI_WEB_API_BASE_URL` | PI Web API server base URL | `None` | No |
|
||||||
| `PI_WEB_API_AUTH_TYPE` | PI Web API authentication type (basic/bearer) | `None` | No |
|
| `PI_WEB_API_AUTH_TYPE` | PI Web API authentication type (basic/bearer) | `None` | No |
|
||||||
| `PI_WEB_API_AUTH_TOKEN` | PI Web API authentication token | `None` | No |
|
| `PI_WEB_API_AUTH_TOKEN` | PI Web API authentication token | `None` | No |
|
||||||
@@ -903,6 +913,12 @@ Legacy MinIO object layout (relative key):
|
|||||||
`training_datasets/{model_name}/{object_prefix}_{timestamp}.parquet` where `object_prefix` is sanitized
|
`training_datasets/{model_name}/{object_prefix}_{timestamp}.parquet` where `object_prefix` is sanitized
|
||||||
(slashes replaced by underscores) to keep a stable model-level directory.
|
(slashes replaced by underscores) to keep a stable model-level directory.
|
||||||
|
|
||||||
|
## OPC UA Communication
|
||||||
|
|
||||||
|
Full reference: **[docs/opc-communication.md](docs/opc-communication.md)** (connection lifecycle, Tier-1 `Bad*` reconnect, connection lock / session readiness, metrics, PostgreSQL confidence **12** vs **14**, tests).
|
||||||
|
|
||||||
|
Implementation plan: [`.cursor/plans/opc_bad_reconnect_ac4c6045.plan.md`](.cursor/plans/opc_bad_reconnect_ac4c6045.plan.md).
|
||||||
|
|
||||||
### OPC Configuration
|
### OPC Configuration
|
||||||
|
|
||||||
For multiple OPC servers, use the `OPC_CONFIG` environment variable:
|
For multiple OPC servers, use the `OPC_CONFIG` environment variable:
|
||||||
@@ -1126,8 +1142,10 @@ laborious/
|
|||||||
- Ensure proper connection pool configuration
|
- Ensure proper connection pool configuration
|
||||||
|
|
||||||
4. **OPC Connection Failures**
|
4. **OPC Connection Failures**
|
||||||
- Verify OPC server is accessible
|
- See [docs/opc-communication.md](docs/opc-communication.md)
|
||||||
|
- Verify OPC server is accessible and `OPC_RECONNECTION_INTERVAL` is appropriate
|
||||||
- Check certificate and key file paths
|
- Check certificate and key file paths
|
||||||
|
- Correlate `opc_write_attempts_total` with `opc_session_*` metrics; count session errors via `prediction_confidence = 14`
|
||||||
- Review OPC server logs for connection issues
|
- Review OPC server logs for connection issues
|
||||||
|
|
||||||
5. **PI Web API Connection Failures**
|
5. **PI Web API Connection Failures**
|
||||||
|
|||||||
160
docs/opc-communication.md
Normal file
160
docs/opc-communication.md
Normal file
@@ -0,0 +1,160 @@
|
|||||||
|
# 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 Temporal activity layer in [`laborious/activities/opc.py`](../laborious/activities/opc.py).
|
||||||
|
|
||||||
|
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*
|
||||||
|
|
||||||
|
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* | `_start_reconnect_on_bad` → `_run_reconnect_on_bad` → `_reconnect_locked()` (respects `reconnection_interval`) |
|
||||||
|
| Write | `write_data()` checks `_session_ready`, validates, then one `get_node` + `write_value` |
|
||||||
|
| Shutdown | `close()` disconnects all repositories |
|
||||||
|
|
||||||
|
### 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` (`asyncio.Lock`) | Held for the entire `disconnect` → `connect` path. Only one connection-maintenance task at a time. |
|
||||||
|
| `_session_ready` (`asyncio.Event`) | Set when a session is ready for writes; cleared before reconnect starts and set again after a successful connect. |
|
||||||
|
|
||||||
|
**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 `_session_ready` is cleared → **fail immediately** (`opc_error_kind=reconnect_in_progress`).
|
||||||
|
2. `validate_connection()` checks `protocol.state` only (read-only).
|
||||||
|
3. Single `get_node` + `write_value` (no retry).
|
||||||
|
|
||||||
|
**Reconnect path (`_run_reconnect_on_bad`):**
|
||||||
|
|
||||||
|
1. `_start_reconnect_on_bad` clears `_session_ready` and schedules the task when the interval allows.
|
||||||
|
2. `async 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 with an optional `asyncio.Semaphore(1)` while keeping the connection lock semantics above.
|
||||||
|
|
||||||
|
**Future threads:** replace `asyncio.Lock` / `Event` with `threading` primitives or route all OPC I/O through one dedicated loop.
|
||||||
|
|
||||||
|
## 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: [`e2e/test_predictions_batch_format_export.py`](../e2e/test_predictions_batch_format_export.py), scenarios in [`e2e/scenarios.md`](../e2e/scenarios.md)
|
||||||
@@ -2,7 +2,14 @@
|
|||||||
Pytest configuration and fixtures for E2E tests.
|
Pytest configuration and fixtures for E2E tests.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import sys
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
# E2E workflows under test do not run ModelAnalysis; stub before Activities import.
|
||||||
|
_model_analysis_module = MagicMock()
|
||||||
|
_model_analysis_module.ModelAnalysis = MagicMock
|
||||||
|
sys.modules.setdefault('sientia', MagicMock())
|
||||||
|
sys.modules.setdefault('sientia.ModelAnalysis', _model_analysis_module)
|
||||||
from io import BytesIO
|
from io import BytesIO
|
||||||
|
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
|
|||||||
@@ -64,7 +64,8 @@ def assert_prediction(
|
|||||||
prediction: float = 0.5,
|
prediction: float = 0.5,
|
||||||
prediction_confidence: int | Decimal = 0,
|
prediction_confidence: int | Decimal = 0,
|
||||||
prediction_status: str = 'Good',
|
prediction_status: str = 'Good',
|
||||||
comments: str = '',
|
comments: str | None = None,
|
||||||
|
comments_contains: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""
|
"""
|
||||||
Assert exactly one prediction row exists for model_id with expected columns.
|
Assert exactly one prediction row exists for model_id with expected columns.
|
||||||
@@ -75,7 +76,8 @@ def assert_prediction(
|
|||||||
prediction: Expected prediction value.
|
prediction: Expected prediction value.
|
||||||
prediction_confidence: Expected confidence (int or Decimal for numeric column).
|
prediction_confidence: Expected confidence (int or Decimal for numeric column).
|
||||||
prediction_status: Expected status string.
|
prediction_status: Expected status string.
|
||||||
comments: Expected comments string.
|
comments: Expected exact comments string (optional).
|
||||||
|
comments_contains: Substring expected in comments when queued (optional).
|
||||||
"""
|
"""
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
@@ -98,7 +100,12 @@ def assert_prediction(
|
|||||||
str(prediction_confidence)
|
str(prediction_confidence)
|
||||||
), f'Expected prediction_confidence={prediction_confidence}, got {row[2]}'
|
), f'Expected prediction_confidence={prediction_confidence}, got {row[2]}'
|
||||||
assert row[3] == prediction_status, f"Expected prediction_status='{prediction_status}', got {row[3]}"
|
assert row[3] == prediction_status, f"Expected prediction_status='{prediction_status}', got {row[3]}"
|
||||||
|
if comments is not None:
|
||||||
assert row[4] == comments, f"Expected comments='{comments}', got {row[4]}"
|
assert row[4] == comments, f"Expected comments='{comments}', got {row[4]}"
|
||||||
|
if comments_contains is not None:
|
||||||
|
assert comments_contains in row[4], (
|
||||||
|
f"Expected comments to contain '{comments_contains}', got {row[4]}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def assert_continue(
|
def assert_continue(
|
||||||
|
|||||||
@@ -495,6 +495,48 @@ These paths do **not** rely on Temporal activity retries for export failures: th
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
#### Scenario 3.2.4: OPC Session / Channel Bad* (Tier-1)
|
||||||
|
**Description**: OPC write fails with a Tier-1 session or channel status (e.g. `BadSessionIdInvalid`) while transport may still appear open on the client
|
||||||
|
|
||||||
|
**Input**:
|
||||||
|
- Valid prediction and OPC output config
|
||||||
|
- Mock or server returning Tier-1 `UaStatusCodeError` on write (no write retry in the same activity)
|
||||||
|
|
||||||
|
**Expected Behavior**:
|
||||||
|
- `write_opc_data` fails forward for affected tags; background reconnect may be scheduled if `OPC_RECONNECTION_INTERVAL` allows
|
||||||
|
- Workflow **completes**
|
||||||
|
- PostgreSQL row uses **`prediction_confidence` 14** and comment prefix `OPC UA session/channel error:` (including OPC status name)
|
||||||
|
- `opc_write_attempts_total` records `result=BadSessionIdInvalid` (or matching status); no second write attempt in the same activity
|
||||||
|
|
||||||
|
**Assertions**:
|
||||||
|
- Workflow completes
|
||||||
|
- `prediction_confidence = 14`
|
||||||
|
- `comments` matches `OPC UA session/channel error:%`
|
||||||
|
- Generic OPC error confidence **12** is not used for this case
|
||||||
|
|
||||||
|
**Reference**: [docs/opc-communication.md](../docs/opc-communication.md), plan `.cursor/plans/opc_bad_reconnect_ac4c6045.plan.md`
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
#### Scenario 3.2.5: OPC Write Blocked During Reconnect
|
||||||
|
**Description**: A write is attempted while the repository is reconnecting (session not ready)
|
||||||
|
|
||||||
|
**Input**:
|
||||||
|
- Valid prediction
|
||||||
|
- Simulated slow reconnect (e.g. delayed `connect`) or concurrent writes where the first triggers reconnect
|
||||||
|
|
||||||
|
**Expected Behavior**:
|
||||||
|
- Second write (or parallel write) is rejected **immediately** when reconnect is in progress or `_session_ready` is cleared — **without** calling `write_value`
|
||||||
|
- No wait/sleep on the write path; no duplicate `connect` from parallel writers (connection lock)
|
||||||
|
- `prediction_confidence = 14`, `comments = OPC UA reconnect in progress` (distinguish from Tier-1 `Bad*` via comment prefix in SQL)
|
||||||
|
|
||||||
|
**Assertions**:
|
||||||
|
- At most one reconnect sequence (`disconnect` + `connect`) for the overlapping window
|
||||||
|
- No write retry after failure
|
||||||
|
- Tests in `test_opc_repository` (unit) and optional e2e in `test_predictions_batch_format_export.py`
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
#### Scenario 3.2.3: PI Web API Partial Write Error
|
#### Scenario 3.2.3: PI Web API Partial Write Error
|
||||||
**Description**: Two prediction tags attempt to be written to PI Web API, but only one succeeds
|
**Description**: Two prediction tags attempt to be written to PI Web API, but only one succeeds
|
||||||
|
|
||||||
|
|||||||
@@ -153,7 +153,6 @@ async def test_scenario_3_1_1_default_prediction_export(
|
|||||||
'addr_1',
|
'addr_1',
|
||||||
0,
|
0,
|
||||||
'float',
|
'float',
|
||||||
ANY,
|
|
||||||
{
|
{
|
||||||
'model_id': 311,
|
'model_id': 311,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
@@ -165,7 +164,6 @@ async def test_scenario_3_1_1_default_prediction_export(
|
|||||||
'addr_2',
|
'addr_2',
|
||||||
2,
|
2,
|
||||||
'float',
|
'float',
|
||||||
ANY,
|
|
||||||
{
|
{
|
||||||
'model_id': 311,
|
'model_id': 311,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
@@ -252,20 +250,28 @@ async def test_scenario_3_1_2_export_with_opc_only(
|
|||||||
opc_write_data = cast(Any, test_activities.opc_repository['1'].write_data)
|
opc_write_data = cast(Any, test_activities.opc_repository['1'].write_data)
|
||||||
opc_write_data.assert_has_calls(
|
opc_write_data.assert_has_calls(
|
||||||
[
|
[
|
||||||
call('addr_1', 0.5, 'float', ANY,
|
call(
|
||||||
|
'addr_1',
|
||||||
|
0.5,
|
||||||
|
'float',
|
||||||
{
|
{
|
||||||
'model_id': 312,
|
'model_id': 312,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
'schedule_name': 'test-schedule',
|
'schedule_name': 'test-schedule',
|
||||||
'workflow_name': 'predictions_batch',
|
'workflow_name': 'predictions_batch',
|
||||||
}),
|
},
|
||||||
call('addr_2', 0, 'float', ANY,
|
),
|
||||||
|
call(
|
||||||
|
'addr_2',
|
||||||
|
0,
|
||||||
|
'float',
|
||||||
{
|
{
|
||||||
'model_id': 312,
|
'model_id': 312,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
'schedule_name': 'test-schedule',
|
'schedule_name': 'test-schedule',
|
||||||
'workflow_name': 'predictions_batch',
|
'workflow_name': 'predictions_batch',
|
||||||
}),
|
},
|
||||||
|
),
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -500,20 +506,28 @@ async def test_scenario_3_1_5_export_without_transformed_data(
|
|||||||
opc_write_data = cast(Any, test_activities.opc_repository['1'].write_data)
|
opc_write_data = cast(Any, test_activities.opc_repository['1'].write_data)
|
||||||
opc_write_data.assert_has_calls(
|
opc_write_data.assert_has_calls(
|
||||||
[
|
[
|
||||||
call('addr_1', 0.5, 'float', ANY,
|
call(
|
||||||
|
'addr_1',
|
||||||
|
0.5,
|
||||||
|
'float',
|
||||||
{
|
{
|
||||||
'model_id': 315,
|
'model_id': 315,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
'schedule_name': 'test-schedule',
|
'schedule_name': 'test-schedule',
|
||||||
'workflow_name': 'predictions_batch',
|
'workflow_name': 'predictions_batch',
|
||||||
}),
|
},
|
||||||
call('addr_2', 0, 'float', ANY,
|
),
|
||||||
|
call(
|
||||||
|
'addr_2',
|
||||||
|
0,
|
||||||
|
'float',
|
||||||
{
|
{
|
||||||
'model_id': 315,
|
'model_id': 315,
|
||||||
'model_name': 'test_model',
|
'model_name': 'test_model',
|
||||||
'schedule_name': 'test-schedule',
|
'schedule_name': 'test-schedule',
|
||||||
'workflow_name': 'predictions_batch',
|
'workflow_name': 'predictions_batch',
|
||||||
}),
|
},
|
||||||
|
),
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -648,6 +662,67 @@ async def test_scenario_3_2_2_opc_write_error(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
@pytest.mark.integration
|
||||||
|
async def test_scenario_3_2_4_opc_session_bad_error(
|
||||||
|
temporal_test_env: WorkflowEnvironment,
|
||||||
|
temporal_worker: Worker,
|
||||||
|
test_activities: Activities,
|
||||||
|
postgres_engine,
|
||||||
|
):
|
||||||
|
"""
|
||||||
|
Scenario 3.2.4: OPC session/channel Tier-1 Bad* (e.g. BadSessionIdInvalid).
|
||||||
|
|
||||||
|
PostgreSQL stores prediction_confidence 14 and a stable session error comment.
|
||||||
|
"""
|
||||||
|
client = temporal_test_env.client
|
||||||
|
|
||||||
|
model_id = 324
|
||||||
|
|
||||||
|
insert_sample_data(postgres_engine, model_id, [23.5, 78.2])
|
||||||
|
|
||||||
|
opc_write_data = cast(Any, test_activities.opc_repository['1'].write_data)
|
||||||
|
opc_write_data.return_value = (
|
||||||
|
False,
|
||||||
|
{
|
||||||
|
'notification_id': 'OPC_WRITE_DATA_ERROR_1',
|
||||||
|
'message': 'BadSessionIdInvalid',
|
||||||
|
'block': 'opc_repository',
|
||||||
|
'level': NotificationLevel.ERROR,
|
||||||
|
'attachment_content': 'BadSessionIdInvalid',
|
||||||
|
'opc_error_kind': 'session_bad',
|
||||||
|
'opc_status': 'BadSessionIdInvalid',
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
input_data = get_base_input_data(model_id)
|
||||||
|
input_data['opc_output_config'] = {
|
||||||
|
'1': {
|
||||||
|
'prediction_tags': {
|
||||||
|
'addr_1': {
|
||||||
|
'data_type': 'float',
|
||||||
|
}
|
||||||
|
},
|
||||||
|
'confidence_tags': {
|
||||||
|
'addr_2': {
|
||||||
|
'data_type': 'float',
|
||||||
|
}
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
input_data['pi_web_api_output_config'] = None
|
||||||
|
|
||||||
|
await start_and_await_workflow(
|
||||||
|
client, PredictionsBatch.run, input_data, make_workflow_id('test-opc-session-bad')
|
||||||
|
)
|
||||||
|
|
||||||
|
assert_prediction(
|
||||||
|
postgres_engine,
|
||||||
|
model_id,
|
||||||
|
prediction_confidence=14,
|
||||||
|
comments='OPC UA session/channel error: BadSessionIdInvalid',
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
@pytest.mark.integration
|
@pytest.mark.integration
|
||||||
|
|||||||
25
inter_arrival.py
Normal file
25
inter_arrival.py
Normal file
@@ -0,0 +1,25 @@
|
|||||||
|
# %%
|
||||||
|
|
||||||
|
# Load logs.txt
|
||||||
|
with open('logs.txt', 'r') as file:
|
||||||
|
lines = file.readlines()
|
||||||
|
|
||||||
|
# %%
|
||||||
|
import re
|
||||||
|
# Grep "inter-arrival_s=number" with regex
|
||||||
|
intervals = []
|
||||||
|
for line in lines:
|
||||||
|
match = re.search(r'inter-arrival_s=([0-9.]+)', line)
|
||||||
|
if match:
|
||||||
|
intervals.append(float(match.group(1)))
|
||||||
|
# %%
|
||||||
|
|
||||||
|
print(intervals)
|
||||||
|
# %%
|
||||||
|
import matplotlib.pyplot as plt
|
||||||
|
plt.plot(intervals)
|
||||||
|
plt.ylabel('Inter-arrival time (s)')
|
||||||
|
plt.xlabel('Sample')
|
||||||
|
plt.title('Inter-arrival time distribution')
|
||||||
|
plt.show()
|
||||||
|
# %%
|
||||||
@@ -15,6 +15,16 @@ with workflow.unsafe.imports_passed_through():
|
|||||||
from laborious.utils.repository.opc_repository import OpcRepository
|
from laborious.utils.repository.opc_repository import OpcRepository
|
||||||
|
|
||||||
OPC_WRITTING_ERROR_CONFIDENCE = 12
|
OPC_WRITTING_ERROR_CONFIDENCE = 12
|
||||||
|
OPC_SESSION_BAD_CONFIDENCE = 14
|
||||||
|
OPC_SESSION_BAD_COMMENT_PREFIX = 'OPC UA session/channel error:'
|
||||||
|
OPC_WRITTING_ERROR_MESSAGE = 'Some data could not be written to OPC servers'
|
||||||
|
OPC_RECONNECT_IN_PROGRESS_COMMENT = 'OPC UA reconnect in progress'
|
||||||
|
OPC_COMMENT_SEPARATOR = ' | '
|
||||||
|
|
||||||
|
|
||||||
|
def _opc_session_bad_comment(opc_status: str | None) -> str:
|
||||||
|
status = opc_status or 'Unknown'
|
||||||
|
return f'{OPC_SESSION_BAD_COMMENT_PREFIX} {status}'
|
||||||
|
|
||||||
|
|
||||||
class OPC(SientiaMonitoring):
|
class OPC(SientiaMonitoring):
|
||||||
@@ -118,29 +128,18 @@ class OPC(SientiaMonitoring):
|
|||||||
data_type: str,
|
data_type: str,
|
||||||
tag_type: str,
|
tag_type: str,
|
||||||
metadata: dict[str, Any],
|
metadata: dict[str, Any],
|
||||||
) -> float | None:
|
) -> tuple[float | None, dict[str, Any] | None]:
|
||||||
"""
|
"""
|
||||||
Write data to a specific OPC server tag with comprehensive error handling.
|
Write data to a specific OPC server tag with comprehensive error handling.
|
||||||
|
|
||||||
This method provides a secure and reliable way to write data to OPC servers
|
Return:
|
||||||
with automatic error handling, notification integration, and detailed logging.
|
tuple[float | None, dict[str, Any] | None]: Response time on success, or
|
||||||
It validates server availability before attempting write operations and
|
(None, error info_data) on repository failure.
|
||||||
provides comprehensive error reporting for operational monitoring.
|
|
||||||
|
|
||||||
Args:
|
|
||||||
- server_id (str): The id of the OPC server.
|
|
||||||
- tag (str): The tag to write to.
|
|
||||||
- data (Any): The data to write.
|
|
||||||
- data_type (str): The data type.
|
|
||||||
- tag_type (str): The tag type.
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
- bool: True if the data was written successfully, False otherwise.
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
try:
|
try:
|
||||||
is_success, info_data = await self.opc_repository[server_id].write_data(
|
is_success, info_data = await self.opc_repository[server_id].write_data(
|
||||||
tag, data, data_type, self.logger, metadata
|
tag, data, data_type, metadata
|
||||||
)
|
)
|
||||||
if not is_success:
|
if not is_success:
|
||||||
await self.send_notification_async(
|
await self.send_notification_async(
|
||||||
@@ -151,8 +150,8 @@ class OPC(SientiaMonitoring):
|
|||||||
level=info_data.get('level', NotificationLevel.ERROR),
|
level=info_data.get('level', NotificationLevel.ERROR),
|
||||||
attachment_content=info_data.get('attachment_content', None),
|
attachment_content=info_data.get('attachment_content', None),
|
||||||
)
|
)
|
||||||
return None
|
return None, info_data
|
||||||
return info_data['response_time']
|
return info_data['response_time'], None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
await self.send_notification_async(
|
await self.send_notification_async(
|
||||||
@@ -205,7 +204,7 @@ class OPC(SientiaMonitoring):
|
|||||||
config: dict[str, Any],
|
config: dict[str, Any],
|
||||||
data: DataFrame,
|
data: DataFrame,
|
||||||
metadata: dict[str, Any],
|
metadata: dict[str, Any],
|
||||||
) -> tuple[bool, dict[str, float | None]]:
|
) -> tuple[bool, dict[str, float | None], bool, str | None, bool]:
|
||||||
"""
|
"""
|
||||||
Manage the writing of prediction and confidence data to OPC server tags.
|
Manage the writing of prediction and confidence data to OPC server tags.
|
||||||
|
|
||||||
@@ -234,10 +233,13 @@ class OPC(SientiaMonitoring):
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
response_times: dict[str, float | None] = {}
|
response_times: dict[str, float | None] = {}
|
||||||
|
session_bad_seen = False
|
||||||
|
session_bad_status: str | None = None
|
||||||
|
reconnect_in_progress_seen = False
|
||||||
|
|
||||||
if 'prediction_tags' in config:
|
if 'prediction_tags' in config:
|
||||||
for tag, tag_config in config['prediction_tags'].items():
|
for tag, tag_config in config['prediction_tags'].items():
|
||||||
response_time = await self.write_data(
|
response_time, error_info = await self.write_data(
|
||||||
server_id=server_id,
|
server_id=server_id,
|
||||||
tag=tag,
|
tag=tag,
|
||||||
data=data.head(1)['prediction'].values[0],
|
data=data.head(1)['prediction'].values[0],
|
||||||
@@ -245,6 +247,13 @@ class OPC(SientiaMonitoring):
|
|||||||
tag_type='prediction',
|
tag_type='prediction',
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
)
|
)
|
||||||
|
if error_info:
|
||||||
|
kind = error_info.get('opc_error_kind')
|
||||||
|
if kind == 'session_bad':
|
||||||
|
session_bad_seen = True
|
||||||
|
session_bad_status = error_info.get('opc_status', session_bad_status)
|
||||||
|
elif kind == 'reconnect_in_progress':
|
||||||
|
reconnect_in_progress_seen = True
|
||||||
if response_time is not None:
|
if response_time is not None:
|
||||||
self.info(
|
self.info(
|
||||||
f'Prediction data written to OPC server {server_id} for tag {tag}.',
|
f'Prediction data written to OPC server {server_id} for tag {tag}.',
|
||||||
@@ -254,7 +263,7 @@ class OPC(SientiaMonitoring):
|
|||||||
|
|
||||||
if 'confidence_tags' in config:
|
if 'confidence_tags' in config:
|
||||||
for tag, tag_config in config['confidence_tags'].items():
|
for tag, tag_config in config['confidence_tags'].items():
|
||||||
response_time = await self.write_data(
|
response_time, error_info = await self.write_data(
|
||||||
server_id=server_id,
|
server_id=server_id,
|
||||||
tag=tag,
|
tag=tag,
|
||||||
data=data.head(1)['prediction_confidence'].values[0],
|
data=data.head(1)['prediction_confidence'].values[0],
|
||||||
@@ -262,6 +271,13 @@ class OPC(SientiaMonitoring):
|
|||||||
tag_type='confidence',
|
tag_type='confidence',
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
)
|
)
|
||||||
|
if error_info:
|
||||||
|
kind = error_info.get('opc_error_kind')
|
||||||
|
if kind == 'session_bad':
|
||||||
|
session_bad_seen = True
|
||||||
|
session_bad_status = error_info.get('opc_status', session_bad_status)
|
||||||
|
elif kind == 'reconnect_in_progress':
|
||||||
|
reconnect_in_progress_seen = True
|
||||||
if response_time is not None:
|
if response_time is not None:
|
||||||
self.info(
|
self.info(
|
||||||
f'Confidence data written to OPC server {server_id} for tag {tag}.',
|
f'Confidence data written to OPC server {server_id} for tag {tag}.',
|
||||||
@@ -271,7 +287,13 @@ class OPC(SientiaMonitoring):
|
|||||||
|
|
||||||
success = None not in response_times.values()
|
success = None not in response_times.values()
|
||||||
|
|
||||||
return success, response_times
|
return (
|
||||||
|
success,
|
||||||
|
response_times,
|
||||||
|
session_bad_seen,
|
||||||
|
session_bad_status,
|
||||||
|
reconnect_in_progress_seen,
|
||||||
|
)
|
||||||
|
|
||||||
@activity.defn(name='write_opc_data')
|
@activity.defn(name='write_opc_data')
|
||||||
async def write_opc_data(
|
async def write_opc_data(
|
||||||
@@ -301,6 +323,9 @@ class OPC(SientiaMonitoring):
|
|||||||
self.info(f'Data to write: {data.size} rows', metadata)
|
self.info(f'Data to write: {data.size} rows', metadata)
|
||||||
|
|
||||||
success = True
|
success = True
|
||||||
|
session_bad_seen = False
|
||||||
|
session_bad_status: str | None = None
|
||||||
|
reconnect_in_progress_seen = False
|
||||||
|
|
||||||
metrics: dict[str, dict[str, float | None]] = {}
|
metrics: dict[str, dict[str, float | None]] = {}
|
||||||
|
|
||||||
@@ -309,22 +334,48 @@ class OPC(SientiaMonitoring):
|
|||||||
success = False
|
success = False
|
||||||
continue
|
continue
|
||||||
|
|
||||||
local_success, local_response_times = await self.manage_output_tags(
|
(
|
||||||
server_id, config, data, metadata
|
local_success,
|
||||||
)
|
local_response_times,
|
||||||
|
local_session_bad,
|
||||||
|
local_status,
|
||||||
|
local_reconnect_in_progress,
|
||||||
|
) = await self.manage_output_tags(server_id, config, data, metadata)
|
||||||
metrics[server_id] = local_response_times
|
metrics[server_id] = local_response_times
|
||||||
local_count = len(local_response_times)
|
local_count = len(local_response_times)
|
||||||
success = success and local_success
|
success = success and local_success
|
||||||
|
if local_session_bad:
|
||||||
|
session_bad_seen = True
|
||||||
|
session_bad_status = local_status or session_bad_status
|
||||||
|
if local_reconnect_in_progress:
|
||||||
|
reconnect_in_progress_seen = True
|
||||||
|
|
||||||
self.info(
|
self.info(
|
||||||
f'Process completed for OPC server {server_id}: {local_count} of {len(config.get("prediction_tags", []))} prediction tags and {len(config.get("confidence_tags", []))} confidence tags',
|
f'Process completed for OPC server {server_id}: {local_count} of {len(config.get("prediction_tags", []))} prediction tags and {len(config.get("confidence_tags", []))} confidence tags',
|
||||||
metadata,
|
metadata,
|
||||||
)
|
)
|
||||||
|
|
||||||
return self.process_confidence(data, success, metadata), metrics
|
return (
|
||||||
|
self.process_confidence(
|
||||||
|
data,
|
||||||
|
success,
|
||||||
|
metadata,
|
||||||
|
session_bad=session_bad_seen,
|
||||||
|
opc_status=session_bad_status,
|
||||||
|
reconnect_in_progress=reconnect_in_progress_seen,
|
||||||
|
),
|
||||||
|
metrics,
|
||||||
|
)
|
||||||
|
|
||||||
def process_confidence(
|
def process_confidence(
|
||||||
self, data: DataFrame, success: bool, metadata: dict[str, Any]
|
self,
|
||||||
|
data: DataFrame,
|
||||||
|
success: bool,
|
||||||
|
metadata: dict[str, Any],
|
||||||
|
*,
|
||||||
|
session_bad: bool = False,
|
||||||
|
opc_status: str | None = None,
|
||||||
|
reconnect_in_progress: bool = False,
|
||||||
) -> dict[Hashable, Any]:
|
) -> dict[Hashable, Any]:
|
||||||
"""
|
"""
|
||||||
Process prediction confidence based on OPC write operation success.
|
Process prediction confidence based on OPC write operation success.
|
||||||
@@ -352,16 +403,26 @@ class OPC(SientiaMonitoring):
|
|||||||
This allows downstream systems to handle data quality appropriately.
|
This allows downstream systems to handle data quality appropriately.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
message = 'Some data could not be written to OPC servers'
|
|
||||||
|
|
||||||
if not success:
|
if not success:
|
||||||
data['prediction_confidence'] = OPC_WRITTING_ERROR_CONFIDENCE
|
comment_parts: list[str] = []
|
||||||
data['comments'] = message
|
confidence = OPC_WRITTING_ERROR_CONFIDENCE
|
||||||
|
|
||||||
|
if session_bad:
|
||||||
|
comment_parts.append(_opc_session_bad_comment(opc_status))
|
||||||
|
confidence = OPC_SESSION_BAD_CONFIDENCE
|
||||||
|
if reconnect_in_progress:
|
||||||
|
comment_parts.append(OPC_RECONNECT_IN_PROGRESS_COMMENT)
|
||||||
|
confidence = OPC_SESSION_BAD_CONFIDENCE
|
||||||
|
if not comment_parts:
|
||||||
|
comment_parts.append(OPC_WRITTING_ERROR_MESSAGE)
|
||||||
|
|
||||||
|
comments = OPC_COMMENT_SEPARATOR.join(comment_parts)
|
||||||
|
data['prediction_confidence'] = confidence
|
||||||
|
data['comments'] = comments
|
||||||
self.debug(
|
self.debug(
|
||||||
f'{message}, setting confidence to {OPC_WRITTING_ERROR_CONFIDENCE}.',
|
f'OPC write issues, confidence={confidence}, comments={comments}',
|
||||||
metadata,
|
metadata,
|
||||||
)
|
)
|
||||||
|
|
||||||
else:
|
else:
|
||||||
self.debug('Data written to OPC servers successfully.', metadata)
|
self.debug('Data written to OPC servers successfully.', metadata)
|
||||||
|
|
||||||
|
|||||||
@@ -89,6 +89,40 @@ OPC_CONNECTION_STATUS = Gauge(
|
|||||||
['pod_id', 'server_name', 'server_url'],
|
['pod_id', 'server_name', 'server_url'],
|
||||||
)
|
)
|
||||||
|
|
||||||
|
_OPC_SESSION_DEBUG_LABELS = ['pod_id', 'server_name', 'runtime', 'opc_server_id', 'session_id']
|
||||||
|
|
||||||
|
OPC_SESSION_CREATED_TOTAL = Counter(
|
||||||
|
'opc_session_created_total',
|
||||||
|
'OPC UA sessions established (after successful connect)',
|
||||||
|
_OPC_SESSION_DEBUG_LABELS,
|
||||||
|
)
|
||||||
|
|
||||||
|
OPC_SESSION_CLOSED_TOTAL = Counter(
|
||||||
|
'opc_session_closed_total',
|
||||||
|
'OPC UA client disconnects completed (session tear-down initiated)',
|
||||||
|
_OPC_SESSION_DEBUG_LABELS,
|
||||||
|
)
|
||||||
|
|
||||||
|
OPC_SESSION_REVISED_TIMEOUT_MS = Gauge(
|
||||||
|
'opc_session_revised_timeout_milliseconds',
|
||||||
|
'Server-revised OPC UA session timeout (RevisedSessionTimeout) in ms after connect',
|
||||||
|
_OPC_SESSION_DEBUG_LABELS,
|
||||||
|
)
|
||||||
|
|
||||||
|
OPC_WRITE_ATTEMPT_LABELS = [*_OPC_SESSION_DEBUG_LABELS, 'model_id', 'model_name', 'result']
|
||||||
|
|
||||||
|
OPC_WRITE_ATTEMPTS_TOTAL = Counter(
|
||||||
|
'opc_write_attempts_total',
|
||||||
|
'OPC UA write attempts with session and outcome (result=OK or exception class name)',
|
||||||
|
OPC_WRITE_ATTEMPT_LABELS,
|
||||||
|
)
|
||||||
|
|
||||||
|
OPC_WRITE_INTER_ARRIVAL_OVER_SESSION_TIMEOUT_TOTAL = Counter(
|
||||||
|
'opc_write_inter_arrival_over_session_timeout_total',
|
||||||
|
'Successful writes where seconds since the previous successful write exceeded RevisedSessionTimeout (ms)',
|
||||||
|
_OPC_SESSION_DEBUG_LABELS,
|
||||||
|
)
|
||||||
|
|
||||||
# ================== Model metrics ==================
|
# ================== Model metrics ==================
|
||||||
|
|
||||||
MODEL_READ_LAG = Histogram(
|
MODEL_READ_LAG = Histogram(
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ from typing import Any
|
|||||||
from asyncua import Client
|
from asyncua import Client
|
||||||
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
||||||
from asyncua.ua import DataValue, Variant, VariantType
|
from asyncua.ua import DataValue, Variant, VariantType
|
||||||
|
from asyncua.ua.uaerrors import UaStatusCodeError
|
||||||
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
|
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
|
||||||
from sientia_do.notifications.models import NotificationLevel
|
from sientia_do.notifications.models import NotificationLevel
|
||||||
from sientia_do.observability.logger import Logger
|
from sientia_do.observability.logger import Logger
|
||||||
@@ -17,6 +18,119 @@ from sientia_do.observability.sientia_monitoring import SientiaMonitoring
|
|||||||
|
|
||||||
from laborious import metrics
|
from laborious import metrics
|
||||||
|
|
||||||
|
# Requested session and secure channel lifetime (ms) before server revision; 10 minutes.
|
||||||
|
OPC_UA_SESSION_AND_CHANNEL_TIMEOUT_MS = 10 * 60 * 1000
|
||||||
|
|
||||||
|
|
||||||
|
class OpcClientAlreadyExistsError(RuntimeError):
|
||||||
|
"""Raised when _create_client is called while self.client is already set."""
|
||||||
|
|
||||||
|
|
||||||
|
class OpcSessionAlreadyConnectedError(RuntimeError):
|
||||||
|
"""Raised when _open_session is called while a UA session is already open."""
|
||||||
|
|
||||||
|
|
||||||
|
class OpcClientNotInitializedError(RuntimeError):
|
||||||
|
"""Raised when _open_session is called before _create_client."""
|
||||||
|
|
||||||
|
|
||||||
|
RECONNECTABLE_OPC_BAD_NAMES: frozenset[str] = frozenset(
|
||||||
|
{
|
||||||
|
'BadSessionIdInvalid',
|
||||||
|
'BadSessionClosed',
|
||||||
|
'BadSessionNotActivated',
|
||||||
|
'BadSecureChannelIdInvalid',
|
||||||
|
'BadSecureChannelClosed',
|
||||||
|
'BadSecureChannelTokenUnknown',
|
||||||
|
'BadTcpSecureChannelUnknown',
|
||||||
|
'BadServerNotConnected',
|
||||||
|
'BadConnectionClosed',
|
||||||
|
'BadDisconnect',
|
||||||
|
'BadConnectionRejected',
|
||||||
|
'BadCommunicationError',
|
||||||
|
'BadRequestInterrupted',
|
||||||
|
'BadUnknownResponse',
|
||||||
|
'BadTimeout',
|
||||||
|
'BadRequestTimeout',
|
||||||
|
'BadSequenceNumberInvalid',
|
||||||
|
'BadSequenceNumberUnknown',
|
||||||
|
'BadSecurityModeInsufficient',
|
||||||
|
'BadRequestHeaderInvalid',
|
||||||
|
'BadInvalidState',
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _opc_authentication_token_str(client: Client | None) -> str:
|
||||||
|
"""
|
||||||
|
Serialize the current OPC UA authentication token (session handle) for logging and metrics.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
str: Token string, or "unknown" if unavailable.
|
||||||
|
"""
|
||||||
|
if client is None:
|
||||||
|
return 'unknown'
|
||||||
|
try:
|
||||||
|
proto = client.uaclient.protocol
|
||||||
|
if proto is None:
|
||||||
|
return 'unknown'
|
||||||
|
tok = getattr(proto, 'authentication_token', None)
|
||||||
|
if tok is None:
|
||||||
|
return 'unknown'
|
||||||
|
return str(tok)
|
||||||
|
except Exception:
|
||||||
|
return 'unknown'
|
||||||
|
|
||||||
|
|
||||||
|
def _opc_status_from_exception(exc: BaseException) -> str:
|
||||||
|
"""
|
||||||
|
Resolve OPC UA status name from an exception, including chained UaStatusCodeError causes.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
exc (BaseException): Raised error from asyncua.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
str: Status class name or generic Python exception name.
|
||||||
|
"""
|
||||||
|
current: BaseException | None = exc
|
||||||
|
while current is not None:
|
||||||
|
if isinstance(current, UaStatusCodeError):
|
||||||
|
return type(current).__name__
|
||||||
|
current = current.__cause__
|
||||||
|
return type(exc).__name__
|
||||||
|
|
||||||
|
|
||||||
|
def is_reconnectable_opcua_bad(exc: BaseException) -> bool:
|
||||||
|
"""
|
||||||
|
Return whether the exception is a Tier-1 OPC UA Bad* that should trigger reconnect.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
exc (BaseException): Raised error from get_node or write_value.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
bool: True if reconnect should be scheduled.
|
||||||
|
"""
|
||||||
|
return _opc_status_from_exception(exc) in RECONNECTABLE_OPC_BAD_NAMES
|
||||||
|
|
||||||
|
|
||||||
|
def _model_labels_from_write_metadata(metadata: dict[str, Any] | None) -> dict[str, str]:
|
||||||
|
"""
|
||||||
|
Extract model_id and model_name from write metadata for Prometheus labels.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
metadata (dict[str, Any] | None): Context passed into write_data; may omit keys.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
dict[str, str]: Labels model_id and model_name, defaulting to "unknown".
|
||||||
|
"""
|
||||||
|
if not metadata:
|
||||||
|
return {'model_id': 'unknown', 'model_name': 'unknown'}
|
||||||
|
return {
|
||||||
|
'model_id': str(metadata.get('model_id', 'unknown')),
|
||||||
|
'model_name': str(metadata.get('model_name', 'unknown')),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
data_type_map = {
|
data_type_map = {
|
||||||
'float': {
|
'float': {
|
||||||
'converter': float,
|
'converter': float,
|
||||||
@@ -63,8 +177,6 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
self.cert_path = cert_path
|
self.cert_path = cert_path
|
||||||
self.private_key_path = private_key_path
|
self.private_key_path = private_key_path
|
||||||
self.server_cert_path = server_cert_path
|
self.server_cert_path = server_cert_path
|
||||||
self.logger = logger
|
|
||||||
self.error_count = 0
|
|
||||||
self.reconnection_interval = reconnection_interval
|
self.reconnection_interval = reconnection_interval
|
||||||
self.last_reconnection_time: None | datetime = None
|
self.last_reconnection_time: None | datetime = None
|
||||||
self.disconnection_interval = 10.0
|
self.disconnection_interval = 10.0
|
||||||
@@ -79,27 +191,63 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
'workflow_name': 'opc_repository',
|
'workflow_name': 'opc_repository',
|
||||||
'schedule_name': '-',
|
'schedule_name': '-',
|
||||||
}
|
}
|
||||||
|
self._last_write_mono: float | None = None
|
||||||
|
self._connection_lock = asyncio.Lock()
|
||||||
|
self._session_ready = asyncio.Event()
|
||||||
|
self._reconnect_task: asyncio.Task[None] | None = None
|
||||||
|
|
||||||
async def set_security(self):
|
def _opc_debug_tags(self, session_id: str) -> dict[str, str]:
|
||||||
|
return {
|
||||||
|
'pod_id': str(getattr(self, 'pod_id', 'unknown')),
|
||||||
|
'server_name': self.server_name,
|
||||||
|
'runtime': str(getattr(self, 'runtime', 'unknown')),
|
||||||
|
'opc_server_id': self.id,
|
||||||
|
'session_id': session_id,
|
||||||
|
}
|
||||||
|
|
||||||
|
def _is_session_open(self) -> bool:
|
||||||
"""
|
"""
|
||||||
Configures the security settings for the OPC UA client.
|
Return whether the asyncua client has an open transport session.
|
||||||
This method sets up the security policy, certificates, and timeouts
|
|
||||||
required for establishing a secure connection with the OPC UA server.
|
Return:
|
||||||
|
bool: True when protocol exists and is not closed.
|
||||||
|
"""
|
||||||
|
if self.client is None:
|
||||||
|
return False
|
||||||
|
try:
|
||||||
|
proto = self.client.uaclient.protocol
|
||||||
|
return proto is not None and proto.state != 'closed'
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
|
|
||||||
|
def _reconnection_window_elapsed(self) -> bool:
|
||||||
|
"""
|
||||||
|
Return whether enough time has passed since the last reconnect attempt.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
bool: True if a new reconnect is allowed.
|
||||||
|
"""
|
||||||
|
if self.last_reconnection_time is None:
|
||||||
|
return True
|
||||||
|
return (
|
||||||
|
datetime.now() - self.last_reconnection_time
|
||||||
|
).total_seconds() > self.reconnection_interval
|
||||||
|
|
||||||
|
def _not_connected_error(self) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
'notification_id': f'OPC_CONNECTION_NOT_READY_{self.id}',
|
||||||
|
'message': f'OPC server {self.id} is not connected',
|
||||||
|
'block': 'opc_repository',
|
||||||
|
'level': NotificationLevel.WARNING,
|
||||||
|
}
|
||||||
|
|
||||||
|
async def set_security(self) -> None:
|
||||||
|
"""
|
||||||
|
Configure certificates and timeouts on the asyncua client.
|
||||||
|
|
||||||
Raises:
|
Raises:
|
||||||
ValueError: If either the certificate path or private key path is not provided.
|
ValueError: If cert paths or client are missing.
|
||||||
Attributes:
|
|
||||||
- cert_path (str): Path to the client's certificate file.
|
|
||||||
- private_key_path (str): Path to the client's private key file.
|
|
||||||
- server_cert_path (str, optional): Path to the server's certificate file.
|
|
||||||
- server_uri (str): The URI of the server to be used as the application URI.
|
|
||||||
- client (opcua.Client): The OPC UA client instance.
|
|
||||||
- logger (logging.Logger): Logger instance for logging information.
|
|
||||||
Security Settings:
|
|
||||||
- Security Policy: Basic256
|
|
||||||
- Secure Channel Timeout: 10,000,000 ms
|
|
||||||
- Session Timeout: 10,000,000 ms
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
if self.cert_path is None or self.private_key_path is None:
|
if self.cert_path is None or self.private_key_path is None:
|
||||||
raise ValueError(
|
raise ValueError(
|
||||||
'Certificate and private key paths must be provided for secure connection.'
|
'Certificate and private key paths must be provided for secure connection.'
|
||||||
@@ -113,91 +261,107 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
raise ValueError('Client must be initialized before setting security')
|
raise ValueError('Client must be initialized before setting security')
|
||||||
|
|
||||||
self.client.application_uri = self.server_uri
|
self.client.application_uri = self.server_uri
|
||||||
self.logger.custom_info('Setting security...', self.metadata)
|
self.info('Setting security...', self.metadata)
|
||||||
await self.client.set_security(
|
await self.client.set_security(
|
||||||
SecurityPolicyBasic256,
|
SecurityPolicyBasic256,
|
||||||
certificate=str(cert),
|
certificate=str(cert),
|
||||||
private_key=str(private_key),
|
private_key=str(private_key),
|
||||||
server_certificate=str(server_cert) if server_cert else None,
|
server_certificate=str(server_cert) if server_cert else None,
|
||||||
)
|
)
|
||||||
self.client.secure_channel_timeout = 10000000
|
self.client.secure_channel_timeout = OPC_UA_SESSION_AND_CHANNEL_TIMEOUT_MS
|
||||||
self.client.session_timeout = 10000000
|
self.client.session_timeout = OPC_UA_SESSION_AND_CHANNEL_TIMEOUT_MS
|
||||||
|
|
||||||
async def connect(self) -> tuple[bool, dict[str, Any]]:
|
async def _create_client(self) -> None:
|
||||||
"""
|
"""
|
||||||
Establishes a connection to the OPC server.
|
Instantiate the asyncua Client and apply security when configured.
|
||||||
This method initializes the OPC client using the provided URL and
|
|
||||||
sets up security if a certificate path is specified. It then
|
Caller must hold _connection_lock. Does not open a UA session.
|
||||||
attempts to connect to the server and logs the connection status.
|
|
||||||
Raises:
|
Raises:
|
||||||
Exception: If the connection to the OPC server fails.
|
OpcClientAlreadyExistsError: If self.client is already set.
|
||||||
"""
|
"""
|
||||||
|
if self.client is not None:
|
||||||
|
raise OpcClientAlreadyExistsError(
|
||||||
|
f'OPC client already exists for server {self.id}; '
|
||||||
|
'call disconnect() before creating a new client'
|
||||||
|
)
|
||||||
|
|
||||||
self.client = Client(self.url, timeout=10, watchdog_intervall=3600000) # type: ignore[attr-defined]
|
self.client = Client(self.url, timeout=10, watchdog_intervall=50) # type: ignore[attr-defined]
|
||||||
|
|
||||||
self.client.name = self.pod_id
|
self.client.name = self.pod_id
|
||||||
self.client.application_name = self.pod_id
|
self.client.application_name = self.pod_id
|
||||||
pod_uri = self.pod_id.replace('-', ':')
|
pod_uri = self.pod_id.replace('-', ':')
|
||||||
self.client.application_uri = pod_uri
|
self.client.application_uri = pod_uri
|
||||||
self.client.product_uri = pod_uri
|
self.client.product_uri = pod_uri
|
||||||
|
|
||||||
if self.cert_path:
|
if self.cert_path:
|
||||||
await self.set_security()
|
await self.set_security()
|
||||||
self.logger.custom_info(
|
|
||||||
f'Starting connection to OPC server {self.id}:{self.server_name}...', self.metadata
|
async def _open_session(self) -> tuple[bool, dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Open the OPC UA session on the existing client.
|
||||||
|
|
||||||
|
Caller must hold _connection_lock.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
OpcClientNotInitializedError: If self.client is None.
|
||||||
|
OpcSessionAlreadyConnectedError: If a session is already open.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
tuple[bool, dict[str, Any]]: Success flag and error payload on connect failure.
|
||||||
|
"""
|
||||||
|
if self.client is None:
|
||||||
|
raise OpcClientNotInitializedError(
|
||||||
|
f'OPC client is not initialized for server {self.id}; '
|
||||||
|
'call _create_client() before opening a session'
|
||||||
|
)
|
||||||
|
if self._is_session_open():
|
||||||
|
raise OpcSessionAlreadyConnectedError(
|
||||||
|
f'OPC session already connected for server {self.id}; '
|
||||||
|
'call disconnect() before connecting again'
|
||||||
)
|
)
|
||||||
return await self.try_connect()
|
|
||||||
|
|
||||||
async def try_connect(self) -> tuple[bool, dict[str, Any]]:
|
|
||||||
"""
|
|
||||||
Attempt to establish connection to the OPC server.
|
|
||||||
|
|
||||||
This method performs the actual connection attempt to the OPC server
|
|
||||||
and handles connection failures with comprehensive error reporting.
|
|
||||||
It updates reconnection timing and provides detailed error information
|
|
||||||
for operational monitoring and debugging.
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
tuple[bool, dict[str, Any]]: Connection result
|
|
||||||
- bool: True if connection successful, False otherwise
|
|
||||||
- dict: Error information if connection failed
|
|
||||||
"""
|
|
||||||
|
|
||||||
tags = {
|
tags = {
|
||||||
'pod_id': self.pod_id,
|
'pod_id': self.pod_id,
|
||||||
'server_name': self.server_name,
|
'server_name': self.server_name,
|
||||||
}
|
}
|
||||||
await self.emit_metric(metrics.OPC_CONNECTIONS_TOTAL, tags)
|
await self.emit_metric(metrics.OPC_CONNECTIONS_TOTAL, tags)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
self.last_reconnection_time = datetime.now()
|
|
||||||
if self.client is None:
|
|
||||||
return False, {
|
|
||||||
'notification_id': f'OPC_CONNECTION_ERROR_{self.id}',
|
|
||||||
'message': 'Client is not initialized',
|
|
||||||
'block': 'opc_repository',
|
|
||||||
'level': NotificationLevel.ERROR,
|
|
||||||
}
|
|
||||||
await self.client.connect()
|
await self.client.connect()
|
||||||
|
|
||||||
|
session_id = _opc_authentication_token_str(self.client)
|
||||||
|
revised_session_timeout_ms = int(self.client.session_timeout)
|
||||||
|
revised_secure_channel_timeout_ms = int(self.client.secure_channel_timeout)
|
||||||
|
self.info(
|
||||||
|
f'OPC new session connected opc_server_id={self.id} session_id={session_id} '
|
||||||
|
f'revised_session_timeout_ms={revised_session_timeout_ms} '
|
||||||
|
f'revised_secure_channel_timeout_ms={revised_secure_channel_timeout_ms}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
await self.emit_metric(
|
||||||
|
metrics.OPC_SESSION_CREATED_TOTAL, self._opc_debug_tags(session_id)
|
||||||
|
)
|
||||||
|
await self.emit_metric(
|
||||||
|
metric_object=metrics.OPC_SESSION_REVISED_TIMEOUT_MS,
|
||||||
|
method='set',
|
||||||
|
tags=self._opc_debug_tags(session_id),
|
||||||
|
value=revised_session_timeout_ms,
|
||||||
|
)
|
||||||
await self.emit_metric(
|
await self.emit_metric(
|
||||||
metric_object=metrics.OPC_CONNECTION_STATUS,
|
metric_object=metrics.OPC_CONNECTION_STATUS,
|
||||||
method='set',
|
method='set',
|
||||||
tags={
|
tags={**tags, 'server_url': self.url},
|
||||||
**tags,
|
|
||||||
'server_url': self.url,
|
|
||||||
},
|
|
||||||
value=1,
|
value=1,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
self._last_write_mono = None
|
||||||
|
self._session_ready.set()
|
||||||
return True, {}
|
return True, {}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.disconnect()
|
await self._disconnect_locked()
|
||||||
|
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
self.logger.custom_error(trace, self.metadata)
|
self.error(trace, self.metadata)
|
||||||
|
|
||||||
await self.emit_metric(metrics.OPC_CONNECTIONS_FAILED, tags)
|
await self.emit_metric(metrics.OPC_CONNECTIONS_FAILED, tags)
|
||||||
|
|
||||||
return False, {
|
return False, {
|
||||||
'notification_id': f'OPC_CONNECTION_ERROR_{self.id}',
|
'notification_id': f'OPC_CONNECTION_ERROR_{self.id}',
|
||||||
'message': f'Failed to connect to OPC server: {e}',
|
'message': f'Failed to connect to OPC server: {e}',
|
||||||
@@ -206,21 +370,45 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
'attachment_content': trace,
|
'attachment_content': trace,
|
||||||
}
|
}
|
||||||
|
|
||||||
async def disconnection_fallback(self) -> list:
|
async def _connect_locked(self) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
|
||||||
Tries 5 times to disconnect from the OPC UA server, with a delay of 100ms x try.
|
|
||||||
"""
|
"""
|
||||||
|
Create the client when absent, then open a UA session.
|
||||||
|
|
||||||
|
Caller must hold _connection_lock.
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
OpcSessionAlreadyConnectedError: If a session is already open.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
tuple[bool, dict[str, Any]]: Result from _open_session on connect failure.
|
||||||
|
"""
|
||||||
|
if self._is_session_open():
|
||||||
|
raise OpcSessionAlreadyConnectedError(
|
||||||
|
f'OPC session already connected for server {self.id}; '
|
||||||
|
'call disconnect() before connecting again'
|
||||||
|
)
|
||||||
|
if self.client is None:
|
||||||
|
await self._create_client()
|
||||||
|
return await self._open_session()
|
||||||
|
|
||||||
|
async def _disconnection_fallback(self) -> list[dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Try up to five times to disconnect from the OPC UA server.
|
||||||
|
"""
|
||||||
assert self.client is not None
|
assert self.client is not None
|
||||||
error_stack = []
|
error_stack: list[dict[str, Any]] = []
|
||||||
for i in range(5):
|
for i in range(5):
|
||||||
try:
|
try:
|
||||||
self.logger.info(f'Disconnecting from OPC UA server, attempt {i + 1} of 5')
|
self.info(
|
||||||
|
f'Disconnecting from OPC UA server, attempt {i + 1} of 5',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
await self.client.disconnect()
|
await self.client.disconnect()
|
||||||
return []
|
return []
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
self.logger.error(
|
self.error(
|
||||||
f'Failed to disconnect from OPC UA server in attempt {i + 1} of 5: {e}'
|
f'Failed to disconnect from OPC UA server in attempt {i + 1} of 5: {e}',
|
||||||
|
self.metadata,
|
||||||
)
|
)
|
||||||
error_stack.append(
|
error_stack.append(
|
||||||
{
|
{
|
||||||
@@ -232,18 +420,26 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
await asyncio.sleep(self.disconnection_interval * i)
|
await asyncio.sleep(self.disconnection_interval * i)
|
||||||
return error_stack
|
return error_stack
|
||||||
|
|
||||||
async def disconnect(self):
|
async def _disconnect_locked(self) -> None:
|
||||||
"""
|
"""
|
||||||
Gracefully disconnect from the OPC server.
|
Tear down the current session and client.
|
||||||
|
|
||||||
This method safely terminates the connection to the OPC server
|
Caller must hold _connection_lock.
|
||||||
and cleans up client resources. It handles disconnection errors
|
|
||||||
gracefully and ensures proper resource cleanup.
|
|
||||||
"""
|
"""
|
||||||
|
self._last_write_mono = None
|
||||||
|
self._session_ready.clear()
|
||||||
|
|
||||||
if self.client is None:
|
if self.client is None:
|
||||||
return
|
return
|
||||||
|
|
||||||
errors = await self.disconnection_fallback()
|
session_id = _opc_authentication_token_str(self.client)
|
||||||
|
self.info(
|
||||||
|
f'OPC disconnecting opc_server_id={self.id} session_id={session_id}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
await self.emit_metric(metrics.OPC_SESSION_CLOSED_TOTAL, self._opc_debug_tags(session_id))
|
||||||
|
|
||||||
|
errors = await self._disconnection_fallback()
|
||||||
if errors:
|
if errors:
|
||||||
await self.send_notification_async(
|
await self.send_notification_async(
|
||||||
metadata=self.metadata,
|
metadata=self.metadata,
|
||||||
@@ -254,7 +450,8 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
attachment_content=json.dumps(errors, indent=4),
|
attachment_content=json.dumps(errors, indent=4),
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
self.logger.warning(f'Disconnected from OPC server {self.id} successfully')
|
self.warning(f'Disconnected from OPC server {self.id} successfully', self.metadata)
|
||||||
|
|
||||||
await self.emit_metric(
|
await self.emit_metric(
|
||||||
metric_object=metrics.OPC_CONNECTION_STATUS,
|
metric_object=metrics.OPC_CONNECTION_STATUS,
|
||||||
method='set',
|
method='set',
|
||||||
@@ -265,186 +462,286 @@ class OpcRepository(SientiaMonitoring):
|
|||||||
},
|
},
|
||||||
value=0,
|
value=0,
|
||||||
)
|
)
|
||||||
|
|
||||||
self.client = None
|
self.client = None
|
||||||
|
|
||||||
|
async def _reconnect_locked(self) -> tuple[bool, dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Close the current session and open a new one.
|
||||||
|
|
||||||
|
Caller must hold _connection_lock. Records last_reconnection_time for interval gating.
|
||||||
|
|
||||||
|
Return:
|
||||||
|
tuple[bool, dict[str, Any]]: Result from _connect_locked after teardown.
|
||||||
|
"""
|
||||||
|
self.last_reconnection_time = datetime.now()
|
||||||
|
await self._disconnect_locked()
|
||||||
|
return await self._connect_locked()
|
||||||
|
|
||||||
|
async def connect(self) -> tuple[bool, dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Open an OPC UA session under the connection lock (worker initialization).
|
||||||
|
"""
|
||||||
|
async with self._connection_lock:
|
||||||
|
self.info(
|
||||||
|
f'Starting connection to OPC server {self.id}:{self.server_name}...',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
return await self._connect_locked()
|
||||||
|
|
||||||
|
async def disconnect(self) -> None:
|
||||||
|
"""
|
||||||
|
Gracefully disconnect from the OPC server under the connection lock.
|
||||||
|
"""
|
||||||
|
async with self._connection_lock:
|
||||||
|
await self._disconnect_locked()
|
||||||
|
|
||||||
async def validate_connection(self) -> tuple[bool, dict[str, Any]]:
|
async def validate_connection(self) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Validate and maintain OPC server connection health.
|
Read-only check that the asyncua protocol is open.
|
||||||
|
|
||||||
This method performs comprehensive connection validation and
|
Caller must ensure _session_ready before writing. Does not connect or reconnect.
|
||||||
implements automatic reconnection logic for production reliability.
|
|
||||||
It handles various connection states and implements intelligent
|
|
||||||
reconnection strategies with error counting and timing controls.
|
|
||||||
|
|
||||||
Connection Validation:
|
Return:
|
||||||
1. Checks client existence and connection state
|
tuple[bool, dict[str, Any]]: (True, {}) when open, otherwise (False, error).
|
||||||
2. Implements error counting with automatic disconnection
|
"""
|
||||||
3. Enforces reconnection timing windows
|
if self._is_session_open():
|
||||||
4. Provides detailed error reporting and notifications
|
return True, {}
|
||||||
|
self.error(f'OPC server {self.id} is not connected', self.metadata)
|
||||||
|
return False, self._not_connected_error()
|
||||||
|
|
||||||
Reconnection Strategy:
|
async def _start_reconnect_on_bad(self, opc_status: str, session_id: str) -> None:
|
||||||
- Error Count Threshold: Disconnects after 5 consecutive errors
|
"""
|
||||||
- Reconnection Window: Enforces minimum intervals between attempts
|
Schedule a background reconnect if interval and task state allow it.
|
||||||
- Automatic Recovery: Attempts reconnection when conditions allow
|
|
||||||
- State Monitoring: Continuously monitors connection health
|
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
None
|
opc_status (str): OPC UA status name that triggered reconnect.
|
||||||
|
session_id (str): Session token before failure.
|
||||||
Returns:
|
|
||||||
tuple[bool, dict[str, Any]]: Connection validation result
|
|
||||||
- bool: True if connection is healthy, False otherwise
|
|
||||||
- dict: Error information if validation fails
|
|
||||||
"""
|
"""
|
||||||
if self.client is None:
|
if not self._reconnection_window_elapsed():
|
||||||
return await self.connect()
|
self.warning(
|
||||||
|
f'OPC reconnect skipped reason=reconnection_window opc_server_id={self.id} '
|
||||||
# if self.error_count > 5: # NOSONAR
|
f'opc_status={opc_status}',
|
||||||
# self.logger.custom_warning(
|
self.metadata,
|
||||||
# f'OPC server {self.id} will be disconnected due to multiple errors', self.metadata
|
|
||||||
# )
|
|
||||||
# try:
|
|
||||||
# await self.disconnect()
|
|
||||||
# except Exception as e:
|
|
||||||
# trace = traceback.format_exc()
|
|
||||||
# self.logger.custom_error(
|
|
||||||
# f'Failed to disconnect from OPC server: {e}', self.metadata
|
|
||||||
# )
|
|
||||||
# self.logger.custom_error(trace, self.metadata)
|
|
||||||
# self.logger.custom_info(
|
|
||||||
# f'Attempting to reconnect to OPC server {self.id}...', self.metadata
|
|
||||||
# )
|
|
||||||
# return await self.connect()
|
|
||||||
|
|
||||||
# Check if client is connected using asyncua's connection state
|
|
||||||
try:
|
|
||||||
if (
|
|
||||||
self.client.uaclient.protocol is None
|
|
||||||
or self.client.uaclient.protocol.state == 'closed'
|
|
||||||
):
|
|
||||||
# OPC server is not connected
|
|
||||||
self.logger.custom_error(f'OPC server {self.id} is not connected', self.metadata)
|
|
||||||
if (
|
|
||||||
self.last_reconnection_time is None
|
|
||||||
or (datetime.now() - self.last_reconnection_time).total_seconds()
|
|
||||||
> self.reconnection_interval
|
|
||||||
):
|
|
||||||
await self.disconnect()
|
|
||||||
self.logger.custom_info(
|
|
||||||
f'Trying to reconnect to OPC server {self.id}...', self.metadata
|
|
||||||
)
|
)
|
||||||
return await self.connect()
|
return
|
||||||
|
if self._reconnect_task is not None and not self._reconnect_task.done():
|
||||||
|
self.warning(
|
||||||
|
f'OPC reconnect skipped reason=in_progress opc_server_id={self.id} '
|
||||||
|
f'opc_status={opc_status}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
return False, {
|
self._session_ready.clear()
|
||||||
'notification_id': f'OPC_CONNECTION_AWAITING_RECONNECTION_WINDOW_{self.id}',
|
self.info(
|
||||||
'message': f'OPC server {self.id} is not connected, waiting for next reconnection window...',
|
f'OPC reconnect scheduled after opc_status={opc_status} opc_server_id={self.id} '
|
||||||
'block': 'opc_repository',
|
f'old_session_id={session_id}',
|
||||||
'level': NotificationLevel.WARNING,
|
self.metadata,
|
||||||
}
|
)
|
||||||
return True, {}
|
self._reconnect_task = asyncio.create_task(
|
||||||
except Exception as e:
|
self._run_reconnect_on_bad(opc_status, session_id)
|
||||||
trace = traceback.format_exc()
|
)
|
||||||
message = f'Failed to validate connection to OPC server: {e}'
|
|
||||||
self.logger.custom_error(message, self.metadata)
|
async def _run_reconnect_on_bad(self, opc_status: str, session_id: str) -> None:
|
||||||
return False, {
|
"""
|
||||||
'notification_id': f'OPC_CONNECTION_CHECK_ERROR_{self.id}',
|
Tear down and re-establish the OPC UA session under the connection lock.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
opc_status (str): OPC UA status that triggered reconnect.
|
||||||
|
session_id (str): Previous session token string.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
async with self._connection_lock:
|
||||||
|
self.info(
|
||||||
|
f'OPC reconnect started opc_status={opc_status} opc_server_id={self.id} '
|
||||||
|
f'old_session_id={session_id}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
await self._reconnect_locked()
|
||||||
|
except Exception:
|
||||||
|
self.error(
|
||||||
|
f'OPC reconnect task failed opc_server_id={self.id} opc_status={opc_status}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
self.error(traceback.format_exc(), self.metadata)
|
||||||
|
|
||||||
|
async def _log_write_inter_arrival(self, session_id: str, node: str) -> None:
|
||||||
|
"""
|
||||||
|
Log elapsed wall time since the previous successful OPC write on this repository.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
session_id (str): Current OPC UA session token string.
|
||||||
|
node (str): Node id written in this operation.
|
||||||
|
"""
|
||||||
|
now = time.monotonic()
|
||||||
|
if self._last_write_mono is not None:
|
||||||
|
delta_s = now - self._last_write_mono
|
||||||
|
self.info(
|
||||||
|
f'OPC write inter-arrival_s={delta_s:.6f} opc_server_id={self.id} '
|
||||||
|
f'session_id={session_id} node={node}',
|
||||||
|
self.metadata,
|
||||||
|
)
|
||||||
|
if self.client is not None:
|
||||||
|
session_timeout_ms = float(self.client.session_timeout)
|
||||||
|
if session_timeout_ms > 0 and delta_s > (session_timeout_ms / 1000.0):
|
||||||
|
await self.emit_metric(
|
||||||
|
metrics.OPC_WRITE_INTER_ARRIVAL_OVER_SESSION_TIMEOUT_TOTAL,
|
||||||
|
self._opc_debug_tags(session_id),
|
||||||
|
)
|
||||||
|
self._last_write_mono = now
|
||||||
|
|
||||||
|
async def _emit_opc_write_metric(
|
||||||
|
self, session_id: str, result: str, metadata: dict[str, Any] | None
|
||||||
|
) -> None:
|
||||||
|
await self.emit_metric(
|
||||||
|
metrics.OPC_WRITE_ATTEMPTS_TOTAL,
|
||||||
|
{
|
||||||
|
**self._opc_debug_tags(session_id),
|
||||||
|
**_model_labels_from_write_metadata(metadata),
|
||||||
|
'result': result,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
def _write_failure_payload(
|
||||||
|
self,
|
||||||
|
notification_id: str,
|
||||||
|
message: str,
|
||||||
|
level: NotificationLevel = NotificationLevel.ERROR,
|
||||||
|
attachment_content: str | None = None,
|
||||||
|
opc_error_kind: str | None = None,
|
||||||
|
opc_status: str | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
payload: dict[str, Any] = {
|
||||||
|
'notification_id': notification_id,
|
||||||
'message': message,
|
'message': message,
|
||||||
'block': 'opc_repository',
|
'block': 'opc_repository',
|
||||||
'level': NotificationLevel.ERROR,
|
'level': level,
|
||||||
'attachment_content': trace,
|
|
||||||
}
|
}
|
||||||
|
if attachment_content is not None:
|
||||||
|
payload['attachment_content'] = attachment_content
|
||||||
|
if opc_error_kind is not None:
|
||||||
|
payload['opc_error_kind'] = opc_error_kind
|
||||||
|
if opc_status is not None:
|
||||||
|
payload['opc_status'] = opc_status
|
||||||
|
return payload
|
||||||
|
|
||||||
async def write_data(
|
async def _handle_tier1_bad(
|
||||||
self, node: str, value: Any, data_type: str, logger: Logger, metadata: dict[str, Any]
|
self,
|
||||||
|
exc: BaseException,
|
||||||
|
session_id: str,
|
||||||
|
node: str,
|
||||||
|
metadata: dict[str, Any],
|
||||||
|
phase: str,
|
||||||
) -> tuple[bool, dict[str, Any]]:
|
) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Write data to OPC server with comprehensive validation and monitoring.
|
Record metrics/logs and schedule reconnect after a Tier-1 Bad* error.
|
||||||
|
|
||||||
This method provides secure and reliable data writing to OPC servers
|
|
||||||
with automatic connection validation, data type conversion, and
|
|
||||||
comprehensive error handling. It implements performance monitoring
|
|
||||||
and metrics collection for operational visibility.
|
|
||||||
|
|
||||||
Data Writing Process:
|
|
||||||
1. Connection validation and automatic reconnection
|
|
||||||
2. Node validation and error handling
|
|
||||||
3. Data type conversion and validation
|
|
||||||
4. OPC data writing with timestamp
|
|
||||||
5. Performance metrics collection
|
|
||||||
6. Error handling and notification
|
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
node (str): OPC node identifier to write data to
|
exc (BaseException): Tier-1 OPC UA error.
|
||||||
value (Any): Data value to write to the OPC node
|
session_id (str): Session token at failure time.
|
||||||
data_type (str): Data type for OPC conversion
|
node (str): Node id being written.
|
||||||
logger (Logger): Logger instance for operation logging
|
metadata (dict[str, Any]): Write context.
|
||||||
metadata (dict[str, Any]): Context metadata for logging and metrics
|
phase (str): get_node or write_value.
|
||||||
|
|
||||||
Returns:
|
Return:
|
||||||
tuple[bool, dict[str, Any]]: Write operation result
|
tuple[bool, dict[str, Any]]: Always (False, error payload).
|
||||||
- bool: True if write successful, False otherwise
|
|
||||||
- dict: Error information if write failed
|
|
||||||
"""
|
"""
|
||||||
|
opc_status = _opc_status_from_exception(exc)
|
||||||
|
trace = traceback.format_exc()
|
||||||
|
self.error(trace, metadata)
|
||||||
|
await self._emit_opc_write_metric(session_id, opc_status, metadata)
|
||||||
|
self.error(
|
||||||
|
f'OPC write failed opc_status={opc_status} opc_server_id={self.id} '
|
||||||
|
f'session_id={session_id} model_id={metadata.get("model_id", "unknown")} '
|
||||||
|
f'model_name={metadata.get("model_name", "unknown")} node={node} phase={phase}',
|
||||||
|
metadata,
|
||||||
|
)
|
||||||
|
await self._start_reconnect_on_bad(opc_status, session_id)
|
||||||
|
return False, self._write_failure_payload(
|
||||||
|
notification_id=f'OPC_WRITE_DATA_ERROR_{self.id}',
|
||||||
|
message=f'Failed to {phase} on OPC server: {exc} | metadata: {metadata}',
|
||||||
|
attachment_content=trace,
|
||||||
|
opc_error_kind='session_bad',
|
||||||
|
opc_status=opc_status,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def write_data(
|
||||||
|
self, node: str, value: Any, data_type: str, metadata: dict[str, Any]
|
||||||
|
) -> tuple[bool, dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Write data to OPC server with a single attempt and Tier-1 Bad* reconnect scheduling.
|
||||||
|
"""
|
||||||
|
if not self._session_ready.is_set():
|
||||||
|
await self._emit_opc_write_metric('unknown', 'ReconnectInProgress', metadata)
|
||||||
|
self.warning(
|
||||||
|
f'OPC write rejected reconnect_in_progress opc_server_id={self.id} '
|
||||||
|
f'model_id={metadata.get("model_id", "unknown")} '
|
||||||
|
f'model_name={metadata.get("model_name", "unknown")}',
|
||||||
|
metadata,
|
||||||
|
)
|
||||||
|
return False, {
|
||||||
|
'notification_id': f'OPC_WRITE_RECONNECT_IN_PROGRESS_{self.id}',
|
||||||
|
'message': f'OPC write skipped: reconnect in progress | metadata: {metadata}',
|
||||||
|
'block': 'opc_repository',
|
||||||
|
'level': NotificationLevel.WARNING,
|
||||||
|
'opc_error_kind': 'reconnect_in_progress',
|
||||||
|
}
|
||||||
|
|
||||||
is_connected, error = await self.validate_connection()
|
is_connected, error = await self.validate_connection()
|
||||||
|
|
||||||
if not is_connected:
|
if not is_connected:
|
||||||
|
await self._emit_opc_write_metric('unknown', 'NotConnected', metadata)
|
||||||
return False, error
|
return False, error
|
||||||
|
|
||||||
start_time = time.time()
|
start_time = time.time()
|
||||||
|
session_id = _opc_authentication_token_str(self.client)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# ignored because self.validate_connection is called before, so we know self.client is not None
|
|
||||||
node_obj = self.client.get_node(node) # type: ignore[union-attr]
|
node_obj = self.client.get_node(node) # type: ignore[union-attr]
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
if is_reconnectable_opcua_bad(e):
|
||||||
|
return await self._handle_tier1_bad(e, session_id, node, metadata, 'get_node')
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
logger.custom_error(trace, metadata.get('schedule_name', 'N/A'))
|
self.error(trace, metadata)
|
||||||
self.error_count += 1
|
await self._emit_opc_write_metric(
|
||||||
return False, {
|
session_id, f'GetNodeError:{type(e).__name__}', metadata
|
||||||
'notification_id': f'OPC_WRITE_GET_NODE_ERROR_{self.id}',
|
)
|
||||||
'message': f'Failed to get node from OPC server: {e} | metadata: {metadata}',
|
return False, self._write_failure_payload(
|
||||||
'block': 'opc_repository',
|
notification_id=f'OPC_WRITE_GET_NODE_ERROR_{self.id}',
|
||||||
'level': NotificationLevel.ERROR,
|
message=f'Failed to get node from OPC server: {e} | metadata: {metadata}',
|
||||||
'attachment_content': trace,
|
attachment_content=trace,
|
||||||
}
|
)
|
||||||
|
|
||||||
if data_type not in data_type_map:
|
if data_type not in data_type_map:
|
||||||
return False, {
|
await self._emit_opc_write_metric(session_id, 'UnsupportedDataType', metadata)
|
||||||
'notification_id': f'OPC_WRITE_DATA_TYPE_ERROR_{self.id}',
|
return False, self._write_failure_payload(
|
||||||
'message': f'Unsupported data type: {data_type} | metadata: {metadata}',
|
notification_id=f'OPC_WRITE_DATA_TYPE_ERROR_{self.id}',
|
||||||
'block': 'opc_repository',
|
message=f'Unsupported data type: {data_type} | metadata: {metadata}',
|
||||||
'level': NotificationLevel.ERROR,
|
)
|
||||||
}
|
|
||||||
|
|
||||||
data = data_type_map[data_type]['converter'](value)
|
data = data_type_map[data_type]['converter'](value)
|
||||||
logger.custom_info(f'Writing {data} - {type(data)} to {node}', metadata)
|
self.info(f'Writing {data} - {type(data)} to {node}', metadata)
|
||||||
# now = datetime.now() # NOSONAR
|
|
||||||
ua_data = DataValue(
|
ua_data = DataValue(
|
||||||
Variant(data, data_type_map[data_type]['opc_type']),
|
Variant(data, data_type_map[data_type]['opc_type']),
|
||||||
# SourceTimestamp=DateTime( # NOSONAR
|
|
||||||
# now.year, now.month, now.day, now.hour, now.minute, now.second, now.microsecond # NOSONAR
|
|
||||||
# ), # NOSONAR
|
|
||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await node_obj.write_value(ua_data)
|
await node_obj.write_value(ua_data)
|
||||||
|
|
||||||
end_time = time.time()
|
end_time = time.time()
|
||||||
response_time = end_time - start_time
|
response_time = end_time - start_time
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
if is_reconnectable_opcua_bad(e):
|
||||||
|
return await self._handle_tier1_bad(e, session_id, node, metadata, 'write_value')
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
logger.custom_error(trace, metadata)
|
self.error(trace, metadata)
|
||||||
self.error_count += 1
|
await self._emit_opc_write_metric(session_id, type(e).__name__, metadata)
|
||||||
return False, {
|
return False, self._write_failure_payload(
|
||||||
'notification_id': f'OPC_WRITE_DATA_ERROR_{self.id}',
|
notification_id=f'OPC_WRITE_DATA_ERROR_{self.id}',
|
||||||
'message': f'Failed to write data to OPC server: {e} | metadata: {metadata}',
|
message=f'Failed to write data to OPC server: {e} | metadata: {metadata}',
|
||||||
'block': 'opc_repository',
|
attachment_content=trace,
|
||||||
'level': NotificationLevel.ERROR,
|
)
|
||||||
'attachment_content': trace,
|
|
||||||
}
|
await self._emit_opc_write_metric(session_id, 'OK', metadata)
|
||||||
self.error_count = 0
|
await self._log_write_inter_arrival(session_id, node)
|
||||||
|
|
||||||
return True, {
|
return True, {
|
||||||
'response_time': response_time,
|
'response_time': response_time,
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ with workflow.unsafe.imports_passed_through():
|
|||||||
)
|
)
|
||||||
from laborious.workflows.sub_workflows.prediction_process import PredictionProcess
|
from laborious.workflows.sub_workflows.prediction_process import PredictionProcess
|
||||||
|
|
||||||
POD_ID = os.getenv('POD_ID')
|
POD_ID = os.getenv('HOSTNAME')
|
||||||
SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091'))
|
SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091'))
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -18,3 +18,4 @@ testcontainers[postgres,minio] # PostgreSQL and MinIO containers for E2E tests
|
|||||||
# Development Tools
|
# Development Tools
|
||||||
ipython>=8.12.0 # Enhanced Python shell
|
ipython>=8.12.0 # Enhanced Python shell
|
||||||
ipdb>=0.13.13 # IPython debugger
|
ipdb>=0.13.13 # IPython debugger
|
||||||
|
ipykernel==6.30.1 # IPython kernel for Jupyter notebooks
|
||||||
|
|||||||
@@ -5,7 +5,15 @@ from pandas import DataFrame
|
|||||||
from pytest import mark
|
from pytest import mark
|
||||||
from sientia_do.notifications.models import NotificationLevel
|
from sientia_do.notifications.models import NotificationLevel
|
||||||
|
|
||||||
from laborious.activities.opc import OPC
|
from laborious.activities.opc import (
|
||||||
|
OPC,
|
||||||
|
OPC_COMMENT_SEPARATOR,
|
||||||
|
OPC_RECONNECT_IN_PROGRESS_COMMENT,
|
||||||
|
OPC_SESSION_BAD_COMMENT_PREFIX,
|
||||||
|
OPC_SESSION_BAD_CONFIDENCE,
|
||||||
|
OPC_WRITTING_ERROR_CONFIDENCE,
|
||||||
|
OPC_WRITTING_ERROR_MESSAGE,
|
||||||
|
)
|
||||||
|
|
||||||
metadata = {
|
metadata = {
|
||||||
'metadata': {
|
'metadata': {
|
||||||
@@ -206,7 +214,7 @@ WRITE_DATA_CASES = [
|
|||||||
async def test_write_data_success(opc, tag, data_type, data):
|
async def test_write_data_success(opc, tag, data_type, data):
|
||||||
opc.opc_repository['server1'].write_data.return_value = (True, {'response_time': 0.1})
|
opc.opc_repository['server1'].write_data.return_value = (True, {'response_time': 0.1})
|
||||||
|
|
||||||
result = await opc.write_data(
|
response_time, error_info = await opc.write_data(
|
||||||
server_id='server1',
|
server_id='server1',
|
||||||
tag=tag,
|
tag=tag,
|
||||||
data=data,
|
data=data,
|
||||||
@@ -214,10 +222,9 @@ async def test_write_data_success(opc, tag, data_type, data):
|
|||||||
tag_type='prediction',
|
tag_type='prediction',
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
)
|
)
|
||||||
assert result == 0.1
|
assert response_time == 0.1
|
||||||
opc.opc_repository['server1'].write_data.assert_called_once_with(
|
assert error_info is None
|
||||||
tag, data, data_type, opc.logger, metadata
|
opc.opc_repository['server1'].write_data.assert_called_once_with(tag, data, data_type, metadata)
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
@@ -233,7 +240,7 @@ async def test_write_data_failed(opc):
|
|||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
result = await opc.write_data(
|
response_time, error_info = await opc.write_data(
|
||||||
server_id='server1',
|
server_id='server1',
|
||||||
tag='tag1',
|
tag='tag1',
|
||||||
data=50,
|
data=50,
|
||||||
@@ -241,7 +248,8 @@ async def test_write_data_failed(opc):
|
|||||||
tag_type='prediction',
|
tag_type='prediction',
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
)
|
)
|
||||||
assert result is None
|
assert response_time is None
|
||||||
|
assert error_info is not None
|
||||||
|
|
||||||
opc.send_notification_async.assert_called_once_with(
|
opc.send_notification_async.assert_called_once_with(
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
@@ -283,7 +291,7 @@ async def test_write_data_exception(opc):
|
|||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_manage_output_tags_success(opc):
|
async def test_manage_output_tags_success(opc):
|
||||||
opc.write_data = AsyncMock(return_value=0.1)
|
opc.write_data = AsyncMock(return_value=(0.1, None))
|
||||||
|
|
||||||
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
||||||
config = {
|
config = {
|
||||||
@@ -291,7 +299,7 @@ async def test_manage_output_tags_success(opc):
|
|||||||
'confidence_tags': {'tag2': {'data_type': 'float'}},
|
'confidence_tags': {'tag2': {'data_type': 'float'}},
|
||||||
}
|
}
|
||||||
|
|
||||||
output_data, opc_metrics = await opc.manage_output_tags(
|
output_data, opc_metrics, _, _, _ = await opc.manage_output_tags(
|
||||||
server_id='server1',
|
server_id='server1',
|
||||||
config=config,
|
config=config,
|
||||||
data=data,
|
data=data,
|
||||||
@@ -323,7 +331,13 @@ async def test_manage_output_tags_success(opc):
|
|||||||
|
|
||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
@mark.parametrize('side_effect', [[0.1, None], [None, 0.2]])
|
@mark.parametrize(
|
||||||
|
'side_effect',
|
||||||
|
[
|
||||||
|
[(0.1, None), (None, {})],
|
||||||
|
[(None, {}), (0.2, None)],
|
||||||
|
],
|
||||||
|
)
|
||||||
async def test_manage_output_tags_failed(opc, side_effect):
|
async def test_manage_output_tags_failed(opc, side_effect):
|
||||||
opc.write_data = AsyncMock(side_effect=side_effect)
|
opc.write_data = AsyncMock(side_effect=side_effect)
|
||||||
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
||||||
@@ -331,14 +345,14 @@ async def test_manage_output_tags_failed(opc, side_effect):
|
|||||||
'prediction_tags': {'tag1': {'data_type': 'float'}},
|
'prediction_tags': {'tag1': {'data_type': 'float'}},
|
||||||
'confidence_tags': {'tag2': {'data_type': 'float'}},
|
'confidence_tags': {'tag2': {'data_type': 'float'}},
|
||||||
}
|
}
|
||||||
output_data, opc_metrics = await opc.manage_output_tags(
|
output_data, opc_metrics, _, _, _ = await opc.manage_output_tags(
|
||||||
server_id='server1',
|
server_id='server1',
|
||||||
config=config,
|
config=config,
|
||||||
data=data,
|
data=data,
|
||||||
metadata=metadata['metadata'],
|
metadata=metadata['metadata'],
|
||||||
)
|
)
|
||||||
assert output_data is False
|
assert output_data is False
|
||||||
assert opc_metrics == {'tag1': side_effect[0], 'tag2': side_effect[1]}
|
assert opc_metrics == {'tag1': side_effect[0][0], 'tag2': side_effect[1][0]}
|
||||||
opc.write_data.assert_has_calls(
|
opc.write_data.assert_has_calls(
|
||||||
[
|
[
|
||||||
call(
|
call(
|
||||||
@@ -363,12 +377,12 @@ async def test_manage_output_tags_failed(opc, side_effect):
|
|||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_manage_output_tags_do_nothing(opc):
|
async def test_manage_output_tags_do_nothing(opc):
|
||||||
opc.write_data = AsyncMock(return_value=0.1)
|
opc.write_data = AsyncMock(return_value=(0.1, None))
|
||||||
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
||||||
config = {
|
config = {
|
||||||
'_invalid_key': {'tag1': {'data_type': 'float'}},
|
'_invalid_key': {'tag1': {'data_type': 'float'}},
|
||||||
}
|
}
|
||||||
output_data, opc_metrics = await opc.manage_output_tags(
|
output_data, opc_metrics, _, _, _ = await opc.manage_output_tags(
|
||||||
server_id='server1',
|
server_id='server1',
|
||||||
config=config,
|
config=config,
|
||||||
data=data,
|
data=data,
|
||||||
@@ -395,7 +409,9 @@ async def test_write_opc_data_success(mock_dataframe, opc):
|
|||||||
}
|
}
|
||||||
|
|
||||||
# Act
|
# Act
|
||||||
opc.manage_output_tags = AsyncMock(return_value=(True, {'tag1': 0.1, 'tag2': 0.2}))
|
opc.manage_output_tags = AsyncMock(
|
||||||
|
return_value=(True, {'tag1': 0.1, 'tag2': 0.2}, False, None, False)
|
||||||
|
)
|
||||||
|
|
||||||
opc.process_confidence = MagicMock(return_value={'data': 'data'})
|
opc.process_confidence = MagicMock(return_value={'data': 'data'})
|
||||||
output_data, opc_metrics = await opc.write_opc_data(input_data)
|
output_data, opc_metrics = await opc.write_opc_data(input_data)
|
||||||
@@ -413,6 +429,9 @@ async def test_write_opc_data_success(mock_dataframe, opc):
|
|||||||
mock_dataframe.return_value,
|
mock_dataframe.return_value,
|
||||||
True,
|
True,
|
||||||
metadata['metadata'],
|
metadata['metadata'],
|
||||||
|
session_bad=False,
|
||||||
|
opc_status=None,
|
||||||
|
reconnect_in_progress=False,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -462,13 +481,130 @@ async def test_write_opc_data_no_validate_server(opc):
|
|||||||
],
|
],
|
||||||
)
|
)
|
||||||
def test_process_confidence(opc, data, success, expected):
|
def test_process_confidence(opc, data, success, expected):
|
||||||
# Act
|
result = opc.process_confidence(data, success, metadata['metadata'])
|
||||||
result = opc.process_confidence(data, success, metadata)
|
|
||||||
|
|
||||||
# Assert
|
|
||||||
assert result['prediction_confidence'][0] == expected
|
assert result['prediction_confidence'][0] == expected
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_confidence_session_bad(opc):
|
||||||
|
data = DataFrame({'prediction_confidence': [0.9]})
|
||||||
|
result = opc.process_confidence(
|
||||||
|
data,
|
||||||
|
False,
|
||||||
|
metadata['metadata'],
|
||||||
|
session_bad=True,
|
||||||
|
opc_status='BadSessionIdInvalid',
|
||||||
|
)
|
||||||
|
assert result['prediction_confidence'][0] == OPC_SESSION_BAD_CONFIDENCE
|
||||||
|
assert result['comments'][0].startswith(OPC_SESSION_BAD_COMMENT_PREFIX)
|
||||||
|
assert 'BadSessionIdInvalid' in result['comments'][0]
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_confidence_generic_failure(opc):
|
||||||
|
data = DataFrame({'prediction_confidence': [0.9]})
|
||||||
|
result = opc.process_confidence(data, False, metadata['metadata'])
|
||||||
|
assert result['prediction_confidence'][0] == OPC_WRITTING_ERROR_CONFIDENCE
|
||||||
|
assert result['comments'][0] == OPC_WRITTING_ERROR_MESSAGE
|
||||||
|
|
||||||
|
|
||||||
|
@mark.asyncio
|
||||||
|
async def test_manage_output_tags_session_bad(opc):
|
||||||
|
opc.write_data = AsyncMock(
|
||||||
|
side_effect=[
|
||||||
|
(
|
||||||
|
None,
|
||||||
|
{
|
||||||
|
'opc_error_kind': 'session_bad',
|
||||||
|
'opc_status': 'BadSessionIdInvalid',
|
||||||
|
'notification_id': 'OPC_WRITE_DATA_ERROR_1',
|
||||||
|
'message': 'bad',
|
||||||
|
'block': 'opc_repository',
|
||||||
|
'level': NotificationLevel.ERROR,
|
||||||
|
},
|
||||||
|
),
|
||||||
|
(0.2, None),
|
||||||
|
]
|
||||||
|
)
|
||||||
|
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
||||||
|
config = {
|
||||||
|
'prediction_tags': {'tag1': {'data_type': 'float'}},
|
||||||
|
'confidence_tags': {'tag2': {'data_type': 'float'}},
|
||||||
|
}
|
||||||
|
(
|
||||||
|
success,
|
||||||
|
metrics,
|
||||||
|
session_bad_seen,
|
||||||
|
opc_status,
|
||||||
|
reconnect_in_progress,
|
||||||
|
) = await opc.manage_output_tags('server1', config, data, metadata['metadata'])
|
||||||
|
assert success is False
|
||||||
|
assert session_bad_seen is True
|
||||||
|
assert reconnect_in_progress is False
|
||||||
|
assert opc_status == 'BadSessionIdInvalid'
|
||||||
|
assert metrics['tag1'] is None
|
||||||
|
assert metrics['tag2'] == 0.2
|
||||||
|
|
||||||
|
|
||||||
|
@mark.asyncio
|
||||||
|
async def test_manage_output_tags_reconnect_in_progress(opc):
|
||||||
|
opc.write_data = AsyncMock(
|
||||||
|
return_value=(
|
||||||
|
None,
|
||||||
|
{
|
||||||
|
'opc_error_kind': 'reconnect_in_progress',
|
||||||
|
'notification_id': 'OPC_WRITE_RECONNECT_IN_PROGRESS_1',
|
||||||
|
'message': 'skipped',
|
||||||
|
'block': 'opc_repository',
|
||||||
|
'level': NotificationLevel.WARNING,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
data = DataFrame({'prediction': [0.75], 'prediction_confidence': [0.95]})
|
||||||
|
config = {'prediction_tags': {'tag1': {'data_type': 'float'}}}
|
||||||
|
(
|
||||||
|
success,
|
||||||
|
metrics,
|
||||||
|
session_bad_seen,
|
||||||
|
opc_status,
|
||||||
|
reconnect_in_progress,
|
||||||
|
) = await opc.manage_output_tags('server1', config, data, metadata['metadata'])
|
||||||
|
assert success is False
|
||||||
|
assert session_bad_seen is False
|
||||||
|
assert reconnect_in_progress is True
|
||||||
|
assert opc_status is None
|
||||||
|
assert metrics['tag1'] is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_confidence_reconnect_in_progress(opc):
|
||||||
|
data = DataFrame({'prediction_confidence': [0.9]})
|
||||||
|
result = opc.process_confidence(
|
||||||
|
data,
|
||||||
|
False,
|
||||||
|
metadata['metadata'],
|
||||||
|
reconnect_in_progress=True,
|
||||||
|
)
|
||||||
|
assert result['prediction_confidence'][0] == OPC_SESSION_BAD_CONFIDENCE
|
||||||
|
assert result['comments'][0] == OPC_RECONNECT_IN_PROGRESS_COMMENT
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_confidence_concatenates_multiple_comments(opc):
|
||||||
|
data = DataFrame({'prediction_confidence': [0.9]})
|
||||||
|
session_comment = f'{OPC_SESSION_BAD_COMMENT_PREFIX} BadSessionIdInvalid'
|
||||||
|
|
||||||
|
result = opc.process_confidence(
|
||||||
|
data,
|
||||||
|
False,
|
||||||
|
metadata['metadata'],
|
||||||
|
session_bad=True,
|
||||||
|
opc_status='BadSessionIdInvalid',
|
||||||
|
reconnect_in_progress=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result['prediction_confidence'][0] == OPC_SESSION_BAD_CONFIDENCE
|
||||||
|
assert result['comments'][0] == OPC_COMMENT_SEPARATOR.join(
|
||||||
|
[session_comment, OPC_RECONNECT_IN_PROGRESS_COMMENT]
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@mark.asyncio
|
@mark.asyncio
|
||||||
async def test_validate_server(opc):
|
async def test_validate_server(opc):
|
||||||
assert await opc.validate_server('server1', metadata) is True
|
assert await opc.validate_server('server1', metadata) is True
|
||||||
|
|||||||
@@ -1,12 +1,20 @@
|
|||||||
|
import asyncio
|
||||||
import json
|
import json
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from unittest.mock import ANY, AsyncMock, MagicMock, Mock, patch
|
from unittest.mock import ANY, AsyncMock, MagicMock, Mock, patch
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
||||||
|
from asyncua.ua.uaerrors import BadNodeIdUnknown, BadSessionIdInvalid
|
||||||
from sientia_do.notifications.models import NotificationLevel
|
from sientia_do.notifications.models import NotificationLevel
|
||||||
|
|
||||||
from laborious.utils.repository.opc_repository import OpcRepository
|
from laborious.utils.repository.opc_repository import (
|
||||||
|
OpcClientAlreadyExistsError,
|
||||||
|
OpcClientNotInitializedError,
|
||||||
|
OpcRepository,
|
||||||
|
OpcSessionAlreadyConnectedError,
|
||||||
|
is_reconnectable_opcua_bad,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
@@ -33,6 +41,11 @@ def opc_repository(mock_logger):
|
|||||||
repository.send_notification = MagicMock()
|
repository.send_notification = MagicMock()
|
||||||
repository.send_notification_async = AsyncMock()
|
repository.send_notification_async = AsyncMock()
|
||||||
repository.emit_metric = AsyncMock()
|
repository.emit_metric = AsyncMock()
|
||||||
|
repository.info = MagicMock()
|
||||||
|
repository.error = MagicMock()
|
||||||
|
repository.warning = MagicMock()
|
||||||
|
repository.debug = MagicMock()
|
||||||
|
repository._session_ready.set()
|
||||||
return repository
|
return repository
|
||||||
|
|
||||||
|
|
||||||
@@ -65,7 +78,6 @@ def test_init(opc_repository):
|
|||||||
assert opc_repository.reconnection_interval == 60
|
assert opc_repository.reconnection_interval == 60
|
||||||
assert opc_repository.client is None
|
assert opc_repository.client is None
|
||||||
assert opc_repository.last_reconnection_time is None
|
assert opc_repository.last_reconnection_time is None
|
||||||
assert opc_repository.error_count == 0
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
@@ -80,8 +92,8 @@ async def test_set_security(opc_repository, mock_client):
|
|||||||
private_key='/path/to/key.pem',
|
private_key='/path/to/key.pem',
|
||||||
server_certificate='/path/to/server_cert.pem',
|
server_certificate='/path/to/server_cert.pem',
|
||||||
)
|
)
|
||||||
assert mock_client.secure_channel_timeout == 10000000
|
assert mock_client.secure_channel_timeout == 600_000
|
||||||
assert mock_client.session_timeout == 10000000
|
assert mock_client.session_timeout == 600_000
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
@@ -106,48 +118,95 @@ async def test_set_security_missing_client(opc_repository):
|
|||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_connect_with_security(opc_repository, mock_client):
|
async def test_connect_with_security(opc_repository, mock_client):
|
||||||
opc_repository.try_connect = AsyncMock(return_value=(True, {}))
|
opc_repository._create_client = AsyncMock()
|
||||||
|
opc_repository._open_session = AsyncMock(return_value=(True, {}))
|
||||||
result = await opc_repository.connect()
|
result = await opc_repository.connect()
|
||||||
|
|
||||||
opc_repository.try_connect.assert_called_once()
|
opc_repository._create_client.assert_called_once()
|
||||||
assert opc_repository.client == mock_client
|
opc_repository._open_session.assert_called_once()
|
||||||
assert result == (True, {})
|
assert result == (True, {})
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_connect_without_security(opc_repository, mock_client):
|
async def test_connect_without_security(opc_repository, mock_client):
|
||||||
opc_repository.cert_path = None
|
opc_repository.cert_path = None
|
||||||
opc_repository.try_connect = AsyncMock(return_value=(True, {}))
|
opc_repository._create_client = AsyncMock()
|
||||||
|
opc_repository._open_session = AsyncMock(return_value=(True, {}))
|
||||||
opc_repository.set_security = AsyncMock()
|
opc_repository.set_security = AsyncMock()
|
||||||
result = await opc_repository.connect()
|
result = await opc_repository.connect()
|
||||||
|
|
||||||
opc_repository.try_connect.assert_called_once()
|
opc_repository._create_client.assert_called_once()
|
||||||
|
opc_repository._open_session.assert_called_once()
|
||||||
opc_repository.set_security.assert_not_called()
|
opc_repository.set_security.assert_not_called()
|
||||||
assert opc_repository.client == mock_client
|
|
||||||
assert result == (True, {})
|
assert result == (True, {})
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_try_connect_success(opc_repository):
|
async def test_connect_raises_when_session_already_open(opc_repository, mock_client):
|
||||||
opc_repository.last_reconnection_time = None
|
opc_repository.client = mock_client
|
||||||
|
proto = MagicMock()
|
||||||
|
proto.state = 'open'
|
||||||
|
mock_client.uaclient = MagicMock(protocol=proto)
|
||||||
|
|
||||||
|
with pytest.raises(OpcSessionAlreadyConnectedError, match='disconnect'):
|
||||||
|
await opc_repository.connect()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_create_client_raises_when_client_exists(opc_repository, mock_client):
|
||||||
|
opc_repository.client = mock_client
|
||||||
|
|
||||||
|
with pytest.raises(OpcClientAlreadyExistsError, match='already exists'):
|
||||||
|
await opc_repository._create_client()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_open_session_success(opc_repository):
|
||||||
|
closed_proto = MagicMock()
|
||||||
|
closed_proto.state = 'closed'
|
||||||
opc_repository.client = AsyncMock()
|
opc_repository.client = AsyncMock()
|
||||||
result = await opc_repository.try_connect()
|
opc_repository.client.uaclient = MagicMock(protocol=closed_proto)
|
||||||
|
opc_repository.client.session_timeout = 600_000
|
||||||
|
opc_repository.client.secure_channel_timeout = 600_000
|
||||||
|
|
||||||
|
open_proto = MagicMock()
|
||||||
|
open_proto.state = 'open'
|
||||||
|
open_proto.authentication_token = 'tok'
|
||||||
|
|
||||||
|
async def connect_side_effect():
|
||||||
|
opc_repository.client.uaclient.protocol = open_proto
|
||||||
|
|
||||||
|
opc_repository.client.connect = AsyncMock(side_effect=connect_side_effect)
|
||||||
|
|
||||||
|
result = await opc_repository._open_session()
|
||||||
|
|
||||||
opc_repository.client.connect.assert_called_once()
|
opc_repository.client.connect.assert_called_once()
|
||||||
assert opc_repository.last_reconnection_time is not None
|
assert opc_repository.last_reconnection_time is None
|
||||||
assert result == (True, {})
|
assert result == (True, {})
|
||||||
|
assert opc_repository._session_ready.is_set()
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_try_connect_fail(opc_repository):
|
async def test_open_session_raises_when_already_connected(opc_repository, mock_client):
|
||||||
opc_repository.last_reconnection_time = None
|
opc_repository.client = mock_client
|
||||||
opc_repository.disconnect = AsyncMock()
|
proto = MagicMock()
|
||||||
|
proto.state = 'open'
|
||||||
|
mock_client.uaclient = MagicMock(protocol=proto)
|
||||||
|
|
||||||
|
with pytest.raises(OpcSessionAlreadyConnectedError, match='disconnect'):
|
||||||
|
await opc_repository._open_session()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_open_session_fail(opc_repository):
|
||||||
|
opc_repository._disconnect_locked = AsyncMock()
|
||||||
opc_repository.client = MagicMock()
|
opc_repository.client = MagicMock()
|
||||||
opc_repository.client.connect.side_effect = Exception('Test error')
|
opc_repository.client.uaclient = MagicMock(protocol=MagicMock(state='closed'))
|
||||||
|
opc_repository.client.connect = AsyncMock(side_effect=Exception('Test error'))
|
||||||
|
|
||||||
is_connected, error_data = await opc_repository.try_connect()
|
is_connected, error_data = await opc_repository._open_session()
|
||||||
|
|
||||||
opc_repository.disconnect.assert_called_once()
|
opc_repository._disconnect_locked.assert_called_once()
|
||||||
opc_repository.client.connect.assert_called_once()
|
opc_repository.client.connect.assert_called_once()
|
||||||
assert is_connected is False
|
assert is_connected is False
|
||||||
assert error_data['notification_id'] == f'OPC_CONNECTION_ERROR_{opc_repository.id}'
|
assert error_data['notification_id'] == f'OPC_CONNECTION_ERROR_{opc_repository.id}'
|
||||||
@@ -158,25 +217,18 @@ async def test_try_connect_fail(opc_repository):
|
|||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_try_connect_no_client(opc_repository):
|
async def test_open_session_raises_when_no_client(opc_repository):
|
||||||
opc_repository.client = None
|
opc_repository.client = None
|
||||||
result = await opc_repository.try_connect()
|
|
||||||
assert result == (
|
with pytest.raises(OpcClientNotInitializedError, match='not initialized'):
|
||||||
False,
|
await opc_repository._open_session()
|
||||||
{
|
|
||||||
'notification_id': f'OPC_CONNECTION_ERROR_{opc_repository.id}',
|
|
||||||
'message': 'Client is not initialized',
|
|
||||||
'block': 'opc_repository',
|
|
||||||
'level': NotificationLevel.ERROR,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_disconnection_fallback_success(opc_repository, mock_client):
|
async def test_disconnection_fallback_success(opc_repository, mock_client):
|
||||||
opc_repository.client = mock_client
|
opc_repository.client = mock_client
|
||||||
mock_client.disconnect.return_value = True
|
mock_client.disconnect.return_value = True
|
||||||
result = await opc_repository.disconnection_fallback()
|
result = await opc_repository._disconnection_fallback()
|
||||||
|
|
||||||
mock_client.disconnect.assert_called_once()
|
mock_client.disconnect.assert_called_once()
|
||||||
assert result == []
|
assert result == []
|
||||||
@@ -186,7 +238,7 @@ async def test_disconnection_fallback_success(opc_repository, mock_client):
|
|||||||
async def test_disconnection_fallback_fail(opc_repository, mock_client):
|
async def test_disconnection_fallback_fail(opc_repository, mock_client):
|
||||||
opc_repository.client = mock_client
|
opc_repository.client = mock_client
|
||||||
mock_client.disconnect.side_effect = Exception('Test error')
|
mock_client.disconnect.side_effect = Exception('Test error')
|
||||||
result = await opc_repository.disconnection_fallback()
|
result = await opc_repository._disconnection_fallback()
|
||||||
assert result == [
|
assert result == [
|
||||||
{'attempt': 1, 'error': 'Test error', 'traceback': ANY},
|
{'attempt': 1, 'error': 'Test error', 'traceback': ANY},
|
||||||
{'attempt': 2, 'error': 'Test error', 'traceback': ANY},
|
{'attempt': 2, 'error': 'Test error', 'traceback': ANY},
|
||||||
@@ -200,10 +252,10 @@ async def test_disconnection_fallback_fail(opc_repository, mock_client):
|
|||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_disconnect(opc_repository, mock_client):
|
async def test_disconnect(opc_repository, mock_client):
|
||||||
opc_repository.client = mock_client
|
opc_repository.client = mock_client
|
||||||
opc_repository.disconnection_fallback = AsyncMock(return_value=[])
|
opc_repository._disconnection_fallback = AsyncMock(return_value=[])
|
||||||
await opc_repository.disconnect()
|
await opc_repository.disconnect()
|
||||||
|
|
||||||
opc_repository.disconnection_fallback.assert_called_once()
|
opc_repository._disconnection_fallback.assert_called_once()
|
||||||
assert opc_repository.client is None
|
assert opc_repository.client is None
|
||||||
|
|
||||||
|
|
||||||
@@ -216,12 +268,12 @@ async def test_disconnect_no_client(opc_repository):
|
|||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_disconnect_error(opc_repository, mock_client):
|
async def test_disconnect_error(opc_repository, mock_client):
|
||||||
opc_repository.client = mock_client
|
opc_repository.client = mock_client
|
||||||
opc_repository.disconnection_fallback = AsyncMock(
|
opc_repository._disconnection_fallback = AsyncMock(
|
||||||
return_value=[{'attempt': 1, 'error': 'Test error', 'traceback': 'text'}]
|
return_value=[{'attempt': 1, 'error': 'Test error', 'traceback': 'text'}]
|
||||||
)
|
)
|
||||||
await opc_repository.disconnect()
|
await opc_repository.disconnect()
|
||||||
|
|
||||||
opc_repository.disconnection_fallback.assert_called_once()
|
opc_repository._disconnection_fallback.assert_called_once()
|
||||||
opc_repository.send_notification_async.assert_called_once_with(
|
opc_repository.send_notification_async.assert_called_once_with(
|
||||||
metadata=opc_repository.metadata,
|
metadata=opc_repository.metadata,
|
||||||
notification_id=f'OPC_DISCONNECTION_ERROR_{opc_repository.id}',
|
notification_id=f'OPC_DISCONNECTION_ERROR_{opc_repository.id}',
|
||||||
@@ -238,91 +290,24 @@ async def test_disconnect_error(opc_repository, mock_client):
|
|||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_validate_connection_none_client(opc_repository):
|
async def test_validate_connection_none_client(opc_repository):
|
||||||
opc_repository.client = None
|
opc_repository.client = None
|
||||||
opc_repository.connect = AsyncMock(return_value=(True, {}))
|
|
||||||
response = await opc_repository.validate_connection()
|
response = await opc_repository.validate_connection()
|
||||||
assert response == (True, {})
|
assert response == (False, opc_repository._not_connected_error())
|
||||||
opc_repository.connect.assert_called_once()
|
|
||||||
|
|
||||||
|
|
||||||
# @pytest.mark.asyncio
|
|
||||||
# async def test_validate_connection_error_count_disconnect_error(opc_repository):
|
|
||||||
# opc_repository.error_count = 6
|
|
||||||
# opc_repository.client = AsyncMock()
|
|
||||||
# opc_repository.disconnect = AsyncMock(side_effect=Exception('Test error'))
|
|
||||||
# opc_repository.connect = AsyncMock(return_value=(True, {}))
|
|
||||||
|
|
||||||
# response = await opc_repository.validate_connection()
|
|
||||||
# assert response == opc_repository.connect.return_value
|
|
||||||
# opc_repository.disconnect.assert_called_once()
|
|
||||||
# opc_repository.connect.assert_called_once()
|
|
||||||
# opc_repository.logger.custom_error.assert_has_calls(
|
|
||||||
# [
|
|
||||||
# call('Failed to disconnect from OPC server: Test error', ANY),
|
|
||||||
# ]
|
|
||||||
# )
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_validate_connection_error_validate_connection_error(opc_repository):
|
async def test_validate_connection_session_not_open(opc_repository):
|
||||||
opc_repository.client = MagicMock(uaclient=Exception('Test error'))
|
|
||||||
opc_repository.error_count = 0
|
|
||||||
|
|
||||||
response = await opc_repository.validate_connection()
|
|
||||||
|
|
||||||
assert response == (
|
|
||||||
False,
|
|
||||||
{
|
|
||||||
'notification_id': f'OPC_CONNECTION_CHECK_ERROR_{opc_repository.id}',
|
|
||||||
'message': "Failed to validate connection to OPC server: 'Exception' object has no attribute 'protocol'",
|
|
||||||
'block': 'opc_repository',
|
|
||||||
'level': NotificationLevel.ERROR,
|
|
||||||
'attachment_content': ANY,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
@patch('laborious.utils.repository.opc_repository.datetime')
|
|
||||||
async def test_validate_connection_lost_not_time_to_reconnect(_mock_datetime, opc_repository):
|
|
||||||
_mock_datetime.now = MagicMock(return_value=datetime(2025, 1, 1, 0, 0, 0))
|
|
||||||
opc_repository.error_count = 0
|
|
||||||
opc_repository.client = MagicMock()
|
opc_repository.client = MagicMock()
|
||||||
opc_repository.client.uaclient.protocol = None
|
opc_repository.client.uaclient.protocol = None
|
||||||
opc_repository.last_reconnection_time = datetime(2025, 1, 1, 0, 0, 0)
|
|
||||||
opc_repository.connect = MagicMock(return_value=(True, {}))
|
|
||||||
|
|
||||||
response = await opc_repository.validate_connection()
|
response = await opc_repository.validate_connection()
|
||||||
opc_repository.connect.assert_not_called()
|
|
||||||
assert response == (
|
|
||||||
False,
|
|
||||||
{
|
|
||||||
'notification_id': f'OPC_CONNECTION_AWAITING_RECONNECTION_WINDOW_{opc_repository.id}',
|
|
||||||
'message': f'OPC server {opc_repository.id} is not connected, waiting for next reconnection window...',
|
|
||||||
'block': 'opc_repository',
|
|
||||||
'level': NotificationLevel.WARNING,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
|
assert response == (False, opc_repository._not_connected_error())
|
||||||
@pytest.mark.asyncio
|
opc_repository.error.assert_called_once()
|
||||||
@patch('laborious.utils.repository.opc_repository.datetime')
|
|
||||||
async def test_validate_connection_lost_time_to_reconnect(mock_datetime, opc_repository):
|
|
||||||
mock_datetime.now = MagicMock(return_value=datetime(2025, 1, 1, 1, 0, 0))
|
|
||||||
opc_repository.error_count = 0
|
|
||||||
opc_repository.client = AsyncMock()
|
|
||||||
opc_repository.client.uaclient.protocol = None
|
|
||||||
opc_repository.last_reconnection_time = datetime(2025, 1, 1, 0, 0, 0)
|
|
||||||
opc_repository.connect = AsyncMock(return_value=(True, {}))
|
|
||||||
|
|
||||||
response = await opc_repository.validate_connection()
|
|
||||||
opc_repository.connect.assert_called_once()
|
|
||||||
assert response == opc_repository.connect.return_value
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_validate_connection_success(opc_repository):
|
async def test_validate_connection_success(opc_repository):
|
||||||
opc_repository.client = MagicMock()
|
opc_repository.client = MagicMock()
|
||||||
opc_repository.error_count = 0
|
|
||||||
opc_repository.client.uaclient.protocol = MagicMock()
|
opc_repository.client.uaclient.protocol = MagicMock()
|
||||||
opc_repository.client.uaclient.protocol.state = 'open'
|
opc_repository.client.uaclient.protocol.state = 'open'
|
||||||
|
|
||||||
@@ -337,9 +322,7 @@ async def test_write_data_validate_connection_do_nothing(opc_repository):
|
|||||||
mock_node = AsyncMock()
|
mock_node = AsyncMock()
|
||||||
opc_repository.client.get_node.return_value = mock_node
|
opc_repository.client.get_node.return_value = mock_node
|
||||||
|
|
||||||
result = await opc_repository.write_data(
|
result = await opc_repository.write_data('ns=2;s=TestNode', 42.0, 'float', metadata['metadata'])
|
||||||
'ns=2;s=TestNode', 42.0, 'float', opc_repository.logger, metadata['metadata']
|
|
||||||
)
|
|
||||||
|
|
||||||
opc_repository.validate_connection.assert_called_once()
|
opc_repository.validate_connection.assert_called_once()
|
||||||
opc_repository.client.get_node.assert_called_once_with('ns=2;s=TestNode')
|
opc_repository.client.get_node.assert_called_once_with('ns=2;s=TestNode')
|
||||||
@@ -350,11 +333,8 @@ async def test_write_data_validate_connection_do_nothing(opc_repository):
|
|||||||
async def test_write_data_validate_connection_failed(opc_repository):
|
async def test_write_data_validate_connection_failed(opc_repository):
|
||||||
opc_repository.validate_connection = AsyncMock(return_value=(False, {}))
|
opc_repository.validate_connection = AsyncMock(return_value=(False, {}))
|
||||||
opc_repository.client = AsyncMock()
|
opc_repository.client = AsyncMock()
|
||||||
opc_repository.error_count = 0
|
|
||||||
|
|
||||||
result = await opc_repository.write_data(
|
result = await opc_repository.write_data('ns=2;s=TestNode', 42.0, 'float', metadata['metadata'])
|
||||||
'ns=2;s=TestNode', 42.0, 'float', opc_repository.logger, metadata['metadata']
|
|
||||||
)
|
|
||||||
|
|
||||||
opc_repository.validate_connection.assert_called_once()
|
opc_repository.validate_connection.assert_called_once()
|
||||||
opc_repository.client.get_node.assert_not_called()
|
opc_repository.client.get_node.assert_not_called()
|
||||||
@@ -365,11 +345,10 @@ async def test_write_data_validate_connection_failed(opc_repository):
|
|||||||
async def test_write_data_get_node_failed(opc_repository):
|
async def test_write_data_get_node_failed(opc_repository):
|
||||||
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
||||||
opc_repository.client = AsyncMock()
|
opc_repository.client = AsyncMock()
|
||||||
opc_repository.error_count = 0
|
|
||||||
opc_repository.client.get_node = MagicMock(side_effect=Exception('Test error'))
|
opc_repository.client.get_node = MagicMock(side_effect=Exception('Test error'))
|
||||||
|
|
||||||
is_success, error_data = await opc_repository.write_data(
|
is_success, error_data = await opc_repository.write_data(
|
||||||
'ns=2;s=TestNode', 42.0, 'float', opc_repository.logger, metadata['metadata']
|
'ns=2;s=TestNode', 42.0, 'float', metadata['metadata']
|
||||||
)
|
)
|
||||||
|
|
||||||
opc_repository.validate_connection.assert_called_once()
|
opc_repository.validate_connection.assert_called_once()
|
||||||
@@ -393,7 +372,7 @@ async def test_write_data_invalid_data_type(opc_repository, mock_client):
|
|||||||
mock_client.get_node = MagicMock(return_value=mock_node)
|
mock_client.get_node = MagicMock(return_value=mock_node)
|
||||||
|
|
||||||
is_success, error_data = await opc_repository.write_data(
|
is_success, error_data = await opc_repository.write_data(
|
||||||
'ns=2;s=TestNode', 42.0, 'invalid_type', opc_repository.logger, metadata['metadata']
|
'ns=2;s=TestNode', 42.0, 'invalid_type', metadata['metadata']
|
||||||
)
|
)
|
||||||
|
|
||||||
opc_repository.validate_connection.assert_called_once()
|
opc_repository.validate_connection.assert_called_once()
|
||||||
@@ -417,9 +396,7 @@ async def test_write_data(opc_repository, mock_client):
|
|||||||
mock_node = AsyncMock()
|
mock_node = AsyncMock()
|
||||||
mock_client.get_node = MagicMock(return_value=mock_node)
|
mock_client.get_node = MagicMock(return_value=mock_node)
|
||||||
|
|
||||||
result = await opc_repository.write_data(
|
result = await opc_repository.write_data('ns=2;s=TestNode', 42.0, 'float', metadata['metadata'])
|
||||||
'ns=2;s=TestNode', 42.0, 'float', opc_repository.logger, metadata['metadata']
|
|
||||||
)
|
|
||||||
|
|
||||||
mock_client.get_node.assert_called_once_with('ns=2;s=TestNode')
|
mock_client.get_node.assert_called_once_with('ns=2;s=TestNode')
|
||||||
mock_node.write_value.assert_called_once()
|
mock_node.write_value.assert_called_once()
|
||||||
@@ -431,12 +408,11 @@ async def test_write_data_write_value_failed(opc_repository, mock_client):
|
|||||||
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
||||||
opc_repository.client = mock_client
|
opc_repository.client = mock_client
|
||||||
mock_node = AsyncMock()
|
mock_node = AsyncMock()
|
||||||
opc_repository.error_count = 0
|
|
||||||
mock_client.get_node = MagicMock(return_value=mock_node)
|
mock_client.get_node = MagicMock(return_value=mock_node)
|
||||||
mock_node.write_value.side_effect = Exception('Test error')
|
mock_node.write_value.side_effect = Exception('Test error')
|
||||||
|
|
||||||
is_success, error_data = await opc_repository.write_data(
|
is_success, error_data = await opc_repository.write_data(
|
||||||
'ns=2;s=TestNode', 42.0, 'float', opc_repository.logger, metadata['metadata']
|
'ns=2;s=TestNode', 42.0, 'float', metadata['metadata']
|
||||||
)
|
)
|
||||||
|
|
||||||
opc_repository.validate_connection.assert_called_once()
|
opc_repository.validate_connection.assert_called_once()
|
||||||
@@ -451,3 +427,108 @@ async def test_write_data_write_value_failed(opc_repository, mock_client):
|
|||||||
assert error_data['block'] == 'opc_repository'
|
assert error_data['block'] == 'opc_repository'
|
||||||
assert error_data['level'] == NotificationLevel.ERROR
|
assert error_data['level'] == NotificationLevel.ERROR
|
||||||
assert error_data['attachment_content'] is not None
|
assert error_data['attachment_content'] is not None
|
||||||
|
|
||||||
|
|
||||||
|
def test_is_reconnectable_opcua_bad():
|
||||||
|
assert is_reconnectable_opcua_bad(BadSessionIdInvalid()) is True
|
||||||
|
assert is_reconnectable_opcua_bad(BadNodeIdUnknown()) is False
|
||||||
|
assert is_reconnectable_opcua_bad(Exception('other')) is False
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_write_data_bad_session_id_invalid_schedules_reconnect(opc_repository, mock_client):
|
||||||
|
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
||||||
|
opc_repository.client = mock_client
|
||||||
|
opc_repository._start_reconnect_on_bad = AsyncMock()
|
||||||
|
mock_node = AsyncMock()
|
||||||
|
mock_client.get_node = MagicMock(return_value=mock_node)
|
||||||
|
mock_node.write_value.side_effect = BadSessionIdInvalid()
|
||||||
|
|
||||||
|
is_success, error_data = await opc_repository.write_data(
|
||||||
|
'ns=2;s=TestNode', 42.0, 'float', metadata['metadata']
|
||||||
|
)
|
||||||
|
|
||||||
|
mock_node.write_value.assert_called_once()
|
||||||
|
opc_repository._start_reconnect_on_bad.assert_called_once()
|
||||||
|
assert is_success is False
|
||||||
|
assert error_data['opc_error_kind'] == 'session_bad'
|
||||||
|
assert error_data['opc_status'] == 'BadSessionIdInvalid'
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_write_data_reconnect_in_progress_immediate(opc_repository):
|
||||||
|
opc_repository._session_ready.clear()
|
||||||
|
opc_repository._reconnect_task = asyncio.create_task(asyncio.sleep(60))
|
||||||
|
opc_repository.validate_connection = AsyncMock()
|
||||||
|
|
||||||
|
is_success, error_data = await opc_repository.write_data(
|
||||||
|
'ns=2;s=TestNode', 42.0, 'float', metadata['metadata']
|
||||||
|
)
|
||||||
|
|
||||||
|
opc_repository._reconnect_task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await opc_repository._reconnect_task
|
||||||
|
opc_repository._reconnect_task = None
|
||||||
|
|
||||||
|
opc_repository.validate_connection.assert_not_called()
|
||||||
|
assert is_success is False
|
||||||
|
assert error_data['opc_error_kind'] == 'reconnect_in_progress'
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_start_reconnect_on_bad_skips_within_interval(opc_repository):
|
||||||
|
opc_repository.last_reconnection_time = datetime.now()
|
||||||
|
opc_repository.reconnection_interval = 3600
|
||||||
|
|
||||||
|
await opc_repository._start_reconnect_on_bad('BadSessionIdInvalid', 'tok')
|
||||||
|
|
||||||
|
assert opc_repository._reconnect_task is None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_parallel_bad_writes_single_reconnect_task(opc_repository, mock_client):
|
||||||
|
opc_repository.validate_connection = AsyncMock(return_value=(True, {}))
|
||||||
|
opc_repository.client = mock_client
|
||||||
|
opc_repository.reconnection_interval = 0
|
||||||
|
opc_repository.last_reconnection_time = None
|
||||||
|
mock_node = AsyncMock()
|
||||||
|
mock_client.get_node = MagicMock(return_value=mock_node)
|
||||||
|
mock_node.write_value.side_effect = BadSessionIdInvalid()
|
||||||
|
|
||||||
|
connect_count = 0
|
||||||
|
|
||||||
|
async def slow_reconnect():
|
||||||
|
nonlocal connect_count
|
||||||
|
connect_count += 1
|
||||||
|
await asyncio.sleep(0.05)
|
||||||
|
opc_repository._session_ready.set()
|
||||||
|
return True, {}
|
||||||
|
|
||||||
|
opc_repository._reconnect_locked = slow_reconnect
|
||||||
|
|
||||||
|
results = await asyncio.gather(
|
||||||
|
opc_repository.write_data('ns=2;s=TestNode', 1.0, 'float', metadata['metadata']),
|
||||||
|
opc_repository.write_data('ns=2;s=TestNode2', 2.0, 'float', metadata['metadata']),
|
||||||
|
)
|
||||||
|
await asyncio.sleep(0.15)
|
||||||
|
|
||||||
|
assert connect_count <= 1
|
||||||
|
assert 1 <= mock_node.write_value.call_count <= 2
|
||||||
|
error_kinds = [r[1].get('opc_error_kind') for r in results]
|
||||||
|
assert error_kinds.count('session_bad') >= 1
|
||||||
|
assert all(k in ('session_bad', 'reconnect_in_progress') for k in error_kinds)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
@patch('laborious.utils.repository.opc_repository.datetime')
|
||||||
|
async def test_reconnect_locked_sets_last_reconnection_time(mock_datetime, opc_repository):
|
||||||
|
mock_datetime.now = MagicMock(return_value=datetime(2025, 1, 1, 12, 0, 0))
|
||||||
|
opc_repository._disconnect_locked = AsyncMock()
|
||||||
|
opc_repository._connect_locked = AsyncMock(return_value=(True, {}))
|
||||||
|
|
||||||
|
result = await opc_repository._reconnect_locked()
|
||||||
|
|
||||||
|
opc_repository._disconnect_locked.assert_called_once()
|
||||||
|
opc_repository._connect_locked.assert_called_once()
|
||||||
|
assert result == (True, {})
|
||||||
|
assert opc_repository.last_reconnection_time == datetime(2025, 1, 1, 12, 0, 0)
|
||||||
|
|||||||
Reference in New Issue
Block a user