mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-08-28 05:27:41 +00:00
feat(types): add RunEvent envelope fields + ActorRef (schema v2)
Adds stage_id, parallel_group_id, parallel_branch_id, tool_call_id, and actor to RunEvent per the v2 concrete-shape proposal. Introduces ActorRef/ActorKind types. Serialization omits absent fields rather than writing null. Stubs StoredEventFields with matching defaults; population in stored_event_fields() follows in a later commit. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
6f354ffc46
commit
4eec9124fa
7 changed files with 268 additions and 62 deletions
|
|
@ -198,8 +198,13 @@ fn run_event(run_id: fabro_types::RunId, node_id: Option<String>, body: EventBod
|
|||
run_id,
|
||||
node_id,
|
||||
node_label: None,
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -90,8 +90,13 @@ fn pr_view_reads_pull_request_from_store_without_pull_request_json() {
|
|||
run_id,
|
||||
node_id: None,
|
||||
node_label: None,
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body: EventBody::PullRequestCreated(PullRequestCreatedProps {
|
||||
pr_url: "https://github.com/fabro-sh/fabro/pull/123".to_string(),
|
||||
pr_number: 123,
|
||||
|
|
|
|||
|
|
@ -611,8 +611,13 @@ mod tests {
|
|||
run_id: fixtures::RUN_1,
|
||||
node_id: node_id.map(ToOwned::to_owned),
|
||||
node_label: None,
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body,
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -49,7 +49,7 @@ pub use run::{
|
|||
RunSubjectProvenance,
|
||||
};
|
||||
pub use run_blob_id::RunBlobId;
|
||||
pub use run_event::{EventBody, RunEvent, RunNoticeLevel};
|
||||
pub use run_event::{ActorKind, ActorRef, EventBody, RunEvent, RunNoticeLevel};
|
||||
pub use run_id::RunId;
|
||||
pub use run_id::fixtures;
|
||||
pub use sandbox_record::SandboxRecord;
|
||||
|
|
|
|||
|
|
@ -27,6 +27,23 @@ pub enum RunNoticeLevel {
|
|||
Error,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ActorKind {
|
||||
User,
|
||||
Agent,
|
||||
System,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct ActorRef {
|
||||
pub kind: ActorKind,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub id: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub display: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct RunEvent {
|
||||
pub id: String,
|
||||
|
|
@ -34,8 +51,13 @@ pub struct RunEvent {
|
|||
pub run_id: RunId,
|
||||
pub node_id: Option<String>,
|
||||
pub node_label: Option<String>,
|
||||
pub stage_id: Option<String>,
|
||||
pub parallel_group_id: Option<String>,
|
||||
pub parallel_branch_id: Option<String>,
|
||||
pub session_id: Option<String>,
|
||||
pub parent_session_id: Option<String>,
|
||||
pub tool_call_id: Option<String>,
|
||||
pub actor: Option<ActorRef>,
|
||||
pub body: EventBody,
|
||||
}
|
||||
|
||||
|
|
@ -271,9 +293,19 @@ struct RunEventRaw {
|
|||
#[serde(default)]
|
||||
node_label: Option<String>,
|
||||
#[serde(default)]
|
||||
stage_id: Option<String>,
|
||||
#[serde(default)]
|
||||
parallel_group_id: Option<String>,
|
||||
#[serde(default)]
|
||||
parallel_branch_id: Option<String>,
|
||||
#[serde(default)]
|
||||
session_id: Option<String>,
|
||||
#[serde(default)]
|
||||
parent_session_id: Option<String>,
|
||||
#[serde(default)]
|
||||
tool_call_id: Option<String>,
|
||||
#[serde(default)]
|
||||
actor: Option<ActorRef>,
|
||||
event: String,
|
||||
#[serde(default = "default_properties")]
|
||||
properties: Value,
|
||||
|
|
@ -283,6 +315,23 @@ fn default_properties() -> Value {
|
|||
Value::Object(Map::new())
|
||||
}
|
||||
|
||||
struct RunEventParts<'a> {
|
||||
id: String,
|
||||
ts: DateTime<Utc>,
|
||||
run_id: RunId,
|
||||
node_id: Option<String>,
|
||||
node_label: Option<String>,
|
||||
stage_id: Option<String>,
|
||||
parallel_group_id: Option<String>,
|
||||
parallel_branch_id: Option<String>,
|
||||
session_id: Option<String>,
|
||||
parent_session_id: Option<String>,
|
||||
tool_call_id: Option<String>,
|
||||
actor: Option<ActorRef>,
|
||||
event: &'a str,
|
||||
properties: &'a Value,
|
||||
}
|
||||
|
||||
impl EventBody {
|
||||
pub fn event_name(&self) -> &str {
|
||||
match self {
|
||||
|
|
@ -524,17 +573,22 @@ fn is_known_event_name(event: &str) -> bool {
|
|||
impl RunEvent {
|
||||
pub fn from_value(value: Value) -> serde_json::Result<Self> {
|
||||
let raw: RunEventRaw = serde_json::from_value(value)?;
|
||||
Self::from_parts(
|
||||
raw.id,
|
||||
raw.ts,
|
||||
raw.run_id,
|
||||
raw.node_id,
|
||||
raw.node_label,
|
||||
raw.session_id,
|
||||
raw.parent_session_id,
|
||||
&raw.event,
|
||||
&raw.properties,
|
||||
)
|
||||
Self::from_parts(RunEventParts {
|
||||
id: raw.id,
|
||||
ts: raw.ts,
|
||||
run_id: raw.run_id,
|
||||
node_id: raw.node_id,
|
||||
node_label: raw.node_label,
|
||||
stage_id: raw.stage_id,
|
||||
parallel_group_id: raw.parallel_group_id,
|
||||
parallel_branch_id: raw.parallel_branch_id,
|
||||
session_id: raw.session_id,
|
||||
parent_session_id: raw.parent_session_id,
|
||||
tool_call_id: raw.tool_call_id,
|
||||
actor: raw.actor,
|
||||
event: &raw.event,
|
||||
properties: &raw.properties,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn from_ref(value: &Value) -> serde_json::Result<Self> {
|
||||
|
|
@ -559,58 +613,78 @@ impl RunEvent {
|
|||
.get("properties")
|
||||
.cloned()
|
||||
.unwrap_or_else(default_properties);
|
||||
Self::from_parts(
|
||||
id.to_string(),
|
||||
let actor = match obj.get("actor") {
|
||||
Some(value) if !value.is_null() => Some(ActorRef::deserialize(value)?),
|
||||
_ => None,
|
||||
};
|
||||
Self::from_parts(RunEventParts {
|
||||
id: id.to_string(),
|
||||
ts,
|
||||
run_id,
|
||||
obj.get("node_id")
|
||||
node_id: obj
|
||||
.get("node_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
obj.get("node_label")
|
||||
node_label: obj
|
||||
.get("node_label")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
obj.get("session_id")
|
||||
stage_id: obj
|
||||
.get("stage_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
obj.get("parent_session_id")
|
||||
parallel_group_id: obj
|
||||
.get("parallel_group_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
parallel_branch_id: obj
|
||||
.get("parallel_branch_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
session_id: obj
|
||||
.get("session_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
parent_session_id: obj
|
||||
.get("parent_session_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
tool_call_id: obj
|
||||
.get("tool_call_id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::to_string),
|
||||
actor,
|
||||
event,
|
||||
&properties,
|
||||
)
|
||||
properties: &properties,
|
||||
})
|
||||
}
|
||||
|
||||
fn from_parts(
|
||||
id: String,
|
||||
ts: DateTime<Utc>,
|
||||
run_id: RunId,
|
||||
node_id: Option<String>,
|
||||
node_label: Option<String>,
|
||||
session_id: Option<String>,
|
||||
parent_session_id: Option<String>,
|
||||
event: &str,
|
||||
properties: &Value,
|
||||
) -> serde_json::Result<Self> {
|
||||
fn from_parts(parts: RunEventParts<'_>) -> serde_json::Result<Self> {
|
||||
let body_payload = json!({
|
||||
"event": event,
|
||||
"properties": properties,
|
||||
"event": parts.event,
|
||||
"properties": parts.properties,
|
||||
});
|
||||
let body: EventBody = match serde_json::from_value(body_payload) {
|
||||
Ok(body) => body,
|
||||
Err(err) if is_known_event_name(event) => return Err(err),
|
||||
Err(err) if is_known_event_name(parts.event) => return Err(err),
|
||||
Err(_) => EventBody::Unknown {
|
||||
name: event.to_string(),
|
||||
properties: properties.clone(),
|
||||
name: parts.event.to_string(),
|
||||
properties: parts.properties.clone(),
|
||||
},
|
||||
};
|
||||
Ok(Self {
|
||||
id,
|
||||
ts,
|
||||
run_id,
|
||||
node_id,
|
||||
node_label,
|
||||
session_id,
|
||||
parent_session_id,
|
||||
id: parts.id,
|
||||
ts: parts.ts,
|
||||
run_id: parts.run_id,
|
||||
node_id: parts.node_id,
|
||||
node_label: parts.node_label,
|
||||
stage_id: parts.stage_id,
|
||||
parallel_group_id: parts.parallel_group_id,
|
||||
parallel_branch_id: parts.parallel_branch_id,
|
||||
session_id: parts.session_id,
|
||||
parent_session_id: parts.parent_session_id,
|
||||
tool_call_id: parts.tool_call_id,
|
||||
actor: parts.actor,
|
||||
body,
|
||||
})
|
||||
}
|
||||
|
|
@ -643,6 +717,27 @@ impl RunEvent {
|
|||
if let Some(value) = &self.node_label {
|
||||
map.insert("node_label".to_string(), Value::String(value.clone()));
|
||||
}
|
||||
if let Some(value) = &self.stage_id {
|
||||
map.insert("stage_id".to_string(), Value::String(value.clone()));
|
||||
}
|
||||
if let Some(value) = &self.parallel_group_id {
|
||||
map.insert(
|
||||
"parallel_group_id".to_string(),
|
||||
Value::String(value.clone()),
|
||||
);
|
||||
}
|
||||
if let Some(value) = &self.parallel_branch_id {
|
||||
map.insert(
|
||||
"parallel_branch_id".to_string(),
|
||||
Value::String(value.clone()),
|
||||
);
|
||||
}
|
||||
if let Some(value) = &self.tool_call_id {
|
||||
map.insert("tool_call_id".to_string(), Value::String(value.clone()));
|
||||
}
|
||||
if let Some(actor) = &self.actor {
|
||||
map.insert("actor".to_string(), serde_json::to_value(actor)?);
|
||||
}
|
||||
map.insert("properties".to_string(), self.body.properties_value()?);
|
||||
Ok(Value::Object(map))
|
||||
}
|
||||
|
|
@ -697,8 +792,13 @@ mod tests {
|
|||
run_id: fixtures::RUN_1,
|
||||
node_id: Some("build".to_string()),
|
||||
node_label: Some("Build".to_string()),
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body: EventBody::StageCompleted(StageCompletedProps {
|
||||
index: 1,
|
||||
duration_ms: 1234,
|
||||
|
|
@ -873,4 +973,90 @@ mod tests {
|
|||
assert_eq!(serialized["event"], value["event"]);
|
||||
assert_eq!(serialized["properties"], value["properties"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn run_event_round_trips_new_envelope_fields() {
|
||||
let value = json!({
|
||||
"id": "evt_envelope",
|
||||
"ts": "2026-04-08T16:21:11.106Z",
|
||||
"run_id": fixtures::RUN_1,
|
||||
"event": "agent.tool.completed",
|
||||
"stage_id": "code@1",
|
||||
"node_id": "code",
|
||||
"node_label": "Code",
|
||||
"parallel_group_id": "code@1",
|
||||
"parallel_branch_id": "code@1:0",
|
||||
"session_id": "ses_child",
|
||||
"parent_session_id": "ses_parent",
|
||||
"tool_call_id": "call_1",
|
||||
"actor": {
|
||||
"kind": "agent",
|
||||
"id": "ses_child",
|
||||
"display": "claude-sonnet"
|
||||
},
|
||||
"properties": {
|
||||
"tool_name": "read_file",
|
||||
"tool_call_id": "call_1",
|
||||
"output": {"summary": "read"},
|
||||
"is_error": false,
|
||||
"visit": 1
|
||||
}
|
||||
});
|
||||
|
||||
let parsed = RunEvent::from_value(value.clone()).unwrap();
|
||||
assert_eq!(parsed.stage_id.as_deref(), Some("code@1"));
|
||||
assert_eq!(parsed.parallel_group_id.as_deref(), Some("code@1"));
|
||||
assert_eq!(parsed.parallel_branch_id.as_deref(), Some("code@1:0"));
|
||||
assert_eq!(parsed.tool_call_id.as_deref(), Some("call_1"));
|
||||
let actor = parsed.actor.as_ref().expect("actor present");
|
||||
assert_eq!(actor.kind, ActorKind::Agent);
|
||||
assert_eq!(actor.id.as_deref(), Some("ses_child"));
|
||||
assert_eq!(actor.display.as_deref(), Some("claude-sonnet"));
|
||||
|
||||
let serialized = parsed.to_value().unwrap();
|
||||
assert_eq!(serialized["stage_id"], value["stage_id"]);
|
||||
assert_eq!(serialized["parallel_group_id"], value["parallel_group_id"]);
|
||||
assert_eq!(
|
||||
serialized["parallel_branch_id"],
|
||||
value["parallel_branch_id"]
|
||||
);
|
||||
assert_eq!(serialized["tool_call_id"], value["tool_call_id"]);
|
||||
assert_eq!(serialized["actor"], value["actor"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn run_event_omits_absent_envelope_fields() {
|
||||
let event = RunEvent {
|
||||
id: "evt_bare".to_string(),
|
||||
ts: DateTime::parse_from_rfc3339("2026-04-04T12:00:00.000Z")
|
||||
.unwrap()
|
||||
.with_timezone(&Utc),
|
||||
run_id: fixtures::RUN_1,
|
||||
node_id: None,
|
||||
node_label: None,
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body: EventBody::RunStarted(RunStartedProps {
|
||||
name: "demo".to_string(),
|
||||
base_branch: None,
|
||||
base_sha: None,
|
||||
run_branch: None,
|
||||
worktree_dir: None,
|
||||
goal: None,
|
||||
}),
|
||||
};
|
||||
|
||||
let serialized = event.to_value().unwrap();
|
||||
let obj = serialized.as_object().unwrap();
|
||||
assert!(!obj.contains_key("stage_id"));
|
||||
assert!(!obj.contains_key("parallel_group_id"));
|
||||
assert!(!obj.contains_key("parallel_branch_id"));
|
||||
assert!(!obj.contains_key("tool_call_id"));
|
||||
assert!(!obj.contains_key("actor"));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1248,12 +1248,17 @@ pub fn event_name(event: &Event) -> &'static str {
|
|||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
#[derive(Debug, Default)]
|
||||
struct StoredEventFields {
|
||||
session_id: Option<String>,
|
||||
parent_session_id: Option<String>,
|
||||
node_id: Option<String>,
|
||||
node_label: Option<String>,
|
||||
stage_id: Option<String>,
|
||||
parallel_group_id: Option<String>,
|
||||
parallel_branch_id: Option<String>,
|
||||
tool_call_id: Option<String>,
|
||||
actor: Option<fabro_types::ActorRef>,
|
||||
}
|
||||
|
||||
fn default_node_label(node_id: Option<&String>, node_label: Option<String>) -> Option<String> {
|
||||
|
|
@ -1285,10 +1290,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields {
|
|||
let node_id = Some(node_id.clone());
|
||||
let node_label = default_node_label(node_id.as_ref(), Some(name.clone()));
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::CheckpointCompleted { node_id, .. }
|
||||
|
|
@ -1306,10 +1310,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields {
|
|||
let node_id = Some(node_id.clone());
|
||||
let node_label = default_node_label(node_id.as_ref(), None);
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::Agent {
|
||||
|
|
@ -1325,15 +1328,15 @@ fn stored_event_fields(event: &Event) -> StoredEventFields {
|
|||
parent_session_id: parent_session_id.clone(),
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::GitCommit { node_id, .. } => {
|
||||
let node_label = default_node_label(node_id.as_ref(), None);
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id: node_id.clone(),
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::ParallelBranchStarted { branch, .. }
|
||||
|
|
@ -1341,10 +1344,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields {
|
|||
let node_id = Some(branch.clone());
|
||||
let node_label = default_node_label(node_id.as_ref(), None);
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::Prompt { stage, .. }
|
||||
|
|
@ -1355,28 +1357,21 @@ fn stored_event_fields(event: &Event) -> StoredEventFields {
|
|||
let node_id = Some(stage.clone());
|
||||
let node_label = default_node_label(node_id.as_ref(), None);
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
Event::StallWatchdogTimeout { node, .. } => {
|
||||
let node_id = Some(node.clone());
|
||||
let node_label = default_node_label(node_id.as_ref(), None);
|
||||
StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id,
|
||||
node_label,
|
||||
..StoredEventFields::default()
|
||||
}
|
||||
}
|
||||
_ => StoredEventFields {
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
node_id: None,
|
||||
node_label: None,
|
||||
},
|
||||
_ => StoredEventFields::default(),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2394,8 +2389,13 @@ pub fn to_run_event_at(run_id: &RunId, event: &Event, ts: chrono::DateTime<Utc>)
|
|||
run_id: *run_id,
|
||||
node_id: fields.node_id,
|
||||
node_label: fields.node_label,
|
||||
stage_id: fields.stage_id,
|
||||
parallel_group_id: fields.parallel_group_id,
|
||||
parallel_branch_id: fields.parallel_branch_id,
|
||||
session_id: fields.session_id,
|
||||
parent_session_id: fields.parent_session_id,
|
||||
tool_call_id: fields.tool_call_id,
|
||||
actor: fields.actor,
|
||||
body,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -193,8 +193,13 @@ mod tests {
|
|||
run_id: fixtures::RUN_1,
|
||||
node_id: None,
|
||||
node_label: None,
|
||||
stage_id: None,
|
||||
parallel_group_id: None,
|
||||
parallel_branch_id: None,
|
||||
session_id: None,
|
||||
parent_session_id: None,
|
||||
tool_call_id: None,
|
||||
actor: None,
|
||||
body: EventBody::RunSubmitted(RunSubmittedProps {
|
||||
reason: None,
|
||||
definition_blob: None,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue