Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/pr.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ jobs:
strategy:
fail-fast: false
matrix:
python-version: ["3.13", "3.14", "3.15"]
python-version: ["3.12", "3.13"]

steps:
- name: Checkout repository
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/pypi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,8 @@ name: Publish to PyPI
on:
release:
types: [released]
push:
branches: [main]
# push:
# branches: [main]
workflow_dispatch:

permissions:
Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ Echodataflow streamlines echosounder data processing by combining [Prefect](http

1. Set up a computing environment using Conda:
```bash
conda create --name echodataflow -c conda-forge python=3.12
conda create --name echodataflow -c conda-forge python=3.13
conda activate echodataflow
```

Expand All @@ -34,7 +34,7 @@ Echodataflow streamlines echosounder data processing by combining [Prefect](http
clone the repo and install it like below:
```bash
git clone https://github.com/echostack-org/echodataflow.git # clone the repo
pip install -e ".[test,lint,docs]" # install in editable mode with dev tools
pip install -e ".[test,lint,docs,mission]" # install in editable mode with dev tools
```

3. Pip install the `segmentation_inference` package that contains a version of the hake segmentation model.
Expand Down
8 changes: 0 additions & 8 deletions src/echodataflow/flows/flows_acoustics.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,14 +40,6 @@
read_or_create_ledger,
)

from echodataflow.utils.processing_ledger import (
get_raw_files_to_process,
initialize_ledger,
mark_raw_completed,
mark_raw_failed,
mark_raw_processing,
resolve_database,
)
from echodataflow.tasks.tasks_acoustics import (
task_create_MVBS,
task_raw2Sv,
Expand Down
27 changes: 18 additions & 9 deletions src/echodataflow/operations/operations_postprocessing.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
MVBS_COLUMNS_POSTPROCESSING,
PREDICTION_COLUMNS_POSTPROCESSING,
filter_time_range,
normalize_ledger_dates,
read_manifest,
write_manifest,
)
Expand Down Expand Up @@ -127,7 +128,9 @@ def build_Sv_ledger(raw_files: pd.DataFrame) -> pd.DataFrame:
invalid = ledger.loc[ledger["timestamp"].isna(), "raw_filename"].tolist()
if invalid:
raise ValueError(f"Could not parse timestamps from raw filenames: {invalid}")
ledger["timestamp"] = pd.to_datetime(ledger["timestamp"], utc=True)
ledger["timestamp"] = pd.to_datetime(ledger["timestamp"], utc=True).astype(
"datetime64[ns, UTC]"
)

ledger["Sv_filename"] = pd.NA
ledger["raw2Sv_status"] = "pending"
Expand Down Expand Up @@ -156,7 +159,10 @@ def build_MVBS_ledger(
if no_data_gap_hours <= 0:
raise ValueError("no_data_gap_hours must be greater than zero")
if ledger_Sv.empty:
return pd.DataFrame(columns=MVBS_COLUMNS_POSTPROCESSING)
return normalize_ledger_dates(
pd.DataFrame(columns=MVBS_COLUMNS_POSTPROCESSING),
["slice_start", "slice_end", "first_ping_time", "last_ping_time"],
)

df_Sv = ledger_Sv.copy()
df_Sv["timestamp"] = pd.to_datetime(df_Sv["timestamp"], utc=True)
Expand Down Expand Up @@ -205,9 +211,9 @@ def build_MVBS_ledger(
}
)
ledger = pd.DataFrame.from_records(records, columns=MVBS_COLUMNS_POSTPROCESSING)
for column in ("first_ping_time", "last_ping_time"):
ledger[column] = pd.to_datetime(ledger[column], utc=True)
return ledger
return normalize_ledger_dates(
ledger, ["slice_start", "slice_end", "first_ping_time", "last_ping_time"]
)


def failure_state(attempt_count: int, max_flow_run_attempts: int) -> tuple[int, str]:
Expand Down Expand Up @@ -277,7 +283,10 @@ def build_prediction_ledger(
) -> pd.DataFrame:
"""Preplan every prediction window and its required MVBS slices."""
if ledger_MVBS.empty:
return pd.DataFrame(columns=PREDICTION_COLUMNS_POSTPROCESSING)
return normalize_ledger_dates(
pd.DataFrame(columns=PREDICTION_COLUMNS_POSTPROCESSING),
["slice_start", "slice_end", "first_ping_time", "last_ping_time"],
)

df_MVBS = ledger_MVBS.copy()
df_MVBS["slice_start"] = pd.to_datetime(df_MVBS["slice_start"], utc=True)
Expand Down Expand Up @@ -317,9 +326,9 @@ def build_prediction_ledger(
}
)
ledger = pd.DataFrame.from_records(records, columns=PREDICTION_COLUMNS_POSTPROCESSING)
for column in ("first_ping_time", "last_ping_time"):
ledger[column] = pd.to_datetime(ledger[column], utc=True)
return ledger
return normalize_ledger_dates(
ledger, ["slice_start", "slice_end", "first_ping_time", "last_ping_time"]
)


