use a2a_rs::domain::{Message, Part, Task, TaskState};
use serde_json::{Value, json};
pub fn send_message_params(
objective: &str,
output_contract: Option<&str>,
message_id: &str,
) -> Value {
send_message_params_cmd(objective, None, output_contract, message_id)
}
pub fn send_message_params_cmd(
objective: &str,
command: Option<&Value>,
output_contract: Option<&str>,
message_id: &str,
) -> Value {
let mut m = Message::user_text(objective.to_string(), message_id.to_string());
if objective.trim().is_empty() {
m.parts.clear();
}
if let Some(env) = command
&& let Ok(d) = serde_json::from_value(json!({ "agentd": env }))
{
m.parts.push(Part::data(d));
}
if let Some(contract) = output_contract.filter(|c| !c.is_empty()) {
m.parts
.push(Part::text(format!("Required output: {contract}")));
}
json!({ "message": serde_json::to_value(&m).unwrap_or(Value::Null) })
}
pub fn task_of(v: &Value) -> Option<Task> {
let body = match v.get("task") {
Some(t) => t,
None => v,
};
serde_json::from_value(body.clone()).ok()
}
pub fn task_id_of(v: &Value) -> String {
task_of(v).map(|t| t.id).unwrap_or_default()
}
pub fn task_state_of(v: &Value) -> TaskState {
task_of(v)
.and_then(|t| t.status.as_option().and_then(|s| s.state.as_known()))
.unwrap_or(TaskState::TASK_STATE_UNSPECIFIED)
}
pub fn is_terminal(state: TaskState) -> bool {
matches!(
state,
TaskState::TASK_STATE_COMPLETED
| TaskState::TASK_STATE_FAILED
| TaskState::TASK_STATE_CANCELED
| TaskState::TASK_STATE_REJECTED
)
}
pub fn artifact_text_of(v: &Value) -> String {
let Some(task) = task_of(v) else {
return String::new();
};
let mut out = String::new();
for artifact in &task.artifacts {
for part in &artifact.parts {
if let Some(a2a_rs::domain::part::Content::Text(t)) = &part.content {
if !out.is_empty() {
out.push('\n');
}
out.push_str(t);
}
}
}
out
}
pub fn describe(state: TaskState) -> &'static str {
match state {
TaskState::TASK_STATE_COMPLETED => "completed",
TaskState::TASK_STATE_FAILED => "failed",
TaskState::TASK_STATE_CANCELED => "canceled",
TaskState::TASK_STATE_REJECTED => "rejected",
TaskState::TASK_STATE_INPUT_REQUIRED => "input-required",
TaskState::TASK_STATE_AUTH_REQUIRED => "auth-required",
TaskState::TASK_STATE_WORKING => "working",
TaskState::TASK_STATE_SUBMITTED => "submitted",
_ => "unspecified",
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_objective_goes_out_as_a_user_message() {
let p = send_message_params("summarise the incident", Some("one paragraph"), "m1");
let m = &p["message"];
assert_eq!(m["role"], "ROLE_USER");
assert_eq!(m["messageId"], "m1");
assert_eq!(m["parts"][0]["text"], "summarise the incident");
assert_eq!(m["parts"][1]["text"], "Required output: one paragraph");
let back: Message = serde_json::from_value(m.clone()).expect("their Message");
assert_eq!(back.parts.len(), 2);
}
#[test]
fn a_peers_reply_is_read_with_the_specs_type() {
let reply = json!({"task": {
"id": "t-9",
"contextId": "c-1",
"status": {"state": "TASK_STATE_COMPLETED", "timestamp": "2026-08-17T14:00:00Z"},
"artifacts": [
{"artifactId": "a1", "parts": [{"text": "first"}]},
{"artifactId": "a2", "parts": [{"text": "second"}]}
]
}});
assert_eq!(task_id_of(&reply), "t-9");
assert_eq!(task_state_of(&reply), TaskState::TASK_STATE_COMPLETED);
assert!(is_terminal(task_state_of(&reply)));
assert_eq!(artifact_text_of(&reply), "first\nsecond");
}
#[test]
fn an_unfinished_or_unreadable_reply_is_not_mistaken_for_an_answer() {
let working = json!({"task": {"id": "t", "contextId": "c", "status": {"state": "TASK_STATE_WORKING"}}});
assert!(!is_terminal(task_state_of(&working)));
assert_eq!(artifact_text_of(&working), "");
let garbage = json!({"nope": true});
assert_eq!(task_state_of(&garbage), TaskState::TASK_STATE_UNSPECIFIED);
assert!(!is_terminal(task_state_of(&garbage)));
}
}