From ea84d9bc726f76c59faef90d384d10f7ff37c95c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 9 Apr 2026 12:09:00 -0400 Subject: [PATCH] refactor(store): flatten EventEnvelope wire shape via serde Replace the hand-written to_wire_value / from_wire_value helpers and the wire_event_envelope_from_generated bridge with #[serde(flatten)] on EventEnvelope.payload. Derived serde now produces and accepts the wire shape natively: { "seq": 42, "id": "...", "event": "...", ... } instead of the nested { "seq": 42, "payload": { ... } } the derive would otherwise emit. #[serde(flatten)] composes fine with the #[serde(transparent)] EventPayload(Value) wrapper, so the inner payload object is merged into the outer map on both sides. - fabro-store/src/types.rs: add #[serde(flatten)]; delete the two wire helpers (33 lines of Value-map poking); update the round-trip test to assert the shape is actually flat. - fabro-server/src/server.rs: sse_event_from_store serializes the envelope directly; api_event_envelope_from_store pipelines to_value into from_value. - fabro-cli/src/server_client.rs: buffer_sse_events parses straight into EventEnvelope via serde_json::from_str; list_run_events uses the existing convert_type helper in place of the deleted wire_event_envelope_from_generated bridge. - fabro-cli tests: helpers that called from_wire_value now call serde_json::from_value. Drops the shape check that from_wire_value used to perform on parse (id/ts/run_id/event must exist as strings): that check extracted run_id from the payload and then validated it against itself, so it only guaranteed presence, not correctness. EventPayload::new(value, expected_run_id) still runs the same check where a caller has a real external run_id to cross-match. Generated code and the OpenAPI allOf(seq, RunEvent) schema are untouched; the wire JSON is byte-identical before and after. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/crates/fabro-cli/src/server_client.rs | 11 +---- lib/crates/fabro-cli/tests/it/cmd/support.rs | 2 +- lib/crates/fabro-cli/tests/it/workflow/mod.rs | 2 +- lib/crates/fabro-server/src/server.rs | 21 ++++----- lib/crates/fabro-store/src/types.rs | 46 +++---------------- 5 files changed, 21 insertions(+), 61 deletions(-) diff --git a/lib/crates/fabro-cli/src/server_client.rs b/lib/crates/fabro-cli/src/server_client.rs index a4aa44822..97b325301 100644 --- a/lib/crates/fabro-cli/src/server_client.rs +++ b/lib/crates/fabro-cli/src/server_client.rs @@ -72,20 +72,13 @@ impl RunAttachEventStream { fn buffer_sse_events(&mut self, finalize: bool) -> Result<()> { for payload in sse::drain_sse_payloads(&mut self.pending_bytes, finalize) { - let value: serde_json::Value = serde_json::from_str(&payload)?; self.buffered_events - .push_back(EventEnvelope::from_wire_value(value)?); + .push_back(serde_json::from_str(&payload)?); } Ok(()) } } -fn wire_event_envelope_from_generated(value: types::EventEnvelope) -> Result { - let value = - serde_json::to_value(value).context("failed to serialize generated EventEnvelope")?; - EventEnvelope::from_wire_value(value).map_err(Into::into) -} - pub(crate) use fabro_store::RunProjection; pub(crate) async fn connect_server(storage_dir: &Path) -> Result { @@ -414,7 +407,7 @@ impl ServerStoreClient { let page_events = parsed .data .into_iter() - .map(wire_event_envelope_from_generated) + .map(convert_type::<_, EventEnvelope>) .collect::>>()?; let next_page_since_seq = page_events.last().map(|event| event.seq.saturating_add(1)); all_events.extend(page_events); diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index d93851add..53e75af42 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -667,7 +667,7 @@ pub(crate) fn run_events(run_dir: &Path) -> Vec { .expect("event list response should contain a data array"); items .into_iter() - .map(EventEnvelope::from_wire_value) + .map(serde_json::from_value) .collect::, _>>() .expect("wire event envelope list should parse") } diff --git a/lib/crates/fabro-cli/tests/it/workflow/mod.rs b/lib/crates/fabro-cli/tests/it/workflow/mod.rs index 6fe7383f1..b07bed9e2 100644 --- a/lib/crates/fabro-cli/tests/it/workflow/mod.rs +++ b/lib/crates/fabro-cli/tests/it/workflow/mod.rs @@ -179,7 +179,7 @@ fn run_events(run_dir: &Path) -> Vec { .expect("event list response should contain a data array"); items .into_iter() - .map(EventEnvelope::from_wire_value) + .map(serde_json::from_value) .collect::, _>>() .expect("wire event envelope list should parse") } diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index acf1fce6a..93f6bfde4 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -1480,8 +1480,7 @@ fn event_matches_run_filter(event: &EventEnvelope, run_filter: Option<&HashSet Option { - let wire = event.to_wire_value().ok()?; - let data = serde_json::to_string(&wire).ok()?; + let data = serde_json::to_string(event).ok()?; let data = redact_jsonl_line(&data); Some(Event::default().data(data)) } @@ -2381,15 +2380,15 @@ fn octet_stream_response(bytes: Bytes) -> Response { #[allow(clippy::result_large_err)] fn api_event_envelope_from_store(event: &EventEnvelope) -> Result { - fn serialize_error(err: impl std::fmt::Display) -> Response { - ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - format!("Failed to serialize stored event: {err}"), - ) - .into_response() - } - let value = event.to_wire_value().map_err(serialize_error)?; - serde_json::from_value(value).map_err(serialize_error) + serde_json::to_value(event) + .and_then(serde_json::from_value) + .map_err(|err| { + ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to serialize stored event: {err}"), + ) + .into_response() + }) } fn clear_live_run_state(run: &mut ManagedRun) { diff --git a/lib/crates/fabro-store/src/types.rs b/lib/crates/fabro-store/src/types.rs index 1c63553b5..4c3d23f75 100644 --- a/lib/crates/fabro-store/src/types.rs +++ b/lib/crates/fabro-store/src/types.rs @@ -83,46 +83,10 @@ impl TryFrom<&EventPayload> for RunEvent { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct EventEnvelope { pub seq: u32, + #[serde(flatten)] pub payload: EventPayload, } -impl EventEnvelope { - pub fn from_wire_value(value: serde_json::Value) -> Result { - let serde_json::Value::Object(mut obj) = value else { - return Err(StoreError::InvalidEvent( - "wire EventEnvelope must be a JSON object".into(), - )); - }; - let seq = obj - .remove("seq") - .and_then(|value| value.as_u64()) - .and_then(|value| u32::try_from(value).ok()) - .ok_or_else(|| { - StoreError::InvalidEvent("wire EventEnvelope missing valid seq".into()) - })?; - let run_id = obj - .get("run_id") - .and_then(|value| value.as_str()) - .ok_or_else(|| StoreError::InvalidEvent("wire EventEnvelope missing run_id".into()))? - .parse() - .map_err(|err| StoreError::InvalidEvent(format!("invalid wire run_id: {err}")))?; - let payload = EventPayload::new(serde_json::Value::Object(obj), &run_id)?; - Ok(Self { seq, payload }) - } - - pub fn to_wire_value(&self) -> Result { - let mut value = self.payload.as_value().clone(); - let map = value.as_object_mut().ok_or_else(|| { - StoreError::InvalidEvent("stored event payload must be a JSON object".into()) - })?; - map.insert( - "seq".to_string(), - serde_json::Value::Number(u64::from(self.seq).into()), - ); - Ok(value) - } -} - #[cfg(test)] mod tests { use chrono::{TimeZone, Utc}; @@ -160,9 +124,13 @@ mod tests { let payload = EventPayload::new(event.to_value().unwrap(), &fixtures::RUN_1).unwrap(); let envelope = EventEnvelope { seq: 7, payload }; - let wire = envelope.to_wire_value().unwrap(); - let parsed = EventEnvelope::from_wire_value(wire).unwrap(); + let wire = serde_json::to_value(&envelope).unwrap(); + assert_eq!(wire["seq"], 7); + assert_eq!(wire["id"], "evt_1"); + assert_eq!(wire["event"], "run.completed"); + assert!(wire.get("payload").is_none(), "wire shape must be flat"); + let parsed: EventEnvelope = serde_json::from_value(wire).unwrap(); assert_eq!(parsed, envelope); } }