def read_or_create_ledger(
Expand Down
20 changes: 13 additions & 7 deletions src/echodataflow/utils/manifests.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,10 +64,20 @@
]


def normalize_ledger_dates(df: pd.DataFrame, date_columns: list[str]) -> pd.DataFrame:
"""Normalize ledger dates in place to UTC with nanosecond storage precision."""
for column in date_columns:
if column in df:
df[column] = pd.to_datetime(df[column], format="mixed", utc=True).astype(
"datetime64[ns, UTC]"
)
return df


def read_manifest(path: Path, columns: list[str], date_columns: list[str]) -> pd.DataFrame:
# Return a schema-correct empty manifest on the first run
if not path.exists():
return pd.DataFrame(columns=columns)
return normalize_ledger_dates(pd.DataFrame(columns=columns), date_columns)
df = pd.read_csv(path, index_col=0)

# Validation: catch missing or unexpected columns
Expand Down Expand Up @@ -98,12 +108,8 @@ def read_manifest(path: Path, columns: list[str], date_columns: list[str]) -> pd
if column in df:
df[column] = df[column].fillna("")

# Set datetime columns to UTC
for column in date_columns:
if column in df:
# Ledger updates can mix legacy naive values with UTC-qualified values
df[column] = pd.to_datetime(df[column], format="mixed", utc=True)
return df
# Ledger updates can mix legacy naive values with UTC-qualified values.
return normalize_ledger_dates(df, date_columns)


def write_manifest(df: pd.DataFrame, path: Path) -> None:
Expand Down
154 changes: 0 additions & 154 deletions tests/test_flow_raw2Sv.py

This file was deleted.

20 changes: 20 additions & 0 deletions tests/test_manifests.py
Original file line number Diff line number Diff line change
Expand Up @@ -249,3 +249,23 @@ def test_filter_time_range_can_exclude_interval_ending_at_exact_start():
)

assert selected["name"].tolist() == ["overlaps"]


@pytest.mark.parametrize(
"stored_value",
[None, "", "2025-06-11T21:47:11+0000", "2025-06-11T21:47:11.580137+0000"],
)
def test_manifest_dates_use_nanosecond_precision(tmp_path, stored_value):
path = tmp_path / "ledger.csv"
columns = ["first_ping_time", "last_ping_time"]
if stored_value is not None:
pd.DataFrame({column: [stored_value] for column in columns}).to_csv(path)

ledger = read_manifest(path, columns, columns)
for column in columns:
assert str(ledger[column].dtype) == "datetime64[ns, UTC]"
if stored_value:
assert ledger.loc[0, column] == pd.Timestamp(stored_value)
ping_time = pd.Timestamp("2025-06-11T21:47:11.580137+0000")
ledger.loc[0, column] = ping_time
assert ledger.loc[0, column] == ping_time