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
68 changes: 50 additions & 18 deletions src/quant_platform_kit/strategy_lifecycle/live_equity.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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(
Expand All @@ -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)
Expand Down Expand Up @@ -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]],
*,
Expand Down Expand Up @@ -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 = (
Expand All @@ -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,
Expand All @@ -602,15 +619,18 @@ 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)
latest_days = _latest_contiguous_observation_days(ordered_days, contract)
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]
Expand All @@ -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)


Expand All @@ -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
Expand Down
44 changes: 38 additions & 6 deletions src/quant_platform_kit/strategy_lifecycle/return_collector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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] = (
Expand Down Expand Up @@ -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():
Expand All @@ -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(
Expand Down Expand Up @@ -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,
Expand Down
47 changes: 46 additions & 1 deletion tests/test_lifecycle_performance_monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
)
Expand All @@ -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()
Loading
Loading