def _parse_log(
structured_log: StructuredLog,
slo_resolver: SloResolver,
log_kind: str,
) -> dict:
"""Parse a structured log DataFrame into the JS data payload."""
df = structured_log.df
success_by_qid = structured_log.query_success()
events = df.sort_values("rel_time_s")
# --- Event type value sets (strings) for filtering ---
query_lifecycle_values = {
e.value for e in EventType.query_lifecycle_types()
}
routing_values = {e.value for e in EventType.routing_types()}
cluster_lifecycle_values = {
e.value for e in EventType.cluster_lifecycle_types()
}
autoscaler_values = {e.value for e in EventType.autoscaler_types()}
# --- Build per-query event timeline ---
query_events: dict[str, list[dict]] = defaultdict(list)
for _, row in events.iterrows():
qid = row.get("query_id")
if pd.isna(qid) or not qid:
continue
et = row["event_type"]
if et not in query_lifecycle_values and et not in routing_values:
continue
# Skip events emitted by the autoscaler's internal counterfactual-replay
# router — those synthetic queries have no execution lifecycle events.
if row.get("source") == "Autoscaler.QueryRouter":
continue
query_events[qid].append(
{
"rel_time_s": float(row["rel_time_s"]),
"event_type": et,
"cluster_name": row.get("cluster_name", ""),
"query_text_id": str(row.get("query_text_id", "")),
"details": row.get("details", {}),
}
)
# --- Reconstruct queries ---
queries = []
for qid, evts in query_events.items():
evts.sort(key=lambda e: e["rel_time_s"])
by_type: dict[str, list[dict]] = defaultdict(list)
for e in evts:
by_type[e["event_type"]].append(e)
# Arrival
arrival_evts = by_type.get(EventType.ARRIVAL.value, [])
arrival_s = arrival_evts[0]["rel_time_s"] if arrival_evts else None
# Execution start (required)
exec_start_evts = by_type.get(EventType.QUERY_EXECUTION_START.value, [])
if not exec_start_evts:
raise ValueError(
f"Query {qid!r} is missing a QUERY_EXECUTION_START event. "
"Check emission sites."
)
exec_start_s = exec_start_evts[0]["rel_time_s"]
# Execution finish (optional — query may have been interrupted)
exec_finish_evts = by_type.get(
EventType.QUERY_EXECUTION_FINISH.value, []
)
exec_finish_s: float | None = (
exec_finish_evts[0]["rel_time_s"] if exec_finish_evts else None
)
# Completion (required)
completion_evts = by_type.get(EventType.COMPLETION.value, [])
if not completion_evts:
raise ValueError(
f"Query {qid!r} is missing a COMPLETION event. "
"Check emission sites."
)
completion_s: float = completion_evts[0]["rel_time_s"]
success: bool | None = success_by_qid.get(qid)
# Cluster name from QUERY_ROUTED or execution events
routed_evts = by_type.get(EventType.QUERY_ROUTED.value, [])
if routed_evts:
cluster_name = routed_evts[0]["cluster_name"]
else:
cluster_name = exec_start_evts[0]["cluster_name"]
# Query text id from any event
query_text_id = ""
for e in evts:
qtid = e.get("query_text_id", "")
if qtid and str(qtid) != "nan":
query_text_id = str(qtid)
break
slo_s = slo_resolver.resolve(query_text_id if query_text_id else None)
rpu = _safe_rpu(cluster_name)
# Use arrival_s if available, otherwise exec_start_s
if arrival_s is None:
arrival_s = exec_start_s
# End-to-end latency: arrival (or exec start as fallback) to completion.
latency_s = completion_s - arrival_s
# For overall bar extent
end_s = completion_s
violates_slo = (not success) or (latency_s > slo_s)
queries.append(
{
"query_id": qid,
"query_text_id": query_text_id,
"cluster_name": cluster_name,
"rpu": rpu,
"arrival_s": arrival_s,
"exec_start_s": exec_start_s,
"exec_finish_s": exec_finish_s,
"completion_s": completion_s,
"start_s": arrival_s,
"end_s": end_s,
"latency_s": latency_s,
"slo_s": slo_s,
"success": success,
"violates_slo": violates_slo,
"state": "completed",
}
)
# --- Cluster lifecycle events ---
cluster_events_list = []
cl_mask = events["event_type"].isin(cluster_lifecycle_values)
for _, row in events[cl_mask].iterrows():
cname = row.get("cluster_name", "")
details = row.get("details", {})
cluster_events_list.append(
{
"rel_time_s": float(row["rel_time_s"]),
"event_type": row["event_type"],
"cluster_name": cname,
"rpu": _safe_rpu(cname),
"reason": details.get("reason"),
}
)
# --- Autoscaler events ---
autoscaler_events = []
as_mask = events["event_type"].isin(autoscaler_values)
for _, row in events[as_mask].iterrows():
details = row.get("details", {})
cname = row.get("cluster_name", "")
autoscaler_events.append(
{
"rel_time_s": float(row["rel_time_s"]),
"event_type": row["event_type"],
"cluster_name": cname,
"rpu": _safe_rpu(cname),
"slo_violation": details.get("slo_violation"),
"cost": details.get("cost"),
"slo_threshold": details.get("slo_threshold"),
}
)
# --- Routing score events (grouped by query_id) ---
routing_scores: dict[str, list[dict]] = defaultdict(list)
rs_mask = events["event_type"] == EventType.ROUTING_SCORE.value
for _, row in events[rs_mask].iterrows():
qid = row.get("query_id", "")
details = row.get("details", {})
routing_scores[qid].append(
{
"rel_time_s": float(row["rel_time_s"]),
"cluster_name": row.get("cluster_name", ""),
"rpu": _safe_rpu(row.get("cluster_name", "")),
"latency_s": details.get("latency_s_for_routing"),
"slo_violation": details.get("slo_violation"),
"cost": details.get("cost"),
}
)
# --- Latency update events ---
latency_update_events = []
lu_mask = events["event_type"] == EventType.LATENCY_UPDATE.value
for _, row in events[lu_mask].iterrows():
details = row.get("details", {})
latency_update_events.append(
{
"rel_time_s": float(row["rel_time_s"]),
"query_id": row.get("query_id", ""),
"cluster_name": row.get("cluster_name", ""),
"old_latency_s": details.get("old_latency_s"),
"latency_s": details.get("latency_s"),
}
)
# --- Run metadata ---
run_meta: dict[str, Any] = {}
rs_rows = events[events["event_type"] == EventType.RUN_START.value]
if not rs_rows.empty:
d = rs_rows.iloc[0].get("details", {})
run_meta = {
"workload_name": d.get("workload_name", ""),
"num_queries": d.get("num_queries"),
"routing_policy": d.get("routing_policy", ""),
"closed_loop": d.get("closed_loop"),
}
# --- Run finish time ---
run_finish_events = []
rf_rows = events[events["event_type"] == EventType.RUN_FINISH.value]
for _, row in rf_rows.iterrows():
run_finish_events.append(
{
"rel_time_s": float(row["rel_time_s"]),
"event_type": EventType.RUN_FINISH.value,
}
)
# --- Arrival times for scrubber ---
arrivals = events[
events["event_type"] == EventType.ARRIVAL.value
].sort_values("rel_time_s")
arrival_times = [float(t) for t in arrivals["rel_time_s"]]
# --- Time range ---
time_range = [
float(events["rel_time_s"].min()),
float(events["rel_time_s"].max()),
]
return {
"kind": log_kind,
"queries": queries,
"cluster_events": cluster_events_list,
"autoscaler_events": autoscaler_events,
"routing_scores": dict(routing_scores),
"latency_update_events": latency_update_events,
"run_meta": run_meta,
"run_finish_events": run_finish_events,
"arrival_times": arrival_times,
"time_range": time_range,
"default_slo_s": slo_resolver.default_slo_s,
}