diff --git a/CHANGELOG.md b/CHANGELOG.md index a159b6cf..c3ea2e64 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,24 @@ adheres to [Semantic Versioning](https://semver.org/). a dict, or every one a tuple). Such a column now scores like the equivalent nested Arrow column. Columns that mix kinds, such as strings with numbers or lists with scalars, are still flagged. +- `StreamingCleaner` no longer raises `TypeError: Cannot compare tz-naive and + tz-aware timestamps` when a datetime column is tz-naive in one batch and + tz-aware in another (including a naive batch followed by strings with mixed + UTC offsets). Running datetime statistics and time-series watermarks compare + in UTC, reading naive values as UTC, and keep the zone each value arrived in. + Time-series mode warns when a timestamp or event-time column changes + awareness between batches, and `state_` marks such columns `tz_mixed`. +- `cdc_profile` no longer reads small numbers such as row numbers `1..100` as + epoch seconds and reports a 56-year-old batch as passing. When the epoch unit + is inferred and the median event time falls before 1990-01-01, it reports an + `event_time_implausible` error and leaves `freshness_seconds` unset, without + evaluating lateness or ordering. An explicit `event_time_unit=` is trusted as + before, and epoch columns in s, ms, us and ns from 1990 on are unchanged. + Numeric epoch event times before 1990-01-01 now need an explicit + `event_time_unit`; without it such a batch fails with + `event_time_implausible`. + Time-series streaming warns about such columns when `timestamp_unit` is not + set. ## [2.1.0] - 2026-09-15 diff --git a/src/freshdata/cdc.py b/src/freshdata/cdc.py index 583ed7e7..b1ff0660 100644 --- a/src/freshdata/cdc.py +++ b/src/freshdata/cdc.py @@ -27,7 +27,14 @@ import pandas as pd -from .streaming._timeseries import TIMESTAMP_UNITS, parse_timestamps, to_timedelta +from .streaming._timeseries import ( + PLAUSIBLE_EPOCH_START, + TIMESTAMP_UNITS, + describe_implausible_epoch, + implausible_epoch, + parse_timestamps, + to_timedelta, +) __all__ = ["CDCDefect", "CDCReport", "cdc_profile"] @@ -42,6 +49,7 @@ "missing_key", "invalid_operation", "missing_event_time", + "event_time_implausible", "replay_risk", ] Level = Literal["info", "warning", "error"] @@ -157,8 +165,11 @@ def _sample_keys( return tuple(str(v) for v in vals) -def _parse_event_time(raw: pd.Series, unit: str | None = None) -> pd.Series: - """Parse the event-time column; unparseable values become ``NaT``. +def _parse_event_time(raw: pd.Series, unit: str | None = None) -> tuple[pd.Series, str | None]: + """Parse the event-time column; return ``(parsed, epoch_unit)``. + + Unparseable values become ``NaT``. ``epoch_unit`` is the unit numeric + values were read in, or ``None`` for string and datetime columns. Datetime columns pass through. Mixed UTC offsets (e.g. across a DST change) or a mix of naive and offset-aware values are normalised to UTC, reading @@ -167,8 +178,7 @@ def _parse_event_time(raw: pd.Series, unit: str | None = None) -> pd.Series: when *unit* is ``None``, as ``fd.clean_timeseries`` does (see :func:`~freshdata.streaming._timeseries.parse_timestamps`). """ - parsed, _ = parse_timestamps(raw, unit) - return parsed + return parse_timestamps(raw, unit) def _to_event_tz(value: object, tz: Any) -> pd.Timestamp: @@ -238,6 +248,30 @@ def _ordering_defects( return {"out_of_order": n_ooo, "late": n_late} +def _watermark_defects( + df: pd.DataFrame, + ts: pd.Series, + key: str | None, + wm: pd.Timestamp, + defects: list[CDCDefect], +) -> dict[str, int]: + """Detect records strictly before an explicit watermark (already in the events' zone).""" + late_mask = ts < wm + n_late = int(late_mask.sum()) + if n_late: + defects.append( + CDCDefect( + kind="late", + level="error", + n_rows=n_late, + rationale=f"event time before the watermark {wm}", + sample_keys=_sample_keys(late_mask, df, key), + details={"watermark": str(wm)}, + ) + ) + return {"out_of_order": 0, "late": n_late} + + def _missing_key_defect(df: pd.DataFrame, key: str, defects: list[CDCDefect]) -> None: """Flag rows whose CDC key is null (a warning; the trust penalty is unchanged).""" null_key_mask = df[key].isna() @@ -255,6 +289,32 @@ def _missing_key_defect(df: pd.DataFrame, key: str, defects: list[CDCDefect]) -> ) +def _implausible_event_time_defect( + raw: pd.Series, name: str, unit: str, n_present: int, defects: list[CDCDefect] +) -> bool: + """Flag numeric event times whose *inferred* epoch unit gives implausible + instants (see :func:`~freshdata.streaming._timeseries.implausible_epoch`).""" + median = implausible_epoch(raw, unit) + if median is None: + return False + defects.append( + CDCDefect( + kind="event_time_implausible", + level="error", + n_rows=n_present, + rationale=f"{name!r} " + + describe_implausible_epoch(median, unit, "event_time_unit") + + "; freshness, lateness and ordering were not evaluated", + details={ + "inferred_unit": unit, + "median": median, + "plausible_from": str(PLAUSIBLE_EPOCH_START.date()), + }, + ) + ) + return True + + def _freshness( ts: pd.Series, now: pd.Timestamp | None, @@ -343,7 +403,11 @@ def cdc_profile( Epoch unit (``"s"``, ``"ms"``, ``"us"`` or ``"ns"``) of a *numeric* ``event_time`` column, such as Debezium's ``ts_ms``. ``None`` (default) infers it from the magnitude of the values, as ``fd.clean_timeseries`` - does. Ignored for string and datetime columns. + does. Ignored for string and datetime columns. When the unit is + inferred and the median value falls before 1990-01-01 in it (row + numbers or offsets such as ``1..100``), an ``event_time_implausible`` + error is reported and freshness, lateness and ordering are not + evaluated; pass the unit explicitly to read such values as epochs. Returns ------- @@ -366,7 +430,7 @@ def cdc_profile( if stale_after_td is not None and stale_after_td < pd.Timedelta(0): raise ValueError(f"stale_after must not be negative, got {stale_after!r}") - ts = _parse_event_time(df[event_time], event_time_unit) + ts, epoch_unit = _parse_event_time(df[event_time], event_time_unit) event_tz = ts.dt.tz now_ts = _to_event_tz(now, event_tz) if now is not None else None @@ -386,23 +450,20 @@ def cdc_profile( ) ) + # 1b) numeric event times that are not plausible epochs (row numbers, offsets). + # Freshness, lateness and ordering in seconds would be meaningless, so they + # are not evaluated, and the error keeps the gate from passing unverified. + implausible = epoch_unit is not None and event_time_unit is None and ( + _implausible_event_time_defect( + df[event_time], event_time, epoch_unit, n_rows - n_missing, defects + ) + ) + # 2) ordering: out-of-order + late (relative to watermark / running watermark). - if watermark is not None: - wm = _to_event_tz(watermark, event_tz) - late_mask = ts < wm - n_late = int(late_mask.sum()) - counts = {"out_of_order": 0, "late": n_late} - if n_late: - defects.append( - CDCDefect( - kind="late", - level="error", - n_rows=n_late, - rationale=f"event time before the watermark {wm}", - sample_keys=_sample_keys(late_mask, df, key), - details={"watermark": str(wm)}, - ) - ) + if implausible: + counts = {"out_of_order": 0, "late": 0} + elif watermark is not None: + counts = _watermark_defects(df, ts, key, _to_event_tz(watermark, event_tz), defects) else: counts = _ordering_defects(df, ts, key, lateness_td, defects) @@ -444,7 +505,7 @@ def cdc_profile( ) # 5) freshness / stale batch. - freshness_seconds = _freshness(ts, now_ts, stale_after_td, defects) + freshness_seconds = None if implausible else _freshness(ts, now_ts, stale_after_td, defects) # 6) replay-risk batch heuristic. A replay re-delivers changes, so at least # one duplicate-key row is required; late rows alone never trigger it. diff --git a/src/freshdata/streaming/_state.py b/src/freshdata/streaming/_state.py index 88e9add7..7e194e0c 100644 --- a/src/freshdata/streaming/_state.py +++ b/src/freshdata/streaming/_state.py @@ -22,6 +22,7 @@ import pandas as pd from pandas.api.types import is_bool_dtype, is_datetime64_any_dtype, is_numeric_dtype +from ..fieldcheck import _as_utc from ._stats import BoundedCounter, ReservoirSampler, Welford #: Cap on the retained trust-score history (keeps the state bounded on long streams). @@ -86,7 +87,10 @@ def __init__(self, name: str, role: str, baseline_dtype: str, first_seen_batch: self._dt_min: pd.Timestamp | None = None self._dt_max: pd.Timestamp | None = None self._dt_prev_max: pd.Timestamp | None = None + self._dt_tz_aware: bool | None = None # awareness of the latest batch self.datetime_ordered = True # until a batch proves otherwise + #: True once both tz-naive and tz-aware batches were seen for this column. + self.datetime_tz_mixed = False @property def seen(self) -> int: @@ -126,11 +130,21 @@ def _update_datetime(self, s: pd.Series) -> None: if nonnull.empty: return batch_min, batch_max = nonnull.min(), nonnull.max() - self._dt_min = batch_min if self._dt_min is None else min(self._dt_min, batch_min) - self._dt_max = batch_max if self._dt_max is None else max(self._dt_max, batch_max) + # Batches may differ in tz-awareness (a naive batch, then an offset-aware + # one). Compare in UTC, reading naive values as UTC, and keep each stored + # bound in the zone it arrived in. + aware = batch_min.tzinfo is not None + if self._dt_tz_aware is not None and aware != self._dt_tz_aware: + self.datetime_tz_mixed = True + self._dt_tz_aware = aware + self._dt_min = (batch_min if self._dt_min is None + else min(self._dt_min, batch_min, key=_as_utc)) + self._dt_max = (batch_max if self._dt_max is None + else max(self._dt_max, batch_max, key=_as_utc)) # Globally ordered iff each batch starts no earlier than the previous ended # and is itself internally non-decreasing. - if self._dt_prev_max is not None and batch_min < self._dt_prev_max: + if (self._dt_prev_max is not None + and _as_utc(batch_min) < _as_utc(self._dt_prev_max)): self.datetime_ordered = False if not nonnull.is_monotonic_increasing: self.datetime_ordered = False @@ -188,6 +202,8 @@ def to_dict(self) -> dict[str, Any]: "max": str(self._dt_max), "ordered": self.datetime_ordered, } + if self.datetime_tz_mixed: + payload["datetime"]["tz_mixed"] = True return payload diff --git a/src/freshdata/streaming/_timeseries.py b/src/freshdata/streaming/_timeseries.py index f9f9ca6a..2bf7f94f 100644 --- a/src/freshdata/streaming/_timeseries.py +++ b/src/freshdata/streaming/_timeseries.py @@ -30,6 +30,7 @@ from .._util import mask_sensitive_value, safe_median from ..config import CleanConfig from ..engine.context import infer_role +from ..fieldcheck import _as_utc from ..report import CleanReport from ..steps.dtypes import COERCED_CELLS_CAP @@ -132,6 +133,46 @@ def infer_epoch_unit(values: pd.Series) -> str: return "ns" +#: Earliest plausible instant for numeric timestamps whose epoch unit is +#: *inferred*. See :func:`implausible_epoch`. +PLAUSIBLE_EPOCH_START = pd.Timestamp("1990-01-01") +_UNITS_PER_SECOND = {"s": 1, "ms": 10**3, "us": 10**6, "ns": 10**9} + + +def implausible_epoch(values: pd.Series, unit: str) -> float | None: + """The median of numeric *values* if, read as epoch *unit*, it is implausible. + + Values are implausible as timestamps when their median falls before + :data:`PLAUSIBLE_EPOCH_START` (1990-01-01 UTC); ``None`` means plausible + (or no values). Row numbers, counters and small offsets such as ``1..100`` + are inferred as seconds and land in 1970, so they are caught, while any + epoch column in s, ms, us or ns for a date from 1990 on is not. The median, + as in :func:`infer_epoch_unit`, keeps a few stray values from deciding + either way. Callers apply this only to an inferred unit: an explicit unit + is trusted. + """ + median = safe_median(values.dropna()) + if pd.isna(median): + return None + start = PLAUSIBLE_EPOCH_START.value // 10**9 * _UNITS_PER_SECOND[unit] + return float(median) if float(median) < start else None + + +def describe_implausible_epoch(median: float, unit: str, option: str) -> str: + """Explain an :func:`implausible_epoch` result, naming the unit *option*.""" + return (f"numeric values have a median of {median:g}, which as epoch '{unit}' " + f"is before {PLAUSIBLE_EPOCH_START.date()}; they look like row numbers " + f"or offsets, not timestamps. Pass {option}= if they are epoch times") + + +def _as_utc_series(values: pd.Series) -> pd.Series: + """Datetimes in UTC, reading naive values as UTC (the series form of + :func:`~freshdata.fieldcheck._as_utc`), so batches of either awareness compare.""" + if values.dt.tz is None: + return values.dt.tz_localize("UTC") + return values.dt.tz_convert("UTC") + + def parse_timestamps(values: pd.Series, unit: str | None = None ) -> tuple[pd.Series, str | None]: """Parse a timestamp column; return ``(parsed, epoch_unit)``. @@ -310,6 +351,8 @@ def __post_init__(self) -> None: # Per-entity watermarks (keyed by the entity-id tuple, or () for a single # stream); each advances monotonically across batches. self._entity_watermarks: dict[tuple, pd.Timestamp] = {} + # Per column: whether its latest non-empty batch was tz-aware. + self._tz_aware: dict[str, bool] = {} self.late_quarantined_total = 0 self.late_dropped_total = 0 self.anomalies_flagged_total = 0 @@ -393,6 +436,8 @@ def _parse_timestamp_column(self, raw: pd.Series, report: CleanReport) -> pd.Ser f"read numeric timestamps as epoch unit '{epoch_unit}'", column=name, count=int(parsed.notna().sum()), risk="low", rationale=f"epoch unit {source}") + self._warn_implausible_epoch(name, raw, epoch_unit, report) + self._note_tz_awareness(name, parsed, report) lost = raw.notna() & parsed.isna() n_lost = int(lost.sum()) if not n_lost: @@ -420,6 +465,33 @@ def _parse_timestamp_column(self, raw: pd.Series, report: CleanReport) -> pd.Ser "in report.coerced_cells.") return parsed + def _warn_implausible_epoch(self, name: str, raw: pd.Series, unit: str, + report: CleanReport) -> None: + """Warn when numeric timestamps read in an *inferred* unit look like row + numbers or offsets (see :func:`implausible_epoch`); they are still parsed.""" + if self.config.timestamp_unit is not None: + return + median = implausible_epoch(raw, unit) + if median is not None: + report.add_warning(f"column '{name}': " + + describe_implausible_epoch(median, unit, "timestamp_unit")) + + def _note_tz_awareness(self, name: str, parsed: pd.Series, + report: CleanReport) -> None: + """Warn when a column's tz-awareness differs from its previous batch.""" + if not parsed.notna().any(): + return + aware = parsed.dt.tz is not None + before = self._tz_aware.get(name) + self._tz_aware[name] = aware + if before is None or before == aware: + return + now, was = ("tz-aware", "tz-naive") if aware else ("tz-naive", "tz-aware") + report.add_warning( + f"column '{name}': timestamps are {now} in this batch but were {was} in " + "the previous batch; naive timestamps are read as UTC when compared " + "across batches") + # -- step 5 (numbered by the spec): watermark-aware late data --------------- def _handle_late_data(self, df: pd.DataFrame, report: CleanReport @@ -431,7 +503,11 @@ def _handle_late_data(self, df: pd.DataFrame, report: CleanReport event_col = cfg.resolved_event_time_column if event_col not in df.columns: return df, None, {} - event_time, _ = parse_timestamps(df[event_col], cfg.timestamp_unit) + event_time, epoch_unit = parse_timestamps(df[event_col], cfg.timestamp_unit) + if event_col != cfg.timestamp_column: # the timestamp column was checked on parse + if epoch_unit is not None: + self._warn_implausible_epoch(str(event_col), df[event_col], epoch_unit, report) + self._note_tz_awareness(str(event_col), event_time, report) late_mask = self._late_mask(df, event_time, lateness) n_late = int(late_mask.sum()) @@ -485,17 +561,23 @@ def _late_mask(self, df: pd.DataFrame, event_time: pd.Series, ekey = () if not keys else (key,) if len(keys) == 1 else tuple(key) et = event_time.loc[idx] start_wm = self._entity_watermarks.get(ekey) - prior_wm = et.cummax().shift(1) + # A batch's awareness may differ from the watermark's (naive, then + # offset-aware), so compare in UTC, reading naive values as UTC. + # Watermarks keep the zone they arrived in. + et_utc = _as_utc_series(et) + prior_wm = et_utc.cummax().shift(1) if start_wm is not None: - prior_wm = prior_wm.fillna(start_wm).clip(lower=start_wm) + start_utc = _as_utc(start_wm) + prior_wm = prior_wm.fillna(start_utc).clip(lower=start_utc) late_mask.loc[idx] = ( - prior_wm.notna() & (et < prior_wm - lateness)).fillna(False) + prior_wm.notna() & (et_utc < prior_wm - lateness)).fillna(False) batch_max = et.max() if pd.notna(batch_max): self._entity_watermarks[ekey] = ( - batch_max if start_wm is None else max(start_wm, batch_max)) + batch_max if start_wm is None + else max(start_wm, batch_max, key=_as_utc)) if self._entity_watermarks: - self.watermark = max(self._entity_watermarks.values()) + self.watermark = max(self._entity_watermarks.values(), key=_as_utc) return late_mask # -- step 4: ordered dedupe ------------------------------------------------- diff --git a/tests/test_cdc_profile_fixes.py b/tests/test_cdc_profile_fixes.py index 3cb9a459..f876dd86 100644 --- a/tests/test_cdc_profile_fixes.py +++ b/tests/test_cdc_profile_fixes.py @@ -157,6 +157,71 @@ def test_invalid_event_time_unit_raises(): fd.cdc_profile(df, event_time="ts", event_time_unit="minutes") +# -- numbers that are not plausible epochs -------------------------------------- + + +def test_row_numbers_are_not_read_as_epoch_seconds(): + df = pd.DataFrame({"id": range(1, 101), "row_no": range(1, 101), "v": 1.0}) + rep = fd.cdc_profile(df, event_time="row_no", now="2026-09-15", stale_after="1h") + assert not rep.passed + assert rep.freshness_seconds is None + assert _counts(rep) == {"event_time_implausible": 100} + (defect,) = rep.defects + assert defect.level == "error" + assert defect.details == {"inferred_unit": "s", "median": 50.5, "plausible_from": "1990-01-01"} + assert "before 1990-01-01" in defect.rationale + assert "event_time_unit=" in defect.rationale + assert "freshness:" not in rep.summary() + assert rep.to_dict()["defects"][0]["kind"] == "event_time_implausible" + + +def test_implausible_event_time_skips_time_checks_but_keeps_key_checks(): + # Out of order, with a null and a duplicate change, as nullable integers. + df = pd.DataFrame( + {"seq": pd.array([3, 1, 2, 2, None], dtype="Int64"), "k": ["a", "a", "b", "b", "c"]} + ) + for watermark in (None, "2026-01-01"): + rep = fd.cdc_profile(df, event_time="seq", key="k", watermark=watermark, now="2026-01-02") + assert _counts(rep) == { + "missing_event_time": 1, + "event_time_implausible": 4, + "duplicate_key": 2, + "replay_risk": 2, + } + assert rep.freshness_seconds is None + + +@pytest.mark.parametrize("scale", [1, 10**3, 10**6, 10**9]) +@pytest.mark.parametrize(("seconds", "implausible"), [(631152000, False), (631151999, True)]) +def test_implausible_epoch_boundary_is_1990_in_every_unit(scale, seconds, implausible): + df = pd.DataFrame({"ts": [seconds * scale]}) + rep = fd.cdc_profile(df, event_time="ts", now="1990-01-02") + assert ("event_time_implausible" in _counts(rep)) is implausible + assert (rep.freshness_seconds is None) is implausible + + +@pytest.mark.parametrize("values", [[0, 0], [-86400, -3600], [0.5, 1.5]]) +def test_zero_negative_and_float_offsets_are_implausible(values): + rep = fd.cdc_profile(pd.DataFrame({"ts": values}), event_time="ts", now="2026-01-01") + assert _counts(rep) == {"event_time_implausible": 2} + + +def test_real_epoch_seconds_column_is_unchanged(): + df = pd.DataFrame({"ts": [1704067200 + 60 * i for i in range(100)]}) + rep = fd.cdc_profile(df, event_time="ts", now="2024-01-01 01:40", stale_after="1h") + assert rep.passed + assert rep.defects == [] + assert rep.freshness_seconds == pytest.approx(60.0) + + +def test_explicit_event_time_unit_trusts_small_numbers(): + df = pd.DataFrame({"row_no": range(1, 101)}) + rep = fd.cdc_profile(df, event_time="row_no", now="1970-01-01 00:02", event_time_unit="s") + assert rep.passed + assert rep.defects == [] + assert rep.freshness_seconds == pytest.approx(20.0) + + # -- #328: replay_risk needs a duplicate key ------------------------------------ diff --git a/tests/test_streaming_state.py b/tests/test_streaming_state.py index 8279792b..af67a960 100644 --- a/tests/test_streaming_state.py +++ b/tests/test_streaming_state.py @@ -97,3 +97,48 @@ def test_state_to_dict_is_json_friendly(): d = state.to_dict() assert d["rows_seen"] == 2 assert "n" in d["columns"] and "c" in d["columns"] + + +def _dt_state(): + return ColumnState("t", "datetime", "datetime64[ns]", 1, + reservoir_size=10, max_categories=8, seed=0) + + +def test_column_state_datetime_survives_tz_awareness_changes(): + naive = pd.Series(pd.to_datetime(["2020-01-01 10:00", "2020-01-01 11:00"])) + aware = pd.Series( + pd.to_datetime(["2020-01-01 12:00", "2020-01-01 13:00"])).dt.tz_localize("UTC") + + naive_then_aware = _dt_state() + naive_then_aware.update(naive) + naive_then_aware.update(aware) + assert naive_then_aware.datetime_ordered + assert naive_then_aware.to_dict()["datetime"] == { + "min": "2020-01-01 10:00:00", "max": "2020-01-01 13:00:00+00:00", + "ordered": True, "tz_mixed": True} + + aware_then_naive = _dt_state() + aware_then_naive.update(aware) + aware_then_naive.update(naive) # starts before the aware batch ended + assert not aware_then_naive.datetime_ordered + assert aware_then_naive.to_dict()["datetime"]["min"] == "2020-01-01 10:00:00" + assert aware_then_naive.to_dict()["datetime"]["max"] == "2020-01-01 13:00:00+00:00" + + +def test_column_state_datetime_compares_mixed_awareness_in_utc(): + cs = _dt_state() + cs.update(pd.Series(pd.to_datetime(["2020-01-01 10:00", "2020-01-01 11:00"]))) + # 15:00+05:30 is 09:30 UTC: earlier than the naive batch, read as UTC. + cs.update(pd.Series([pd.Timestamp("2020-01-01 15:00", tz="Asia/Kolkata")])) + d = cs.to_dict()["datetime"] + assert d["min"] == "2020-01-01 15:00:00+05:30" + assert d["max"] == "2020-01-01 11:00:00" + assert not cs.datetime_ordered + + +def test_column_state_single_awareness_has_no_tz_mixed_flag(): + cs = _dt_state() + cs.update(pd.Series(pd.to_datetime(["2020-01-01"])).dt.tz_localize("UTC")) + cs.update(pd.Series([pd.Timestamp("2020-01-02", tz="Asia/Kolkata")])) + assert not cs.datetime_tz_mixed + assert "tz_mixed" not in cs.to_dict()["datetime"] diff --git a/tests/test_streaming_timeseries.py b/tests/test_streaming_timeseries.py index 23885a05..98cd146f 100644 --- a/tests/test_streaming_timeseries.py +++ b/tests/test_streaming_timeseries.py @@ -459,6 +459,114 @@ def test_epoch_event_time_column_drives_late_data(): assert any(a.step == "late_data" and a.count == 1 for a in report.actions) +def _implausible_warnings(report): + return [w for w in report.warnings if "before 1990-01-01" in w] + + +def test_row_number_timestamps_warn_unless_unit_is_explicit(): + df = pd.DataFrame({"ts": range(1, 11), "v": [float(i) for i in range(10)]}) + _, report = _ts_clean(df) + (warning,) = _implausible_warnings(report) + assert "column 'ts'" in warning and "timestamp_unit=" in warning + _, report = _ts_clean(df, timestamp_unit="s") + assert not _implausible_warnings(report) + + +@pytest.mark.parametrize("scale", [1, 10**3, 10**6, 10**9]) +def test_real_epoch_timestamps_do_not_warn(scale): + df = pd.DataFrame({"ts": [(1704067200 + 3600 * i) * scale for i in range(4)], + "v": [1.0, 2.0, 3.0, 4.0]}) + _, report = _ts_clean(df) + assert not _implausible_warnings(report) + + +def test_row_number_event_time_column_warns(): + df = pd.DataFrame({"ts": pd.date_range("2024", periods=3, freq="h"), + "ev": [1, 2, 3], "v": [1.0, 2.0, 3.0]}) + _, report = _ts_clean(df, event_time_column="ev", allowed_lateness="1m") + (warning,) = _implausible_warnings(report) + assert "column 'ev'" in warning + + +# -- batches of mixed tz-awareness ---------------------------------------------- + +_TZ_TIMES = ["2026-01-01 10:00", "2026-01-01 09:00", "2026-01-01 11:00"] + + +def _tz_batch(hours=0, tz=None): + ts = pd.Series(pd.to_datetime(_TZ_TIMES)) + pd.Timedelta(hours=hours) + if tz is not None: + ts = ts.dt.tz_localize("UTC").dt.tz_convert(tz) + return pd.DataFrame({"id": [1, 2, 3], "ts": ts, "v": [1.0, 2.0, 3.0]}) + + +def _tz_stream(*batches): + cfg = TimeSeriesCleanConfig(timestamp_column="ts", allowed_lateness="30min") + cleaner = StreamingCleaner(time_series_config=cfg, verbose=False) + return cleaner, [cleaner.clean_batch(b) for b in batches] + + +def _awareness_warnings(report): + return [w for w in report.warnings if "column 'ts': timestamps are" in w] + + +@pytest.mark.parametrize(("first_tz", "second_tz", "watermark"), [ + (None, "UTC", "2026-01-01T13:00:00+00:00"), + ("UTC", None, "2026-01-01T13:00:00"), +]) +def test_stream_survives_tz_awareness_change_between_batches(first_tz, second_tz, watermark): + cleaner, [(out1, rep1), (out2, rep2)] = _tz_stream( + _tz_batch(tz=first_tz), _tz_batch(hours=2, tz=second_tz)) + # 09:00 is late in batch 1, and 11:00 behind 12:00 is late in batch 2. + assert (len(out1), len(out2)) == (2, 2) + assert not _awareness_warnings(rep1) + (warning,) = _awareness_warnings(rep2) + now, was = ("tz-aware", "tz-naive") if second_tz else ("tz-naive", "tz-aware") + assert f"are {now} in this batch but were {was}" in warning + assert "read as UTC" in warning + final = cleaner.finalize().streaming["time_series"] + assert final["late_quarantined_total"] == 2 + assert final["watermark"] == watermark + assert cleaner.state_["columns"]["ts"]["datetime"]["tz_mixed"] is True + + +def test_watermark_compares_naive_then_aware_batches_in_utc(): + first = pd.DataFrame({"ts": pd.to_datetime(["2026-01-01 10:00", "2026-01-01 11:00"]), + "v": [1.0, 2.0]}) + # 16:15 and 15:45 at +05:30 are 10:45 and 10:15 UTC; the cutoff is 10:30 UTC. + second = pd.DataFrame({"ts": pd.to_datetime(["2026-01-01 16:15", "2026-01-01 15:45"]) + .tz_localize("Asia/Kolkata"), "v": [3.0, 4.0]}) + cleaner, [_, (out2, _)] = _tz_stream(first, second) + assert out2["ts"].tolist() == [pd.Timestamp("2026-01-01 16:15", tz="Asia/Kolkata")] + assert cleaner.finalize().streaming["time_series"]["watermark"] == "2026-01-01T11:00:00" + + +def test_watermark_compares_aware_then_naive_batches_in_utc(): + first = pd.DataFrame({"ts": [pd.Timestamp("2026-01-01 16:30", tz="Asia/Kolkata")], + "v": [1.0]}) # 11:00 UTC + second = pd.DataFrame({"ts": pd.to_datetime(["2026-01-01 10:45", "2026-01-01 10:15"]), + "v": [2.0, 3.0]}) + cleaner, [_, (out2, _)] = _tz_stream(first, second) + assert out2["ts"].tolist() == [pd.Timestamp("2026-01-01 10:45")] + assert (cleaner.finalize().streaming["time_series"]["watermark"] + == "2026-01-01T16:30:00+05:30") + + +def test_stream_survives_naive_batch_then_mixed_offset_strings(): + second = pd.DataFrame({ + "id": [1, 2, 3], + # 11:15, 10:00 (naive, read as UTC: late) and 11:30 UTC. + "ts": ["2026-01-01T11:15:00+00:00", "2026-01-01 10:00:00", "2026-01-01T17:00:00+05:30"], + "v": [1.0, 2.0, 3.0]}) + cleaner, [_, (out2, rep2)] = _tz_stream(_tz_batch(), second) + assert out2["ts"].tolist() == [pd.Timestamp("2026-01-01 11:15", tz="UTC"), + pd.Timestamp("2026-01-01 11:30", tz="UTC")] + assert len(_awareness_warnings(rep2)) == 1 + final = cleaner.finalize().streaming["time_series"] + assert final["late_quarantined_total"] == 2 + assert final["watermark"] == "2026-01-01T11:30:00+00:00" + + # -- anomaly config and MAD (#290, #291) ---------------------------------------- def test_anomaly_window_size_one_is_rejected():