diff --git a/src/quant_platform_kit/strategy_lifecycle/live_equity.py b/src/quant_platform_kit/strategy_lifecycle/live_equity.py index 01a89577..11a1643e 100644 --- a/src/quant_platform_kit/strategy_lifecycle/live_equity.py +++ b/src/quant_platform_kit/strategy_lifecycle/live_equity.py @@ -234,9 +234,10 @@ def extract_external_cash_flow(payload: Mapping[str, Any] | None) -> float | Non Producers should write one of ``net_external_cash_flow`` or ``external_cash_flow`` in account currency: deposits are positive and withdrawals are negative. Internal cash sweeps, realized PnL and broker - cash balances must not be supplied here. A missing field means zero; a - present but invalid field returns ``None`` so the affected daily return is - excluded instead of being misreported. + cash balances must not be supplied here. A missing field is unknown and + returns ``None``. Only an explicit finite value, including zero, is a + confirmed flow. A present but invalid field also returns ``None`` so the + affected daily return is excluded instead of being misreported. """ if not isinstance(payload, Mapping): return None @@ -254,7 +255,7 @@ def extract_external_cash_flow(payload: Mapping[str, Any] | None) -> float | Non for key in _EXTERNAL_CASH_FLOW_KEYS: if key in candidate: return _as_finite_number(candidate.get(key)) - return 0.0 + return None def cash_flow_adjusted_return( @@ -268,9 +269,10 @@ def cash_flow_adjusted_return( This is a daily time-weighted-return-compatible calculation: ``(ending_equity - signed_external_flow) / previous_equity - 1``. It is exact when flows occur at the end of the observation period, which is the - only timing available in persisted daily run records. ``None`` represents - insufficient or impossible evidence and is intentionally not coerced to a - zero return. + only timing available in persisted daily run records. The default flow of + zero is an explicit research assumption by the caller, not evidence that a + live record omitted its cash-flow field. ``None`` represents insufficient + or impossible evidence and is intentionally not coerced to a zero return. """ start = _as_float(previous_equity) end = _as_float(ending_equity) @@ -498,6 +500,13 @@ def _empty_live_return_result(status: str, detail: str = "") -> LiveReturnSeries return LiveReturnSeriesResult(series=series, status=status, detail=detail) +def _retain_unknown_cash_flow(detail: str, *, unknown: bool) -> str: + """Keep the original reason and still name an earlier unknown cash flow.""" + if not unknown or "invalid_cash_flow" in detail: + return detail + return f"{detail};invalid_cash_flow" if detail else "invalid_cash_flow" + + def live_run_records_to_return_series_result( records: Sequence[Mapping[str, Any]], *, @@ -574,6 +583,8 @@ def live_run_records_to_return_series_result( points.append((recorded_at, equity, cash_flow)) if len(points) < 2: + if invalid_cash_flow_dates: + return _empty_live_return_result("insufficient_observations", "invalid_cash_flow") return _empty_live_return_result("insufficient_observations") frame = ( @@ -582,18 +593,24 @@ def live_run_records_to_return_series_result( .groupby("date", sort=True, as_index=False) .agg({"equity": "last", "external_cash_flow": "sum"}) ) + unknown_cash_flow = bool(invalid_cash_flow_dates) if invalid_cash_flow_dates: # A plain return series cannot preserve separate comparable segments. - # Keep only the latest segment so downstream metrics cannot compound - # valid returns from opposite sides of an unknown cash-flow interval. + # Keep only observations after the latest unknown flow so metrics cannot + # compound returns from opposite sides of that interval. frame = frame[frame["date"] > max(invalid_cash_flow_dates)] if len(frame) < 2: + if unknown_cash_flow: + return _empty_live_return_result("insufficient_observations", "invalid_cash_flow") return _empty_live_return_result("insufficient_observations") span_start = _as_calendar_date(frame["date"].iloc[0]) span_end = _as_calendar_date(frame["date"].iloc[-1]) if span_start is None or span_end is None: - return _empty_live_return_result("insufficient_observations") + return _empty_live_return_result( + "insufficient_observations", + _retain_unknown_cash_flow("", unknown=unknown_cash_flow), + ) ready, readiness_detail = exchange_holiday_calendar_readiness( contract, span_start=span_start, @@ -602,7 +619,10 @@ def live_run_records_to_return_series_result( if not ready: # Unknown / synthetic / uncovered calendars are not computable. Do not # silently return a short weekday-approximated segment as success. - return _empty_live_return_result("incomplete_calendar", readiness_detail) + return _empty_live_return_result( + "incomplete_calendar", + _retain_unknown_cash_flow(readiness_detail, unknown=unknown_cash_flow), + ) frame = frame.sort_values("date", kind="stable").set_index("date") ordered_days = list(frame.index) @@ -610,7 +630,7 @@ def live_run_records_to_return_series_result( if len(latest_days) < 2: return _empty_live_return_result( "incomplete_observation_gap", - "no_contiguous_session_pair", + _retain_unknown_cash_flow("no_contiguous_session_pair", unknown=unknown_cash_flow), ) truncated = latest_days != ordered_days frame = frame.loc[latest_days] @@ -628,15 +648,25 @@ def live_run_records_to_return_series_result( return_points.append((as_of, adjusted_return)) previous_equity = current_equity if not return_points: - return _empty_live_return_result("insufficient_observations") + return _empty_live_return_result( + "insufficient_observations", + _retain_unknown_cash_flow("", unknown=unknown_cash_flow), + ) returns = pd.Series( (value for _, value in return_points), index=pd.Index((as_of for as_of, _ in return_points), name="date"), dtype=float, ) returns.name = "live_return" - status = "truncated_after_observation_gap" if truncated else "ok" - detail = "latest_contiguous_segment" if truncated else readiness_detail + if unknown_cash_flow: + status = "truncated_after_invalid_cash_flow" + detail = "latest_contiguous_segment" + elif truncated: + status = "truncated_after_observation_gap" + detail = "latest_contiguous_segment" + else: + status = "ok" + detail = readiness_detail return LiveReturnSeriesResult(series=returns.astype(float), status=status, detail=detail) @@ -651,9 +681,11 @@ def live_run_records_to_return_series( Multiple records from the same day use the final equity observation and accumulate their declared external flows. This prevents a pure deposit or withdrawal from becoming a spurious gain or loss in lifecycle monitoring. - If a declared flow is invalid, only the latest comparable segment after - that date is returned because a plain Series cannot preserve segment - boundaries for downstream compounding. + A missing external-flow field is unknown, not a confirmed zero. Invalid or + unknown flows keep only the latest comparable segment after that date, + with status ``truncated_after_invalid_cash_flow`` when a segment remains. + A plain Series cannot preserve segment boundaries for downstream + compounding. For exchange calendars, missing verified holiday source/coverage yields an empty series with status ``incomplete_calendar`` (see diff --git a/src/quant_platform_kit/strategy_lifecycle/return_collector.py b/src/quant_platform_kit/strategy_lifecycle/return_collector.py index e2decd74..65d61afb 100644 --- a/src/quant_platform_kit/strategy_lifecycle/return_collector.py +++ b/src/quant_platform_kit/strategy_lifecycle/return_collector.py @@ -172,7 +172,14 @@ def collect_from_live_runs_result( if derived.status == "ok" and not derived.series.empty: series_by_profile[profile] = derived.series continue - if derived.status == "truncated_after_observation_gap" and not derived.series.empty: + if ( + derived.status + in { + "truncated_after_observation_gap", + "truncated_after_invalid_cash_flow", + } + and not derived.series.empty + ): # Explicit truncation: usable latest segment, not silent full success. series_by_profile[profile] = derived.series incomplete_by_profile[profile] = ( @@ -222,6 +229,8 @@ def _merge_return_series( self, existing: Mapping[str, pd.Series], incoming: Mapping[str, pd.Series], + *, + incomplete_by_profile: Mapping[str, str] | None = None, ) -> dict[str, pd.Series]: merged = dict(existing) for profile, series in incoming.items(): @@ -230,17 +239,34 @@ def _merge_return_series( continue if series.empty: continue - # Prefer the observation_status already attached to the higher-priority - # (CSV/research) series when both sources contribute. + incoming_status = str(getattr(series, "attrs", {}).get("observation_status") or "") + if incoming_status == "truncated_after_invalid_cash_flow": + # The live segment already stops at the unknown flow. Do not + # stitch earlier CSV returns back onto it or relabel it ok. + merged[profile] = series + continue preferred_status = str( getattr(merged[profile], "attrs", {}).get("observation_status") - or getattr(series, "attrs", {}).get("observation_status") + or incoming_status or "ok" ) combined = pd.concat([merged[profile], series]).sort_index() combined = combined[~combined.index.duplicated(keep="last")] combined.attrs["observation_status"] = preferred_status merged[profile] = combined + for profile, reason in (incomplete_by_profile or {}).items(): + text = str(reason or "") + if "invalid_cash_flow" not in text: + continue + live = incoming.get(profile) + if live is not None and not live.empty: + continue + current = merged.get(profile) + if current is None or current.empty: + continue + stamped = current.copy() + stamped.attrs["observation_status"] = "truncated_after_invalid_cash_flow" + merged[profile] = stamped return merged def collect( @@ -300,12 +326,18 @@ def collect( for profile, series in live_outcome.series_by_profile.items(): stamped = series.copy() reason = str(live_outcome.incomplete_by_profile.get(profile) or "") - if reason.startswith("truncated_after_observation_gap"): + if reason.startswith("truncated_after_invalid_cash_flow"): + stamped.attrs["observation_status"] = "truncated_after_invalid_cash_flow" + elif reason.startswith("truncated_after_observation_gap"): stamped.attrs["observation_status"] = "truncated_after_observation_gap" else: stamped.attrs["observation_status"] = "ok" live_series[profile] = stamped - return self._merge_return_series(all_strategies, live_series) + return self._merge_return_series( + all_strategies, + live_series, + incomplete_by_profile=live_outcome.incomplete_by_profile, + ) def collect_benchmark( self, diff --git a/tests/test_lifecycle_performance_monitor.py b/tests/test_lifecycle_performance_monitor.py index b38ed0a5..68acfb1c 100644 --- a/tests/test_lifecycle_performance_monitor.py +++ b/tests/test_lifecycle_performance_monitor.py @@ -436,7 +436,7 @@ def test_run_monitor_live_truncation_stamps_status_not_full_history(self) -> Non "domain": "crypto", "recorded_at": day, "record_kind": "execution", - "execution_result": {"total_equity": equity}, + "execution_result": {"total_equity": equity, "external_cash_flow": 0.0}, }, stream_id="stream-1", ) @@ -460,6 +460,51 @@ def test_run_monitor_live_truncation_stamps_status_not_full_history(self) -> Non # Only post-gap contiguous points are usable (not the pre-gap day). self.assertLessEqual(snapshots[0].windows[2].observation_count, 3) + def test_run_monitor_unknown_cash_flow_truncation_is_not_complete(self) -> None: + revision = "j" * 40 + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + store = PerformanceStore(local_root=root / "store") + rows = ( + ("2026-09-07T20:00:00Z", 100.0, 0.0), + ("2026-09-08T20:00:00Z", 200.0, None), + ("2026-09-09T20:00:00Z", 200.0, 0.0), + ("2026-09-10T20:00:00Z", 202.0, 0.0), + ) + for day, equity, flow in rows: + execution_result: dict[str, object] = {"total_equity": equity} + if flow is not None: + execution_result["external_cash_flow"] = flow + store.save_live_run_record( + "crypto_live_pool_rotation", + "crypto", + { + "strategy_profile": "crypto_live_pool_rotation", + "domain": "crypto", + "recorded_at": day, + "record_kind": "execution", + "execution_result": execution_result, + }, + stream_id="stream-1", + ) + snapshots = run_monitor( + "crypto", + strategy_profile="crypto_live_pool_rotation", + collector=ReturnCollector(projects_root=root, store=store), + store=store, + live_stream_id="stream-1", + windows=(2,), + min_observations=1, + source_revision=revision, + ) + self.assertEqual(len(snapshots), 1) + self.assertEqual( + snapshots[0].observation_status, + "truncated_after_invalid_cash_flow", + ) + self.assertEqual(snapshots[0].windows[2].observation_count, 1) + self.assertAlmostEqual(snapshots[0].latest_return or 0.0, 202.0 / 200.0 - 1.0) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_live_return_collector.py b/tests/test_live_return_collector.py index 60a2aa40..9170095c 100644 --- a/tests/test_live_return_collector.py +++ b/tests/test_live_return_collector.py @@ -164,8 +164,8 @@ def test_extract_equity_from_nested_execution_result(self) -> None: def test_live_run_records_to_return_series(self) -> None: series = live_run_records_to_return_series( [ - {"recorded_at": "2026-07-07T10:00:00+00:00", "total_equity": 100.0}, - {"recorded_at": "2026-07-08T10:00:00+00:00", "total_equity": 101.0}, + {"recorded_at": "2026-07-07T10:00:00+00:00", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-08T10:00:00+00:00", "total_equity": 101.0, "external_cash_flow": 0.0}, ] ) self.assertEqual(len(series), 1) @@ -178,9 +178,131 @@ def test_external_cash_flow_extraction_is_signed_and_does_not_use_cash_balances( ), -50.0, ) - self.assertEqual(extract_external_cash_flow({"cash_balance": 500.0}), 0.0) + self.assertIsNone(extract_external_cash_flow({"cash_balance": 500.0})) + self.assertIsNone(extract_external_cash_flow({"total_equity": 200.0})) + self.assertEqual(extract_external_cash_flow({"external_cash_flow": 0}), 0.0) + self.assertEqual(extract_external_cash_flow({"net_external_cash_flow": "0.0"}), 0.0) self.assertIsNone(extract_external_cash_flow({"net_external_cash_flow": "invalid"})) + def test_math_helper_default_remains_explicit_zero_flow_assumption(self) -> None: + self.assertAlmostEqual(cash_flow_adjusted_return(100.0, 200.0), 1.0) + + def test_missing_cash_flow_from_100_to_200_is_not_a_trusted_return(self) -> None: + result = live_run_records_to_return_series_result( + [ + {"recorded_at": "2026-09-08T20:00:00Z", "total_equity": 100.0}, + {"recorded_at": "2026-09-09T20:00:00Z", "total_equity": 200.0}, + ] + ) + self.assertEqual(result.status, "insufficient_observations") + self.assertNotEqual(result.status, "ok") + self.assertTrue(result.series.empty) + + def test_confirmed_zero_flow_reports_the_equity_change(self) -> None: + result = live_run_records_to_return_series_result( + [ + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 100.0, + "external_cash_flow": 0.0, + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 200.0, + "external_cash_flow": 0.0, + }, + ] + ) + self.assertEqual(result.status, "ok") + self.assertAlmostEqual(float(result.series.iloc[0]), 1.0) + + def test_live_deposit_and_withdrawal_are_not_profit_or_loss(self) -> None: + deposit = live_run_records_to_return_series_result( + [ + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 100.0, + "external_cash_flow": 0.0, + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 200.0, + "external_cash_flow": 100.0, + }, + ] + ) + withdrawal = live_run_records_to_return_series_result( + [ + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 200.0, + "external_cash_flow": 0.0, + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 100.0, + "external_cash_flow": -100.0, + }, + ] + ) + self.assertEqual(deposit.status, "ok") + self.assertAlmostEqual(float(deposit.series.iloc[0]), 0.0) + self.assertEqual(withdrawal.status, "ok") + self.assertAlmostEqual(float(withdrawal.series.iloc[0]), 0.0) + + def test_missing_flow_positions_do_not_bridge_or_look_complete(self) -> None: + def row(day: str, equity: float, flow: float | None = 0.0, omit: bool = False) -> dict: + payload: dict = {"recorded_at": day, "total_equity": equity} + if not omit: + payload["external_cash_flow"] = flow + return payload + + first_missing = live_run_records_to_return_series_result( + [ + row("2026-09-08T20:00:00Z", 100.0, omit=True), + row("2026-09-09T20:00:00Z", 110.0, 0.0), + row("2026-09-10T20:00:00Z", 121.0, 0.0), + ] + ) + self.assertEqual(first_missing.status, "truncated_after_invalid_cash_flow") + self.assertEqual(list(first_missing.series.index), [pd.Timestamp("2026-09-10")]) + self.assertAlmostEqual(float(first_missing.series.iloc[0]), 121.0 / 110.0 - 1.0) + + mixed = live_run_records_to_return_series_result( + [ + row("2026-09-07T20:00:00Z", 100.0, 0.0), + row("2026-09-08T09:00:00Z", 150.0, 50.0), + row("2026-09-08T20:00:00Z", 200.0, omit=True), + row("2026-09-09T20:00:00Z", 200.0, 0.0), + row("2026-09-10T20:00:00Z", 202.0, 0.0), + ] + ) + self.assertEqual(mixed.status, "truncated_after_invalid_cash_flow") + self.assertEqual(list(mixed.series.index), [pd.Timestamp("2026-09-10")]) + self.assertAlmostEqual(float(mixed.series.iloc[0]), 0.01) + + middle = live_run_records_to_return_series_result( + [ + row("2026-09-07T20:00:00Z", 100.0, 0.0), + row("2026-09-08T20:00:00Z", 200.0, omit=True), + row("2026-09-09T20:00:00Z", 200.0, 0.0), + row("2026-09-10T20:00:00Z", 202.0, 0.0), + ] + ) + self.assertEqual(middle.status, "truncated_after_invalid_cash_flow") + self.assertEqual(list(middle.series.index), [pd.Timestamp("2026-09-10")]) + self.assertAlmostEqual(float(middle.series.iloc[0]), 0.01) + + trailing = live_run_records_to_return_series_result( + [ + row("2026-09-08T20:00:00Z", 100.0, 0.0), + row("2026-09-09T20:00:00Z", 110.0, 0.0), + row("2026-09-10T20:00:00Z", 200.0, omit=True), + ] + ) + self.assertEqual(trailing.status, "insufficient_observations") + self.assertTrue(trailing.series.empty) + def test_cash_flow_adjusted_return_is_invariant_to_pure_deposits_and_withdrawals(self) -> None: self.assertAlmostEqual( cash_flow_adjusted_return(100.0, 200.0, net_external_cash_flow=100.0) or 0.0, @@ -199,7 +321,7 @@ def test_cash_flow_adjusted_return_is_invariant_to_pure_deposits_and_withdrawals def test_live_returns_aggregate_same_day_cash_flows_before_adjustment(self) -> None: series = live_run_records_to_return_series( [ - {"recorded_at": "2026-07-07T10:00:00+00:00", "total_equity": 100.0}, + {"recorded_at": "2026-07-07T10:00:00+00:00", "total_equity": 100.0, "external_cash_flow": 0.0}, { "recorded_at": "2026-07-08T09:00:00+00:00", "total_equity": 150.0, @@ -213,6 +335,7 @@ def test_live_returns_aggregate_same_day_cash_flows_before_adjustment(self) -> N { "recorded_at": "2026-07-09T10:00:00+00:00", "total_equity": 220.0, + "external_cash_flow": 0.0, }, ] ) @@ -299,12 +422,12 @@ def test_count_consecutive_losses_trailing_only(self) -> None: def test_consecutive_losses_from_live_run_records(self) -> None: streak = consecutive_losses_from_live_run_records( [ - {"recorded_at": "2026-07-01T10:00:00+00:00", "total_equity": 100.0}, - {"recorded_at": "2026-07-02T10:00:00+00:00", "total_equity": 99.0}, - {"recorded_at": "2026-07-03T10:00:00+00:00", "total_equity": 97.0}, - {"recorded_at": "2026-07-04T10:00:00+00:00", "total_equity": 98.0}, - {"recorded_at": "2026-07-05T10:00:00+00:00", "total_equity": 96.0}, - {"recorded_at": "2026-07-06T10:00:00+00:00", "total_equity": 95.0}, + {"recorded_at": "2026-07-01T10:00:00+00:00", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-02T10:00:00+00:00", "total_equity": 99.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-03T10:00:00+00:00", "total_equity": 97.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-04T10:00:00+00:00", "total_equity": 98.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-05T10:00:00+00:00", "total_equity": 96.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-07-06T10:00:00+00:00", "total_equity": 95.0, "external_cash_flow": 0.0}, ] ) # returns: -1%, -2.02%, +1.03%, -2.04%, -1.04% → trailing streak 2 @@ -326,7 +449,7 @@ def test_resolve_consecutive_losses_from_store(self) -> None: "domain": "us_equity", "recorded_at": day, "record_kind": "execution", - "execution_result": {"total_equity": equity}, + "execution_result": {"total_equity": equity, "external_cash_flow": 0.0}, }, ) self.assertEqual( @@ -365,7 +488,7 @@ def test_stamp_consecutive_losses_on_snapshot(self) -> None: "domain": "us_equity", "recorded_at": day, "record_kind": "execution", - "execution_result": {"total_equity": equity}, + "execution_result": {"total_equity": equity, "external_cash_flow": 0.0}, }, ) snapshot = PortfolioSnapshot( @@ -542,7 +665,7 @@ def test_collect_merges_live_run_returns(self) -> None: "domain": "us_equity", "recorded_at": day, "record_kind": "execution", - "execution_result": {"total_equity": equity}, + "execution_result": {"total_equity": equity, "external_cash_flow": 0.0}, }, stream_id="schwab", ) @@ -588,20 +711,233 @@ def list_live_run_records(self, domain: str) -> list[dict[str, object]]: self.domain = domain return rows - series = ReturnCollector(store=Store()).collect_from_live_runs( - "us_equity", - stream_id="offline-account", - )["audit_case"] + with tempfile.TemporaryDirectory() as tmp: + collector = ReturnCollector(store=Store(), projects_root=Path(tmp)) + outcome = collector.collect_from_live_runs_result( + "us_equity", + stream_id="offline-account", + ) + series = outcome.series_by_profile["audit_case"] + self.assertIn( + "truncated_after_invalid_cash_flow", + outcome.incomplete_by_profile["audit_case"], + ) + collected = collector.collect( + "us_equity", + live_stream_id="offline-account", + )["audit_case"] + self.assertEqual( + collected.attrs["observation_status"], + "truncated_after_invalid_cash_flow", + ) - self.assertEqual(list(series.index), [pd.Timestamp("2026-09-10")]) - self.assertAlmostEqual(float(series.iloc[0]), 0.01) - self.assertAlmostEqual(compute_window_metrics(series).total_return, 0.01) + self.assertEqual(list(series.index), [pd.Timestamp("2026-09-10")]) + self.assertAlmostEqual(float(series.iloc[0]), 0.01) + self.assertAlmostEqual(compute_window_metrics(series).total_return, 0.01) + + def test_collect_missing_cash_flow_is_not_ok(self) -> None: + rows = [ + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 100.0, + "strategy_profile": "missing_flow", + "lifecycle_stream_id": "offline-account", + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 200.0, + "strategy_profile": "missing_flow", + "lifecycle_stream_id": "offline-account", + }, + ] + + class Store: + def list_live_run_records(self, domain: str) -> list[dict[str, object]]: + return rows + + with tempfile.TemporaryDirectory() as tmp: + outcome = ReturnCollector(store=Store(), projects_root=Path(tmp)).collect_from_live_runs_result( + "us_equity", + stream_id="offline-account", + ) + self.assertNotIn("missing_flow", outcome.series_by_profile) + self.assertIn("missing_flow", outcome.incomplete_by_profile) + self.assertNotIn("ok", outcome.incomplete_by_profile["missing_flow"]) + self.assertIn("insufficient_observations", outcome.incomplete_by_profile["missing_flow"]) + self.assertIn("invalid_cash_flow", outcome.incomplete_by_profile["missing_flow"]) + + def test_csv_cannot_bridge_or_relabel_unknown_live_cash_flow(self) -> None: + bridge = "bridge_case" + empty = "empty_case" + stream = "offline-account" + rows = [ + { + "recorded_at": "2026-09-07T20:00:00Z", + "total_equity": 100.0, + "external_cash_flow": 0.0, + "strategy_profile": bridge, + "lifecycle_stream_id": stream, + }, + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 200.0, + "strategy_profile": bridge, + "lifecycle_stream_id": stream, + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 200.0, + "external_cash_flow": 0.0, + "strategy_profile": bridge, + "lifecycle_stream_id": stream, + }, + { + "recorded_at": "2026-09-10T20:00:00Z", + "total_equity": 202.0, + "external_cash_flow": 0.0, + "strategy_profile": bridge, + "lifecycle_stream_id": stream, + }, + { + "recorded_at": "2026-09-08T20:00:00Z", + "total_equity": 100.0, + "strategy_profile": empty, + "lifecycle_stream_id": stream, + }, + { + "recorded_at": "2026-09-09T20:00:00Z", + "total_equity": 200.0, + "strategy_profile": empty, + "lifecycle_stream_id": stream, + }, + ] + + class Store: + def list_live_run_records(self, domain: str) -> list[dict[str, object]]: + return rows + + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + pd.DataFrame( + {"as_of": ["2026-09-07"], bridge: [0.01], empty: [0.01]} + ).to_csv(root / "portfolio_and_tracker_returns.csv", index=False) + collector = ReturnCollector( + artifact_roots={"crypto": root}, + projects_root=root, + store=Store(), + ) + outcome = collector.collect_from_live_runs_result("crypto", stream_id=stream) + collected = collector.collect("crypto", live_stream_id=stream) + + self.assertIn("truncated_after_invalid_cash_flow", outcome.incomplete_by_profile[bridge]) + self.assertEqual(list(outcome.series_by_profile[bridge].index), [pd.Timestamp("2026-09-10")]) + self.assertNotIn(empty, outcome.series_by_profile) + self.assertIn("invalid_cash_flow", outcome.incomplete_by_profile[empty]) + + bridge_series = collected[bridge] + self.assertEqual(list(bridge_series.index), [pd.Timestamp("2026-09-10")]) + self.assertNotIn(pd.Timestamp("2026-09-07"), list(bridge_series.index)) + self.assertAlmostEqual(float(bridge_series.iloc[0]), 202.0 / 200.0 - 1.0) + self.assertEqual(bridge_series.attrs["observation_status"], "truncated_after_invalid_cash_flow") + + empty_series = collected[empty] + self.assertEqual(list(empty_series.index), [pd.Timestamp("2026-09-07")]) + self.assertAlmostEqual(float(empty_series.iloc[0]), 0.01) + self.assertEqual(empty_series.attrs["observation_status"], "truncated_after_invalid_cash_flow") + self.assertNotEqual(empty_series.attrs["observation_status"], "ok") + + def test_unknown_cash_flow_survives_gap_and_calendar_rejection(self) -> None: + gap_profile = "gap_after_unknown" + calendar_profile = "calendar_after_unknown" + stream = "offline-account" + + def live_row(day: str, equity: float, profile: str, flow: float | None) -> dict[str, object]: + row: dict[str, object] = { + "recorded_at": day, + "total_equity": equity, + "strategy_profile": profile, + "lifecycle_stream_id": stream, + } + if flow is not None: + row["external_cash_flow"] = flow + return row + + rows = [ + live_row("2026-09-07T20:00:00Z", 100.0, gap_profile, None), + live_row("2026-09-08T20:00:00Z", 200.0, gap_profile, 0.0), + live_row("2026-09-10T20:00:00Z", 202.0, gap_profile, 0.0), + live_row("2025-12-29T20:00:00Z", 100.0, calendar_profile, None), + live_row("2025-12-30T20:00:00Z", 200.0, calendar_profile, 0.0), + live_row("2025-12-31T20:00:00Z", 202.0, calendar_profile, 0.0), + ] + + class SplitStore: + def list_live_run_records(self, domain: str) -> list[dict[str, object]]: + profile = gap_profile if domain == "crypto" else calendar_profile + return [row for row in rows if row["strategy_profile"] == profile] + + gap_result = live_run_records_to_return_series_result( + [row for row in rows if row["strategy_profile"] == gap_profile], + domain="crypto", + ) + self.assertEqual(gap_result.status, "incomplete_observation_gap") + self.assertIn("no_contiguous_session_pair", gap_result.detail) + self.assertIn("invalid_cash_flow", gap_result.detail) + self.assertTrue(gap_result.series.empty) + + calendar_result = live_run_records_to_return_series_result( + [row for row in rows if row["strategy_profile"] == calendar_profile], + domain="us_equity", + ) + self.assertEqual(calendar_result.status, "incomplete_calendar") + self.assertIn("coverage_exceeded", calendar_result.detail) + self.assertIn("invalid_cash_flow", calendar_result.detail) + self.assertTrue(calendar_result.series.empty) + + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + crypto_root = root / "crypto" + us_root = root / "us" + crypto_root.mkdir() + us_root.mkdir() + pd.DataFrame({"as_of": ["2026-09-06"], gap_profile: [0.01]}).to_csv( + crypto_root / "portfolio_and_tracker_returns.csv", index=False + ) + pd.DataFrame({"as_of": ["2025-12-28"], calendar_profile: [0.01]}).to_csv( + us_root / "portfolio_and_tracker_returns.csv", index=False + ) + collector = ReturnCollector( + artifact_roots={"crypto": crypto_root, "us_equity": us_root}, + projects_root=root, + store=SplitStore(), + ) + crypto_outcome = collector.collect_from_live_runs_result("crypto", stream_id=stream) + us_outcome = collector.collect_from_live_runs_result("us_equity", stream_id=stream) + crypto_collected = collector.collect("crypto", live_stream_id=stream) + us_collected = collector.collect("us_equity", live_stream_id=stream) + + self.assertIn("no_contiguous_session_pair", crypto_outcome.incomplete_by_profile[gap_profile]) + self.assertIn("invalid_cash_flow", crypto_outcome.incomplete_by_profile[gap_profile]) + self.assertNotIn(gap_profile, crypto_outcome.series_by_profile) + gap_series = crypto_collected[gap_profile] + self.assertEqual(list(gap_series.index), [pd.Timestamp("2026-09-06")]) + self.assertEqual(gap_series.attrs["observation_status"], "truncated_after_invalid_cash_flow") + + self.assertIn("coverage_exceeded", us_outcome.incomplete_by_profile[calendar_profile]) + self.assertIn("invalid_cash_flow", us_outcome.incomplete_by_profile[calendar_profile]) + self.assertNotIn(calendar_profile, us_outcome.series_by_profile) + calendar_series = us_collected[calendar_profile] + self.assertEqual(list(calendar_series.index), [pd.Timestamp("2025-12-28")]) + self.assertEqual( + calendar_series.attrs["observation_status"], + "truncated_after_invalid_cash_flow", + ) def test_missing_us_trading_day_does_not_become_single_day_return(self) -> None: result = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-09-16T20:00:00Z", "total_equity": 102.0}, + {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-09-16T20:00:00Z", "total_equity": 102.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -613,8 +949,8 @@ def test_missing_us_trading_day_does_not_become_single_day_return(self) -> None: def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> None: weekend = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-11T20:00:00Z", "total_equity": 100.0}, # Fri - {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 101.0}, # Mon + {"recorded_at": "2026-09-11T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, # Fri + {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, # Mon ], domain="us_equity", ) @@ -624,8 +960,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No # Labor Day 2026-09-07 is a published full-day closure, not a gap. labor = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0}, + {"recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -635,8 +971,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No # Half-day 2026-11-27 remains an expected session after Thanksgiving. half_day_gap = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-11-25T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-11-30T20:00:00Z", "total_equity": 101.0}, + {"recorded_at": "2026-11-25T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-11-30T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -645,8 +981,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No half_day_ok = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-11-25T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-11-27T20:00:00Z", "total_equity": 101.0}, + {"recorded_at": "2026-11-25T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-11-27T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -655,8 +991,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No # HKEX Lunar New Year full closures; 2026-02-16 half-day still required. hk_holiday = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-02-16T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-02-20T20:00:00Z", "total_equity": 101.0}, + {"recorded_at": "2026-02-16T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-02-20T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, ], domain="hk_equity", ) @@ -664,8 +1000,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No out_of_range = live_run_records_to_return_series_result( [ - {"recorded_at": "2025-12-30T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2025-12-31T20:00:00Z", "total_equity": 101.0}, + {"recorded_at": "2025-12-30T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2025-12-31T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -674,8 +1010,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No future = live_run_records_to_return_series_result( [ - {"recorded_at": "2027-01-04T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2027-01-05T20:00:00Z", "total_equity": 101.0}, + {"recorded_at": "2027-01-04T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2027-01-05T20:00:00Z", "total_equity": 101.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -683,8 +1019,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No synthetic = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0}, + {"recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0, "external_cash_flow": 0.0}, ], observation_contract=ReturnObservationContract( calendar_id="XNYS", @@ -701,8 +1037,8 @@ def test_published_2026_calendars_weekend_holiday_halfday_and_bounds(self) -> No def test_crypto_weekend_gap_does_not_bridge_natural_days(self) -> None: result = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-12T20:00:00Z", "total_equity": 100.0}, # Sat - {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 102.0}, # Mon, missing Sun + {"recorded_at": "2026-09-12T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, # Sat + {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 102.0, "external_cash_flow": 0.0}, # Mon, missing Sun ], domain="crypto", ) @@ -712,9 +1048,9 @@ def test_crypto_weekend_gap_does_not_bridge_natural_days(self) -> None: def test_gap_keeps_only_latest_contiguous_segment(self) -> None: result = live_run_records_to_return_series_result( [ - {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 100.0}, - {"recorded_at": "2026-09-16T20:00:00Z", "total_equity": 102.0}, - {"recorded_at": "2026-09-17T20:00:00Z", "total_equity": 103.0}, + {"recorded_at": "2026-09-14T20:00:00Z", "total_equity": 100.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-09-16T20:00:00Z", "total_equity": 102.0, "external_cash_flow": 0.0}, + {"recorded_at": "2026-09-17T20:00:00Z", "total_equity": 103.0, "external_cash_flow": 0.0}, ], domain="us_equity", ) @@ -727,12 +1063,14 @@ def test_collect_from_live_runs_refuses_us_trading_day_gap(self) -> None: { "recorded_at": "2026-09-14T20:00:00Z", "total_equity": 100.0, + "external_cash_flow": 0.0, "strategy_profile": "gap_case", "lifecycle_stream_id": "offline-account", }, { "recorded_at": "2026-09-16T20:00:00Z", "total_equity": 102.0, + "external_cash_flow": 0.0, "strategy_profile": "gap_case", "lifecycle_stream_id": "offline-account", }, @@ -756,12 +1094,14 @@ def test_collect_rejects_synthetic_holiday_source_as_real(self) -> None: { "recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, { "recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, @@ -787,12 +1127,14 @@ def test_collect_default_uses_published_2026_labor_day(self) -> None: { "recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, { "recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, @@ -815,12 +1157,14 @@ def test_custom_coverage_overlay_is_not_swallowed_by_defaults(self) -> None: { "recorded_at": "2026-09-04T20:00:00Z", "total_equity": 100.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, { "recorded_at": "2026-09-08T20:00:00Z", "total_equity": 102.0, + "external_cash_flow": 0.0, "strategy_profile": "holiday_case", "lifecycle_stream_id": "offline-account", }, @@ -858,7 +1202,7 @@ def test_collect_refuses_to_merge_multiple_account_streams(self) -> None: "domain": "us_equity", "recorded_at": day, "record_kind": "execution", - "execution_result": {"total_equity": equity}, + "execution_result": {"total_equity": equity, "external_cash_flow": 0.0}, }, stream_id=stream_id, )