diff --git a/config.yml b/config.yml new file mode 100644 index 00000000000..2d454a66221 --- /dev/null +++ b/config.yml @@ -0,0 +1,13 @@ +config_path: config.yml +controls: + - max: 0.1 + min: 0.0 + name: my_control + perturbation_magnitude: 0.01 + variables: + - {initial_guess: 0.1, name: test} +model: + realizations: [0] + realizations_weights: [1.0] +objective_functions: + - {name: my_objective} diff --git a/tests/everest/entry_points/test_everexport.py b/everest_output/logs/everest-log-config-yml-2026-08-27T0844+0200.txt similarity index 100% rename from tests/everest/entry_points/test_everexport.py rename to everest_output/logs/everest-log-config-yml-2026-08-27T0844+0200.txt diff --git a/everest_output/logs/everest-log-config-yml-2026-08-27T0849+0200.txt b/everest_output/logs/everest-log-config-yml-2026-08-27T0849+0200.txt new file mode 100644 index 00000000000..60560579430 --- /dev/null +++ b/everest_output/logs/everest-log-config-yml-2026-08-27T0849+0200.txt @@ -0,0 +1,2 @@ +2026-08-27 08:49:56,002 - everest - MainThread - DEBUG - No definitions node found in configuration file +2026-08-27 08:49:56,104 - everest - MainThread - DEBUG - No definitions node found in configuration file diff --git a/everest_output/logs/everest-log-config-yml-2026-08-27T0856+0200.txt b/everest_output/logs/everest-log-config-yml-2026-08-27T0856+0200.txt new file mode 100644 index 00000000000..583075620c5 --- /dev/null +++ b/everest_output/logs/everest-log-config-yml-2026-08-27T0856+0200.txt @@ -0,0 +1,2 @@ +2026-08-27 08:56:16,747 - everest - MainThread - DEBUG - No definitions node found in configuration file +2026-08-27 08:56:16,851 - everest - MainThread - DEBUG - No definitions node found in configuration file diff --git a/src/ert/gui/experiments/experiment_client.py b/src/ert/gui/experiments/experiment_client.py deleted file mode 100644 index 59d38d08272..00000000000 --- a/src/ert/gui/experiments/experiment_client.py +++ /dev/null @@ -1,161 +0,0 @@ -from __future__ import annotations - -import logging -import queue -import ssl -import time -import traceback -from base64 import b64encode -from http import HTTPStatus -from pathlib import Path - -import requests -from pydantic import ValidationError -from requests import HTTPError -from websockets.exceptions import ConnectionClosedError -from websockets.sync.client import connect - -from _ert.threading import ErtThread -from ert.ensemble_evaluator import EvaluatorServerConfig -from ert.run_models import RunModelAPI -from ert.run_models.event import StatusEvents, status_event_from_json -from everest.strings import EverEndpoints - -logger = logging.getLogger(__name__) - - -class ExperimentClient: - def __init__( - self, - experiment_id: str, - url: str, - cert_file: str, - username: str, - password: str, - ssl_context: ssl.SSLContext, - ) -> None: - self._experiment_id = experiment_id - self._url = url - self._cert = cert_file - self._username = username - self._password = password - self._ssl_context = ssl_context - - self._is_alive = False - self._start_time: int | None = None - - def _http_get(self, endpoint: str) -> requests.Response: - return requests.get( - f"{self._url}/{endpoint}", - verify=self._cert, - auth=(self._username, self._password), - proxies={"http": None, "https": None}, # type: ignore - ) - - def _http_post(self, endpoint: str) -> requests.Response: - return requests.post( - f"{self._url}/{endpoint}", - verify=self._cert, - auth=(self._username, self._password), - proxies={"http": None, "https": None}, # type: ignore - ) - - @property - def config(self) -> dict[str, str]: - return self._http_get( - f"{EverEndpoints.CONFIG_PATH}/{self._experiment_id}" - ).json() - - @property - def credentials(self) -> str: - return b64encode(f"{self._username}:{self._password}".encode()).decode() - - def setup_event_queue_from_ws_endpoint( - self, - refresh_interval: float = 0.01, - open_timeout: float = 30, - websocket_recv_timeout: float = 1.0, - ) -> tuple[queue.SimpleQueue[StatusEvents], ErtThread]: - event_queue: queue.SimpleQueue[StatusEvents] = queue.SimpleQueue() - - def passthrough_ws_events() -> None: - try: # ruff: ignore[too-many-statements-in-try-clause] - with connect( - self._url.replace("https://", "wss://") - + f"/{EverEndpoints.EVENTS}/{self._experiment_id}", - ssl=self._ssl_context, - open_timeout=open_timeout, - additional_headers={"Authorization": f"Basic {self.credentials}"}, - ) as websocket: - while not self._is_alive: - try: - message = websocket.recv(timeout=websocket_recv_timeout) - except TimeoutError: - message = None - if message: - try: - event = status_event_from_json(message) - event_queue.put(event) - except ValidationError as e: - logger.error( - "Error when processing event %s", exc_info=e - ) - - time.sleep(refresh_interval) - except ConnectionClosedError: - logger.debug("Connection closed by server") - except Exception: - logger.debug(traceback.format_exc()) - - monitor_thread = ErtThread( - name="everest_gui_event_monitor", - target=passthrough_ws_events, - daemon=True, - ) - - return event_queue, monitor_thread - - def create_run_model_api(self) -> RunModelAPI: - def start_fn( - evaluator_server_config: EvaluatorServerConfig, - *, - rerun_failed_realizations: bool = False, - ) -> None: - pass - - return RunModelAPI( - experiment_name=Path(self.config["config_path"]).name, - supports_rerunning_failed_realizations=False, - start_simulations_thread=start_fn, - cancel=self.stop, - has_failed_realizations=lambda: False, - ) - - def stop(self) -> None: - try: - response = self._http_post(EverEndpoints.STOP) - except requests.exceptions.ConnectionError as e: - logger.error( - "Connection error when cancelling EVEREST " - f"experiment: {''.join(traceback.format_exception(e))}" - ) - print("Failed to cancel experiment") - return - except HTTPError as e: - logger.error( - "HTTP error when cancelling EVEREST " - f"experiment: {''.join(traceback.format_exception(e))}" - ) - print("Failed to cancel experiment") - return - if response.status_code == 200: - logger.info("Cancelled experiment from EVEREST") - print("Successfully cancelled experiment") - else: - logger.error( - f"Failed to cancel EVEREST experiment: " - f"POST @ {self._url}/{EverEndpoints.STOP}, " - f"server responded with status {response.status_code}: " - f"{HTTPStatus(response.status_code).phrase}" - ) - print("Failed to cancel experiment") diff --git a/src/ert/services/ert_client.py b/src/ert/services/ert_client.py index 8d738f9259c..f3dcce0a0ce 100644 --- a/src/ert/services/ert_client.py +++ b/src/ert/services/ert_client.py @@ -2,7 +2,13 @@ import io import json +import logging +import queue +import ssl import threading +import time +import traceback +from base64 import b64encode from collections import OrderedDict from collections.abc import Callable from copy import deepcopy @@ -15,6 +21,16 @@ import numpy as np import numpy.typing as npt import pandas as pd +from pydantic import ValidationError +from websockets.exceptions import ConnectionClosedError +from websockets.sync.client import connect + +from _ert.threading import ErtThread +from ert.run_models.event import ( + StatusEvents, + status_event_from_json, +) +from everest.strings import EverEndpoints from .shared_client import ErtClientConnectionInfo, Methods, SharedClient @@ -24,6 +40,8 @@ _PARQUET = {"accept": "application/x-parquet"} _EXPERIMENT_SERVER = "/experiment_server" +logger = logging.getLogger(__name__) + def _escape(value: str) -> str: """Keys may contain slashes, and the server decodes the path segment once.""" @@ -141,7 +159,7 @@ def version(self) -> str: return str(self._get("/version").json()) def experiments(self) -> list[dict[str, Any]]: - return list(self._get("/experiments").json()) + return self._get("/experiments").json() def ensemble(self, ensemble_id: str) -> dict[str, Any]: return dict(self._get(f"/ensembles/{ensemble_id}").json()) @@ -231,34 +249,106 @@ def experiment_server_is_running(self) -> bool: return response.status_code == httpx.codes.OK def experiment_ids(self) -> list[str]: - response = self._experiment_server_get("experiments") + response = self._experiment_server_get(EverEndpoints.EXPERIMENTS) return list(response.json()["experiment_ids"]) def experiment_status(self, experiment_id: str) -> dict[str, Any]: - return dict(self._experiment_server_get(f"status/{experiment_id}").json()) + return dict( + self._experiment_server_get( + f"{EverEndpoints.STATUS}/{experiment_id}" + ).json() + ) - def experiment_config_path(self, experiment_id: str) -> dict[str, Any]: - return dict(self._experiment_server_get(f"config_path/{experiment_id}").json()) + def experiment_config(self, experiment_id: str) -> dict[str, str]: + return self._experiment_server_get( + f"{EverEndpoints.CONFIG_PATH}/{experiment_id}" + ).json() def experiment_start_time(self, experiment_id: str) -> int: - return int(self._experiment_server_get(f"start_time/{experiment_id}").text) + return int( + self._experiment_server_get( + f"{EverEndpoints.START_TIME}/{experiment_id}" + ).text + ) + + def setup_event_queue_from_ws_endpoint( + self, + experiment_id: str, + refresh_interval: float = 0.01, + open_timeout: float = 30, + websocket_recv_timeout: float = 1.0, + ) -> tuple[queue.SimpleQueue[StatusEvents], ErtThread]: + """Return a queue of experiment events and the thread that fills it. + + The caller owns the thread and must start it. + """ + event_queue: queue.SimpleQueue[StatusEvents] = queue.SimpleQueue() + + url = ( + self.conn_info.base_url.replace("https://", "wss://") + + f"{_EXPERIMENT_SERVER}/{EverEndpoints.EVENTS}/{experiment_id}" + ) + username, password = self._auth + credentials = b64encode(f"{username}:{password}".encode()).decode() + + def passthrough_ws_events() -> None: + try: # ruff: ignore[too-many-statements-in-try-clause] + with connect( + url, + ssl=self._ssl_context, + open_timeout=open_timeout, + additional_headers={"Authorization": f"Basic {credentials}"}, + ) as websocket: + while True: + try: + message = websocket.recv(timeout=websocket_recv_timeout) + except TimeoutError: + message = None + if message: + try: + event_queue.put(status_event_from_json(message)) + except ValidationError as e: + logger.error( + "Error when processing event %s", exc_info=e + ) + + time.sleep(refresh_interval) + except ConnectionClosedError: + logger.debug("Connection closed by server") + except Exception: + logger.debug(traceback.format_exc()) + + monitor_thread = ErtThread( + name="ert_storage_api_event_monitor", + target=passthrough_ws_events, + daemon=True, + ) + + return event_queue, monitor_thread def start_experiment(self, config: dict[str, Any]) -> str: response = self._request( "POST", - f"{_EXPERIMENT_SERVER}/start_experiment", + f"{_EXPERIMENT_SERVER}/{EverEndpoints.START_EXPERIMENT}", auth=self._auth, json=config, ) return str(_checked(response).json()["experiment_id"]) - def stop_experiment_server(self) -> None: - _checked(self._request("POST", f"{_EXPERIMENT_SERVER}/stop", auth=self._auth)) + def stop_experiment_server(self) -> bool: + return ( + self._request( + "POST", + f"{_EXPERIMENT_SERVER}/{EverEndpoints.STOP}", + auth=self._auth, + ).status_code + == 200 + ) def runpath_exists(self, paths: list[str]) -> bool: response = self._request( "POST", - f"{_EXPERIMENT_SERVER}/runpath", + f"{_EXPERIMENT_SERVER}/{EverEndpoints.RUNPATH}", auth=self._auth, json={"paths": paths}, ) @@ -266,6 +356,13 @@ def runpath_exists(self, paths: list[str]) -> bool: # <-------------- Internals --------------> + @property + def _ssl_context(self) -> ssl.SSLContext | None: + cert = self._client.conn_info.cert + if not isinstance(cert, str): + return None + return ssl.create_default_context(cafile=cert) + @property def _auth(self) -> tuple[str, str]: """Experiment-server routes authenticate with HTTP Basic, not the token diff --git a/src/everest/bin/everest_script.py b/src/everest/bin/everest_script.py index 094f33cc69f..25e4a474c91 100755 --- a/src/everest/bin/everest_script.py +++ b/src/everest/bin/everest_script.py @@ -22,7 +22,6 @@ from ert.utils import makedirs_if_needed from everest.config import EverestConfig, ServerConfig from everest.detached import ( - start_experiment, start_server, wait_for_server, ) @@ -164,7 +163,7 @@ def _build_args_parser() -> argparse.ArgumentParser: async def run_everest(options: argparse.Namespace) -> None: try: - ErtClient.for_project( + ErtClient.get_client( Path(ServerConfig.get_session_dir(options.config.output_dir)), connect_timeout=1, ) @@ -228,7 +227,7 @@ async def directory_is_nonempty(path: Path) -> bool: print("Waiting for server ...") logger.debug("Waiting for response from everserver") wait_start_time: float = time.monotonic() - client = ErtClient.for_project( + client = ErtClient.get_client( Path(ServerConfig.get_session_dir(options.config.output_dir)) ) wait_for_server(client, timeout=600) @@ -238,11 +237,7 @@ async def directory_is_nonempty(path: Path) -> bool: f"waiting for {time.monotonic() - wait_start_time:g} seconds. " "Starting experiment" ) - - experiment_id = start_experiment( - server_context=ServerConfig.get_server_context_from_conn_info(client.conn_info), - config=options.config, - ) + experiment_id = client.start_experiment(options.config) # blocks until the run is finished if options.gui: diff --git a/src/everest/bin/kill_script.py b/src/everest/bin/kill_script.py index 6ec3ca94a7d..e474c40fcf0 100755 --- a/src/everest/bin/kill_script.py +++ b/src/everest/bin/kill_script.py @@ -15,7 +15,7 @@ from ert.services.ert_client import ErtClient from everest.bin.utils import setup_logging from everest.config import EverestConfig, ServerConfig -from everest.detached import stop_server, wait_for_server_to_stop +from everest.detached import wait_for_server_to_stop from everest.util import version_info logger = logging.getLogger(__name__) @@ -74,7 +74,7 @@ def _handle_keyboard_interrupt(signal: int, _: Any, *, after: bool = False) -> N def kill_everest(options: argparse.Namespace) -> None: try: - client = ErtClient.for_project( + client = ErtClient.get_client( Path(ServerConfig.get_session_dir(options.config.output_dir)), connect_timeout=1, ) @@ -86,7 +86,7 @@ def kill_everest(options: argparse.Namespace) -> None: print("Server is not running.") return - stopping = stop_server(server_context) + stopping = client.stop_experiment_server() if threading.current_thread() is threading.main_thread(): signal.signal(signal.SIGINT, partial(_handle_keyboard_interrupt, after=True)) diff --git a/src/everest/bin/monitor_script.py b/src/everest/bin/monitor_script.py index 009e904920c..c4eacef3e3b 100755 --- a/src/everest/bin/monitor_script.py +++ b/src/everest/bin/monitor_script.py @@ -10,7 +10,6 @@ from ert.services.ert_client import ErtClient from ert.storage import ErtStorageException, ExperimentState from everest.config import EverestConfig, ServerConfig -from everest.detached.client import get_experiments from .utils import ( ArgParseFormatter, @@ -78,13 +77,13 @@ def _build_args_parser() -> argparse.ArgumentParser: def monitor_everest(options: argparse.Namespace) -> None: config: EverestConfig = options.config try: # ruff: ignore[too-many-statements-in-try-clause] - client = ErtClient.for_project( + client = ErtClient.get_client( Path(ServerConfig.get_session_dir(config.output_dir)), connect_timeout=1 ) server_context = ServerConfig.get_server_context_from_conn_info( client.conn_info ) - experiment_id = get_experiments(server_context)[-1] + experiment_id = client.experiment_ids()[-1] run_detached_monitor(server_context=server_context, experiment_id=experiment_id) try: diff --git a/src/everest/bin/utils.py b/src/everest/bin/utils.py index a11adeafc98..2ed2718313d 100644 --- a/src/everest/bin/utils.py +++ b/src/everest/bin/utils.py @@ -34,10 +34,9 @@ from everest.config.server_config import ServerConfig from everest.detached import ( server_is_running, - start_monitor, - stop_server, wait_for_server_to_stop, ) +from everest.detached.client import start_monitor from everest.strings import EVEREST, OPT_PROGRESS_ID, SIM_PROGRESS_ID from everest.util import format_list @@ -106,14 +105,14 @@ def handle_keyboard_interrupt(signum: int, _: Any, options: argparse.Namespace) "The optimization will be stopped and the program will exit..." ) try: - client = ErtClient.for_project( + client = ErtClient.get_client( Path(ServerConfig.get_session_dir(options.config.output_dir)) ) server_context = ServerConfig.get_server_context_from_conn_info( client.conn_info ) if server_is_running(*server_context): - stop_server(server_context) + client.stop_experiment_server() wait_for_server_to_stop(server_context, timeout=10) except TimeoutError: diff --git a/src/everest/detached/__init__.py b/src/everest/detached/__init__.py index 1224682a9d7..0a216bc8ff6 100644 --- a/src/everest/detached/__init__.py +++ b/src/everest/detached/__init__.py @@ -2,24 +2,18 @@ from .client import ( PROXY, - get_experiments, server_is_running, - start_experiment, start_monitor, start_server, - stop_server, wait_for_server, wait_for_server_to_stop, ) __all__ = [ "PROXY", - "get_experiments", "server_is_running", - "start_experiment", "start_monitor", "start_server", - "stop_server", "wait_for_server", "wait_for_server_to_stop", ] diff --git a/src/everest/detached/client.py b/src/everest/detached/client.py index b62edc2ceae..7a45259d094 100644 --- a/src/everest/detached/client.py +++ b/src/everest/detached/client.py @@ -1,6 +1,5 @@ import asyncio import logging -import re import ssl import time import traceback @@ -93,55 +92,6 @@ def stop_server( return False -def get_experiments( - server_context: tuple[str, str, tuple[str, str]], - retries: int = 5, -) -> list[str]: - url, cert, auth = server_context - for retry in range(retries): - try: - response = requests.get( - f"{url}/{EverEndpoints.EXPERIMENTS}", - verify=cert, - auth=auth, - proxies=PROXY, - ) - response.raise_for_status() - return response.json()["experiment_ids"] - except Exception: - logger.debug(traceback.format_exc()) - time.sleep(retry) - raise RuntimeError("Failed to get experiment_ids") - - -def start_experiment( - server_context: tuple[str, str, tuple[str, str]], - config: EverestConfig, - retries: int = 5, -) -> str: - url, cert, auth = server_context - for retry in range(retries): - try: - start_endpoint = f"{url}/{EverEndpoints.START_EXPERIMENT}" - response = requests.post( - start_endpoint, - verify=cert, - auth=auth, - proxies=PROXY, # type: ignore - json=config.to_dict(), - ) - response.raise_for_status() - return response.json()["experiment_id"] - except Exception: - logger.debug(traceback.format_exc()) - time.sleep(retry) - raise RuntimeError("Failed to start experiment") - - -def extract_errors_from_file(path: str) -> list[str]: - return re.findall(r"(Error \w+.*)", Path(path).read_text(encoding="utf-8")) - - def wait_for_server(api: ErtClient, timeout: float) -> None: """ Waits until the everest server has started. Polls diff --git a/src/everest/detached/everserver.py b/src/everest/detached/everserver.py index 740219a6b65..75490efc54f 100644 --- a/src/everest/detached/everserver.py +++ b/src/everest/detached/everserver.py @@ -19,7 +19,6 @@ from ert.trace import tracer from ert.utils import makedirs_if_needed from everest.config import ServerConfig -from everest.detached import get_experiments from everest.strings import ( DEFAULT_LOGGING_FORMAT, OPTIMIZATION_LOG_DIR, @@ -151,9 +150,7 @@ def main() -> None: client = ErtClient.for_project(Path(server_path)) done = False while not done: - experiment_ids = get_experiments( - ServerConfig.get_server_context_from_conn_info(client.conn_info) - ) + experiment_ids = client.experiment_ids() active = [ ExperimentStatus( **client.experiment_status(experiment_id) diff --git a/src/everest/gui/main_window.py b/src/everest/gui/main_window.py index 55fb8dc322b..2999a27f445 100644 --- a/src/everest/gui/main_window.py +++ b/src/everest/gui/main_window.py @@ -1,6 +1,5 @@ from __future__ import annotations -import ssl from pathlib import Path from PyQt6.QtCore import pyqtSignal as Signal @@ -10,13 +9,14 @@ QMainWindow, ) +from ert.ensemble_evaluator.config import EvaluatorServerConfig from ert.gui.ertnotifier import ErtNotifier from ert.gui.experiments import RunDialog -from ert.gui.experiments.experiment_client import ExperimentClient from ert.plugins import ErtPluginManager +from ert.run_models.run_model import RunModelAPI from ert.services.ert_client import ErtClient from everest.config import ServerConfig -from everest.detached import get_experiments, wait_for_server +from everest.detached import wait_for_server class EverestMainWindow(QMainWindow): @@ -43,40 +43,40 @@ def __init__( self.setCentralWidget(self.central_widget) def run(self) -> None: - client = ErtClient.for_project( + client = ErtClient.get_client( Path(ServerConfig.get_session_dir(self.output_dir)) ) wait_for_server(client, 60) - server_context = ServerConfig.get_server_context_from_conn_info( - client.conn_info + experiment_id = client.experiment_ids()[-1] + config = client.experiment_config(experiment_id) + + config_filename = Path(config["config_path"]).name + self.setWindowTitle(f"EVEREST - {config_filename}") + + def start_fn( + evaluator_server_config: EvaluatorServerConfig, + *, + rerun_failed_realizations: bool = False, + ) -> None: + pass + + run_model_api = RunModelAPI( + experiment_name=config_filename, + supports_rerunning_failed_realizations=False, + start_simulations_thread=start_fn, + cancel=client.stop_experiment_server, # type: ignore + has_failed_realizations=lambda: False, ) - url, cert, auth = server_context - - ssl_context = ssl.create_default_context() - ssl_context.load_verify_locations(cafile=cert) - username, password = auth - - exp_client = ExperimentClient( - experiment_id=get_experiments(server_context)[-1], - url=url, - cert_file=cert, - username=username, - password=password, - ssl_context=ssl_context, - ) - - config = exp_client.config - title = Path(config["config_path"]).name - self.setWindowTitle(f"EVEREST - {title}") - - run_model_api = exp_client.create_run_model_api() - event_queue, event_monitor_thread = exp_client.setup_event_queue_from_ws_endpoint( - refresh_interval=0.02, open_timeout=40, websocket_recv_timeout=1.0 + event_queue, event_monitor_thread = client.setup_event_queue_from_ws_endpoint( + experiment_id=experiment_id, + refresh_interval=0.02, + open_timeout=40, + websocket_recv_timeout=1.0, ) run_dialog = RunDialog( - title=title, + title=config_filename, run_model_api=run_model_api, event_queue=event_queue, notifier=ErtNotifier(), diff --git a/tests/everest/entry_points/config.yml b/tests/everest/entry_points/config.yml new file mode 100644 index 00000000000..2d454a66221 --- /dev/null +++ b/tests/everest/entry_points/config.yml @@ -0,0 +1,13 @@ +config_path: config.yml +controls: + - max: 0.1 + min: 0.0 + name: my_control + perturbation_magnitude: 0.01 + variables: + - {initial_guess: 0.1, name: test} +model: + realizations: [0] + realizations_weights: [1.0] +objective_functions: + - {name: my_objective} diff --git a/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0852+0200.txt b/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0852+0200.txt new file mode 100644 index 00000000000..c75eaea0db6 --- /dev/null +++ b/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0852+0200.txt @@ -0,0 +1,2 @@ +2026-08-27 08:52:29,168 - everest - MainThread - DEBUG - No definitions node found in configuration file +2026-08-27 08:52:29,271 - everest - MainThread - DEBUG - No definitions node found in configuration file diff --git a/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0854+0200.txt b/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0854+0200.txt new file mode 100644 index 00000000000..ef47a98ebd4 --- /dev/null +++ b/tests/everest/entry_points/everest_output/logs/everest-log-config-yml-2026-08-27T0854+0200.txt @@ -0,0 +1,2 @@ +2026-08-27 08:54:37,578 - everest - MainThread - DEBUG - No definitions node found in configuration file +2026-08-27 08:54:37,690 - everest - MainThread - DEBUG - No definitions node found in configuration file diff --git a/tests/everest/entry_points/test_everest_entry.py b/tests/everest/entry_points/test_everest_entry.py index ecd46f52ff2..386cef9c00e 100644 --- a/tests/everest/entry_points/test_everest_entry.py +++ b/tests/everest/entry_points/test_everest_entry.py @@ -1,7 +1,7 @@ import logging import tempfile from pathlib import Path -from unittest.mock import MagicMock, patch +from unittest.mock import DEFAULT, MagicMock, patch import pytest @@ -27,11 +27,9 @@ def raise_system_error(*args, **kwargs): @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( "everest.bin.everest_script.ErtClient", - **{"for_project.side_effect": [TimeoutError(), MagicMock()]}, + **{"for_project.side_effect": [TimeoutError(), DEFAULT]}, ) -@patch("everest.bin.everest_script.start_experiment") def test_everest_entry_debug( - start_experiment_mock, everest_script_api_mock, get_server_context_from_conn_info_mock, start_server_mock, @@ -59,9 +57,9 @@ def test_everest_entry_debug( start_server_mock.assert_called_once() wait_for_server_mock.assert_called_once() start_monitor_mock.assert_called_once() - start_experiment_mock.assert_called_once() + everest_script_api_mock.for_project.return_value.start_experiment.assert_called_once() assert everest_script_api_mock.for_project.call_count == 2 - assert get_server_context_from_conn_info_mock.call_count == 2 + get_server_context_from_conn_info_mock.assert_called_once() # the config file itself is dumped at DEBUG level assert '"controls"' in logstream @@ -76,11 +74,9 @@ def test_everest_entry_debug( @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( "everest.bin.everest_script.ErtClient", - **{"for_project.side_effect": [TimeoutError(), MagicMock()]}, + **{"for_project.side_effect": [TimeoutError(), DEFAULT]}, ) -@patch("everest.bin.everest_script.start_experiment") def test_everest_entry( - start_experiment_mock, everest_script_api_mock, get_server_context_from_conn_info_mock, start_server_mock, @@ -97,24 +93,23 @@ def test_everest_entry( start_server_mock.assert_called_once() wait_for_server_mock.assert_called_once() start_monitor_mock.assert_called_once() - start_experiment_mock.assert_called_once() + everest_script_api_mock.for_project.return_value.start_experiment.assert_called_once() assert everest_script_api_mock.for_project.call_count == 2 - assert get_server_context_from_conn_info_mock.call_count == 2 + get_server_context_from_conn_info_mock.assert_called_once() @patch("everest.bin.everest_script.run_detached_monitor") @patch("everest.bin.everest_script.wait_for_server") @patch("everest.bin.everest_script.start_server") -@patch("everest.bin.everest_script.start_experiment") @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( "everest.bin.everest_script.ErtClient", **{ "for_project.side_effect": [ TimeoutError(), - MagicMock(), + DEFAULT, TimeoutError(), - MagicMock(), + DEFAULT, ] }, ) @@ -126,7 +121,6 @@ def test_everest_entry_detached_already_run( kill_script_api_mock, everest_script_api_mock, get_server_context_from_conn_info_mock, - start_experiment_mock, start_server_mock, wait_for_server_mock, start_monitor_mock, @@ -139,6 +133,9 @@ def test_everest_entry_detached_already_run( Path("config.yml").touch() config = everest_config_with_defaults(config_path="./config.yml") config.write_to_file("config.yml") + start_experiment_mock = ( + everest_script_api_mock.for_project.return_value.start_experiment + ) # start a new run everest_entry(["config.yml"]) @@ -200,13 +197,14 @@ def test_everest_entry_detached_already_run_monitor( @patch("everest.bin.everest_script.run_detached_monitor") @patch("everest.bin.everest_script.wait_for_server") @patch("everest.bin.everest_script.start_server") -@patch("everest.bin.kill_script.stop_server", return_value=True) @patch("everest.bin.kill_script.wait_for_server_to_stop") -@patch("everest.bin.kill_script.ErtClient") +@patch( + "everest.bin.kill_script.ErtClient", + **{"for_project.return_value.stop_experiment_server.return_value": True}, +) def test_everest_entry_detached_running( kill_api_mock, wait_for_server_to_stop_mock, - stop_server_mock, start_server_mock, wait_for_server_mock, start_monitor_mock, @@ -219,6 +217,7 @@ def test_everest_entry_detached_running( Path("config.yml").touch() config = everest_config_with_defaults(config_path="./config.yml") config.write_to_file("config.yml") + stop_server_mock = kill_api_mock.for_project.return_value.stop_experiment_server # can't start a new run if one is already running with capture_streams() as (out, _): @@ -252,12 +251,11 @@ def test_everest_entry_detached_running( @patch("everest.bin.monitor_script.run_detached_monitor") @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( - "everest.bin.monitor_script.get_experiments", return_value=["test-experiment-id"] + "everest.bin.monitor_script.ErtClient", + **{"for_project.return_value.experiment_ids.return_value": ["test-experiment-id"]}, ) -@patch("everest.bin.monitor_script.ErtClient") def test_everest_entry_detached_running_monitor( monitor_script_api_mock, - get_experiments_mock, get_server_context_from_conn_info_mock, start_monitor_mock, change_to_tmpdir, @@ -274,7 +272,7 @@ def test_everest_entry_detached_running_monitor( start_monitor_mock.assert_called_once() monitor_script_api_mock.for_project.assert_called_once() get_server_context_from_conn_info_mock.assert_called_once() - get_experiments_mock.assert_called_once() + monitor_script_api_mock.for_project.return_value.experiment_ids.assert_called_once() @patch("everest.bin.monitor_script.run_detached_monitor") @@ -318,16 +316,14 @@ def mock_ssl(monkeypatch): ) @patch("everest.bin.everest_script.wait_for_server") @patch("everest.bin.everest_script.start_server") -@patch("everest.bin.everest_script.start_experiment") @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( "everest.bin.everest_script.ErtClient", - **{"for_project.side_effect": [TimeoutError(), MagicMock()]}, + **{"for_project.side_effect": [TimeoutError(), DEFAULT]}, ) def test_exception_raised_when_server_run_fails( everest_script_api_mock, get_server_context_from_conn_info_mock, - start_experiment_mock, start_server_mock, wait_for_server_mock, start_monitor_mock, @@ -347,12 +343,11 @@ def test_exception_raised_when_server_run_fails( ) @patch("everest.config.ServerConfig.get_server_context_from_conn_info") @patch( - "everest.bin.monitor_script.get_experiments", return_value=["test-experiment-id"] + "everest.bin.monitor_script.ErtClient", + **{"for_project.return_value.experiment_ids.return_value": ["test-experiment-id"]}, ) -@patch("everest.bin.monitor_script.ErtClient") def test_exception_raised_when_server_run_fails_monitor( monitor_script_api_mock, - get_experiments_mock, get_server_context_from_conn_info_mock, start_monitor_mock, change_to_tmpdir, @@ -418,7 +413,7 @@ def test_that_run_everest_prints_where_it_runs( with ( patch( "everest.bin.everest_script.ErtClient", - **{"for_project.side_effect": [TimeoutError(), MagicMock()]}, + **{"for_project.side_effect": [TimeoutError(), DEFAULT]}, ), patch( "everest.config.ServerConfig.get_server_context_from_conn_info", @@ -426,7 +421,6 @@ def test_that_run_everest_prints_where_it_runs( ), patch("everest.bin.everest_script.start_server"), patch("everest.bin.everest_script.wait_for_server"), - patch("everest.bin.everest_script.start_experiment"), ): everest_entry(["config.yml"]) diff --git a/tests/everest/test_detached.py b/tests/everest/test_detached.py index 9b12fb4c8e2..09de6ba1754 100644 --- a/tests/everest/test_detached.py +++ b/tests/everest/test_detached.py @@ -35,7 +35,6 @@ PROXY, server_is_running, start_server, - stop_server, wait_for_server, wait_for_server_to_stop, ) @@ -85,7 +84,7 @@ async def test_https_requests(change_to_tmpdir): *ServerConfig.get_server_context_from_conn_info(client.conn_info) ) server_context = ServerConfig.get_server_context_from_conn_info(client.conn_info) - if stop_server(server_context): + if client.stop_experiment_server(): wait_for_server_to_stop(server_context, 240) assert not server_is_running(*server_context) diff --git a/tests/everest/test_everest_client.py b/tests/everest/test_everest_client.py deleted file mode 100644 index 192f75610f0..00000000000 --- a/tests/everest/test_everest_client.py +++ /dev/null @@ -1,262 +0,0 @@ -import logging -import ssl -import sys -import threading -import warnings -from pathlib import Path - -import pytest -import requests -import uvicorn -import yaml -from fastapi import FastAPI -from starlette.responses import Response - -from ert.gui.experiments.experiment_client import ExperimentClient -from ert.run_models.event import EverestBatchResultEvent, EverestStatusEvent -from ert.services import ErtClient -from ert.shared import find_available_socket -from everest.bin.everest_script import everest_entry -from everest.config import EverestConfig, ServerConfig -from everest.detached import get_experiments, server_is_running -from everest.strings import EverEndpoints -from tests.ert.utils import wait_until - - -@pytest.fixture -def client_server_mock() -> tuple[FastAPI, threading.Thread, ExperimentClient]: - server_app = FastAPI() - host = "127.0.0.1" - port = find_available_socket(host, range(5000, 5800)).getsockname()[1] - server_url = f"http://{host}:{port}" - - @server_app.get("alive") - def alive(): - return Response("Hello", status_code=200) - - server = uvicorn.Server( - uvicorn.Config(server_app, host=host, port=port, log_level="info") - ) - - server_thread = threading.Thread( - target=server.run, - daemon=True, - ) - - everest_client = ExperimentClient( - experiment_id="test_experiment_id", - url=server_url, - cert_file="N/A", - username="", - password="", - ssl_context=ssl.create_default_context(), - ) - - def wait_until_alive(timeout=60, sleep_between_retries=1) -> None: - def ping_server() -> bool: - try: - requests.get( - f"{server_url}/alive", - verify="N/A", - auth=("", ""), - proxies={"http": None, "https": None}, # type: ignore - ) - except requests.exceptions.ConnectionError: - return False - else: - return True - - # These warnings emitted by uvicorn, which is still using legacy - # websockets. This is a known issue, and does not cause problems in the - # main code. (see: https://github.com/encode/uvicorn/discussions/2476) - # Hence we ignore them in the tests. Potentially, this may be - # removed when this is resolved within uvicorn. - with warnings.catch_warnings(): - warnings.filterwarnings("ignore", message="websockets.legacy is deprecated") - warnings.filterwarnings( - "ignore", - message="websockets.server.WebSocketServerProtocol is deprecated", - ) - wait_until(ping_server, timeout=timeout, interval=sleep_between_retries) - - yield server_app, server_thread, everest_client, wait_until_alive - - if server_thread.is_alive(): - server.should_exit = True - server_thread.join() - - -@pytest.mark.slow -@pytest.mark.flaky(rerun=2) -def test_that_stop_invokes_correct_endpoint( - caplog, client_server_mock: tuple[FastAPI, threading.Thread, ExperimentClient] -): - server_app, server_thread, client, wait_until_alive = client_server_mock - - @server_app.post(f"/{EverEndpoints.STOP}") - def stop(): - return Response("STOP..", 200) - - server_thread.start() - wait_until_alive() - - with caplog.at_level(logging.INFO): - client.stop() - - assert "Cancelled experiment from EVEREST" in caplog.messages - server_thread.should_exit = True - - -@pytest.mark.slow -def test_that_stop_errors_on_non_ok_httpcode( - caplog, client_server_mock: tuple[FastAPI, threading.Thread, ExperimentClient] -): - server_app, server_thread, client, wait_until_alive = client_server_mock - - @server_app.post(f"/{EverEndpoints.STOP}") - def stop(): - return Response("STOP..", 505) - - server_thread.start() - wait_until_alive() - - with caplog.at_level(logging.ERROR): - client.stop() - - assert any( - "Failed to cancel EVEREST experiment" in m - and "server responded with status 505" in m - for m in caplog.messages - ) - - -def test_that_stop_errors_on_server_down( - caplog, client_server_mock: tuple[FastAPI, threading.Thread, ExperimentClient] -): - _, _, client, _ = client_server_mock - - with caplog.at_level(logging.ERROR): - client.stop() - - assert any( - "Connection error when cancelling EVEREST experiment" in m - for m in caplog.messages - ) - - -@pytest.mark.slow -def test_that_stop_errors_on_server_up_but_endpoint_down( - caplog, client_server_mock: tuple[FastAPI, threading.Thread, ExperimentClient] -): - _, server_thread, client, wait_until_alive = client_server_mock - - server_thread.start() - wait_until_alive() - - with caplog.at_level(logging.ERROR): - client.stop() - - assert any( - "server responded with status 404: Not Found" in m for m in caplog.messages - ) - - -@pytest.mark.skip_mac_ci -@pytest.mark.slow -@pytest.mark.xdist_group("math_func/config_minimal.yml") -@pytest.mark.flaky(rerun=3) -@pytest.mark.skipif( - sys.version_info[0:3] == (3, 13, 6), reason="Fails on Python 3.13.6" -) -def test_that_multiple_everest_clients_can_connect_to_server( - cached_example, change_to_tmpdir -): - # We use a cached run for the reference list of received events - path, config_file, _, server_events_list = cached_example( - "math_func/config_minimal.yml" - ) - - config_path = Path(path) / config_file - config_content = yaml.safe_load(config_path.read_text(encoding="utf-8")) - config_content["simulator"] = {"queue_system": {"name": "local", "max_running": 2}} - config_path.write_text( - yaml.dump(config_content, default_flow_style=False), encoding="utf-8" - ) - - ever_config = EverestConfig.load_file(config_path) - - # Run the case through everserver - everest_main_thread = threading.Thread( - target=everest_entry, args=[[str(config_path)]] - ) - - everest_main_thread.start() - api = ErtClient.for_project( - Path(ServerConfig.get_session_dir(ever_config.output_dir)) - ) - - def everserver_is_running(): - return server_is_running( - *ServerConfig.get_server_context_from_conn_info(api.conn_info) - ) - - wait_until(everserver_is_running, interval=1, timeout=300) - - server_context = ServerConfig.get_server_context_from_conn_info(api.conn_info) - url, cert, auth = server_context - - ssl_context = ssl.create_default_context() - ssl_context.load_verify_locations(cafile=cert) - username, password = auth - - client_event_queues = [] - monitor_threads = [] - for _ in range(5): - client = ExperimentClient( - experiment_id=get_experiments(server_context)[-1], - url=url, - cert_file=cert, - username=username, - password=password, - ssl_context=ssl_context, - ) - - # Connect to the websockets endpoint - client_event_queue, monitor_thread = client.setup_event_queue_from_ws_endpoint() - client_event_queues.append(client_event_queue) - monitor_threads.append(monitor_thread) - monitor_thread.start() - - # Wait until the server has finished running the simulation - everest_main_thread.join() - for _thread in monitor_threads: - if _thread.is_alive(): - _thread.join(timeout=5) - - # Expect all the clients to hold the same events - client_event_lists = [] - for event_queue in client_event_queues: - event_list = [] - while not event_queue.empty(): - event_list.append(event_queue.get()) - - client_event_lists.append(event_list) - - first = client_event_lists[0] - assert all(first == other for other in client_event_lists[1:]) - - everest_event_types = (EverestStatusEvent, EverestBatchResultEvent) - - first_everevents = [ - e.event_type for e in first if isinstance(e, everest_event_types) - ] - assert len(first_everevents) > 0 - - server_everevents = [ - e.event_type for e in server_events_list if isinstance(e, everest_event_types) - ] - assert len(server_everevents) > 0 - - # Compare only everest events, as the events from the forward model - # are (at time of writing) not deterministic enough to expect equality - assert first_everevents == server_everevents diff --git a/tests/everest/test_everest_output.py b/tests/everest/test_everest_output.py index c4b13e9adf7..07e652c70af 100644 --- a/tests/everest/test_everest_output.py +++ b/tests/everest/test_everest_output.py @@ -34,9 +34,7 @@ def test_that_one_experiment_creates_one_ensemble_per_batch(cached_example): @patch("everest.bin.everest_script.run_detached_monitor") @patch("everest.bin.everest_script.wait_for_server") @patch("everest.bin.everest_script.start_server") -@patch("everest.bin.everest_script.start_experiment") def test_save_running_config( - mock_start_experiment, mock_start_server, mock_wait_for_server, mock_run_detached_monitor, diff --git a/tests/everest/test_everserver.py b/tests/everest/test_everserver.py index 92645eaafc7..fe17491d33c 100644 --- a/tests/everest/test_everserver.py +++ b/tests/everest/test_everserver.py @@ -28,7 +28,6 @@ from everest.config import EverestConfig, ServerConfig from everest.detached import ( everserver, - start_experiment, start_server, wait_for_server, ) @@ -72,10 +71,7 @@ async def server_running(): driver = await start_server(config, logging.DEBUG) api = ErtClient.for_project(Path(ServerConfig.get_session_dir(config.output_dir))) wait_for_server(api, 120) - start_experiment( - server_context=ServerConfig.get_server_context_from_conn_info(api.conn_info), - config=config, - ) + api.start_experiment(config) await server_running()