use serde::{Deserialize, Serialize};
use std::borrow::Borrow;
use std::fmt;
use std::str::FromStr;
use uuid::Uuid;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct RunId(Uuid);
impl RunId {
pub fn new() -> Self {
Self(Uuid::new_v4())
}
pub fn as_uuid(&self) -> &Uuid {
&self.0
}
pub(crate) fn for_external_delivery(
mob_id: &MobId,
flow_id: &FlowId,
idempotency_key: &str,
) -> Self {
Self(Uuid::new_v5(
&Uuid::NAMESPACE_URL,
format!("meerkat:mob:{mob_id}:flow:{flow_id}:delivery:{idempotency_key}").as_bytes(),
))
}
}
impl Default for RunId {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for RunId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl FromStr for RunId {
type Err = uuid::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Ok(Self(Uuid::parse_str(s)?))
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct PlacedSpawnId(Uuid);
impl PlacedSpawnId {
pub fn new() -> Self {
Self(Uuid::new_v4())
}
pub fn as_uuid(&self) -> &Uuid {
&self.0
}
}
impl Default for PlacedSpawnId {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for PlacedSpawnId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl FromStr for PlacedSpawnId {
type Err = uuid::Error;
fn from_str(value: &str) -> Result<Self, Self::Err> {
Ok(Self(Uuid::parse_str(value)?))
}
}
macro_rules! string_newtype {
($(#[$meta:meta])* $name:ident) => {
$(#[$meta])*
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct $name(String);
impl $name {
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for $name {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl From<String> for $name {
fn from(value: String) -> Self {
Self(value)
}
}
impl From<&str> for $name {
fn from(value: &str) -> Self {
Self(value.to_owned())
}
}
impl Borrow<str> for $name {
fn borrow(&self) -> &str {
&self.0
}
}
impl Borrow<String> for $name {
fn borrow(&self) -> &String {
&self.0
}
}
impl AsRef<str> for $name {
fn as_ref(&self) -> &str {
&self.0
}
}
impl PartialEq<String> for $name {
fn eq(&self, other: &String) -> bool {
&self.0 == other
}
}
impl PartialEq<&String> for $name {
fn eq(&self, other: &&String) -> bool {
&self.0 == *other
}
}
impl PartialEq<str> for $name {
fn eq(&self, other: &str) -> bool {
self.0.as_str() == other
}
}
impl PartialEq<&str> for $name {
fn eq(&self, other: &&str) -> bool {
self.0.as_str() == *other
}
}
};
}
string_newtype!(
MobId
);
string_newtype!(
FlowId
);
string_newtype!(
StepId
);
string_newtype!(
BranchId
);
string_newtype!(
ProfileName
);
string_newtype!(
FrameId
);
string_newtype!(
LoopInstanceId
);
string_newtype!(
FlowNodeId
);
string_newtype!(
LoopId
);
string_newtype!(
AgentIdentity
);
string_newtype!(
RespawnTopologyPeerId
);
impl AgentIdentity {
pub fn objective_lead_for_session(session_id: &meerkat_core::SessionId) -> Self {
Self::from(format!("objective-lead-session:{session_id}"))
}
pub(crate) fn is_system_reserved(&self) -> bool {
self.as_str()
.starts_with(crate::runtime::FLOW_MEMBER_ID_PREFIX)
}
pub(crate) fn is_flow_member_namespace(&self) -> bool {
self.is_system_reserved()
}
pub(crate) fn flow_system_provenance() -> Self {
Self::from(crate::run::FLOW_RUN_PROVENANCE_AGENT_ID)
}
pub(crate) fn is_flow_system_provenance(&self) -> bool {
self.as_str() == crate::run::FLOW_RUN_PROVENANCE_AGENT_ID
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct Generation(u64);
impl Generation {
pub const INITIAL: Self = Self(0);
pub const fn new(value: u64) -> Self {
Self(value)
}
pub const fn get(self) -> u64 {
self.0
}
pub const fn next(self) -> Option<Self> {
match self.0.checked_add(1) {
Some(value) => Some(Self(value)),
None => None,
}
}
}
impl fmt::Display for Generation {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub struct AgentRuntimeId {
pub identity: AgentIdentity,
pub generation: Generation,
}
impl AgentRuntimeId {
pub fn new(identity: AgentIdentity, generation: Generation) -> Self {
Self {
identity,
generation,
}
}
pub fn initial(identity: AgentIdentity) -> Self {
Self {
identity,
generation: Generation::INITIAL,
}
}
}
impl fmt::Display for AgentRuntimeId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}:{}", self.identity, self.generation.get())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct FenceToken(u64);
impl FenceToken {
pub const fn new(value: u64) -> Self {
Self(value)
}
pub const fn get(self) -> u64 {
self.0
}
}
impl fmt::Display for FenceToken {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "fence:{}", self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct WorkRef(Uuid);
impl WorkRef {
pub fn new() -> Self {
Self(Uuid::new_v4())
}
pub fn as_uuid(&self) -> &Uuid {
&self.0
}
pub(crate) fn for_external_delivery(
mob_id: &MobId,
member_id: &AgentIdentity,
idempotency_key: &str,
) -> Self {
Self(Uuid::new_v5(
&Uuid::NAMESPACE_URL,
format!("meerkat:mob:{mob_id}:member:{member_id}:delivery:{idempotency_key}")
.as_bytes(),
))
}
}
impl Default for WorkRef {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for WorkRef {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl FromStr for WorkRef {
type Err = uuid::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Ok(Self(Uuid::parse_str(s)?))
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkSpec {
pub content: meerkat_core::types::ContentInput,
pub origin: WorkOrigin,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_prompt: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub injected_context: Vec<meerkat_core::types::ContentInput>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transient_turn_context: Option<meerkat_core::lifecycle::run_primitive::TurnRequestContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub interaction_id: Option<meerkat_core::interaction::InteractionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub objective_id: Option<meerkat_core::interaction::ObjectiveId>,
}
impl WorkSpec {
pub fn new(content: impl Into<meerkat_core::types::ContentInput>, origin: WorkOrigin) -> Self {
Self {
content: content.into(),
origin,
system_prompt: None,
injected_context: Vec::new(),
transient_turn_context: None,
interaction_id: None,
objective_id: None,
}
}
#[must_use]
pub fn with_system_prompt(mut self, system_prompt: impl Into<String>) -> Self {
self.system_prompt = Some(system_prompt.into());
self
}
#[must_use]
pub fn with_objective_id(
mut self,
objective_id: meerkat_core::interaction::ObjectiveId,
) -> Self {
self.objective_id = Some(objective_id);
self
}
pub fn with_injected_context(
mut self,
injected_context: Vec<meerkat_core::types::ContentInput>,
) -> Self {
self.injected_context = injected_context;
self
}
#[must_use]
pub fn with_transient_turn_context(
mut self,
context: meerkat_core::lifecycle::run_primitive::TurnRequestContext,
) -> Self {
self.transient_turn_context = Some(context);
self
}
pub fn with_interaction_id(
mut self,
interaction_id: meerkat_core::interaction::InteractionId,
) -> Self {
self.interaction_id = Some(interaction_id);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum WorkOrigin {
External,
Internal,
}
impl WorkOrigin {
pub const fn as_str(self) -> &'static str {
match self {
WorkOrigin::External => "External",
WorkOrigin::Internal => "Internal",
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn work_spec_interaction_id_is_serde_additive() {
let json = serde_json::to_value(WorkSpec::new(
"do the thing".to_string(),
WorkOrigin::Internal,
))
.unwrap();
assert!(
json.get("interaction_id").is_none(),
"None must not serialize (serde-additive)"
);
let spec: WorkSpec = serde_json::from_value(json).unwrap();
assert!(spec.interaction_id.is_none());
let id = meerkat_core::interaction::InteractionId(Uuid::new_v4());
let keyed =
WorkSpec::new("do the thing".to_string(), WorkOrigin::External).with_interaction_id(id);
let json = serde_json::to_value(&keyed).unwrap();
let parsed: WorkSpec = serde_json::from_value(json).unwrap();
assert_eq!(parsed.interaction_id, Some(id));
}
#[test]
fn work_spec_system_prompt_is_per_turn_and_serde_additive() {
let absent = WorkSpec::new("first turn", WorkOrigin::Internal);
let absent_json = serde_json::to_value(&absent).unwrap();
assert!(absent_json.get("system_prompt").is_none());
let present = WorkSpec::new("second turn", WorkOrigin::Internal)
.with_system_prompt("updated instructions");
let json = serde_json::to_value(&present).unwrap();
assert_eq!(json["system_prompt"], "updated instructions");
let decoded: WorkSpec = serde_json::from_value(json).unwrap();
assert_eq!(
decoded.system_prompt.as_deref(),
Some("updated instructions")
);
}
#[test]
fn test_run_id_roundtrip_json() {
let run_id = RunId::new();
let encoded = serde_json::to_string(&run_id).unwrap();
let decoded: RunId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, run_id);
}
#[test]
fn test_run_id_roundtrip_parse_display() {
let run_id = RunId::new();
let rendered = run_id.to_string();
let reparsed = RunId::from_str(&rendered).unwrap();
assert_eq!(reparsed, run_id);
}
#[test]
fn external_delivery_run_id_is_stable_and_flow_scoped() {
let mob_id = MobId::from("ops");
let key = "schedule:one";
let first = RunId::for_external_delivery(&mob_id, &FlowId::from("alpha"), key);
assert_eq!(
first,
RunId::for_external_delivery(&mob_id, &FlowId::from("alpha"), key)
);
assert_ne!(
first,
RunId::for_external_delivery(&mob_id, &FlowId::from("beta"), key)
);
}
#[test]
fn test_flow_id_roundtrip_json() {
let id = FlowId::from("flow-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: FlowId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_step_id_roundtrip_json() {
let id = StepId::from("step-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: StepId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_branch_id_roundtrip_json() {
let id = BranchId::from("branch-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: BranchId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_frame_id_roundtrip_json() {
let id = FrameId::from("frame-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: FrameId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_loop_instance_id_roundtrip_json() {
let id = LoopInstanceId::from("loop-instance-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: LoopInstanceId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_flow_node_id_roundtrip_json() {
let id = FlowNodeId::from("node-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: FlowNodeId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_loop_id_roundtrip_json() {
let id = LoopId::from("loop-a");
let encoded = serde_json::to_string(&id).unwrap();
let decoded: LoopId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_existing_ids_roundtrip() {
let mob = MobId::from("mob-a");
let profile = ProfileName::from("lead");
assert_eq!(
serde_json::from_str::<MobId>(&serde_json::to_string(&mob).unwrap()).unwrap(),
mob
);
assert_eq!(
serde_json::from_str::<ProfileName>(&serde_json::to_string(&profile).unwrap()).unwrap(),
profile
);
}
#[test]
fn test_agent_identity_roundtrip_json() {
let id = AgentIdentity::from("researcher");
let encoded = serde_json::to_string(&id).unwrap();
assert_eq!(encoded, "\"researcher\"");
let decoded: AgentIdentity = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, id);
}
#[test]
fn test_agent_identity_display() {
let id = AgentIdentity::from("lead-agent");
assert_eq!(id.to_string(), "lead-agent");
assert_eq!(id.as_str(), "lead-agent");
}
#[test]
fn test_generation_roundtrip_json() {
let generation = Generation::new(42);
let encoded = serde_json::to_string(&generation).unwrap();
assert_eq!(encoded, "42");
let decoded: Generation = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, generation);
}
#[test]
fn test_generation_initial_and_next() {
assert_eq!(Generation::INITIAL.get(), 0);
assert_eq!(Generation::INITIAL.next().expect("generation 1").get(), 1);
assert_eq!(Generation::new(5).next().expect("generation 6").get(), 6);
assert_eq!(Generation::new(u64::MAX).next(), None);
}
#[test]
fn test_generation_ordering() {
assert!(Generation::new(0) < Generation::new(1));
assert!(Generation::new(1) < Generation::new(100));
}
#[test]
fn test_agent_runtime_id_roundtrip_json() {
let rid = AgentRuntimeId::new(AgentIdentity::from("worker"), Generation::new(3));
let encoded = serde_json::to_string(&rid).unwrap();
let decoded: AgentRuntimeId = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, rid);
}
#[test]
fn test_agent_runtime_id_initial() {
let rid = AgentRuntimeId::initial(AgentIdentity::from("worker"));
assert_eq!(rid.identity, AgentIdentity::from("worker"));
assert_eq!(rid.generation, Generation::INITIAL);
}
#[test]
fn test_agent_runtime_id_display() {
let rid = AgentRuntimeId::new(AgentIdentity::from("coder"), Generation::new(2));
assert_eq!(rid.to_string(), "coder:2");
}
#[test]
fn test_fence_token_roundtrip_json() {
let ft = FenceToken::new(99);
let encoded = serde_json::to_string(&ft).unwrap();
assert_eq!(encoded, "99");
let decoded: FenceToken = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, ft);
}
#[test]
fn test_fence_token_display() {
assert_eq!(FenceToken::new(7).to_string(), "fence:7");
}
#[test]
fn test_fence_token_ordering() {
assert!(FenceToken::new(1) < FenceToken::new(2));
}
#[test]
fn test_work_ref_roundtrip_json() {
let wr = WorkRef::new();
let encoded = serde_json::to_string(&wr).unwrap();
let decoded: WorkRef = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, wr);
}
#[test]
fn test_work_ref_roundtrip_parse_display() {
let wr = WorkRef::new();
let rendered = wr.to_string();
let reparsed = WorkRef::from_str(&rendered).unwrap();
assert_eq!(reparsed, wr);
}
#[test]
fn external_delivery_work_ref_is_stable_and_member_scoped() {
let mob_id = MobId::from("ops");
let key = "schedule:one";
let first = WorkRef::for_external_delivery(&mob_id, &AgentIdentity::from("alpha"), key);
assert_eq!(
first,
WorkRef::for_external_delivery(&mob_id, &AgentIdentity::from("alpha"), key)
);
assert_ne!(
first,
WorkRef::for_external_delivery(&mob_id, &AgentIdentity::from("beta"), key)
);
}
#[test]
fn test_work_spec_roundtrip_json() {
let spec = WorkSpec::new("do something".to_owned(), WorkOrigin::External);
let encoded = serde_json::to_string(&spec).unwrap();
let decoded: WorkSpec = serde_json::from_str(&encoded).unwrap();
assert_eq!(
decoded.content,
meerkat_core::types::ContentInput::from("do something".to_string()),
);
assert_eq!(decoded.origin, WorkOrigin::External);
}
#[test]
fn test_work_origin_variants_roundtrip_json() {
for origin in [WorkOrigin::External, WorkOrigin::Internal] {
let encoded = serde_json::to_string(&origin).unwrap();
let decoded: WorkOrigin = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, origin);
}
}
#[test]
fn test_work_spec_injected_context_serde_default_and_omission() {
let decoded: WorkSpec =
serde_json::from_str(r#"{"content":"do something","origin":"External"}"#).unwrap();
assert!(decoded.injected_context.is_empty());
let empty = WorkSpec::new("do something".to_owned(), WorkOrigin::External);
let value = serde_json::to_value(&empty).unwrap();
assert!(value.get("injected_context").is_none());
let spec = WorkSpec::new("do something".to_owned(), WorkOrigin::External)
.with_injected_context(vec![
meerkat_core::types::ContentInput::Text("ambient alpha".to_string()),
meerkat_core::types::ContentInput::Text("ambient beta".to_string()),
]);
let encoded = serde_json::to_string(&spec).unwrap();
let decoded: WorkSpec = serde_json::from_str(&encoded).unwrap();
assert_eq!(
decoded.injected_context,
vec![
meerkat_core::types::ContentInput::Text("ambient alpha".to_string()),
meerkat_core::types::ContentInput::Text("ambient beta".to_string()),
],
);
}
#[test]
fn test_work_spec_internal_origin() {
let spec = WorkSpec::new("coordinate".to_owned(), WorkOrigin::Internal);
assert_eq!(spec.origin, WorkOrigin::Internal);
assert_eq!(
spec.content,
meerkat_core::types::ContentInput::from("coordinate".to_string()),
);
}
#[test]
fn test_work_spec_accepts_multimodal_content() {
let image_block = meerkat_core::types::ContentBlock::Image {
media_type: "image/png".to_string(),
data: meerkat_core::ImageData::Inline {
data: "iVBORw0KGgo=".to_string(),
},
};
let content = meerkat_core::types::ContentInput::Blocks(vec![
meerkat_core::types::ContentBlock::Text {
text: "analyse this".to_string(),
},
image_block.clone(),
]);
let spec = WorkSpec::new(content.clone(), WorkOrigin::External);
assert_eq!(spec.content, content);
}
}