Merge feature/event-causality (Workstream B)
This commit is contained in:
commit
0384ea1354
9 changed files with 419 additions and 8 deletions
|
|
@ -43,6 +43,8 @@ impl EventNormalizer {
|
|||
event: raw.kind.clone(),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}],
|
||||
// `Manual` is never constructed as a `RawEvent::source` in this
|
||||
|
|
@ -156,6 +158,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.device.{}", verb),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Udev,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"id": id,
|
||||
"device": "unknown",
|
||||
|
|
@ -186,36 +190,48 @@ impl EventNormalizer {
|
|||
event: "bread.workspace.changed".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}],
|
||||
"createworkspace" => vec![BreadEvent {
|
||||
event: "bread.workspace.created".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({ "workspace": data }),
|
||||
}],
|
||||
"destroyworkspace" => vec![BreadEvent {
|
||||
event: "bread.workspace.destroyed".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({ "workspace": data }),
|
||||
}],
|
||||
"monitoradded" => vec![BreadEvent {
|
||||
event: "bread.monitor.connected".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({ "name": data }),
|
||||
}],
|
||||
"monitorremoved" => vec![BreadEvent {
|
||||
event: "bread.monitor.disconnected".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({ "name": data }),
|
||||
}],
|
||||
"activewindow" => vec![BreadEvent {
|
||||
event: "bread.window.focus.changed".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}],
|
||||
"activewindowv2" => {
|
||||
|
|
@ -224,6 +240,8 @@ impl EventNormalizer {
|
|||
event: "bread.window.focused".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"address": fields.first().unwrap_or(&"")
|
||||
}),
|
||||
|
|
@ -235,6 +253,8 @@ impl EventNormalizer {
|
|||
event: "bread.window.opened".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"address": fields.first().unwrap_or(&""),
|
||||
"workspace": fields.get(1).unwrap_or(&""),
|
||||
|
|
@ -249,6 +269,8 @@ impl EventNormalizer {
|
|||
event: "bread.window.closed".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({ "address": fields.first().unwrap_or(&"") }),
|
||||
}]
|
||||
}
|
||||
|
|
@ -258,6 +280,8 @@ impl EventNormalizer {
|
|||
event: "bread.window.moved".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"address": fields.first().unwrap_or(&""),
|
||||
"workspace": fields.get(1).unwrap_or(&""),
|
||||
|
|
@ -268,6 +292,8 @@ impl EventNormalizer {
|
|||
event: "bread.hyprland.event".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Hyprland,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}],
|
||||
}
|
||||
|
|
@ -285,6 +311,8 @@ impl EventNormalizer {
|
|||
},
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Power,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
});
|
||||
}
|
||||
|
|
@ -307,6 +335,8 @@ impl EventNormalizer {
|
|||
event: event.to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Power,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
});
|
||||
}
|
||||
|
|
@ -317,6 +347,8 @@ impl EventNormalizer {
|
|||
event: "bread.power.changed".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Power,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
});
|
||||
}
|
||||
|
|
@ -352,6 +384,8 @@ impl EventNormalizer {
|
|||
event: "bread.device.connected".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Bluetooth,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"id": path,
|
||||
"device": "unknown",
|
||||
|
|
@ -365,6 +399,8 @@ impl EventNormalizer {
|
|||
event: "bread.device.disconnected".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Bluetooth,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"id": path,
|
||||
"device": "unknown",
|
||||
|
|
@ -378,6 +414,8 @@ impl EventNormalizer {
|
|||
event: "bread.bluetooth.device.paired".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Bluetooth,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"id": path,
|
||||
"name": name,
|
||||
|
|
@ -390,6 +428,8 @@ impl EventNormalizer {
|
|||
event: "bread.bluetooth.device.unpaired".to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Bluetooth,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: json!({
|
||||
"id": path,
|
||||
"address": address,
|
||||
|
|
@ -434,6 +474,8 @@ impl EventNormalizer {
|
|||
event: name.to_string(),
|
||||
timestamp: raw.timestamp,
|
||||
source: AdapterSource::Network,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data,
|
||||
}]
|
||||
}
|
||||
|
|
@ -448,6 +490,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.terminal.{}", raw.kind),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -457,6 +501,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.remote.{}", raw.kind),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -466,6 +512,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.git.{}", raw.kind),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -475,6 +523,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.project.{}", raw.kind),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -487,6 +537,8 @@ impl EventNormalizer {
|
|||
event: format!("bread.service.{suffix}"),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -503,6 +555,8 @@ impl EventNormalizer {
|
|||
event,
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
@ -525,6 +579,8 @@ impl EventNormalizer {
|
|||
event: raw.kind.clone(),
|
||||
timestamp: raw.timestamp,
|
||||
source: raw.source.clone(),
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
data: raw.payload.clone(),
|
||||
}]
|
||||
}
|
||||
|
|
|
|||
|
|
@ -656,6 +656,8 @@ mod tests {
|
|||
timestamp: 0,
|
||||
source: AdapterSource::System,
|
||||
data,
|
||||
id: bread_shared::new_event_id(),
|
||||
caused_by: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -211,6 +211,19 @@ struct LuaEngine {
|
|||
next_sub_id: Arc<AtomicU64>,
|
||||
next_timer_id: Arc<AtomicU64>,
|
||||
current_module: Arc<Mutex<Option<String>>>,
|
||||
/// The `id` of the `BreadEvent` whose subscriber callback is currently
|
||||
/// executing, if any. Set immediately before invoking a handler in
|
||||
/// `handle_event` and restored (not merely cleared) immediately after,
|
||||
/// so nested/reentrant dispatch and a single event fan-out to multiple
|
||||
/// subscriptions both see the correct parent id. Read synchronously by
|
||||
/// `bread.emit()`'s binding to populate the outgoing event's
|
||||
/// `caused_by`. The Lua engine runs single-threaded/cooperatively (one
|
||||
/// engine instance processes `LuaMessage`s serially on a dedicated
|
||||
/// thread), so a callback invocation and any `bread.emit()` inside it
|
||||
/// happen synchronously within one `handle_event` call — no additional
|
||||
/// synchronization beyond the existing `Mutex` (mirroring
|
||||
/// `current_module` above) is needed.
|
||||
current_dispatch_id: Arc<Mutex<Option<String>>>,
|
||||
modules: Arc<Mutex<HashMap<String, ModuleInfo>>>,
|
||||
module_decls: Arc<Mutex<HashMap<String, ModuleDecl>>>,
|
||||
module_order: Arc<Mutex<Vec<String>>>,
|
||||
|
|
@ -240,6 +253,7 @@ impl LuaEngine {
|
|||
next_sub_id: Arc::new(AtomicU64::new(1)),
|
||||
next_timer_id: Arc::new(AtomicU64::new(1)),
|
||||
current_module: Arc::new(Mutex::new(None)),
|
||||
current_dispatch_id: Arc::new(Mutex::new(None)),
|
||||
modules: Arc::new(Mutex::new(HashMap::new())),
|
||||
module_decls: Arc::new(Mutex::new(HashMap::new())),
|
||||
module_order: Arc::new(Mutex::new(Vec::new())),
|
||||
|
|
@ -438,6 +452,7 @@ impl LuaEngine {
|
|||
bread.set("off", off_fn)?;
|
||||
|
||||
let emit_tx = self.emit_tx.clone();
|
||||
let current_dispatch_id = self.current_dispatch_id.clone();
|
||||
let emit_fn =
|
||||
self.lua
|
||||
.create_function(move |lua, (event_name, payload): (String, Value)| {
|
||||
|
|
@ -447,8 +462,19 @@ impl LuaEngine {
|
|||
.from_value::<serde_json::Value>(other)
|
||||
.unwrap_or_else(|_| serde_json::json!({})),
|
||||
};
|
||||
// Daemon-internal emit — same trusted path as adapter/IPC
|
||||
// construction, tagged System. If this runs synchronously
|
||||
// inside a subscriber callback (see `handle_event`), thread
|
||||
// the currently-dispatching event's id through as
|
||||
// `caused_by` so the causality chain can be reconstructed.
|
||||
let caused_by = current_dispatch_id
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.clone();
|
||||
let mut event = BreadEvent::new(event_name, AdapterSource::System, data);
|
||||
event.caused_by = caused_by;
|
||||
emit_tx
|
||||
.send(BreadEvent::new(event_name, AdapterSource::System, data))
|
||||
.send(event)
|
||||
.map_err(|_| LuaError::external("event channel closed"))?;
|
||||
Ok(())
|
||||
})?;
|
||||
|
|
@ -1489,6 +1515,11 @@ impl LuaEngine {
|
|||
}
|
||||
|
||||
self.set_current_module(module.clone());
|
||||
// Every subscriber invocation triggered by this event should see the
|
||||
// same parent id, and a `bread.emit()` call made synchronously
|
||||
// inside the callback should attribute its `caused_by` to *this*
|
||||
// event, not whatever was dispatching (if anything) before it.
|
||||
let previous_dispatch_id = self.set_current_dispatch_id(Some(event.id.clone()));
|
||||
let result = match kind {
|
||||
HandlerKind::Event => {
|
||||
let event_value = json_to_lua(&self.lua, &event)?;
|
||||
|
|
@ -1502,6 +1533,9 @@ impl LuaEngine {
|
|||
callback.call::<_, ()>((new_lua, old_lua))
|
||||
}
|
||||
};
|
||||
// Restore rather than clear, so this correctly unwinds if dispatch
|
||||
// is ever reentrant (e.g. a callback that pumps messages itself).
|
||||
self.set_current_dispatch_id(previous_dispatch_id);
|
||||
self.set_current_module(None);
|
||||
|
||||
if let Err(err) = result {
|
||||
|
|
@ -1678,6 +1712,18 @@ impl LuaEngine {
|
|||
}
|
||||
}
|
||||
|
||||
/// Set the "currently dispatching" event id, returning whatever id was
|
||||
/// there before. Callers restore the previous value (rather than
|
||||
/// clearing to `None`) when the handler invocation finishes, so nested/
|
||||
/// reentrant dispatch unwinds correctly — see `current_dispatch_id`'s
|
||||
/// doc comment on the struct definition.
|
||||
fn set_current_dispatch_id(&self, id: Option<String>) -> Option<String> {
|
||||
match self.current_dispatch_id.lock() {
|
||||
Ok(mut guard) => std::mem::replace(&mut *guard, id),
|
||||
Err(poisoned) => std::mem::replace(&mut *poisoned.into_inner(), id),
|
||||
}
|
||||
}
|
||||
|
||||
fn cancel_all_timers(&self) {
|
||||
if let Ok(mut map) = self.timers.lock() {
|
||||
for (_, entry) in map.drain() {
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
use std::collections::HashMap;
|
||||
use std::fs;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::{Child, Command, Stdio};
|
||||
|
|
@ -675,6 +676,147 @@ async fn events_stream_receives_emitted_events() -> Result<()> {
|
|||
Ok(())
|
||||
}
|
||||
|
||||
/// A chain of 3 handlers, each reacting to the previous one's emitted event
|
||||
/// entirely inside the Lua runtime (no test-side pumping between links):
|
||||
/// module A handles the external trigger and emits X, module B subscribes
|
||||
/// to X and emits Y, module C subscribes to Y and emits Z. This exercises
|
||||
/// the "current dispatch id" mechanism in `breadd::lua::LuaEngine` end to
|
||||
/// end: `caused_by` must thread through every hop of the chain, not just a
|
||||
/// single emit.
|
||||
#[tokio::test]
|
||||
async fn event_causality_chain_threads_caused_by_across_handlers() -> Result<()> {
|
||||
let harness = TestHarness::spawn_with_init(
|
||||
r#"
|
||||
-- module A: reacts to the external trigger, emits X
|
||||
bread.on("bread.chain.trigger", function(event)
|
||||
bread.emit("bread.chain.x", {})
|
||||
end)
|
||||
|
||||
-- module B: reacts to X, emits Y
|
||||
bread.on("bread.chain.x", function(event)
|
||||
bread.emit("bread.chain.y", {})
|
||||
end)
|
||||
|
||||
-- module C: reacts to Y, emits Z
|
||||
bread.on("bread.chain.y", function(event)
|
||||
bread.emit("bread.chain.z", {})
|
||||
end)
|
||||
"#,
|
||||
)?;
|
||||
harness.wait_until_ready().await?;
|
||||
|
||||
let stream = UnixStream::connect(harness.socket_path()).await?;
|
||||
let (read_half, mut write_half) = stream.into_split();
|
||||
let subscribe = json!({
|
||||
"id": "sub-chain",
|
||||
"method": "events.subscribe",
|
||||
"params": { "filter": "bread.chain.*" }
|
||||
});
|
||||
write_half
|
||||
.write_all(format!("{}\n", serde_json::to_string(&subscribe)?).as_bytes())
|
||||
.await?;
|
||||
|
||||
let mut reader = BufReader::new(read_half).lines();
|
||||
let ack = reader
|
||||
.next_line()
|
||||
.await?
|
||||
.ok_or_else(|| anyhow!("missing subscribe ack"))?;
|
||||
let ack_json: Value = serde_json::from_str(&ack)?;
|
||||
assert_eq!(
|
||||
ack_json
|
||||
.get("result")
|
||||
.and_then(|v| v.get("subscribed"))
|
||||
.and_then(Value::as_bool),
|
||||
Some(true)
|
||||
);
|
||||
|
||||
// The trigger itself comes in over IPC's unsourced `emit`, outside any
|
||||
// Lua handler — its `caused_by` must be None. Everything downstream
|
||||
// (X, Y, Z) is emitted by `bread.emit()` from inside a running handler.
|
||||
harness
|
||||
.send_request("emit", json!({ "event": "bread.chain.trigger", "data": {} }))
|
||||
.await?;
|
||||
|
||||
let mut events: HashMap<String, Value> = HashMap::new();
|
||||
let deadline = Instant::now() + Duration::from_secs(5);
|
||||
while events.len() < 4 && Instant::now() < deadline {
|
||||
let Some(line) = reader.next_line().await? else {
|
||||
break;
|
||||
};
|
||||
let event: Value = serde_json::from_str(&line)?;
|
||||
let name = event
|
||||
.get("event")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or_default()
|
||||
.to_string();
|
||||
if name.starts_with("bread.chain.") {
|
||||
events.insert(name, event);
|
||||
}
|
||||
}
|
||||
|
||||
let trigger = events
|
||||
.get("bread.chain.trigger")
|
||||
.expect("missing trigger event");
|
||||
let x = events.get("bread.chain.x").expect("missing X event");
|
||||
let y = events.get("bread.chain.y").expect("missing Y event");
|
||||
let z = events.get("bread.chain.z").expect("missing Z event");
|
||||
|
||||
let trigger_id = trigger
|
||||
.get("id")
|
||||
.and_then(Value::as_str)
|
||||
.expect("trigger event missing id")
|
||||
.to_string();
|
||||
let x_id = x
|
||||
.get("id")
|
||||
.and_then(Value::as_str)
|
||||
.expect("X event missing id")
|
||||
.to_string();
|
||||
let y_id = y
|
||||
.get("id")
|
||||
.and_then(Value::as_str)
|
||||
.expect("Y event missing id")
|
||||
.to_string();
|
||||
let z_id = z
|
||||
.get("id")
|
||||
.and_then(Value::as_str)
|
||||
.expect("Z event missing id")
|
||||
.to_string();
|
||||
|
||||
assert_eq!(
|
||||
trigger.get("caused_by").and_then(Value::as_str),
|
||||
None,
|
||||
"IPC-originated trigger event should have no caused_by"
|
||||
);
|
||||
assert_eq!(
|
||||
x.get("caused_by").and_then(Value::as_str),
|
||||
Some(trigger_id.as_str()),
|
||||
"X should be caused_by the trigger event's id"
|
||||
);
|
||||
assert_eq!(
|
||||
y.get("caused_by").and_then(Value::as_str),
|
||||
Some(x_id.as_str()),
|
||||
"Y should be caused_by X's id, not the original trigger"
|
||||
);
|
||||
assert_eq!(
|
||||
z.get("caused_by").and_then(Value::as_str),
|
||||
Some(y_id.as_str()),
|
||||
"Z should be caused_by Y's id"
|
||||
);
|
||||
|
||||
let ids: std::collections::HashSet<&str> =
|
||||
[trigger_id.as_str(), x_id.as_str(), y_id.as_str(), z_id.as_str()]
|
||||
.into_iter()
|
||||
.collect();
|
||||
assert_eq!(
|
||||
ids.len(),
|
||||
4,
|
||||
"every event in the chain should have a distinct id"
|
||||
);
|
||||
|
||||
harness.shutdown();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn workflow_reaches_done_via_wait_any_happy_path() -> Result<()> {
|
||||
let harness = TestHarness::spawn_with_init(
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue