diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 4bb32578..e17d241a 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -2,10 +2,14 @@ name: Integration Test on: pull_request: - types: [ready_for_review] + types: [opened, synchronize, reopened] branches: - main +concurrency: + group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + jobs: docker: timeout-minutes: 30 @@ -21,7 +25,7 @@ jobs: version: latest - name: Start Flight Blender and dependencies - run: docker compose --env-file .env.tests -f docker-compose.fb.yml up -d + run: docker compose --env-file .env.tests -f docker-compose.fb.yml up -d --wait working-directory: ./tests - name: Install uv @@ -34,6 +38,7 @@ jobs: run: uv run openutm-verify --debug --config config/default.yaml - uses: actions/upload-artifact@v4 + if: always() with: name: test-reports path: reports/ diff --git a/config/default.yaml b/config/default.yaml index 130120f3..49223fe3 100644 --- a/config/default.yaml +++ b/config/default.yaml @@ -37,18 +37,19 @@ data_files: # List of test scenario IDs to execute scenarios: - # "F1_happy_path": - # telemetry: "config/bern/telemetry_f1.json" - # "F2_contingent_path": - # telemetry: "config/bern/telemetry_f2.json" - # "F3_non_conforming_path": - # telemetry: "config/bern/telemetry_f3.json" - # "F5_non_conforming_path": + F1_happy_path: + telemetry: "config/bern/telemetry_f1.json" + F2_contingent_path: + telemetry: "config/bern/telemetry_f2.json" + F3_non_conforming_path: + telemetry: "config/bern/telemetry_f3.json" + # F5_non_conforming_path: # telemetry: "config/bern/telemetry_f5.json" - # "opensky_live_data": - # "add_flight_declaration": - # "geo_fence_upload": - # "sdsp_track_heartbeat": + opensky_live_data: + add_flight_declaration: + geo_fence_upload: + # sdsp_track: + # sdsp_heartbeat: openutm_sim_air_traffic_data: # Reporting configuration diff --git a/src/openutm_verification/cli/__init__.py b/src/openutm_verification/cli/__init__.py index 95a170ee..07b98c4b 100644 --- a/src/openutm_verification/cli/__init__.py +++ b/src/openutm_verification/cli/__init__.py @@ -2,6 +2,7 @@ Command Line Interface for OpenUTM Verification Tool. """ +import sys from datetime import datetime, timezone from pathlib import Path @@ -41,12 +42,13 @@ def main(): log_file = setup_logging(output_dir, base_filename, config.reporting.formats, args.debug) # Run verification scenarios - run_verification_scenarios(config, args.config) + failed = run_verification_scenarios(config, args.config) if log_file: from loguru import logger logger.info(f"Log file saved to: {log_file}") + sys.exit(failed) if __name__ == "__main__": diff --git a/src/openutm_verification/core/clients/flight_blender/flight_blender_client.py b/src/openutm_verification/core/clients/flight_blender/flight_blender_client.py index 9641de1b..bcd74e92 100644 --- a/src/openutm_verification/core/clients/flight_blender/flight_blender_client.py +++ b/src/openutm_verification/core/clients/flight_blender/flight_blender_client.py @@ -1,6 +1,7 @@ import json import time import uuid +from contextlib import contextmanager from dataclasses import asdict from typing import Any, Dict, List, Optional @@ -10,8 +11,12 @@ from openutm_verification.core.clients.flight_blender.base_client import ( BaseBlenderAPIClient, ) +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.core.execution.scenario_runner import scenario_step -from openutm_verification.core.reporting.reporting_models import Status, StepResult +from openutm_verification.core.reporting.reporting_models import ( + Status, + StepResult, +) from openutm_verification.models import ( FlightBlenderError, HeartbeatMessage, @@ -75,6 +80,8 @@ def __init__(self, base_url: str, credentials: Dict[str, Any], request_timeout: self.latest_geo_fence_id: Optional[str] = None # Context: store the most recently created flight declaration id for teardown/steps self.latest_flight_declaration_id: Optional[str] = None + # Context: store the generated telemetry states for the current scenario + self.telemetry_states: Optional[List[Dict[str, Any]]] = None logger.debug(f"Initialized FlightBlenderClient with base_url={base_url}, request_timeout={request_timeout}") def __exit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None: @@ -89,11 +96,10 @@ def __exit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None: return super().__exit__(exc_type, exc_val, exc_tb) @scenario_step("Upload Geo Fence") - def upload_geo_fence(self, operation_id: Optional[str] = None, filename: Optional[str] = None) -> Dict[str, Any]: + def upload_geo_fence(self, filename: Optional[str] = None) -> Dict[str, Any]: """Upload an Area-of-Interest (Geo Fence) to Flight Blender. Args: - operation_id: Not used for geo-fence upload (included for API consistency). filename: Path to the GeoJSON file containing the geo-fence definition. Returns: @@ -121,12 +127,9 @@ def upload_geo_fence(self, operation_id: Optional[str] = None, filename: Optiona return body @scenario_step("Get Geo Fence") - def get_geo_fence(self, operation_id: Optional[str] = None) -> Dict[str, Any]: + def get_geo_fence(self) -> Dict[str, Any]: """Retrieve the details of the most recently uploaded geo-fence. - Args: - operation_id: Not used for geo-fence retrieval (included for API consistency). - Returns: The JSON response from the API containing geo-fence details, or a dict indicating skip if no geo-fence ID is available. @@ -237,13 +240,12 @@ def upload_flight_declaration(self, declaration: str | Any) -> Dict[str, Any]: return response_json @scenario_step("Update Operation State") - def update_operation_state(self, operation_id: str, new_state: OperationState, duration_seconds: int = 0) -> Dict[str, Any]: + def update_operation_state(self, new_state: OperationState, duration_seconds: int = 0) -> Dict[str, Any]: """Update the state of a flight operation. Posts the new state and optionally waits for the specified duration. Args: - operation_id: The ID of the operation to update. new_state: The new OperationState to set. duration_seconds: Optional seconds to sleep after update (default 0). @@ -253,12 +255,12 @@ def update_operation_state(self, operation_id: str, new_state: OperationState, d Raises: FlightBlenderError: If the update request fails. """ - endpoint = f"/flight_declaration_ops/flight_declaration_state/{operation_id}" - logger.debug(f"Updating operation {operation_id} to state {new_state.name}") + endpoint = f"/flight_declaration_ops/flight_declaration_state/{self.latest_flight_declaration_id}" + logger.debug(f"Updating operation {self.latest_flight_declaration_id} to state {new_state.name}") payload = {"state": new_state.value, "submitted_by": "hh@auth.com"} response = self.put(endpoint, json=payload) - logger.info(f"Operation state updated for {operation_id} to {new_state.name}") + logger.info(f"Operation state updated for {self.latest_flight_declaration_id} to {new_state.name}") if duration_seconds > 0: logger.debug(f"Sleeping for {duration_seconds} seconds after state update") time.sleep(duration_seconds) @@ -281,11 +283,10 @@ def _load_telemetry_file(self, filename: str) -> List[Dict[str, Any]]: rid_json = json.loads(rid_json_file.read()) return rid_json["current_states"] - def _submit_telemetry_states_impl(self, operation_id: str, states: List[Dict[str, Any]], duration_seconds: int = 0) -> Optional[Dict[str, Any]]: + def _submit_telemetry_states_impl(self, states: List[Dict[str, Any]], duration_seconds: int = 0) -> Optional[Dict[str, Any]]: """Internal implementation for submitting telemetry states. Args: - operation_id: The ID of the operation for telemetry submission. states: List of telemetry state dictionaries. duration_seconds: Optional maximum duration in seconds to submit telemetry (default 0 for unlimited). @@ -296,9 +297,9 @@ def _submit_telemetry_states_impl(self, operation_id: str, states: List[Dict[str FlightBlenderError: If maximum waiting time is exceeded due to rate limits. """ endpoint = "/flight_stream/set_telemetry" - logger.debug(f"Submitting telemetry for operation {operation_id}") + logger.debug(f"Submitting telemetry for operation {self.latest_flight_declaration_id}") - rid_operator_details = _create_rid_operator_details(operation_id) + rid_operator_details = _create_rid_operator_details(self.latest_flight_declaration_id) last_response = None maximum_waiting_time = 10.0 @@ -337,14 +338,13 @@ def _submit_telemetry_states_impl(self, operation_id: str, states: List[Dict[str return last_response @scenario_step("Submit Telemetry (from file)") - def submit_telemetry_from_file(self, operation_id: str, filename: str, duration_seconds: int = 0) -> Optional[Dict[str, Any]]: + def submit_telemetry_from_file(self, filename: str, duration_seconds: int = 0) -> Optional[Dict[str, Any]]: """Submit telemetry data for a flight operation. Loads telemetry states from file and submits them sequentially, with optional duration limiting and error handling for rate limits. Args: - operation_id: The ID of the operation for telemetry submission. filename: Path to the JSON file containing telemetry data. duration_seconds: Optional maximum duration in seconds to submit telemetry (default 0 for unlimited). @@ -355,7 +355,7 @@ def submit_telemetry_from_file(self, operation_id: str, filename: str, duration_ FlightBlenderError: If maximum waiting time is exceeded due to rate limits. """ states = self._load_telemetry_file(filename) - return self._submit_telemetry_states_impl(operation_id, states, duration_seconds) + return self._submit_telemetry_states_impl(states, duration_seconds) @scenario_step("Wait X seconds") def wait_x_seconds(self, wait_time_seconds: int = 5) -> str: @@ -366,15 +366,14 @@ def wait_x_seconds(self, wait_time_seconds: int = 5) -> str: return f"Waited for Flight Blender to process {wait_time_seconds} seconds." @scenario_step("Submit Telemetry") - def submit_telemetry(self, operation_id: str, states: List[Dict[str, Any]], duration_seconds: int = 0) -> Optional[Dict[str, Any]]: + def submit_telemetry(self, states: Optional[List[Dict[str, Any]]] = None, duration_seconds: int = 0) -> Optional[Dict[str, Any]]: """Submit telemetry data for a flight operation from in-memory states. Submits telemetry states sequentially from the provided list, with optional duration limiting and error handling for rate limits. Args: - operation_id: The ID of the operation for telemetry submission. - states: List of telemetry state dictionaries. + states: List of telemetry state dictionaries. If None, uses the generated telemetry states from context. duration_seconds: Optional maximum duration in seconds to submit telemetry (default 0 for unlimited). Returns: @@ -383,12 +382,15 @@ def submit_telemetry(self, operation_id: str, states: List[Dict[str, Any]], dura Raises: FlightBlenderError: If maximum waiting time is exceeded due to rate limits. """ - return self._submit_telemetry_states_impl(operation_id, states, duration_seconds) + telemetry_states = states or self.telemetry_states + if telemetry_states is None: + raise ValueError("Telemetry states are required and could not be resolved from context.") + + return self._submit_telemetry_states_impl(telemetry_states, duration_seconds) @scenario_step("Check Operation State") def check_operation_state( self, - operation_id: str, expected_state: OperationState, duration_seconds: int = 0, ) -> str: @@ -398,30 +400,27 @@ def check_operation_state( and returns a success status. Args: - operation_id: The ID of the operation to check. expected_state: The expected OperationState. duration_seconds: Seconds to wait for processing. Returns: A dictionary with the check result. """ - logger.info(f"Checking operation state for {operation_id} (simulated)...") + logger.info(f"Checking operation state for {self.latest_flight_declaration_id} (simulated)...") logger.info(f"Waiting for {duration_seconds} seconds for Flight Blender to process state...") time.sleep(duration_seconds) - logger.info(f"Flight state check for {operation_id} completed (simulated).") + logger.info(f"Flight state check for {self.latest_flight_declaration_id} completed (simulated).") return f"Waited for Flight Blender to process {expected_state} state." @scenario_step("Check Operation State Connected") def check_operation_state_connected( self, - operation_id: str, expected_state: OperationState, duration_seconds: int = 0, ) -> Dict[str, Any]: """Check the operation state by polling the API until the expected state is reached. Args: - operation_id: The ID of the operation to check. expected_state: The expected OperationState. duration_seconds: Maximum seconds to poll for the state. @@ -431,36 +430,36 @@ def check_operation_state_connected( Raises: FlightBlenderError: If the expected state is not reached within the timeout. """ - endpoint = f"/flight_declaration_ops/flight_declaration/{operation_id}" - logger.info(f"Checking operation state for {operation_id}, expecting {expected_state.name}") + endpoint = f"/flight_declaration_ops/flight_declaration/{self.latest_flight_declaration_id}" + logger.info(f"Checking operation state for {self.latest_flight_declaration_id}, expecting {expected_state.name}") start_time = time.time() while time.time() - start_time < duration_seconds: response = self.get(endpoint) data = response.json() current_state_value = data.get("state") - logger.debug(f"Current state for {operation_id}: {current_state_value}") + logger.debug(f"Current state for {self.latest_flight_declaration_id}: {current_state_value}") if current_state_value == expected_state.value: - logger.info(f"Operation {operation_id} reached expected state {expected_state.name}") + logger.info(f"Operation {self.latest_flight_declaration_id} reached expected state {expected_state.name}") return data time.sleep(1) - logger.error(f"Operation {operation_id} did not reach expected state {expected_state.name} within {duration_seconds} seconds") - raise FlightBlenderError(f"Operation {operation_id} did not reach expected state {expected_state.name} within {duration_seconds} seconds") + logger.error( + f"Operation {self.latest_flight_declaration_id} did not reach expected state {expected_state.name} within {duration_seconds} seconds" + ) + raise FlightBlenderError( + f"Operation {self.latest_flight_declaration_id} did not reach expected state {expected_state.name} within {duration_seconds} seconds" + ) @scenario_step("Delete Flight Declaration") - def delete_flight_declaration(self, operation_id: Optional[str] = None) -> Dict[str, Any]: + def delete_flight_declaration(self) -> Dict[str, Any]: """Delete a flight declaration by ID. - Args: - operation_id: Optional ID of the flight declaration to delete. If not provided, - uses the latest uploaded flight declaration ID. - Returns: A dictionary with deletion status, including whether it was successful. """ - op_id = operation_id or self.latest_flight_declaration_id + op_id = self.latest_flight_declaration_id if not op_id: logger.warning("No flight declaration ID available for deletion") return { @@ -600,6 +599,7 @@ def initialize_heartbeat_websocket_connection(self, session_id: str) -> Any: ws = self.create_websocket_connection(endpoint=endpoint) return ws + @scenario_step("Verify SDSP Track") def initialize_verify_sdsp_track( self, expected_heartbeat_interval_seconds: int, @@ -672,6 +672,7 @@ def initialize_verify_sdsp_track( duration=duration, ) + @scenario_step("Verify SDSP Heartbeat") def initialize_verify_sdsp_heartbeat( self, expected_heartbeat_interval_seconds: int, @@ -746,3 +747,31 @@ def initialize_verify_sdsp_heartbeat( def close_heartbeat_websocket_connection(self, ws_connection: Any) -> None: ws_connection.close() + + @scenario_step("Setup Flight Declaration") + def setup_flight_declaration(self, flight_declaration_path: str, telemetry_path: str) -> None: + """Generates data and uploads flight declaration.""" + from openutm_verification.scenarios.common import ( + generate_flight_declaration, + generate_telemetry, + ) + + flight_declaration = generate_flight_declaration(flight_declaration_path) + telemetry_states = generate_telemetry(telemetry_path) + + self.telemetry_states = telemetry_states + + upload_result = self.upload_flight_declaration(flight_declaration) + + if upload_result.status == Status.FAIL: + logger.error(f"Flight declaration upload failed: {upload_result}") + raise FlightBlenderError("Failed to upload flight declaration during setup_flight_declaration") + + @contextmanager + def flight_declaration(self, data_files: DataFiles): + """Context manager to setup and teardown a flight operation based on scenario config.""" + self.setup_flight_declaration(data_files.flight_declaration, data_files.telemetry) + try: + yield + finally: + self.delete_flight_declaration() diff --git a/src/openutm_verification/core/execution/dependencies.py b/src/openutm_verification/core/execution/dependencies.py index f9666df6..eaec477d 100644 --- a/src/openutm_verification/core/execution/dependencies.py +++ b/src/openutm_verification/core/execution/dependencies.py @@ -8,7 +8,7 @@ from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient from openutm_verification.core.clients.opensky.base_client import create_opensky_settings from openutm_verification.core.clients.opensky.opensky_client import OpenSkyClient -from openutm_verification.core.execution.config_models import AppConfig, ScenarioId, get_settings +from openutm_verification.core.execution.config_models import AppConfig, DataFiles, ScenarioId, get_settings from openutm_verification.core.execution.dependency_resolution import CONTEXT, dependency from openutm_verification.core.reporting.reporting_models import ScenarioResult from openutm_verification.scenarios.registry import SCENARIO_REGISTRY @@ -46,6 +46,23 @@ def scenario_id() -> Generator[ScenarioId, None, None]: yield CONTEXT.get()["scenario_id"] +@dependency(DataFiles) +def data_files(scenario_id: ScenarioId) -> Generator[DataFiles, None, None]: + """Provides data files configuration for dependency injection. + + Returns: + An instance of DataFiles. + """ + config = get_settings() + scenario_config = config.scenarios.get(scenario_id) or config.data_files + data = DataFiles( + telemetry=scenario_config.telemetry or config.data_files.telemetry, + flight_declaration=scenario_config.flight_declaration or config.data_files.flight_declaration, + geo_fence=scenario_config.geo_fence or config.data_files.geo_fence, + ) + yield data + + @dependency(AppConfig) def app_config() -> Generator[AppConfig, None, None]: """Provides the application configuration for dependency injection. diff --git a/src/openutm_verification/core/execution/execution.py b/src/openutm_verification/core/execution/execution.py index 84ab3754..bc3da4bc 100644 --- a/src/openutm_verification/core/execution/execution.py +++ b/src/openutm_verification/core/execution/execution.py @@ -113,3 +113,4 @@ def run_verification_scenarios(config: AppConfig, config_path: Path): base_filename = f"report_{run_timestamp.strftime('%Y-%m-%dT%H-%M-%SZ')}" generate_reports(report_data, config.reporting, base_filename) + return failed_scenarios diff --git a/src/openutm_verification/core/execution/scenario_runner.py b/src/openutm_verification/core/execution/scenario_runner.py index 42495f99..97f1402a 100644 --- a/src/openutm_verification/core/execution/scenario_runner.py +++ b/src/openutm_verification/core/execution/scenario_runner.py @@ -1,6 +1,8 @@ +import contextvars import time +from dataclasses import dataclass, field from functools import wraps -from typing import Any, Callable +from typing import Any, Callable, List, Optional, ParamSpec, Protocol, TypeVar, cast, overload from loguru import logger @@ -8,11 +10,64 @@ from openutm_verification.core.reporting.reporting_models import Status, StepResult from openutm_verification.models import FlightBlenderError +T = TypeVar("T") +P = ParamSpec("P") +R = TypeVar("R", bound=StepResult[Any]) -def scenario_step(step_name: str) -> Callable: - def decorator(func: Callable) -> Callable: + +@dataclass +class ScenarioState: + steps: List[StepResult[Any]] = field(default_factory=list) + active: bool = False + + +_scenario_state: contextvars.ContextVar[Optional[ScenarioState]] = contextvars.ContextVar("scenario_state", default=None) + + +class ScenarioContext: + def __init__(self): + self._token = None + self._state: Optional[ScenarioState] = None + + def __enter__(self): + self._state = ScenarioState(active=True) + self._token = _scenario_state.set(self._state) + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + if self._state: + self._state.active = False + if self._token: + _scenario_state.reset(self._token) + + @classmethod + def add_result(cls, result: StepResult[Any]) -> None: + state = _scenario_state.get() + if state and state.active: + state.steps.append(result) + + @property + def steps(self) -> List[StepResult[Any]]: + if self._state: + return self._state.steps + state = _scenario_state.get() + return state.steps if state else [] + + +class StepDecorator(Protocol): + @overload + def __call__(self, func: Callable[P, R]) -> Callable[P, R]: ... + + @overload + def __call__(self, func: Callable[P, T]) -> Callable[P, StepResult[T]]: ... + + def __call__(self, func: Callable[P, Any]) -> Callable[P, Any]: ... + + +def scenario_step(step_name: str) -> StepDecorator: + def decorator(func: Callable[P, Any]) -> Callable[P, StepResult[Any]]: @wraps(func) - def wrapper(*args: Any, **kwargs: Any) -> StepResult: + def wrapper(*args: P.args, **kwargs: P.kwargs) -> StepResult[Any]: logger.info("-" * 50) logger.info(f"Executing step: '{step_name}'...") start_time = time.time() @@ -20,26 +75,37 @@ def wrapper(*args: Any, **kwargs: Any) -> StepResult: result = func(*args, **kwargs) duration = time.time() - start_time logger.info(f"Step '{step_name}' successful in {duration:.2f} seconds.") - return StepResult(name=step_name, status=Status.PASS, duration=duration, details=result) + + if isinstance(result, StepResult): + step_result = result + else: + step_result = StepResult(name=step_name, status=Status.PASS, duration=duration, details=result) + + ScenarioContext.add_result(step_result) + return step_result except (FlightBlenderError, OpenSkyError) as e: duration = time.time() - start_time logger.error(f"Step '{step_name}' failed after {duration:.2f} seconds: {e}") - return StepResult( + step_result = StepResult( name=step_name, status=Status.FAIL, duration=duration, error_message=str(e), ) + ScenarioContext.add_result(step_result) + return step_result except Exception as e: duration = time.time() - start_time logger.error(f"Step '{step_name}' encountered an unexpected error after {duration:.2f} seconds: {e}") - return StepResult( + step_result = StepResult( name=step_name, status=Status.FAIL, duration=duration, error_message=f"Unexpected error: {e}", ) + ScenarioContext.add_result(step_result) + return step_result return wrapper - return decorator + return cast(StepDecorator, decorator) diff --git a/src/openutm_verification/core/reporting/reporting_models.py b/src/openutm_verification/core/reporting/reporting_models.py index d428108b..a9996f1f 100644 --- a/src/openutm_verification/core/reporting/reporting_models.py +++ b/src/openutm_verification/core/reporting/reporting_models.py @@ -3,7 +3,7 @@ """ from enum import StrEnum -from typing import Any, Dict, List, Optional +from typing import Any, Dict, Generic, List, Optional, TypeVar from pydantic import BaseModel @@ -17,13 +17,16 @@ class Status(StrEnum): FAIL = "FAIL" -class StepResult(BaseModel): +T = TypeVar("T") + + +class StepResult(BaseModel, Generic[T]): """Data model for a single step within a scenario.""" name: str status: Status duration: float - details: Optional[Any] = None + details: T = None # type: ignore error_message: Optional[str] = None @@ -33,7 +36,7 @@ class ScenarioResult(BaseModel): name: str status: Status duration_seconds: float - steps: List[StepResult] + steps: List[StepResult[Any]] error_message: Optional[str] = None flight_declaration_filename: Optional[str] = None telemetry_filename: Optional[str] = None diff --git a/src/openutm_verification/scenarios/common.py b/src/openutm_verification/scenarios/common.py index 49e36171..4734150d 100644 --- a/src/openutm_verification/scenarios/common.py +++ b/src/openutm_verification/scenarios/common.py @@ -1,23 +1,9 @@ import json -from functools import partial from pathlib import Path -from typing import Any, List, cast +from typing import Any, List from loguru import logger -from openutm_verification.core.clients.air_traffic.air_traffic_client import ( - AirTrafficClient, -) -from openutm_verification.core.clients.flight_blender.flight_blender_client import ( - FlightBlenderClient, -) -from openutm_verification.core.clients.opensky.opensky_client import OpenSkyClient -from openutm_verification.core.execution.config_models import config -from openutm_verification.core.reporting.reporting_models import ( - ScenarioResult, - Status, - StepResult, -) from openutm_verification.simulator.flight_declaration import FlightDeclarationGenerator from openutm_verification.simulator.geo_json_telemetry import GeoJSONFlightsSimulator from openutm_verification.simulator.models.flight_data_types import ( @@ -27,109 +13,7 @@ DEFAULT_TELEMETRY_DURATION = 30 # seconds -def _callable_name(func_like: Any) -> str: - """Best-effort name extraction for partials or callables.""" - target = getattr(func_like, "func", func_like) - return getattr(target, "__name__", "") - - -def _redact_fetch_details(res: StepResult) -> tuple[StepResult, Any | None]: - """Normalize fetch result details to only include count and extract observations. - - Returns the possibly-modified StepResult and the extracted observations (or None). - """ - if res.status != Status.PASS: - return res, None - - details = res.details - observations: Any | None = None - if isinstance(details, dict) and "observations" in details: - observations = details.get("observations") - try: - res.details = {"count": len(observations or [])} - except TypeError: - res.details = {"count": 0} - elif isinstance(details, list): - observations = details - res.details = {"count": len(details)} - else: - res.details = {"count": 0} - observations = None - return res, observations - - -def _run_submit_airtraffic_flow(steps: list[partial[Any]]) -> List[StepResult]: - """Execute OpenSky flow steps returning the list of StepResults.""" - - def _execute_step(step_func: partial[Any], current_observations: Any | None) -> tuple[StepResult, Any | None]: - name = _callable_name(step_func) - # Submit step consumes observations - if "submit_air_traffic" in name: - if current_observations: - return ( - step_func(observations=current_observations), - current_observations, - ) - return ( - StepResult( - name="Submit Air Traffic (skipped)", - status=Status.PASS, - duration=0.0, - details="No observations to submit", - ), - current_observations, - ) - - # Generic execution (fetch steps usually have opensky_client already bound via partial) - res: StepResult = step_func() - if "fetch" in name or "generate" in name: - res, observations = _redact_fetch_details(res) - return res, observations - return res, current_observations - - results: List[StepResult] = [] - observations: Any | None = None - for step in steps: - step_result, observations = _execute_step(step, observations) - results.append(step_result) - if step_result.status == Status.FAIL: - break - return results - - -def _run_declaration_flow( - fb_client: FlightBlenderClient, - flight_declaration: Any, - telemetry_states: List[Any], - steps: list[partial[Any]], -) -> List[StepResult]: - """Execute standard declaration + steps + teardown using generated data and return StepResults.""" - upload_result = cast(StepResult, fb_client.upload_flight_declaration(flight_declaration)) - if upload_result.status == Status.FAIL or upload_result.details is None: - # Return early with failure - return [upload_result] - - operation_id = upload_result.details["id"] - - all_steps: List[StepResult] = [upload_result] - for step_func in steps: - kwargs = {} - if "submit_telemetry" in _callable_name(step_func): - kwargs["states"] = telemetry_states - elif "submit_telemetry_from_file" in _callable_name(step_func): - # TODO: read from file - kwargs["states"] = telemetry_states - step_result: StepResult = step_func(operation_id, **kwargs) - all_steps.append(step_result) - if step_result.status == Status.FAIL: - break - - teardown_result: StepResult = cast(StepResult, fb_client.delete_flight_declaration(operation_id)) - all_steps.append(teardown_result) - return all_steps - - -def _generate_flight_declaration(config_path: str) -> Any: +def generate_flight_declaration(config_path: str) -> Any: """Generate a flight declaration from the config file at the given path.""" try: generator = FlightDeclarationGenerator(bounds_path=Path(config_path)) @@ -139,7 +23,7 @@ def _generate_flight_declaration(config_path: str) -> Any: raise -def _generate_telemetry(config_path: str, duration: int = DEFAULT_TELEMETRY_DURATION) -> List[Any]: +def generate_telemetry(config_path: str, duration: int = DEFAULT_TELEMETRY_DURATION) -> List[Any]: """Generate telemetry states from the GeoJSON config file at the given path.""" try: logger.debug(f"Generating telemetry states from {config_path} for duration {duration} seconds") @@ -155,230 +39,6 @@ def _generate_telemetry(config_path: str, duration: int = DEFAULT_TELEMETRY_DURA logger.error(f"Failed to generate telemetry states from {config_path}: {e}") raise - -def run_sdsp_scenario_template( - scenario_id: str, - *, - fb_client: FlightBlenderClient | None = None, - steps: list[partial[Any]], -) -> ScenarioResult: - step_results: List[StepResult] = [] - - for step_func in steps: - logger.debug(f"Executing step: {step_func}") - params = step_func.keywords - logger.debug(f"Parameters in the step: {params}") - step_result: StepResult = step_func() - step_results.append(step_result) - logger.debug(f"Step result: {step_result}") - - for result in step_results: - if type(result) is not StepResult: - logger.error(f"Invalid step result type: {type(result)}") - # print(result) - final_status = Status.PASS if all(s.status == Status.PASS for s in step_results) else Status.FAIL - total_duration = sum(s.duration for s in step_results) - - return ScenarioResult( - name=scenario_id, - status=final_status, - duration_seconds=total_duration, - steps=step_results, - ) - - -def run_air_traffic_scenario_template( - scenario_id: str, - *, - fb_client: FlightBlenderClient | None = None, - air_traffic_client: AirTrafficClient | None = None, - steps: list[partial[Any]], -) -> ScenarioResult: - step_results: List[StepResult] = [] - - single_or_multiple_sensors = config.air_traffic_simulator_settings.single_or_multiple_sensors - - def _execute_step(step_func: partial[Any], current_observations: Any | None) -> tuple[StepResult, Any | None]: - name = _callable_name(step_func) - # Submit step consumes observations - if "submit_simulated_air_traffic" in name: - if current_observations: - return ( - step_func( - observations=current_observations, - single_or_multiple_sensors=single_or_multiple_sensors, - ), - current_observations, - ) - return ( - StepResult( - name="Submit Simulated Air Traffic (skipped)", - status=Status.PASS, - duration=0.0, - details="No observations to submit", - ), - current_observations, - ) - - res: StepResult = step_func() - if "fetch" in name or "generate" in name: - res, observations = _redact_fetch_details(res) - return res, observations - return res, current_observations - - step_results: List[StepResult] = [] - if air_traffic_client is not None and fb_client is not None: - observations: Any | None = None - for step in steps: - step_result, observations = _execute_step(step, observations) - step_results.append(step_result) - if step_result.status == Status.FAIL: - break - - final_status = Status.PASS if all(s.status == Status.PASS for s in step_results) else Status.FAIL - total_duration = sum(s.duration for s in step_results) - - return ScenarioResult( - name=scenario_id, - status=final_status, - duration_seconds=total_duration, - steps=step_results, - ) - - -def run_scenario_template( - scenario_id: str, - *, - fb_client: FlightBlenderClient | None = None, - air_traffic_client: AirTrafficClient | None = None, - opensky_client: OpenSkyClient | None = None, - steps: list[partial[Any]], - duration: int = DEFAULT_TELEMETRY_DURATION, -) -> ScenarioResult: - """Unified scenario runner supporting multiple client combinations. - - Supported flows: - 1. Declaration flow: fb_client only (generates flight declaration + telemetry) - 2. OpenSky flow: fb_client + opensky_client (fetches live data) - 3. Air traffic simulation: fb_client + air_traffic_client (generates simulated data) - 4. Declaration + OpenSky: fb_client + opensky_client (declaration flow) - 5. Declaration + Air traffic: fb_client + air_traffic_client (declaration flow) - """ - - # Log which clients are active - active_clients = [] - if fb_client: - active_clients.append(f"FB={fb_client.__class__.__name__}") - if opensky_client: - active_clients.append(f"OpenSky={opensky_client.__class__.__name__}") - if air_traffic_client: - active_clients.append(f"AirTraffic={air_traffic_client.__class__.__name__}") - - logger.debug(f"Active clients for '{scenario_id}': {', '.join(active_clients) if active_clients else 'None'}") - - # Determine scenario type based on client combination - is_declaration_flow = fb_client is not None and air_traffic_client is None and opensky_client is None - is_opensky_flow = opensky_client is not None - is_air_traffic_flow = air_traffic_client is not None - - # Only generate flight declaration/telemetry data for declaration flows - flight_declaration = None - telemetry_states = None - - if is_declaration_flow: - # Get scenario-specific config paths, falling back to global defaults - scenario_config = config.scenarios.get(scenario_id) - if scenario_config is None: - scenario_config = config.data_files - - telemetry_path = scenario_config.telemetry or config.data_files.telemetry - flight_declaration_path = scenario_config.flight_declaration or config.data_files.flight_declaration - - if not telemetry_path or not flight_declaration_path: - error_msg = ( - f"Declaration flow for '{scenario_id}' missing required config paths: " - f"telemetry={telemetry_path}, flight_declaration={flight_declaration_path}" - ) - logger.error(error_msg) - return ScenarioResult( - name=scenario_id, - status=Status.FAIL, - duration_seconds=0, - steps=[], - error_message="Missing configuration paths for data generation.", - ) - - try: - flight_declaration = _generate_flight_declaration(flight_declaration_path) - telemetry_states = _generate_telemetry(telemetry_path, duration=duration) - logger.info(f"Generated flight declaration and {len(telemetry_states)} telemetry states") - except Exception as e: - logger.error(f"Failed to generate data for scenario '{scenario_id}': {e}") - return ScenarioResult( - name=scenario_id, - status=Status.FAIL, - duration_seconds=0, - steps=[], - error_message=f"Data generation failed: {e}", - ) - - # Route to appropriate flow based on client combination - step_results: List[StepResult] = [] - - if is_declaration_flow: - # Standard declaration flow: upload -> steps -> teardown - logger.debug(f"Running declaration flow for '{scenario_id}'") - # Type narrowing: These are guaranteed non-None by is_declaration_flow logic - assert fb_client is not None, "fb_client must be set for declaration flow" - assert flight_declaration is not None, "flight_declaration must be generated for declaration flow" - assert telemetry_states is not None, "telemetry_states must be generated for declaration flow" - step_results = _run_declaration_flow( - fb_client=fb_client, - flight_declaration=flight_declaration, - telemetry_states=telemetry_states, - steps=steps, - ) - elif is_opensky_flow or is_air_traffic_flow: - # Data-fetching flows: fetch/generate -> submit - flow_type = "OpenSky" if is_opensky_flow else "Air Traffic" - logger.debug(f"Running {flow_type} submission flow for '{scenario_id}'") - step_results = _run_submit_airtraffic_flow(steps) - else: - # No valid clients provided - logger.error(f"Scenario '{scenario_id}' has no valid client configuration.") - return ScenarioResult( - name=scenario_id, - status=Status.FAIL, - duration_seconds=0, - steps=[], - error_message="No valid client configuration provided (need fb_client, opensky_client, or air_traffic_client).", - ) - - final_status = Status.PASS if all(s.status == Status.PASS for s in step_results) else Status.FAIL - total_duration = sum(s.duration for s in step_results) - - return ScenarioResult( - name=scenario_id, - status=final_status, - duration_seconds=total_duration, - steps=step_results, - flight_declaration_data=flight_declaration, - telemetry_data=telemetry_states, - ) - - -def get_telemetry_path(telemetry_filename: str) -> str: - """Helper to get the full path to a telemetry file.""" - parent_dir = Path(__file__).parent.resolve() - return str(parent_dir / f"../assets/rid_samples/{telemetry_filename}") - - -def get_flight_declaration_path(flight_declaration_filename: str) -> str: - """Helper to get the full path to a flight declaration file.""" - parent_dir = Path(__file__).parent.resolve() - return str(parent_dir / f"../assets/flight_declarations_samples/{flight_declaration_filename}") - - def get_geo_fence_path(geo_fence_filename: str) -> str: """Helper to get the full path to a geo-fence file.""" parent_dir = Path(__file__).parent.resolve() diff --git a/src/openutm_verification/scenarios/registry.py b/src/openutm_verification/scenarios/registry.py index 66d43983..e4944a20 100644 --- a/src/openutm_verification/scenarios/registry.py +++ b/src/openutm_verification/scenarios/registry.py @@ -13,10 +13,52 @@ def run_my_scenario(client, scenario_id): # ... """ +from functools import wraps +from typing import Any, Callable, ParamSpec, TypeVar + +from loguru import logger + +from openutm_verification.core.execution.scenario_runner import ScenarioContext +from openutm_verification.core.reporting.reporting_models import ( + ScenarioResult, + Status, +) + SCENARIO_REGISTRY = {} +T = TypeVar("T") +P = ParamSpec("P") + + +def _run_scenario_simple(scenario_id: str, func: Callable, args, kwargs) -> ScenarioResult: + """Runs a scenario without auto-setup.""" + try: + with ScenarioContext() as ctx: + result = func(*args, **kwargs) + + if isinstance(result, ScenarioResult): + return result + steps = ctx.steps -def register_scenario(scenario_id: str): + final_status = ( + Status.PASS if all(s.status == Status.PASS for s in steps) else Status.FAIL + ) + total_duration = sum(s.duration for s in steps) + return ScenarioResult( + name=scenario_id, + status=final_status, + duration_seconds=total_duration, + steps=steps, + ) + + except Exception as e: + logger.error(f"Scenario '{scenario_id}' failed: {e}") + raise e + + +def register_scenario( + scenario_id: str, +) -> Callable[[Callable[P, Any]], Callable[P, ScenarioResult]]: """ A decorator to register a test scenario function. @@ -25,10 +67,15 @@ def register_scenario(scenario_id: str): This ID is used in the configuration file. """ - def decorator(func): + def decorator(func: Callable[P, Any]) -> Callable[P, ScenarioResult]: if scenario_id in SCENARIO_REGISTRY: raise ValueError(f"Scenario with ID '{scenario_id}' is already registered.") - SCENARIO_REGISTRY[scenario_id] = func - return func + + @wraps(func) + def wrapper(*args: P.args, **kwargs: P.kwargs) -> ScenarioResult: + return _run_scenario_simple(scenario_id, func, args, kwargs) + + SCENARIO_REGISTRY[scenario_id] = wrapper + return wrapper return decorator diff --git a/src/openutm_verification/scenarios/test_add_flight_declaration.py b/src/openutm_verification/scenarios/test_add_flight_declaration.py index 0370f49a..cd453bc3 100644 --- a/src/openutm_verification/scenarios/test_add_flight_declaration.py +++ b/src/openutm_verification/scenarios/test_add_flight_declaration.py @@ -1,15 +1,11 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.models import OperationState -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("add_flight_declaration") -def test_add_flight_declaration(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_add_flight_declaration(fb_client: FlightBlenderClient, data_files: DataFiles) -> None: """Runs the add flight declaration scenario. This scenario replicates the behavior of the add_flight_declaration.py importer: @@ -21,19 +17,12 @@ def test_add_flight_declaration(fb_client: FlightBlenderClient, scenario_id: Sce Args: fb_client: The FlightBlenderClient instance for API interaction. - scenario_id: The unique name of the scenario being run. + data_files: The DataFiles instance containing file paths for telemetry, flight declaration, and geo-fence. Returns: A ScenarioResult object containing the results of the scenario execution. """ - steps = [ - partial(fb_client.update_operation_state, new_state=OperationState.ACTIVATED, duration_seconds=20), - partial(fb_client.submit_telemetry, duration_seconds=30), - partial(fb_client.update_operation_state, new_state=OperationState.ENDED), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) + with fb_client.flight_declaration(data_files): + fb_client.update_operation_state(new_state=OperationState.ACTIVATED, duration_seconds=20) + fb_client.submit_telemetry(duration_seconds=30) + fb_client.update_operation_state(new_state=OperationState.ENDED) diff --git a/src/openutm_verification/scenarios/test_airtraffic_data_openutm_sim.py b/src/openutm_verification/scenarios/test_airtraffic_data_openutm_sim.py index 38175df3..26ba0169 100644 --- a/src/openutm_verification/scenarios/test_airtraffic_data_openutm_sim.py +++ b/src/openutm_verification/scenarios/test_airtraffic_data_openutm_sim.py @@ -1,14 +1,9 @@ -from functools import partial - from openutm_verification.core.clients.air_traffic.air_traffic_client import ( AirTrafficClient, ) from openutm_verification.core.clients.flight_blender.flight_blender_client import ( FlightBlenderClient, ) -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult -from openutm_verification.scenarios.common import run_air_traffic_scenario_template from openutm_verification.scenarios.registry import register_scenario @@ -16,21 +11,11 @@ def test_openutm_sim_air_traffic_data( fb_client: FlightBlenderClient, air_traffic_client: AirTrafficClient, - scenario_id: ScenarioId, -) -> ScenarioResult: +) -> None: """Generate simulated air traffic data using OpenSky client and submit to Flight Blender using template. The OpenSky client is provided by the caller; this function focuses on orchestration only. """ - - steps = [ - partial(air_traffic_client.generate_simulated_air_traffic_data), - partial(fb_client.submit_simulated_air_traffic), - ] - - return run_air_traffic_scenario_template( - fb_client=fb_client, - air_traffic_client=air_traffic_client, - scenario_id=scenario_id, - steps=steps, - ) + step_result = air_traffic_client.generate_simulated_air_traffic_data() + observations = step_result.details + fb_client.submit_simulated_air_traffic(observations=observations) diff --git a/src/openutm_verification/scenarios/test_f1_flow.py b/src/openutm_verification/scenarios/test_f1_flow.py index 8564e17d..915537f1 100644 --- a/src/openutm_verification/scenarios/test_f1_flow.py +++ b/src/openutm_verification/scenarios/test_f1_flow.py @@ -1,15 +1,11 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.models import OperationState -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("F1_happy_path") -def test_f1_happy_path(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_f1_happy_path(fb_client: FlightBlenderClient, data_files: DataFiles): """Runs the F1 happy path scenario. This scenario simulates a complete, successful flight operation: @@ -24,14 +20,7 @@ def test_f1_happy_path(fb_client: FlightBlenderClient, scenario_id: ScenarioId) Returns: A ScenarioResult object containing the results of the scenario execution. """ - steps = [ - partial(fb_client.update_operation_state, new_state=OperationState.ACTIVATED), - partial(fb_client.submit_telemetry, duration_seconds=30), - partial(fb_client.update_operation_state, new_state=OperationState.ENDED), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) + with fb_client.flight_declaration(data_files): + fb_client.update_operation_state(new_state=OperationState.ACTIVATED) + fb_client.submit_telemetry(duration_seconds=30) + fb_client.update_operation_state(new_state=OperationState.ENDED) diff --git a/src/openutm_verification/scenarios/test_f2_flow.py b/src/openutm_verification/scenarios/test_f2_flow.py index c89028b5..51e855c9 100644 --- a/src/openutm_verification/scenarios/test_f2_flow.py +++ b/src/openutm_verification/scenarios/test_f2_flow.py @@ -1,15 +1,11 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.models import OperationState -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("F2_contingent_path") -def test_f2_contingent_path(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_f2_contingent_path(fb_client: FlightBlenderClient, data_files: DataFiles): """Runs the F2 contingent path scenario. This scenario simulates a flight operation that enters a contingent state: @@ -25,15 +21,8 @@ def test_f2_contingent_path(fb_client: FlightBlenderClient, scenario_id: Scenari Returns: A ScenarioResult object containing the results of the scenario execution. """ - steps = [ - partial(fb_client.update_operation_state, new_state=OperationState.ACTIVATED), - partial(fb_client.submit_telemetry, duration_seconds=10), - partial(fb_client.update_operation_state, new_state=OperationState.CONTINGENT, duration_seconds=7), - partial(fb_client.update_operation_state, new_state=OperationState.ENDED), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) + with fb_client.flight_declaration(data_files): + fb_client.update_operation_state(new_state=OperationState.ACTIVATED) + fb_client.submit_telemetry(duration_seconds=10) + fb_client.update_operation_state(new_state=OperationState.CONTINGENT, duration_seconds=7) + fb_client.update_operation_state(new_state=OperationState.ENDED) diff --git a/src/openutm_verification/scenarios/test_f3_flow.py b/src/openutm_verification/scenarios/test_f3_flow.py index 71583822..c126a983 100644 --- a/src/openutm_verification/scenarios/test_f3_flow.py +++ b/src/openutm_verification/scenarios/test_f3_flow.py @@ -1,15 +1,11 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.models import OperationState -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("F3_non_conforming_path") -def test_f3_non_conforming_path(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_f3_non_conforming_path(fb_client: FlightBlenderClient, data_files: DataFiles): """Runs the F3 non-conforming path scenario. This scenario simulates a flight that deviates from its declared flight plan, @@ -21,20 +17,13 @@ def test_f3_non_conforming_path(fb_client: FlightBlenderClient, scenario_id: Sce Args: fb_client: The FlightBlenderClient instance for API interaction. - scenario_id: The unique name of the scenario being run. + data_files: The DataFiles instance containing file paths for telemetry, flight declaration, and geo-fence. Returns: A ScenarioResult object containing the results of the scenario execution. """ - steps = [ - partial(fb_client.update_operation_state, new_state=OperationState.ACTIVATED), - partial(fb_client.submit_telemetry, duration_seconds=20), - partial(fb_client.check_operation_state, expected_state=OperationState.NONCONFORMING, duration_seconds=5), - partial(fb_client.update_operation_state, new_state=OperationState.ENDED), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) + with fb_client.flight_declaration(data_files): + fb_client.update_operation_state(new_state=OperationState.ACTIVATED) + fb_client.submit_telemetry(duration_seconds=20) + fb_client.check_operation_state(expected_state=OperationState.NONCONFORMING, duration_seconds=5) + fb_client.update_operation_state(new_state=OperationState.ENDED) diff --git a/src/openutm_verification/scenarios/test_f5_flow.py b/src/openutm_verification/scenarios/test_f5_flow.py index 034d59b6..fe910bce 100644 --- a/src/openutm_verification/scenarios/test_f5_flow.py +++ b/src/openutm_verification/scenarios/test_f5_flow.py @@ -1,25 +1,14 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult +from openutm_verification.core.execution.config_models import DataFiles from openutm_verification.models import OperationState -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("F5_non_conforming_path") -def test_f5_non_conforming_contingent_path(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: - steps = [ - partial(fb_client.update_operation_state, new_state=OperationState.ACTIVATED), - partial(fb_client.submit_telemetry, duration_seconds=20), - partial(fb_client.check_operation_state_connected, expected_state=OperationState.NONCONFORMING, duration_seconds=5), - partial(fb_client.update_operation_state, new_state=OperationState.CONTINGENT), - partial(fb_client.update_operation_state, new_state=OperationState.ENDED), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) +def test_f5_non_conforming_contingent_path(fb_client: FlightBlenderClient, data_files: DataFiles) -> None: + with fb_client.flight_declaration(data_files): + fb_client.update_operation_state(new_state=OperationState.ACTIVATED) + fb_client.submit_telemetry(duration_seconds=20) + fb_client.check_operation_state_connected(expected_state=OperationState.NONCONFORMING, duration_seconds=5) + fb_client.update_operation_state(new_state=OperationState.CONTINGENT) + fb_client.update_operation_state(new_state=OperationState.ENDED) diff --git a/src/openutm_verification/scenarios/test_geo_fence_upload.py b/src/openutm_verification/scenarios/test_geo_fence_upload.py index 48c071fa..999e1233 100644 --- a/src/openutm_verification/scenarios/test_geo_fence_upload.py +++ b/src/openutm_verification/scenarios/test_geo_fence_upload.py @@ -1,22 +1,10 @@ -from functools import partial - from openutm_verification.core.clients.flight_blender.flight_blender_client import FlightBlenderClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult -from openutm_verification.scenarios.common import get_geo_fence_path, run_scenario_template +from openutm_verification.scenarios.common import get_geo_fence_path from openutm_verification.scenarios.registry import register_scenario @register_scenario("geo_fence_upload") -def test_geo_fence_upload(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_geo_fence_upload(fb_client: FlightBlenderClient) -> None: """Upload a geo-fence (Area of Interest) and then delete it (teardown).""" - steps = [ - partial(fb_client.upload_geo_fence, filename=get_geo_fence_path("geo_fence.geojson")), - partial(fb_client.get_geo_fence), - ] - - return run_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, - ) + fb_client.upload_geo_fence(filename=get_geo_fence_path("geo_fence.geojson")) + fb_client.get_geo_fence() diff --git a/src/openutm_verification/scenarios/test_opensky_live_data.py b/src/openutm_verification/scenarios/test_opensky_live_data.py index 41d03690..26bc63b2 100644 --- a/src/openutm_verification/scenarios/test_opensky_live_data.py +++ b/src/openutm_verification/scenarios/test_opensky_live_data.py @@ -1,5 +1,4 @@ import time -from functools import partial from loguru import logger @@ -7,14 +6,13 @@ FlightBlenderClient, ) from openutm_verification.core.clients.opensky.opensky_client import OpenSkyClient -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult, Status -from openutm_verification.scenarios.common import run_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("opensky_live_data") -def test_opensky_live_data(fb_client: FlightBlenderClient, opensky_client: OpenSkyClient, scenario_id: ScenarioId) -> ScenarioResult: +def test_opensky_live_data( + fb_client: FlightBlenderClient, opensky_client: OpenSkyClient +) -> None: """Fetch live flight data from OpenSky and submit to Flight Blender using template. The OpenSky client is provided by the caller; this function focuses on orchestration only. @@ -24,36 +22,15 @@ def test_opensky_live_data(fb_client: FlightBlenderClient, opensky_client: OpenS iteration_count = 5 # total number of iterations wait_time = 3 # seconds to sleep between iterations - aggregated_steps = [] - overall_status = Status.PASS - total_duration = 0.0 - for i in range(iteration_count): logger.info(f"OpenSky iteration {i + 1}/{iteration_count}") - steps = [ - partial(opensky_client.fetch_data), - partial(fb_client.submit_air_traffic), - ] - - result = run_scenario_template( - fb_client=fb_client, - opensky_client=opensky_client, - scenario_id=f"{scenario_id} (iter {i + 1})", - steps=steps, - ) - - aggregated_steps.extend(result.steps) - total_duration += result.duration_seconds - if result.status == Status.FAIL: - overall_status = Status.FAIL + + step_result = opensky_client.fetch_data() + observations = step_result.details + + if observations: + fb_client.submit_air_traffic(observations=observations) if i < iteration_count - 1: logger.info(f"Waiting {wait_time} seconds before next iteration...") time.sleep(wait_time) - - return ScenarioResult( - name=scenario_id, - status=overall_status, - duration_seconds=total_duration, - steps=aggregated_steps, - ) diff --git a/src/openutm_verification/scenarios/test_sdsp_heartbeat.py b/src/openutm_verification/scenarios/test_sdsp_heartbeat.py index 8711547a..5e7d3b19 100644 --- a/src/openutm_verification/scenarios/test_sdsp_heartbeat.py +++ b/src/openutm_verification/scenarios/test_sdsp_heartbeat.py @@ -1,52 +1,38 @@ import uuid -from functools import partial from loguru import logger from openutm_verification.core.clients.flight_blender.flight_blender_client import ( FlightBlenderClient, ) -from openutm_verification.core.execution.config_models import ScenarioId -from openutm_verification.core.reporting.reporting_models import ScenarioResult from openutm_verification.models import SDSPSessionAction -from openutm_verification.scenarios.common import run_sdsp_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("sdsp_heartbeat") -def sdsp_heartbeat(fb_client: FlightBlenderClient, scenario_id: ScenarioId) -> ScenarioResult: +def sdsp_heartbeat(fb_client: FlightBlenderClient): """Runs the SDSP heartbeat scenario. This scenario """ session_id = str(uuid.uuid4()) logger.info(f"Starting SDSP heartbeat scenario with session ID: {session_id}") - steps = [ - partial( - fb_client.start_stop_sdsp_session, - action=SDSPSessionAction.START, - session_id=session_id, - ), - # Wait for some time to simulate heartbeat period - partial(fb_client.wait_x_seconds, wait_time_seconds=2), - partial( - fb_client.initialize_verify_sdsp_heartbeat, - session_id=session_id, - expected_heartbeat_interval_seconds=1, - expected_heartbeat_count=3, - ), - partial( - fb_client.wait_x_seconds, - wait_time_seconds=5, - ), - partial( - fb_client.start_stop_sdsp_session, - action=SDSPSessionAction.STOP, - session_id=session_id, - ), - ] - - return run_sdsp_scenario_template( - fb_client=fb_client, - scenario_id=scenario_id, - steps=steps, + + fb_client.start_stop_sdsp_session( + action=SDSPSessionAction.START, + session_id=session_id, + ) + # Wait for some time to simulate heartbeat period + fb_client.wait_x_seconds(wait_time_seconds=2) + + fb_client.initialize_verify_sdsp_heartbeat( + session_id=session_id, + expected_heartbeat_interval_seconds=1, + expected_heartbeat_count=3, + ) + + fb_client.wait_x_seconds(wait_time_seconds=5) + + fb_client.start_stop_sdsp_session( + action=SDSPSessionAction.STOP, + session_id=session_id, ) diff --git a/src/openutm_verification/scenarios/test_sdsp_track.py b/src/openutm_verification/scenarios/test_sdsp_track.py index 058ca0b0..e7d411a5 100644 --- a/src/openutm_verification/scenarios/test_sdsp_track.py +++ b/src/openutm_verification/scenarios/test_sdsp_track.py @@ -1,51 +1,38 @@ import uuid -from functools import partial from loguru import logger from openutm_verification.core.clients.flight_blender.flight_blender_client import ( FlightBlenderClient, ) -from openutm_verification.core.reporting.reporting_models import ScenarioResult from openutm_verification.models import SDSPSessionAction -from openutm_verification.scenarios.common import run_sdsp_scenario_template from openutm_verification.scenarios.registry import register_scenario @register_scenario("sdsp_track") -def sdsp_track(fb_client: FlightBlenderClient, scenario_name: str) -> ScenarioResult: +def sdsp_track(fb_client: FlightBlenderClient): """Runs the SDSP track scenario. This scenario """ session_id = str(uuid.uuid4()) logger.info(f"Starting SDSP track scenario with session ID: {session_id}") - steps = [ - partial( - fb_client.start_stop_sdsp_session, - action=SDSPSessionAction.START, - session_id=session_id, - ), - # Wait for some time to simulate track period - partial(fb_client.wait_x_seconds, wait_time_seconds=2), - partial( - fb_client.initialize_verify_sdsp_track, - session_id=session_id, - expected_heartbeat_interval_seconds=1, - expected_heartbeat_count=3, - ), - partial( - fb_client.wait_x_seconds, - wait_time_seconds=5, - ), - partial( - fb_client.start_stop_sdsp_session, - action=SDSPSessionAction.STOP, - session_id=session_id, - ), - ] - - return run_sdsp_scenario_template( - fb_client=fb_client, - scenario_name=scenario_name, - steps=steps, + + fb_client.start_stop_sdsp_session( + action=SDSPSessionAction.START, + session_id=session_id, + ) + # Wait for some time to simulate track period + fb_client.wait_x_seconds(wait_time_seconds=2) + + fb_client.initialize_verify_sdsp_track( + session_id=session_id, + expected_heartbeat_interval_seconds=1, + expected_heartbeat_count=3, + ) + + fb_client.wait_x_seconds(wait_time_seconds=5) + + fb_client.start_stop_sdsp_session( + action=SDSPSessionAction.STOP, + session_id=session_id, ) diff --git a/tests/docker-compose.fb.yml b/tests/docker-compose.fb.yml index ccdc91c1..092bcb61 100644 --- a/tests/docker-compose.fb.yml +++ b/tests/docker-compose.fb.yml @@ -38,6 +38,12 @@ services: depends_on: - redis-blender - db-blender + healthcheck: + test: ["CMD-SHELL", "python3 -c \"import urllib.request, json, sys; sys.exit(0 if json.load(urllib.request.urlopen('http://127.0.0.1:8000/ping'))['message'] == 'pong' else 1)\""] + interval: 10s + timeout: 10s + retries: 5 + start_period: 10s flight-blender-celery: