mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-11 22:53:00 +00:00
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) <noreply@anthropic.com>
This commit is contained in:
parent
747ccc0383
commit
ea84d9bc72
5 changed files with 21 additions and 61 deletions
|
|
@ -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<EventEnvelope> {
|
||||
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<ServerStoreClient> {
|
||||
|
|
@ -414,7 +407,7 @@ impl ServerStoreClient {
|
|||
let page_events = parsed
|
||||
.data
|
||||
.into_iter()
|
||||
.map(wire_event_envelope_from_generated)
|
||||
.map(convert_type::<_, EventEnvelope>)
|
||||
.collect::<Result<Vec<EventEnvelope>>>()?;
|
||||
let next_page_since_seq = page_events.last().map(|event| event.seq.saturating_add(1));
|
||||
all_events.extend(page_events);
|
||||
|
|
|
|||
|
|
@ -667,7 +667,7 @@ pub(crate) fn run_events(run_dir: &Path) -> Vec<EventEnvelope> {
|
|||
.expect("event list response should contain a data array");
|
||||
items
|
||||
.into_iter()
|
||||
.map(EventEnvelope::from_wire_value)
|
||||
.map(serde_json::from_value)
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.expect("wire event envelope list should parse")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -179,7 +179,7 @@ fn run_events(run_dir: &Path) -> Vec<EventEnvelope> {
|
|||
.expect("event list response should contain a data array");
|
||||
items
|
||||
.into_iter()
|
||||
.map(EventEnvelope::from_wire_value)
|
||||
.map(serde_json::from_value)
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.expect("wire event envelope list should parse")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1480,8 +1480,7 @@ fn event_matches_run_filter(event: &EventEnvelope, run_filter: Option<&HashSet<R
|
|||
}
|
||||
|
||||
fn sse_event_from_store(event: &EventEnvelope) -> Option<Event> {
|
||||
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<ApiEventEnvelope, Response> {
|
||||
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) {
|
||||
|
|
|
|||
|
|
@ -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<Self> {
|
||||
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<serde_json::Value> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue