use a2a_rs::domain::{
Artifact, Message, Part, Role, Task as WireTask, TaskArtifactUpdateEvent, TaskState,
TaskStatus, TaskStatusUpdateEvent,
};
use buffa::MessageField;
use buffa_types::google::protobuf::{Struct, Timestamp};
use serde_json::{Value, json};
use crate::a2a::tasks::{State, Task};
pub fn stamp(ms: u64) -> Timestamp {
Timestamp {
seconds: (ms / 1000) as i64,
nanos: ((ms % 1000) * 1_000_000) as i32,
..Default::default()
}
}
pub fn timestamp_string(ms: u64) -> String {
serde_json::to_value(stamp(ms))
.ok()
.and_then(|v| v.as_str().map(str::to_string))
.unwrap_or_default()
}
impl State {
pub fn to_wire(self) -> TaskState {
match self {
State::Submitted => TaskState::TASK_STATE_SUBMITTED,
State::Working => TaskState::TASK_STATE_WORKING,
State::InputRequired => TaskState::TASK_STATE_INPUT_REQUIRED,
State::Completed => TaskState::TASK_STATE_COMPLETED,
State::Failed => TaskState::TASK_STATE_FAILED,
State::Canceled => TaskState::TASK_STATE_CANCELED,
State::Rejected => TaskState::TASK_STATE_REJECTED,
}
}
}
fn metadata(v: Value) -> MessageField<Struct> {
match serde_json::from_value::<Struct>(v) {
Ok(s) => MessageField::some(s),
Err(_) => MessageField::none(),
}
}
pub fn agent_message(task_id: &str, context_id: &str, text: &str) -> Message {
let mut m = Message::agent_text(text.to_string(), format!("{task_id}.status"));
m.task_id = task_id.to_string();
m.context_id = context_id.to_string();
m
}
fn status_of(t: &Task) -> TaskStatus {
let message = t
.message
.as_deref()
.map(|m| agent_message(&t.id, &t.context_id, m));
let mut s = TaskStatus::new(t.state.to_wire(), message);
s.timestamp = MessageField::some(stamp(t.updated));
s
}
fn agentd_metadata(t: &Task) -> Value {
let mut m = json!({
"agentd/link": t.link,
"agentd/created": timestamp_string(t.created),
});
if let Some(p) = &t.principal {
m["agentd/principal"] = json!(p);
}
if let Some(sch) = &t.ask_schema {
m["agentd/ask_schema"] = sch.clone();
}
if !t.history.is_empty() {
let history: Vec<Value> = t
.history
.iter()
.map(|h| {
let mut h = h.clone();
if let Some(ms) = h.get("ts").and_then(Value::as_u64) {
h["ts"] = json!(timestamp_string(ms));
}
h
})
.collect();
m["agentd/statusHistory"] = json!(history);
}
m
}
fn artifacts_of(t: &Task) -> Vec<Artifact> {
let mut out = Vec::new();
if let Some(r) = &t.result {
let text = match r {
Value::String(s) => s.clone(),
other => other.to_string(),
};
out.push(Artifact {
artifact_id: format!("{}.result", t.id),
parts: vec![Part::text(text)],
..Default::default()
});
}
for a in &t.artifacts {
out.push(Artifact {
artifact_id: a.clone(),
..Default::default()
});
}
out
}
pub fn result_artifact(t: &Task) -> Option<Artifact> {
artifacts_of(t)
.into_iter()
.next()
.filter(|a| !a.parts.is_empty())
}
pub fn task(t: &Task) -> WireTask {
let mut w = WireTask::new(t.id.clone(), t.context_id.clone());
w.status = MessageField::some(status_of(t));
w.artifacts = artifacts_of(t);
w.metadata = metadata(agentd_metadata(t));
w
}
pub fn task_summary(t: &Task) -> WireTask {
let mut w = WireTask::new(t.id.clone(), t.context_id.clone());
w.status = MessageField::some(status_of(t));
w.metadata = metadata(agentd_metadata(t));
w
}
pub fn status_event(
task_id: &str,
context_id: &str,
state: TaskState,
message: Option<&str>,
at_ms: u64,
) -> TaskStatusUpdateEvent {
let mut s = TaskStatus::new(
state,
message.map(|m| agent_message(task_id, context_id, m)),
);
s.timestamp = MessageField::some(stamp(at_ms));
TaskStatusUpdateEvent {
task_id: task_id.to_string(),
context_id: context_id.to_string(),
kind: "status-update".to_string(),
status: s,
metadata: None,
}
}
pub fn artifact_event(
task_id: &str,
context_id: &str,
artifact: Artifact,
last_chunk: bool,
) -> TaskArtifactUpdateEvent {
TaskArtifactUpdateEvent {
task_id: task_id.to_string(),
context_id: context_id.to_string(),
kind: "artifact-update".to_string(),
artifact,
append: None,
last_chunk: Some(last_chunk),
metadata: None,
}
}
pub fn message_text(m: &Message) -> String {
let mut out = String::new();
for p in &m.parts {
if let Some(a2a_rs::domain::part::Content::Text(t)) = &p.content {
if !out.is_empty() {
out.push('\n');
}
out.push_str(t);
}
}
out
}
pub fn command(m: &Message) -> Option<(String, Value)> {
for p in &m.parts {
let Some(a2a_rs::domain::part::Content::Data(d)) = &p.content else {
continue;
};
let Ok(v) = serde_json::to_value(d) else {
continue;
};
let inner = v.get("data").unwrap_or(&v);
let Some(env) = inner.get("agentd") else {
continue;
};
if let Some(op) = env.get("op").and_then(Value::as_str) {
return Some((op.to_string(), env.clone()));
}
}
None
}
pub fn is_from_caller(m: &Message) -> bool {
m.role.as_known() != Some(Role::ROLE_AGENT)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::a2a::tasks::Link;
#[test]
fn the_projection_is_proto3_json() {
let mut t = Task::new(
"task-1",
"ctx-1",
Some("user:a"),
Link::Run { id: "r1".into() },
);
t.set_result(json!("the answer"));
t.transition(State::Completed, Some("done".into()));
let v = serde_json::to_value(task(&t)).expect("serialize");
assert_eq!(v["id"], "task-1");
assert_eq!(v["contextId"], "ctx-1");
assert_eq!(v["status"]["state"], "TASK_STATE_COMPLETED");
assert_eq!(v["status"]["message"]["role"], "ROLE_AGENT");
assert_eq!(v["status"]["message"]["taskId"], "task-1");
assert!(
v["status"]["timestamp"]
.as_str()
.is_some_and(|s| s.ends_with('Z')),
"timestamps are RFC 3339: {v}"
);
assert_eq!(v["artifacts"][0]["artifactId"], "task-1.result");
assert_eq!(v["artifacts"][0]["parts"][0]["text"], "the answer");
assert_eq!(v["metadata"]["agentd/principal"], "user:a");
assert!(v["history"].is_null(), "history is repeated Message: {v}");
let s = serde_json::to_value(task_summary(&t)).expect("serialize");
assert_eq!(s["status"]["state"], v["status"]["state"]);
assert!(s["state"].is_null());
assert!(s["artifacts"].is_null());
}
#[test]
fn a_task_we_emit_is_a_task_we_can_read_back() {
let t = Task::new("t", "c", None, Link::Turn { ctx: "c".into() });
let v = serde_json::to_value(task(&t)).unwrap();
let back: WireTask = serde_json::from_value(v).expect("round trip through their type");
assert_eq!(back.id, "t");
assert_eq!(
back.status.as_option().unwrap().state.as_known(),
Some(TaskState::TASK_STATE_SUBMITTED)
);
}
#[test]
fn text_and_commands_are_read_from_their_parts() {
let mut m = Message::user_text("please".into(), "m1".into());
assert_eq!(message_text(&m), "please");
assert!(command(&m).is_none());
assert!(is_from_caller(&m));
m.parts.push(Part::data(
serde_json::from_value(json!({"agentd": {"op": "status"}})).unwrap(),
));
let (op, env) = command(&m).expect("a command DataPart");
assert_eq!(op, "status");
assert_eq!(env["op"], "status");
assert_eq!(message_text(&m), "please");
}
}