diff --git a/.github/workflows/test_monitoring.yaml b/.github/workflows/test_monitoring.yaml index 058534c..a9b5127 100644 --- a/.github/workflows/test_monitoring.yaml +++ b/.github/workflows/test_monitoring.yaml @@ -17,12 +17,12 @@ jobs: - name: Checkout recipes uses: actions/checkout@v4 - - name: Checkout echodataflow monitoring PR + - name: Checkout echodataflow uses: actions/checkout@v4 with: repository: echostack-org/echodataflow - ref: refs/pull/174/head path: echodataflow + fetch-depth: 0 - name: Set up Python uses: actions/setup-python@v5 diff --git a/recipes_nrt/deploy/deploy_spasso_iwcps_2026.yaml b/recipes_nrt/deploy/deploy_spasso_iwcps_2026.yaml new file mode 100644 index 0000000..309c5a4 --- /dev/null +++ b/recipes_nrt/deploy/deploy_spasso_iwcps_2026.yaml @@ -0,0 +1,6 @@ +flow_start_time: + +flows: + fetch_spasso: + deployment_name: fetch-spasso-iwcps-2026 + interval: 15 \ No newline at end of file diff --git a/recipes_nrt/params/params_spasso_iwcps_2026.yaml b/recipes_nrt/params/params_spasso_iwcps_2026.yaml new file mode 100644 index 0000000..894e936 --- /dev/null +++ b/recipes_nrt/params/params_spasso_iwcps_2026.yaml @@ -0,0 +1,132 @@ +flows: + fetch_spasso: + host: "download host" + port: 2221 + remote_path: "remote path" + + path_main: "path/to/spasso" + + product_patterns: + - "Copernicus_PHY" + - "Copernicus_SST_L4" + - "Copernicus_SSS_L4" + - "Copernicus_CHL_L4" + - "FTLE_Copernicus_PHY" + - "OW_Copernicus_PHY" + - "KE_Copernicus_PHY" + - "LLADV_Copernicus_PHY" + + username_secret: "spasso-username" + password_secret: "spasso-password" + + task_retries: 3 + task_retry_delay_seconds: 30 + + +navigation: + port: "serial port" + baudrate: 4800 + path_output: "path/to/navigation" + + +dashboard: + title: "App title" + + spasso_path: "path/to/spasso" + navigation_path: "path/to/navigation" + cartopy_path: "path/to/cartopy" + + map_radius_km: 500 + navigation_history_hours: 6 + refresh_seconds: 60 + + plot: + width: 520 + height: 400 + + +products: + sst: + title: "Sea Surface Temperature" + tab: "Environmental" + patterns: + - "*_Copernicus_SST_L4.nc" + variables: + - "analysed_sst" + cmap: "Turbo" + label: "SST (°C)" + transform: "kelvin_to_celsius" + clim: [10.0, 20.0] + + sss: + title: "Sea Surface Salinity" + tab: "Environmental" + patterns: + - "*_Copernicus_SSS_L4.nc" + variables: + - "sos" + cmap: "Viridis" + label: "SSS" + clim: [30.0, 34.0] + + chla: + title: "Chlorophyll-a" + tab: "Environmental" + patterns: + - "*_Copernicus_CHL_L4.nc" + variables: + - "CHL" + cmap: "Viridis" + label: "Chlorophyll-a (mg m⁻³)" + clim: [0.0, 5.0] + + ftle: + title: "Finite Time Lyapunov Exponent" + tab: "Currents / Diagnostics" + patterns: + - "*_FTLE_Copernicus_PHY.nc" + variables: + - "ftle" + cmap: "Greys" + label: "FTLE" + zero_min: true + + ke: + title: "Kinetic Energy" + tab: "Currents / Diagnostics" + patterns: + - "*_KE_Copernicus_PHY.nc" + variables: + - "ke" + cmap: "Viridis" + label: "KE" + zero_min: true + + ow: + title: "Okubo-Weiss Parameter" + tab: "Currents / Diagnostics" + patterns: + - "*_OW_Copernicus_PHY.nc" + variables: + - "ow" + cmap: "RdBu_r" + label: "OW" + + +currents: + title: "Surface Geostrophic Currents" + tab: "Currents / Diagnostics" + + pattern: "????????_Copernicus_PHY.nc" + + u_variable: "ugos" + v_variable: "vgos" + + cmap: "Viridis" + label: "Current speed (m/s)" + zero_min: true + + vectors: + enabled: true + step: 2 + scale: 0.5 \ No newline at end of file diff --git a/scripts/watch_ship_navigation.py b/scripts/watch_ship_navigation.py new file mode 100644 index 0000000..ab25978 --- /dev/null +++ b/scripts/watch_ship_navigation.py @@ -0,0 +1,228 @@ +"""Continuously record ship navigation from an NMEA 0183 serial stream.""" + +from __future__ import annotations + +import argparse +from datetime import datetime, timezone +from pathlib import Path + +import pandas as pd +import serial +import yaml +import time + +def nmea_to_decimal(value: str, hemisphere: str) -> float: + """Convert an NMEA coordinate to decimal degrees.""" + raw = float(value) + + degrees = int(raw // 100) + minutes = raw - degrees * 100 + + decimal = degrees + minutes / 60.0 + + if hemisphere in {"S", "W"}: + decimal *= -1 + + return decimal + + +def parse_rmc(sentence: str) -> dict[str, object] | None: + """Parse a valid GPRMC sentence.""" + if not sentence.startswith("$GPRMC"): + return None + + fields = sentence.split(",") + + if len(fields) < 10 or fields[2] != "A": + return None + + time_string = fields[1] + latitude = fields[3] + latitude_hemisphere = fields[4] + longitude = fields[5] + longitude_hemisphere = fields[6] + speed_knots = fields[7] + course_deg = fields[8] + date_string = fields[9] + + if not all( + [ + time_string, + date_string, + latitude, + latitude_hemisphere, + longitude, + longitude_hemisphere, + ] + ): + return None + + timestamp = datetime.strptime( + date_string + time_string.split(".")[0], + "%d%m%y%H%M%S", + ).replace(tzinfo=timezone.utc) + + return { + "timestamp_utc": timestamp, + "latitude": nmea_to_decimal(latitude, latitude_hemisphere), + "longitude": nmea_to_decimal(longitude, longitude_hemisphere), + "speed_knots": float(speed_knots) if speed_knots else None, + "course_deg": float(course_deg) if course_deg else None, + } + + +def write_parquet_chunk( + records: list[dict[str, object]], + output_dir: Path, +) -> Path: + """Write one completed navigation minute to Parquet.""" + if not records: + raise ValueError("Cannot write an empty navigation chunk.") + + output_dir.mkdir(parents=True, exist_ok=True) + + first_timestamp = records[0]["timestamp_utc"] + + filename = ( + f"ship_navigation_" + f"{first_timestamp:%Y%m%d_%H%M}.parquet" + ) + + output_path = output_dir / filename + + dataframe = pd.DataFrame(records) + dataframe.to_parquet(output_path, index=False) + + return output_path + + +def load_navigation_config(config_path: Path) -> dict: + """Load navigation configuration from the cruise params YAML.""" + with config_path.open("r", encoding="utf-8") as file: + config = yaml.safe_load(file) + + if "navigation" not in config: + raise ValueError( + f"No 'navigation' section found in {config_path}" + ) + + return config["navigation"] + + +def watch_navigation( + port: str, + baudrate: int, + output_dir: Path, +) -> None: + """Continuously read GPS positions and write one Parquet file per minute.""" + print(f"Reading navigation from {port} @ {baudrate} baud") + print(f"Writing navigation to {output_dir}") + print("Press Ctrl+C to stop.") + + records: list[dict[str, object]] = [] + current_minute = None + + try: + while True: + try: + with serial.Serial( + port=port, + baudrate=baudrate, + timeout=1, + ) as stream: + print(f"Connected to {port}.") + + while True: + sentence = ( + stream.readline() + .decode("ascii", errors="ignore") + .strip() + ) + + navigation = parse_rmc(sentence) + + if navigation is None: + continue + + timestamp = navigation["timestamp_utc"] + minute = timestamp.replace( + second=0, + microsecond=0, + ) + + if current_minute is None: + current_minute = minute + + # A new UTC minute has started. + if minute != current_minute: + output_path = write_parquet_chunk( + records, + output_dir, + ) + + print( + f"Saved {len(records)} positions -> " + f"{output_path.name}" + ) + + records = [] + current_minute = minute + + records.append(navigation) + + print( + f"{timestamp.isoformat()} | " + f"lat={navigation['latitude']:.6f} | " + f"lon={navigation['longitude']:.6f} | " + f"speed={navigation['speed_knots']} kn | " + f"course={navigation['course_deg']}°" + ) + + except serial.SerialException as exc: + print( + f"Serial connection lost: {exc}\n" + f"Retrying {port} in 5 seconds..." + ) + time.sleep(5) + + except KeyboardInterrupt: + print("\nStopping navigation reader.") + + # Preserve the unfinished minute when stopped manually. + if records: + output_path = write_parquet_chunk( + records, + output_dir, + ) + + print( + f"Saved final {len(records)} positions -> " + f"{output_path.name}" + ) + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Record ship navigation from an NMEA serial stream." + ) + + parser.add_argument( + "--config", + type=Path, + required=True, + help="Path to the cruise parameter YAML.", + ) + + args = parser.parse_args() + + navigation_config = load_navigation_config(args.config) + + watch_navigation( + port=navigation_config["port"], + baudrate=int(navigation_config["baudrate"]), + output_dir=Path(navigation_config["path_output"]), + ) + + +if __name__ == "__main__": + main() \ No newline at end of file