use std::collections::{BTreeMap, BTreeSet};
use std::fmt;
use meerkat_contracts::wire::{
PortableDefinitionExtract, PortableProfile, PortableSystemPrompt, WireAuthBindingRef,
WireMobRuntimeMode, WireOpaqueJson, WireResolvedToolAccessPolicy, WireTrustedPeerIdentity,
};
use meerkat_core::lifecycle::InputId;
use meerkat_core::ops::OperationId;
use meerkat_core::{
BudgetLimits, ContentInput, Session, SessionGeneration, SessionId, SessionLineageId, ToolName,
};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::ids::{AgentIdentity, MobId, ProfileName};
pub const IDENTITY_INTENT_SCHEMA_VERSION: u32 = 1;
pub const IDENTITY_LEASE_SCHEMA_VERSION: u32 = 1;
pub const IDENTITY_OPERATION_RECEIPT_SCHEMA_VERSION: u32 = 1;
pub const IDENTITY_INTENT_MAX_ENCODED_BYTES: usize = 4 * 1024 * 1024;
pub const IDENTITY_LEASE_MAX_TTL_MS: u64 = 30_000;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct IdentityDeclarationScopeId(String);
impl IdentityDeclarationScopeId {
pub fn new(value: impl Into<String>) -> Result<Self, IdentityIntentError> {
let value = value.into();
validate_text("identity_declaration_scope", &value)?;
Ok(Self(value))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredSessionTarget {
pub session_id: SessionId,
pub lineage_id: SessionLineageId,
pub lineage_generation: SessionGeneration,
pub authority_policy: DesiredSessionAuthorityPolicy,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentitySessionStoreAuthority {
session_id: SessionId,
store_revision: u64,
token: IdentitySessionStoreAuthorityToken,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "profile", rename_all = "snake_case", deny_unknown_fields)]
enum IdentitySessionStoreAuthorityToken {
WholeBlobV1 { blob_sha256: String },
HeadCanonicalV1 { committed_head_token: String },
}
impl IdentitySessionStoreAuthority {
pub(crate) fn from_runtime_authority(
authority: meerkat_runtime::RuntimeSessionAuthority,
) -> Self {
match authority {
meerkat_runtime::RuntimeSessionAuthority::WholeBlob(authority) => Self {
session_id: authority.session_id().clone(),
store_revision: authority.store_revision(),
token: IdentitySessionStoreAuthorityToken::WholeBlobV1 {
blob_sha256: authority.blob_sha256().to_string(),
},
},
meerkat_runtime::RuntimeSessionAuthority::HeadCanonical(authority) => Self {
session_id: authority.session_id().clone(),
store_revision: authority.store_revision(),
token: IdentitySessionStoreAuthorityToken::HeadCanonicalV1 {
committed_head_token: authority.committed_head_token().to_string(),
},
},
}
}
#[cfg(test)]
pub(crate) fn whole_blob_for_test(
session_id: SessionId,
store_revision: u64,
blob_sha256: impl Into<String>,
) -> Self {
Self {
session_id,
store_revision,
token: IdentitySessionStoreAuthorityToken::WholeBlobV1 {
blob_sha256: blob_sha256.into(),
},
}
}
#[must_use]
pub fn session_id(&self) -> &SessionId {
&self.session_id
}
#[must_use]
pub const fn profile(&self) -> meerkat_runtime::RuntimeSessionPersistenceProfile {
match &self.token {
IdentitySessionStoreAuthorityToken::WholeBlobV1 { .. } => {
meerkat_runtime::RuntimeSessionPersistenceProfile::WholeBlobV1
}
IdentitySessionStoreAuthorityToken::HeadCanonicalV1 { .. } => {
meerkat_runtime::RuntimeSessionPersistenceProfile::HeadCanonicalV1
}
}
}
#[must_use]
pub const fn store_revision(&self) -> u64 {
self.store_revision
}
#[must_use]
pub fn token(&self) -> &str {
match &self.token {
IdentitySessionStoreAuthorityToken::WholeBlobV1 { blob_sha256 } => blob_sha256,
IdentitySessionStoreAuthorityToken::HeadCanonicalV1 {
committed_head_token,
} => committed_head_token,
}
}
pub(crate) fn validate(&self) -> Result<(), IdentityIntentError> {
if self.session_id.0.is_nil() || self.store_revision == 0 {
return Err(IdentityIntentError::InvalidSessionStoreAuthority);
}
let valid_token = match &self.token {
IdentitySessionStoreAuthorityToken::WholeBlobV1 { blob_sha256 } => {
has_prefixed_sha256(blob_sha256, &["row-sha256:"])
}
IdentitySessionStoreAuthorityToken::HeadCanonicalV1 {
committed_head_token,
} => has_prefixed_sha256(committed_head_token, &["head-v5-sha256:"]),
};
if !valid_token {
return Err(IdentityIntentError::InvalidSessionStoreAuthority);
}
Ok(())
}
pub(crate) fn observation_version(&self) -> Result<String, IdentityIntentError> {
#[derive(Serialize)]
struct ObservationVersionMaterial<'a> {
domain: &'static str,
authority: &'a IdentitySessionStoreAuthority,
}
self.validate()?;
let bytes = serde_json::to_vec(&ObservationVersionMaterial {
domain: "meerkat.identity.session_store_authority.v1",
authority: self,
})
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
Ok(sha256_digest(&bytes))
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DesiredSessionAuthorityPolicy {
#[default]
CreateIfAbsent,
RequireExisting,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "execution", rename_all = "snake_case", deny_unknown_fields)]
pub enum DesiredExecution {
ControllingSession,
AnyBoundHostSession,
PlacedSession {
host_id: String,
},
External {
address: DesiredExternalAddress,
identity: WireTrustedPeerIdentity,
},
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(try_from = "String", into = "String")]
pub struct DesiredExternalAddress(String);
impl DesiredExternalAddress {
pub fn parse(value: impl AsRef<str>) -> Result<Self, IdentityIntentError> {
let value = value.as_ref();
let parsed = url::Url::parse(value)
.map_err(|error| IdentityIntentError::InvalidExternalAddress(error.to_string()))?;
if parsed.scheme() != "tcp"
|| !parsed.username().is_empty()
|| parsed.password().is_some()
|| parsed.query().is_some()
|| parsed.fragment().is_some()
|| !(parsed.path().is_empty() || parsed.path() == "/")
{
return Err(IdentityIntentError::InvalidExternalAddress(
"external address must be tcp://host:port with no credentials, path, query, or fragment"
.to_string(),
));
}
let host = parsed.host_str().ok_or_else(|| {
IdentityIntentError::InvalidExternalAddress("external address has no host".to_string())
})?;
let port = parsed.port().ok_or_else(|| {
IdentityIntentError::InvalidExternalAddress("external address has no port".to_string())
})?;
let canonical_host = if host.contains(':') {
format!("[{host}]")
} else {
host.to_string()
};
Ok(Self(format!("tcp://{canonical_host}:{port}")))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl TryFrom<String> for DesiredExternalAddress {
type Error = IdentityIntentError;
fn try_from(value: String) -> Result<Self, Self::Error> {
Self::parse(value)
}
}
impl From<DesiredExternalAddress> for String {
fn from(value: DesiredExternalAddress) -> Self {
value.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredMemberOverlay {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<WireOpaqueJson>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub labels: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
pub system_prompt: PortableSystemPrompt,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tool_access_policy: Option<WireResolvedToolAccessPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub budget_limits: Option<BudgetLimits>,
pub runtime_mode: WireMobRuntimeMode,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredLocalCallbackTool {
pub name: ToolName,
pub description: String,
pub input_schema: serde_json::Value,
}
impl DesiredLocalCallbackTool {
pub fn new(
name: impl Into<ToolName>,
description: impl Into<String>,
input_schema: serde_json::Value,
) -> Result<Self, IdentityIntentError> {
let value = Self {
name: name.into(),
description: description.into(),
input_schema,
};
value.validate()?;
Ok(value)
}
pub fn validate(&self) -> Result<(), IdentityIntentError> {
validate_text("local_callback_tool_name", self.name.as_str())?;
validate_text("local_callback_tool_description", &self.description)?;
if !self.input_schema.is_object() {
return Err(IdentityIntentError::InvalidMemberMaterial(format!(
"local callback tool '{}' input schema must be a JSON object",
self.name
)));
}
jsonschema::validator_for(&self.input_schema).map_err(|error| {
IdentityIntentError::InvalidMemberMaterial(format!(
"local callback tool '{}' input schema is invalid: {error}",
self.name
))
})?;
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredMemberMaterial {
pub profile_name: ProfileName,
pub profile: PortableProfile,
pub definition_extract: PortableDefinitionExtract,
pub overlay: DesiredMemberOverlay,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub required_env_keys: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub required_local_callback_tools: Vec<DesiredLocalCallbackTool>,
pub execution: DesiredExecution,
}
impl DesiredMemberMaterial {
pub fn validate(&self) -> Result<(), IdentityIntentError> {
validate_text("profile_name", self.profile_name.as_str())?;
validate_execution(&self.execution)?;
validate_string_set("required_env_key", &self.required_env_keys)?;
validate_required_local_callback_tools(&self.required_local_callback_tools, true)?;
validate_string_set(
"definition_profile_name",
&self.definition_extract.profile_names,
)?;
if !self
.definition_extract
.profile_names
.iter()
.any(|name| name == self.profile_name.as_str())
{
return Err(IdentityIntentError::InvalidMemberMaterial(
"definition extract does not contain the selected profile name".to_string(),
));
}
if self.profile.runtime_mode != self.overlay.runtime_mode {
return Err(IdentityIntentError::InvalidMemberMaterial(
"profile and overlay runtime modes differ".to_string(),
));
}
if matches!(self.execution, DesiredExecution::External { .. })
&& !matches!(self.overlay.runtime_mode, WireMobRuntimeMode::TurnDriven)
{
return Err(IdentityIntentError::InvalidMemberMaterial(
"external execution requires turn-driven runtime mode".to_string(),
));
}
let mut declared_skills = BTreeSet::new();
for skill in &self.profile.skills {
validate_text("profile_skill", skill)?;
if !declared_skills.insert(skill.as_str()) {
return Err(IdentityIntentError::InvalidMemberMaterial(format!(
"profile repeats skill '{skill}'"
)));
}
}
let extracted_skills = self
.definition_extract
.skills
.keys()
.map(String::as_str)
.collect::<BTreeSet<_>>();
if declared_skills != extracted_skills {
return Err(IdentityIntentError::InvalidMemberMaterial(
"profile skills do not exactly match the definition extract".to_string(),
));
}
if let Some(policy) = &self.overlay.tool_access_policy {
match policy {
WireResolvedToolAccessPolicy::AllowList(names)
| WireResolvedToolAccessPolicy::DenyList(names) => {
validate_string_set("tool_access_policy_name", names)?;
}
}
}
let rehydrated = crate::portable_profile::rehydrate_portable_profile(&self.profile)
.map_err(IdentityIntentError::InvalidMemberMaterial)?;
let projected = crate::portable_profile::project_portable_profile(
&rehydrated,
rehydrated.runtime_mode,
&self.definition_extract.models,
self.profile_name.as_str(),
self.profile_name.as_str(),
Vec::new(),
)
.map_err(IdentityIntentError::InvalidMemberMaterial)?;
if projected != self.profile {
return Err(IdentityIntentError::InvalidMemberMaterial(
"portable profile does not round-trip exactly".to_string(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredInitialDelivery {
pub delivery_generation: u64,
pub delivery_id: InputId,
pub message_digest: String,
pub message: ContentInput,
}
impl DesiredInitialDelivery {
pub fn new(
delivery_generation: u64,
delivery_id: InputId,
message: ContentInput,
) -> Result<Self, IdentityIntentError> {
let message_digest = canonical_initial_message_digest(&message)?;
let value = Self {
delivery_generation,
delivery_id,
message_digest,
message,
};
value.validate()?;
Ok(value)
}
pub fn validate(&self) -> Result<(), IdentityIntentError> {
if self.delivery_generation == 0
|| self.delivery_id.0.is_nil()
|| self.message_digest != canonical_initial_message_digest(&self.message)?
{
return Err(IdentityIntentError::InvalidInitialDelivery);
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredMemberSpec {
pub material: DesiredMemberMaterial,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub initial_delivery: Option<DesiredInitialDelivery>,
}
impl DesiredMemberSpec {
pub fn validate(&self) -> Result<(), IdentityIntentError> {
self.material.validate()?;
if let Some(delivery) = &self.initial_delivery {
delivery.validate()?;
}
Ok(())
}
#[must_use]
pub fn execution(&self) -> &DesiredExecution {
&self.material.execution
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityProfileMemberDeclaration {
pub profile_name: ProfileName,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile_override: Option<PortableProfile>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_override: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub external_addressable_override: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<WireOpaqueJson>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub labels: Option<BTreeMap<String, String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub additional_instructions: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub system_prompt_override: Option<PortableSystemPrompt>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tool_access_policy: Option<WireResolvedToolAccessPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_binding: Option<WireAuthBindingRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub budget_limits: Option<BudgetLimits>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_mode: Option<WireMobRuntimeMode>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub required_env_keys: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub required_local_callback_tools: Vec<DesiredLocalCallbackTool>,
pub execution: DesiredExecution,
}
impl IdentityProfileMemberDeclaration {
pub fn validate(&self) -> Result<(), IdentityIntentError> {
validate_text("profile_name", self.profile_name.as_str())?;
if let Some(model) = &self.model_override {
validate_text("identity_model_override", model)?;
}
validate_string_set("required_env_key", &self.required_env_keys)?;
validate_required_local_callback_tools(&self.required_local_callback_tools, false)?;
validate_execution(&self.execution)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct DesiredIdentityEdge {
pub a: AgentIdentity,
pub b: AgentIdentity,
}
impl DesiredIdentityEdge {
pub fn new(left: AgentIdentity, right: AgentIdentity) -> Result<Self, IdentityIntentError> {
if left == right {
return Err(IdentityIntentError::SelfEdge(left));
}
let (a, b) = if left < right {
(left, right)
} else {
(right, left)
};
Ok(Self { a, b })
}
#[must_use]
pub fn owner(&self) -> &AgentIdentity {
&self.a
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(
tag = "desired_presence",
rename_all = "snake_case",
deny_unknown_fields
)]
pub enum IdentityIntent {
Present {
identity: AgentIdentity,
session: DesiredSessionTarget,
member: Box<DesiredMemberSpec>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
owned_wiring: BTreeSet<DesiredIdentityEdge>,
},
Absent {
identity: AgentIdentity,
},
}
impl IdentityIntent {
#[must_use]
pub fn identity(&self) -> &AgentIdentity {
match self {
Self::Present { identity, .. } | Self::Absent { identity } => identity,
}
}
pub fn validate(&self) -> Result<(), IdentityIntentError> {
validate_identity(self.identity())?;
if let Self::Present {
identity,
session,
member,
owned_wiring,
} = self
{
if session.session_id.0.is_nil() {
return Err(IdentityIntentError::NilSessionId);
}
if matches!(
session.authority_policy,
DesiredSessionAuthorityPolicy::CreateIfAbsent
) && session.lineage_generation != SessionGeneration::INITIAL
{
return Err(IdentityIntentError::CreateRequiresInitialGeneration);
}
SessionLineageId::new(session.lineage_id.as_str().to_string()).map_err(|_| {
IdentityIntentError::InvalidText {
field: "session_lineage_id",
}
})?;
member.validate()?;
for edge in owned_wiring {
validate_identity(&edge.a)?;
validate_identity(&edge.b)?;
if edge.a >= edge.b {
return Err(IdentityIntentError::NonCanonicalEdge(edge.clone()));
}
if edge.owner() != identity {
return Err(IdentityIntentError::EdgeOwnedByDifferentIdentity {
identity: identity.clone(),
edge: edge.clone(),
});
}
}
}
let encoded = serde_json::to_vec(self)
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
if encoded.len() > IDENTITY_INTENT_MAX_ENCODED_BYTES {
return Err(IdentityIntentError::TooLarge {
actual: encoded.len(),
maximum: IDENTITY_INTENT_MAX_ENCODED_BYTES,
});
}
Ok(())
}
pub fn digest(&self) -> Result<String, IdentityIntentError> {
self.validate()?;
let encoded = serde_json::to_vec(self)
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
Ok(sha256_digest(&encoded))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "cleanup", rename_all = "snake_case", deny_unknown_fields)]
pub enum IdentityRetirementPlan {
NoKnownRealization,
Targets {
session: DesiredSessionTarget,
execution: DesiredExecution,
incident_wiring: BTreeSet<DesiredIdentityEdge>,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityIntentRecord {
pub schema_version: u32,
pub mob_id: MobId,
pub intent_revision: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub declaration_scope: Option<IdentityDeclarationScopeId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub declaration_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tombstone_generation: Option<u64>,
#[serde(default)]
pub initial_delivery_generation_highwater: u64,
pub retirement_plan: IdentityRetirementPlan,
pub intent_digest: String,
pub authority_digest: String,
pub intent: IdentityIntent,
}
impl IdentityIntentRecord {
pub fn validate(&self) -> Result<(), IdentityIntentError> {
if self.schema_version != IDENTITY_INTENT_SCHEMA_VERSION {
return Err(IdentityIntentError::UnsupportedSchemaVersion {
record: "identity_intent",
version: self.schema_version,
});
}
validate_text("mob_id", self.mob_id.as_str())?;
if self.intent_revision == 0 {
return Err(IdentityIntentError::ZeroIntentRevision);
}
match (&self.declaration_scope, self.declaration_revision) {
(Some(scope), Some(revision)) if revision > 0 => {
validate_text("identity_declaration_scope", scope.as_str())?;
}
(None, None) => {}
_ => return Err(IdentityIntentError::InvalidDeclarationRevision),
}
if self.tombstone_generation == Some(0)
|| (matches!(&self.intent, IdentityIntent::Absent { .. })
&& self.tombstone_generation.is_none())
{
return Err(IdentityIntentError::InvalidTombstoneGeneration);
}
match &self.intent {
IdentityIntent::Present {
member,
session: _,
identity: _,
owned_wiring: _,
} => {
if let Some(delivery) = &member.initial_delivery {
if delivery.delivery_generation != self.initial_delivery_generation_highwater {
return Err(IdentityIntentError::InvalidInitialDelivery);
}
}
validate_retirement_plan(
&self.retirement_plan,
self.intent.identity(),
Some(&self.intent),
)?;
}
IdentityIntent::Absent { .. } => {
validate_retirement_plan(&self.retirement_plan, self.intent.identity(), None)?;
}
}
let digest = self.intent.digest()?;
if digest != self.intent_digest {
return Err(IdentityIntentError::DigestMismatch);
}
if self.canonical_authority_digest()? != self.authority_digest {
return Err(IdentityIntentError::DigestMismatch);
}
Ok(())
}
pub fn canonical_authority_digest(&self) -> Result<String, IdentityIntentError> {
#[derive(Serialize)]
struct AuthorityMaterial<'a> {
domain: &'static str,
schema_version: u32,
mob_id: &'a MobId,
intent_revision: u64,
declaration_scope: &'a Option<IdentityDeclarationScopeId>,
declaration_revision: Option<u64>,
tombstone_generation: Option<u64>,
initial_delivery_generation_highwater: u64,
retirement_plan: &'a IdentityRetirementPlan,
intent_digest: &'a str,
intent: &'a IdentityIntent,
}
let bytes = serde_json::to_vec(&AuthorityMaterial {
domain: "meerkat.identity.intent_authority.v1",
schema_version: self.schema_version,
mob_id: &self.mob_id,
intent_revision: self.intent_revision,
declaration_scope: &self.declaration_scope,
declaration_revision: self.declaration_revision,
tombstone_generation: self.tombstone_generation,
initial_delivery_generation_highwater: self.initial_delivery_generation_highwater,
retirement_plan: &self.retirement_plan,
intent_digest: &self.intent_digest,
intent: &self.intent,
})
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
Ok(sha256_digest(&bytes))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityLeaseClaim {
pub holder_id: String,
pub incarnation_id: String,
pub epoch: u64,
pub renewed_at_ms: u64,
pub expires_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityLeaseRecord {
pub schema_version: u32,
pub epoch_highwater: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub active: Option<IdentityLeaseClaim>,
}
impl IdentityLeaseRecord {
pub fn validate(&self) -> Result<(), IdentityIntentError> {
if self.schema_version != IDENTITY_LEASE_SCHEMA_VERSION {
return Err(IdentityIntentError::UnsupportedSchemaVersion {
record: "identity_lease",
version: self.schema_version,
});
}
if let Some(active) = &self.active {
validate_text("holder_id", &active.holder_id)?;
validate_text("incarnation_id", &active.incarnation_id)?;
if active.epoch == 0 || active.epoch != self.epoch_highwater {
return Err(IdentityIntentError::InvalidLeaseEpoch);
}
active
.expires_at_ms
.checked_sub(active.renewed_at_ms)
.filter(|ttl_ms| *ttl_ms > 0 && *ttl_ms <= IDENTITY_LEASE_MAX_TTL_MS)
.ok_or(IdentityIntentError::InvalidLeaseLifetime)?;
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityLeaseClaimOutcome {
Acquired(IdentityLeaseClaim),
Renewed(IdentityLeaseClaim),
HeldByOther(IdentityLeaseClaim),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityOperationKind {
SessionCreationConsumed,
RetirementProven,
ExternalBinding,
InitialDelivery,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "subject", rename_all = "snake_case", deny_unknown_fields)]
pub enum IdentityOperationSubject {
Identity { identity: AgentIdentity },
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "slot", rename_all = "snake_case", deny_unknown_fields)]
pub enum IdentityOperationSlot {
SessionCreationConsumed {
tombstone_generation: u64,
session_id: SessionId,
lineage_id: SessionLineageId,
lineage_generation: SessionGeneration,
},
RetirementProven {
tombstone_generation: u64,
},
ExternalBinding {
tombstone_generation: u64,
remote_signing_identity: WireTrustedPeerIdentity,
controller_signing_identity: WireTrustedPeerIdentity,
},
InitialDelivery {
tombstone_generation: u64,
session_id: SessionId,
lineage_id: SessionLineageId,
lineage_generation: SessionGeneration,
delivery_generation: u64,
},
}
impl IdentityOperationSlot {
#[must_use]
pub const fn kind(&self) -> IdentityOperationKind {
match self {
Self::SessionCreationConsumed { .. } => IdentityOperationKind::SessionCreationConsumed,
Self::RetirementProven { .. } => IdentityOperationKind::RetirementProven,
Self::ExternalBinding { .. } => IdentityOperationKind::ExternalBinding,
Self::InitialDelivery { .. } => IdentityOperationKind::InitialDelivery,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "operation", rename_all = "snake_case", deny_unknown_fields)]
pub enum IdentityOperationReceiptPayload {
SessionCreationConsumed {
authority: IdentitySessionStoreAuthority,
},
RetirementProven {
absent_authority_digest: String,
},
ExternalBinding {
expected_address: DesiredExternalAddress,
expected_identity: WireTrustedPeerIdentity,
expected_controller_identity: WireTrustedPeerIdentity,
ceremony_id: OperationId,
},
InitialDelivery {
delivery_generation: u64,
delivery_id: InputId,
message_digest: String,
},
}
impl IdentityOperationReceiptPayload {
#[must_use]
pub const fn kind(&self) -> IdentityOperationKind {
match self {
Self::SessionCreationConsumed { .. } => IdentityOperationKind::SessionCreationConsumed,
Self::RetirementProven { .. } => IdentityOperationKind::RetirementProven,
Self::ExternalBinding { .. } => IdentityOperationKind::ExternalBinding,
Self::InitialDelivery { .. } => IdentityOperationKind::InitialDelivery,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityOperationReceipt {
pub schema_version: u32,
pub mob_id: MobId,
pub subject: IdentityOperationSubject,
pub effect_kind: IdentityOperationKind,
pub slot: IdentityOperationSlot,
pub receipt_id: OperationId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub intent_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub intent_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub intent_authority_digest: Option<String>,
pub tombstone_generation: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub audit_lease_epoch: Option<u64>,
pub request_digest: String,
pub payload: IdentityOperationReceiptPayload,
}
impl IdentityOperationReceipt {
pub fn canonical_request_digest(&self) -> Result<String, IdentityIntentError> {
#[derive(Serialize)]
struct DigestMaterial<'a> {
domain: &'static str,
mob_id: &'a MobId,
subject: &'a IdentityOperationSubject,
effect_kind: IdentityOperationKind,
payload: ReceiptRequestMaterial<'a>,
}
#[derive(Serialize)]
#[serde(tag = "operation", rename_all = "snake_case")]
enum ReceiptRequestMaterial<'a> {
SessionCreationConsumed {
tombstone_generation: u64,
authority: &'a IdentitySessionStoreAuthority,
},
RetirementProven {
tombstone_generation: u64,
absent_authority_digest: &'a str,
},
ExternalBinding {
tombstone_generation: u64,
expected_address: &'a DesiredExternalAddress,
expected_identity: &'a WireTrustedPeerIdentity,
expected_controller_identity: &'a WireTrustedPeerIdentity,
ceremony_id: &'a OperationId,
},
InitialDelivery {
tombstone_generation: u64,
session_id: &'a SessionId,
lineage_id: &'a SessionLineageId,
lineage_generation: SessionGeneration,
delivery_generation: u64,
delivery_id: &'a InputId,
message_digest: &'a str,
},
}
let payload = match &self.payload {
IdentityOperationReceiptPayload::SessionCreationConsumed { authority } => {
ReceiptRequestMaterial::SessionCreationConsumed {
tombstone_generation: self.tombstone_generation.unwrap_or(0),
authority,
}
}
IdentityOperationReceiptPayload::RetirementProven {
absent_authority_digest,
} => ReceiptRequestMaterial::RetirementProven {
tombstone_generation: self.tombstone_generation.unwrap_or(0),
absent_authority_digest,
},
IdentityOperationReceiptPayload::ExternalBinding {
expected_address,
expected_identity,
expected_controller_identity,
ceremony_id,
} => ReceiptRequestMaterial::ExternalBinding {
tombstone_generation: self.tombstone_generation.unwrap_or(0),
expected_address,
expected_identity,
expected_controller_identity,
ceremony_id,
},
IdentityOperationReceiptPayload::InitialDelivery {
delivery_generation,
delivery_id,
message_digest,
} => {
let IdentityOperationSlot::InitialDelivery {
session_id,
lineage_id,
lineage_generation,
..
} = &self.slot
else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
ReceiptRequestMaterial::InitialDelivery {
tombstone_generation: self.tombstone_generation.unwrap_or(0),
session_id,
lineage_id,
lineage_generation: *lineage_generation,
delivery_generation: *delivery_generation,
delivery_id,
message_digest,
}
}
};
let bytes = serde_json::to_vec(&DigestMaterial {
domain: "meerkat.identity.operation_receipt.request.v1",
mob_id: &self.mob_id,
subject: &self.subject,
effect_kind: self.effect_kind,
payload,
})
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
Ok(sha256_digest(&bytes))
}
pub fn validate(&self) -> Result<(), IdentityIntentError> {
if self.schema_version != IDENTITY_OPERATION_RECEIPT_SCHEMA_VERSION {
return Err(IdentityIntentError::UnsupportedSchemaVersion {
record: "identity_operation_receipt",
version: self.schema_version,
});
}
validate_text("mob_id", self.mob_id.as_str())?;
let IdentityOperationSubject::Identity { identity } = &self.subject;
validate_identity(identity)?;
if self.receipt_id.0.is_nil() {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
if self.tombstone_generation == Some(0) || self.audit_lease_epoch == Some(0) {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
if self.effect_kind != self.slot.kind() || self.effect_kind != self.payload.kind() {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
let normalized_tombstone = self.tombstone_generation.unwrap_or(0);
match &self.payload {
IdentityOperationReceiptPayload::SessionCreationConsumed { authority } => {
validate_identity_receipt_authority(self)?;
let IdentityOperationSlot::SessionCreationConsumed {
tombstone_generation,
session_id,
lineage_id,
lineage_generation: _,
} = &self.slot
else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
authority.validate()?;
if *tombstone_generation != normalized_tombstone
|| authority.session_id() != session_id
|| lineage_id.as_str().is_empty()
{
return Err(IdentityIntentError::InvalidOperationReceipt);
}
}
IdentityOperationReceiptPayload::RetirementProven {
absent_authority_digest,
} => {
validate_identity_receipt_authority(self)?;
let IdentityOperationSlot::RetirementProven {
tombstone_generation,
} = &self.slot
else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
if *tombstone_generation == 0 || *tombstone_generation != normalized_tombstone {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
validate_sha256_digest(absent_authority_digest)?;
if Some(absent_authority_digest) != self.intent_authority_digest.as_ref() {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
}
IdentityOperationReceiptPayload::ExternalBinding {
expected_address: _,
expected_identity,
expected_controller_identity,
ceremony_id,
} => {
validate_identity_receipt_authority(self)?;
let IdentityOperationSlot::ExternalBinding {
tombstone_generation,
remote_signing_identity,
controller_signing_identity,
} = &self.slot
else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
if *tombstone_generation != normalized_tombstone
|| remote_signing_identity != expected_identity
|| controller_signing_identity != expected_controller_identity
|| ceremony_id.0.is_nil()
{
return Err(IdentityIntentError::InvalidOperationReceipt);
}
expected_identity.resolve().map_err(|error| {
IdentityIntentError::InvalidExternalIdentity(error.to_string())
})?;
expected_controller_identity.resolve().map_err(|error| {
IdentityIntentError::InvalidExternalIdentity(error.to_string())
})?;
}
IdentityOperationReceiptPayload::InitialDelivery {
delivery_generation,
delivery_id,
message_digest,
} => {
validate_identity_receipt_authority(self)?;
let IdentityOperationSlot::InitialDelivery {
tombstone_generation,
delivery_generation: slot_generation,
..
} = &self.slot
else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
if *tombstone_generation != normalized_tombstone
|| *delivery_generation == 0
|| delivery_generation != slot_generation
|| delivery_id.0.is_nil()
{
return Err(IdentityIntentError::InvalidOperationReceipt);
}
validate_sha256_digest(message_digest)?;
}
}
if self.canonical_request_digest()? != self.request_digest {
return Err(IdentityIntentError::DigestMismatch);
}
Ok(())
}
}
fn validate_identity_receipt_authority(
receipt: &IdentityOperationReceipt,
) -> Result<(), IdentityIntentError> {
if receipt.intent_revision.is_none_or(|revision| revision == 0) {
return Err(IdentityIntentError::InvalidOperationReceipt);
}
let Some(intent_digest) = &receipt.intent_digest else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
let Some(authority_digest) = &receipt.intent_authority_digest else {
return Err(IdentityIntentError::InvalidOperationReceipt);
};
validate_sha256_digest(intent_digest)?;
validate_sha256_digest(authority_digest)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityStoredObservation<T> {
Missing,
Valid(T),
Unsupported {
evidence_digest: String,
detail: String,
},
Malformed {
evidence_digest: String,
detail: String,
},
}
impl<T> IdentityStoredObservation<T> {
pub fn validate_evidence(&self) -> Result<(), IdentityIntentError> {
match self {
Self::Unsupported {
evidence_digest, ..
}
| Self::Malformed {
evidence_digest, ..
} => validate_sha256_digest(evidence_digest),
Self::Missing | Self::Valid(_) => Ok(()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityOperationReceiptInsertOutcome {
Inserted(IdentityOperationReceipt),
ExistingExact(IdentityOperationReceipt),
Conflict(IdentityOperationReceipt),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityAuthorityCondition {
Unavailable,
Missing,
Malformed,
PresentCreateIfAbsent,
PresentRequireExisting,
Absent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityLeaseCondition {
Unavailable,
Missing,
Malformed,
HeldByCurrentIncarnation,
HeldByOtherLiveIncarnation,
HeldByExpiredIncarnation,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityResourceCondition {
Unavailable,
Missing,
Matching,
Divergent,
Malformed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentitySessionCondition {
Unavailable,
Missing,
Matching,
RecoverableDivergence,
AmbiguousDivergence,
Malformed,
IrrecoverablyCorrupt,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityReceiptCondition {
NotRequired,
Unavailable,
Missing,
Matching,
Conflicting,
Malformed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityExternalTrustCondition {
NotRequired,
Unavailable,
Matching,
Absent,
Contradictory,
Indeterminate,
Malformed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityExternalCeremonyCondition {
NotRequired,
FreshAvailable,
TemporarilyUnavailable,
AwaitFresh,
SpentOrUnknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityInitialDeliveryCondition {
NotRequired,
Unavailable,
ProvenAbsent,
AcceptedPendingExact,
CommittedExact,
ContentOnlyMatch,
OperationCollision,
Contradictory,
Indeterminate,
Malformed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityReconcileFacts {
pub intent: IdentityAuthorityCondition,
pub lease: IdentityLeaseCondition,
pub external_binding_required: bool,
pub initial_delivery_required: bool,
pub session_creation_receipt: IdentityReceiptCondition,
pub retirement_receipt: IdentityReceiptCondition,
pub session: IdentitySessionCondition,
pub runtime: IdentityResourceCondition,
pub member: IdentityResourceCondition,
pub external_binding_receipt: IdentityReceiptCondition,
pub external_trust: IdentityExternalTrustCondition,
pub external_ceremony: IdentityExternalCeremonyCondition,
pub initial_delivery_receipt: IdentityReceiptCondition,
pub initial_delivery: IdentityInitialDeliveryCondition,
pub wiring: IdentityResourceCondition,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityReconcileDecision {
Backoff,
RepairBlocked,
AcquireLease,
AwaitLease,
SealRetirementProven,
SealSessionCreationConsumed,
EnsureSessionAuthority,
EnsureRuntimeRegistration,
AwaitExternalBindingCeremony,
EnsureExternalBindingReceipt,
EnsureExternalBinding,
EnsureMemberMaterialization,
EnsureInitialDeliveryReceipt,
EnsureInitialDelivery,
AwaitInitialDelivery,
ReconcileWiring,
RetireMemberMaterialization,
RetireRuntimeRegistration,
ReleaseSessionAuthority,
Converged,
Tombstoned,
Quarantined,
}
#[must_use]
pub fn classify_identity_reconciliation(
facts: IdentityReconcileFacts,
) -> IdentityReconcileDecision {
crate::machines::mob_machine::generated_identity_reconcile_decision(facts)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityTargetObservationVersion {
InsertIfAbsent,
Absent {
absence_version: String,
},
Version {
version: String,
},
}
impl IdentityTargetObservationVersion {
fn validate(&self) -> Result<(), IdentityIntentError> {
let value = match self {
Self::InsertIfAbsent => return Ok(()),
Self::Absent { absence_version } => absence_version,
Self::Version { version } => version,
};
if value.is_empty() {
return Err(IdentityIntentError::InvalidObservationVersion);
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityResourceObservation {
Unavailable {
detail: String,
},
Missing {
absence_version: String,
},
Matching {
version: String,
},
Divergent {
version: String,
detail: String,
},
Malformed {
observed_version: Option<String>,
detail: String,
},
}
impl IdentityResourceObservation {
#[must_use]
pub const fn condition(&self) -> IdentityResourceCondition {
match self {
Self::Unavailable { .. } => IdentityResourceCondition::Unavailable,
Self::Missing { .. } => IdentityResourceCondition::Missing,
Self::Matching { .. } => IdentityResourceCondition::Matching,
Self::Divergent { .. } => IdentityResourceCondition::Divergent,
Self::Malformed { .. } => IdentityResourceCondition::Malformed,
}
}
pub fn target_precondition(
&self,
) -> Result<Option<IdentityTargetObservationVersion>, IdentityIntentError> {
let precondition = match self {
Self::Missing { absence_version } => Some(IdentityTargetObservationVersion::Absent {
absence_version: absence_version.clone(),
}),
Self::Matching { version } | Self::Divergent { version, .. } => {
Some(IdentityTargetObservationVersion::Version {
version: version.clone(),
})
}
Self::Unavailable { .. } | Self::Malformed { .. } => None,
};
if let Some(precondition) = &precondition {
precondition.validate()?;
}
Ok(precondition)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum IdentitySessionObservationState {
Unavailable {
detail: String,
},
Missing {
absence_version: String,
},
Matching {
authority: IdentitySessionStoreAuthority,
},
AmbiguousDivergence {
evidence_digest: String,
target: IdentityTargetObservationVersion,
detail: String,
},
Malformed {
evidence_digest: String,
version: String,
detail: String,
},
MalformedUnversioned {
detail: String,
},
IrrecoverablyCorrupt {
evidence_digest: String,
version: String,
detail: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IdentitySessionObservation {
state: IdentitySessionObservationState,
}
impl IdentitySessionObservation {
pub(crate) fn matching(
desired: &DesiredSessionTarget,
observed: &Session,
authority: IdentitySessionStoreAuthority,
) -> Result<Self, IdentityIntentError> {
validate_session_target(desired)?;
authority.validate()?;
if observed.id() != &desired.session_id || authority.session_id() != &desired.session_id {
return Err(IdentityIntentError::SessionStoreAuthorityMismatch);
}
Ok(Self {
state: IdentitySessionObservationState::Matching { authority },
})
}
pub fn missing(absence_version: String) -> Result<Self, IdentityIntentError> {
IdentityTargetObservationVersion::Absent {
absence_version: absence_version.clone(),
}
.validate()?;
Ok(Self {
state: IdentitySessionObservationState::Missing { absence_version },
})
}
#[must_use]
pub fn unavailable(detail: impl Into<String>) -> Self {
Self {
state: IdentitySessionObservationState::Unavailable {
detail: detail.into(),
},
}
}
pub fn ambiguous_divergence(
evidence_digest: String,
target: IdentityTargetObservationVersion,
detail: impl Into<String>,
) -> Result<Self, IdentityIntentError> {
validate_sha256_digest(&evidence_digest)?;
target.validate()?;
Ok(Self {
state: IdentitySessionObservationState::AmbiguousDivergence {
evidence_digest,
target,
detail: detail.into(),
},
})
}
pub fn malformed(
evidence_digest: String,
version: String,
detail: impl Into<String>,
) -> Result<Self, IdentityIntentError> {
validate_sha256_digest(&evidence_digest)?;
IdentityTargetObservationVersion::Version {
version: version.clone(),
}
.validate()?;
Ok(Self {
state: IdentitySessionObservationState::Malformed {
evidence_digest,
version,
detail: detail.into(),
},
})
}
pub fn malformed_unversioned(detail: impl Into<String>) -> Result<Self, IdentityIntentError> {
let detail = detail.into();
validate_text("identity_session_malformed_detail", &detail)?;
Ok(Self {
state: IdentitySessionObservationState::MalformedUnversioned { detail },
})
}
#[must_use]
pub fn malformed_unversioned_detail(&self) -> Option<&str> {
match &self.state {
IdentitySessionObservationState::MalformedUnversioned { detail } => Some(detail),
_ => None,
}
}
pub fn irrecoverably_corrupt(
evidence_digest: String,
version: String,
detail: impl Into<String>,
) -> Result<Self, IdentityIntentError> {
validate_sha256_digest(&evidence_digest)?;
IdentityTargetObservationVersion::Version {
version: version.clone(),
}
.validate()?;
Ok(Self {
state: IdentitySessionObservationState::IrrecoverablyCorrupt {
evidence_digest,
version,
detail: detail.into(),
},
})
}
#[must_use]
pub const fn condition(&self) -> IdentitySessionCondition {
match &self.state {
IdentitySessionObservationState::Unavailable { .. } => {
IdentitySessionCondition::Unavailable
}
IdentitySessionObservationState::Missing { .. } => IdentitySessionCondition::Missing,
IdentitySessionObservationState::Matching { .. } => IdentitySessionCondition::Matching,
IdentitySessionObservationState::AmbiguousDivergence { .. } => {
IdentitySessionCondition::AmbiguousDivergence
}
IdentitySessionObservationState::Malformed { .. } => {
IdentitySessionCondition::Malformed
}
IdentitySessionObservationState::MalformedUnversioned { .. } => {
IdentitySessionCondition::Malformed
}
IdentitySessionObservationState::IrrecoverablyCorrupt { .. } => {
IdentitySessionCondition::IrrecoverablyCorrupt
}
}
}
pub fn target_precondition(
&self,
) -> Result<Option<IdentityTargetObservationVersion>, IdentityIntentError> {
let target = match &self.state {
IdentitySessionObservationState::Missing { absence_version } => {
Some(IdentityTargetObservationVersion::Absent {
absence_version: absence_version.clone(),
})
}
IdentitySessionObservationState::Matching { authority } => {
Some(IdentityTargetObservationVersion::Version {
version: authority.observation_version()?,
})
}
IdentitySessionObservationState::Malformed { version, .. }
| IdentitySessionObservationState::IrrecoverablyCorrupt { version, .. } => {
Some(IdentityTargetObservationVersion::Version {
version: version.clone(),
})
}
IdentitySessionObservationState::AmbiguousDivergence { target, .. } => {
Some(target.clone())
}
IdentitySessionObservationState::Unavailable { .. }
| IdentitySessionObservationState::MalformedUnversioned { .. } => None,
};
if let Some(target) = &target {
target.validate()?;
}
Ok(target)
}
#[must_use]
pub fn store_authority(&self) -> Option<&IdentitySessionStoreAuthority> {
match &self.state {
IdentitySessionObservationState::Matching { authority } => Some(authority),
IdentitySessionObservationState::Unavailable { .. }
| IdentitySessionObservationState::Missing { .. }
| IdentitySessionObservationState::AmbiguousDivergence { .. }
| IdentitySessionObservationState::Malformed { .. }
| IdentitySessionObservationState::MalformedUnversioned { .. }
| IdentitySessionObservationState::IrrecoverablyCorrupt { .. } => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IdentityActuatorTarget {
Session,
Runtime,
SessionCreationReceipt,
RetirementReceipt,
ExternalBindingReceipt,
ExternalBinding,
Member,
InitialDeliveryReceipt,
InitialDelivery,
Wiring,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IdentityActuationPermit {
pub mob_id: MobId,
pub identity: AgentIdentity,
pub target: IdentityActuatorTarget,
pub intent_revision: u64,
pub intent_digest: String,
pub intent_authority_digest: String,
pub lease_epoch: u64,
pub lease_holder_id: String,
pub lease_incarnation_id: String,
pub lease_expires_at_ms: u64,
pub target_observation: IdentityTargetObservationVersion,
}
impl IdentityActuationPermit {
pub fn validate_for_write(&self, observed_at_ms: u64) -> Result<(), IdentityIntentError> {
validate_text("mob_id", self.mob_id.as_str())?;
validate_identity(&self.identity)?;
validate_text("lease_holder_id", &self.lease_holder_id)?;
validate_text("lease_incarnation_id", &self.lease_incarnation_id)?;
if self.intent_revision == 0 || self.lease_epoch == 0 {
return Err(IdentityIntentError::InvalidActuationPermit);
}
validate_sha256_digest(&self.intent_digest)?;
validate_sha256_digest(&self.intent_authority_digest)?;
if observed_at_ms >= self.lease_expires_at_ms {
return Err(IdentityIntentError::ExpiredActuationPermit);
}
self.target_observation.validate()?;
let receipt_target = matches!(
self.target,
IdentityActuatorTarget::SessionCreationReceipt
| IdentityActuatorTarget::RetirementReceipt
| IdentityActuatorTarget::ExternalBindingReceipt
| IdentityActuatorTarget::InitialDeliveryReceipt
);
if receipt_target
!= matches!(
self.target_observation,
IdentityTargetObservationVersion::InsertIfAbsent
)
{
return Err(IdentityIntentError::InvalidActuationPermit);
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdentityConvergenceCondition {
Pending,
Reconciling,
Converged,
Backoff,
RepairBlocked,
Quarantined,
Tombstoned,
Suspended,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IdentityConvergenceStatus {
pub identity: AgentIdentity,
pub intent_revision: Option<u64>,
pub lease_epoch: Option<u64>,
pub decision: Option<IdentityReconcileDecision>,
pub observed_at_ms: u64,
pub detail: Option<String>,
}
impl IdentityConvergenceStatus {
#[must_use]
pub const fn condition(&self) -> IdentityConvergenceCondition {
match self.decision {
None => IdentityConvergenceCondition::Pending,
Some(IdentityReconcileDecision::Backoff) => IdentityConvergenceCondition::Backoff,
Some(IdentityReconcileDecision::RepairBlocked) => {
IdentityConvergenceCondition::RepairBlocked
}
Some(
IdentityReconcileDecision::AwaitLease
| IdentityReconcileDecision::AwaitExternalBindingCeremony
| IdentityReconcileDecision::AwaitInitialDelivery,
) => IdentityConvergenceCondition::Suspended,
Some(IdentityReconcileDecision::Converged) => IdentityConvergenceCondition::Converged,
Some(IdentityReconcileDecision::Tombstoned) => IdentityConvergenceCondition::Tombstoned,
Some(IdentityReconcileDecision::Quarantined) => {
IdentityConvergenceCondition::Quarantined
}
Some(
IdentityReconcileDecision::AcquireLease
| IdentityReconcileDecision::SealRetirementProven
| IdentityReconcileDecision::SealSessionCreationConsumed
| IdentityReconcileDecision::EnsureSessionAuthority
| IdentityReconcileDecision::EnsureRuntimeRegistration
| IdentityReconcileDecision::EnsureExternalBindingReceipt
| IdentityReconcileDecision::EnsureExternalBinding
| IdentityReconcileDecision::EnsureMemberMaterialization
| IdentityReconcileDecision::EnsureInitialDeliveryReceipt
| IdentityReconcileDecision::EnsureInitialDelivery
| IdentityReconcileDecision::ReconcileWiring
| IdentityReconcileDecision::RetireMemberMaterialization
| IdentityReconcileDecision::RetireRuntimeRegistration
| IdentityReconcileDecision::ReleaseSessionAuthority,
) => IdentityConvergenceCondition::Reconciling,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum IdentityIntentError {
InvalidText {
field: &'static str,
},
NilSessionId,
CreateRequiresInitialGeneration,
InvalidExternalAddress(String),
InvalidExternalIdentity(String),
SelfEdge(AgentIdentity),
NonCanonicalEdge(DesiredIdentityEdge),
EdgeOwnedByDifferentIdentity {
identity: AgentIdentity,
edge: DesiredIdentityEdge,
},
TooLarge {
actual: usize,
maximum: usize,
},
Serialization(String),
UnsupportedSchemaVersion {
record: &'static str,
version: u32,
},
ZeroIntentRevision,
InvalidDeclarationRevision,
InvalidMemberMaterial(String),
DigestMismatch,
InvalidLeaseEpoch,
InvalidLeaseLifetime,
InvalidTombstoneGeneration,
InvalidInitialDelivery,
InvalidRetirementPlan,
InvalidOperationReceipt,
InvalidObservationVersion,
InvalidActuationPermit,
ExpiredActuationPermit,
InvalidSessionStoreAuthority,
SessionStoreAuthorityMismatch,
CounterExhausted {
counter: &'static str,
},
}
impl fmt::Display for IdentityIntentError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::InvalidText { field } => {
write!(formatter, "{field} must be nonempty canonical text")
}
Self::NilSessionId => formatter.write_str("identity intent session id must be non-nil"),
Self::CreateRequiresInitialGeneration => formatter
.write_str("CreateIfAbsent requires the initial session lineage generation"),
Self::InvalidExternalAddress(detail) => {
write!(formatter, "invalid desired external address: {detail}")
}
Self::InvalidExternalIdentity(detail) => {
write!(
formatter,
"invalid desired external signing identity: {detail}"
)
}
Self::SelfEdge(identity) => {
write!(formatter, "identity '{identity}' cannot wire to itself")
}
Self::NonCanonicalEdge(edge) => {
write!(formatter, "desired wiring edge is not canonical: {edge:?}")
}
Self::EdgeOwnedByDifferentIdentity { identity, edge } => write!(
formatter,
"identity '{identity}' cannot own desired wiring edge {edge:?}; owner is '{}'",
edge.owner()
),
Self::TooLarge { actual, maximum } => write!(
formatter,
"identity intent is {actual} bytes; maximum is {maximum}"
),
Self::Serialization(detail) => {
write!(formatter, "identity intent serialization failed: {detail}")
}
Self::UnsupportedSchemaVersion { record, version } => {
write!(formatter, "unsupported {record} schema version {version}")
}
Self::ZeroIntentRevision => {
formatter.write_str("identity intent revision must be nonzero")
}
Self::InvalidDeclarationRevision => formatter.write_str(
"identity declaration scope and nonzero revision must be present together",
),
Self::InvalidMemberMaterial(detail) => {
write!(formatter, "invalid desired member material: {detail}")
}
Self::DigestMismatch => {
formatter.write_str("identity intent digest does not match content")
}
Self::InvalidLeaseEpoch => formatter
.write_str("identity lease active epoch must equal a nonzero epoch highwater"),
Self::InvalidLeaseLifetime => write!(
formatter,
"identity lease lifetime must be between 1 and {IDENTITY_LEASE_MAX_TTL_MS}ms"
),
Self::InvalidTombstoneGeneration => formatter.write_str(
"absent identity intent requires a nonzero monotonic tombstone generation",
),
Self::InvalidInitialDelivery => formatter.write_str(
"initial delivery must have a nonzero generation, stable input id, and matching message digest",
),
Self::InvalidRetirementPlan => formatter.write_str(
"identity retirement plan must retain the store-sealed cleanup targets",
),
Self::InvalidOperationReceipt => {
formatter.write_str("identity operation receipt is internally incoherent")
}
Self::InvalidObservationVersion => {
formatter.write_str("identity target observation version must be nonempty")
}
Self::InvalidActuationPermit => {
formatter.write_str("identity actuation permit is internally incoherent")
}
Self::ExpiredActuationPermit => {
formatter.write_str("identity actuation permit lease has expired")
}
Self::InvalidSessionStoreAuthority => formatter
.write_str("identity session store authority is internally incoherent"),
Self::SessionStoreAuthorityMismatch => formatter.write_str(
"store-issued session authority does not match the observed desired session",
),
Self::CounterExhausted { counter } => {
write!(formatter, "identity {counter} counter exhausted")
}
}
}
}
impl std::error::Error for IdentityIntentError {}
fn validate_identity(identity: &AgentIdentity) -> Result<(), IdentityIntentError> {
validate_text("identity", identity.as_str())
}
fn validate_text(field: &'static str, value: &str) -> Result<(), IdentityIntentError> {
if value.is_empty() || value.trim() != value {
Err(IdentityIntentError::InvalidText { field })
} else {
Ok(())
}
}
fn sha256_digest(bytes: &[u8]) -> String {
format!("sha256:{:x}", Sha256::digest(bytes))
}
fn canonical_initial_message_digest(message: &ContentInput) -> Result<String, IdentityIntentError> {
#[derive(Serialize)]
struct DigestMaterial<'a> {
domain: &'static str,
message: &'a ContentInput,
}
let bytes = serde_json::to_vec(&DigestMaterial {
domain: "meerkat.identity.initial_delivery.message.v1",
message,
})
.map_err(|error| IdentityIntentError::Serialization(error.to_string()))?;
Ok(sha256_digest(&bytes))
}
fn validate_sha256_digest(value: &str) -> Result<(), IdentityIntentError> {
let Some(hex) = value.strip_prefix("sha256:") else {
return Err(IdentityIntentError::DigestMismatch);
};
if hex.len() != 64
|| !hex
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
{
return Err(IdentityIntentError::DigestMismatch);
}
Ok(())
}
fn has_prefixed_sha256(value: &str, prefixes: &[&str]) -> bool {
prefixes.iter().any(|prefix| {
value.strip_prefix(prefix).is_some_and(|hex| {
hex.len() == 64
&& hex
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
})
})
}
fn validate_session_target(session: &DesiredSessionTarget) -> Result<(), IdentityIntentError> {
if session.session_id.0.is_nil() {
return Err(IdentityIntentError::NilSessionId);
}
SessionLineageId::new(session.lineage_id.as_str().to_string()).map_err(|_| {
IdentityIntentError::InvalidText {
field: "session_lineage_id",
}
})?;
if matches!(
session.authority_policy,
DesiredSessionAuthorityPolicy::CreateIfAbsent
) && session.lineage_generation != SessionGeneration::INITIAL
{
return Err(IdentityIntentError::CreateRequiresInitialGeneration);
}
Ok(())
}
fn validate_retirement_plan(
plan: &IdentityRetirementPlan,
identity: &AgentIdentity,
current_intent: Option<&IdentityIntent>,
) -> Result<(), IdentityIntentError> {
let IdentityRetirementPlan::Targets {
session,
execution,
incident_wiring,
} = plan
else {
return if current_intent.is_none() {
Ok(())
} else {
Err(IdentityIntentError::InvalidRetirementPlan)
};
};
validate_session_target(session)?;
validate_execution(execution)?;
if incident_wiring
.iter()
.any(|edge| &edge.a != identity && &edge.b != identity)
{
return Err(IdentityIntentError::InvalidRetirementPlan);
}
if let Some(IdentityIntent::Present {
session: desired_session,
member,
owned_wiring,
..
}) = current_intent
{
if session != desired_session
|| execution != member.execution()
|| !owned_wiring.is_subset(incident_wiring)
{
return Err(IdentityIntentError::InvalidRetirementPlan);
}
}
Ok(())
}
fn validate_string_set(field: &'static str, values: &[String]) -> Result<(), IdentityIntentError> {
let mut prior = None;
for value in values {
validate_text(field, value)?;
if prior.is_some_and(|prior: &String| prior >= value) {
return Err(IdentityIntentError::InvalidText { field });
}
prior = Some(value);
}
Ok(())
}
fn validate_required_local_callback_tools(
tools: &[DesiredLocalCallbackTool],
require_canonical_order: bool,
) -> Result<(), IdentityIntentError> {
let mut names = BTreeSet::new();
let mut prior_name = None;
for tool in tools {
tool.validate()?;
if !names.insert(tool.name.as_str()) {
return Err(IdentityIntentError::InvalidMemberMaterial(format!(
"local callback tool '{}' is declared more than once",
tool.name
)));
}
if require_canonical_order
&& prior_name.is_some_and(|prior: &str| prior >= tool.name.as_str())
{
return Err(IdentityIntentError::InvalidMemberMaterial(
"local callback tools are not in canonical name order".to_string(),
));
}
prior_name = Some(tool.name.as_str());
}
Ok(())
}
fn validate_execution(execution: &DesiredExecution) -> Result<(), IdentityIntentError> {
match execution {
DesiredExecution::External {
address: _,
identity,
} => {
identity
.resolve()
.map_err(|error| IdentityIntentError::InvalidExternalIdentity(error.to_string()))?;
}
DesiredExecution::PlacedSession { host_id } => validate_text("host_id", host_id)?,
DesiredExecution::ControllingSession | DesiredExecution::AnyBoundHostSession => {}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn matching_facts() -> IdentityReconcileFacts {
IdentityReconcileFacts {
intent: IdentityAuthorityCondition::PresentRequireExisting,
lease: IdentityLeaseCondition::HeldByCurrentIncarnation,
external_binding_required: false,
initial_delivery_required: false,
session_creation_receipt: IdentityReceiptCondition::NotRequired,
retirement_receipt: IdentityReceiptCondition::NotRequired,
session: IdentitySessionCondition::Matching,
runtime: IdentityResourceCondition::Matching,
member: IdentityResourceCondition::Matching,
external_binding_receipt: IdentityReceiptCondition::NotRequired,
external_trust: IdentityExternalTrustCondition::NotRequired,
external_ceremony: IdentityExternalCeremonyCondition::NotRequired,
initial_delivery_receipt: IdentityReceiptCondition::NotRequired,
initial_delivery: IdentityInitialDeliveryCondition::NotRequired,
wiring: IdentityResourceCondition::Matching,
}
}
fn generated_transition_decision(
authority: &mut crate::machines::mob_machine::MobMachineAuthority,
facts: IdentityReconcileFacts,
) -> IdentityReconcileDecision {
let transition = crate::machines::mob_machine::MobMachineMutator::apply(
authority,
crate::machines::mob_machine::identity_reconciliation_input(facts),
)
.expect("generated identity reconciliation transition must be total");
let effects = transition.into_effects();
match effects.as_slice() {
[
crate::machines::mob_machine::MobMachineEffect::IdentityReconciliationClassified {
decision,
},
] => *decision,
other => panic!(
"generated identity reconciliation transition emitted unexpected effects: {other:?}"
),
}
}
#[test]
fn unversioned_session_corruption_is_malformed_without_write_authority() {
let observation = IdentitySessionObservation::malformed_unversioned(
"persisted Session failed typed decoding",
)
.expect("canonical corruption detail should construct");
assert_eq!(observation.condition(), IdentitySessionCondition::Malformed);
assert_eq!(observation.target_precondition().unwrap(), None);
assert_eq!(observation.store_authority(), None);
assert_eq!(
observation.malformed_unversioned_detail(),
Some("persisted Session failed typed decoding")
);
assert_eq!(
IdentitySessionObservation::unavailable("transport failure")
.malformed_unversioned_detail(),
None
);
}
#[test]
fn unversioned_session_corruption_rejects_noncanonical_detail() {
for detail in ["", " ", " leading", "trailing "] {
assert!(matches!(
IdentitySessionObservation::malformed_unversioned(detail),
Err(IdentityIntentError::InvalidText {
field: "identity_session_malformed_detail"
})
));
}
}
fn session_creation_receipt(
session_id: SessionId,
authority: IdentitySessionStoreAuthority,
) -> IdentityOperationReceipt {
let mut receipt = IdentityOperationReceipt {
schema_version: IDENTITY_OPERATION_RECEIPT_SCHEMA_VERSION,
mob_id: MobId::from("authority-receipt-mob"),
subject: IdentityOperationSubject::Identity {
identity: AgentIdentity::from("authority-receipt-member"),
},
effect_kind: IdentityOperationKind::SessionCreationConsumed,
slot: IdentityOperationSlot::SessionCreationConsumed {
tombstone_generation: 0,
lineage_id: SessionLineageId::for_session(&session_id),
lineage_generation: SessionGeneration::INITIAL,
session_id,
},
receipt_id: OperationId::new(),
intent_revision: Some(1),
intent_digest: Some(
"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
.to_string(),
),
intent_authority_digest: Some(
"sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
.to_string(),
),
tombstone_generation: None,
audit_lease_epoch: Some(1),
request_digest:
"sha256:0000000000000000000000000000000000000000000000000000000000000000"
.to_string(),
payload: IdentityOperationReceiptPayload::SessionCreationConsumed { authority },
};
receipt.request_digest = receipt.canonical_request_digest().unwrap();
receipt
}
#[test]
fn session_creation_receipt_carries_only_exact_store_authority() {
let session_id = SessionId::new();
let authority = IdentitySessionStoreAuthority::whole_blob_for_test(
session_id.clone(),
7,
"row-sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
);
let receipt = session_creation_receipt(session_id.clone(), authority.clone());
receipt.validate().unwrap();
let encoded = serde_json::to_value(&receipt).unwrap();
assert_eq!(
encoded.pointer("/payload/authority/store_revision"),
Some(&serde_json::json!(7)),
);
assert_eq!(
encoded.pointer("/payload/authority/token/profile"),
Some(&serde_json::json!("whole_blob_v1")),
);
assert!(
encoded.pointer("/payload/authority/checkpoint").is_none(),
"receipt authority must not retain Session-owned checkpoint vocabulary",
);
let wrong_authority = IdentitySessionStoreAuthority::whole_blob_for_test(
SessionId::new(),
8,
"row-sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
);
let wrong = session_creation_receipt(session_id, wrong_authority);
assert!(matches!(
wrong.validate(),
Err(IdentityIntentError::InvalidOperationReceipt),
));
}
#[test]
fn unversioned_session_corruption_reuses_strict_malformed_condition_serde() {
let observation =
IdentitySessionObservation::malformed_unversioned("unsupported persisted schema")
.unwrap();
let encoded = serde_json::to_string(&observation.condition()).unwrap();
assert_eq!(encoded, r#""malformed""#);
assert_eq!(
serde_json::from_str::<IdentitySessionCondition>(&encoded).unwrap(),
IdentitySessionCondition::Malformed
);
assert!(
serde_json::from_str::<IdentitySessionCondition>(r#""malformed_unversioned""#).is_err(),
"the observation shape must not extend generated classifier vocabulary"
);
}
#[test]
fn classifier_is_total_over_the_core_observation_product() {
let mut authority = crate::machines::mob_machine::MobMachineAuthority::new();
let initial_state = authority.state().clone();
let intents = [
IdentityAuthorityCondition::Unavailable,
IdentityAuthorityCondition::Missing,
IdentityAuthorityCondition::Malformed,
IdentityAuthorityCondition::PresentCreateIfAbsent,
IdentityAuthorityCondition::PresentRequireExisting,
IdentityAuthorityCondition::Absent,
];
let leases = [
IdentityLeaseCondition::Unavailable,
IdentityLeaseCondition::Missing,
IdentityLeaseCondition::Malformed,
IdentityLeaseCondition::HeldByCurrentIncarnation,
IdentityLeaseCondition::HeldByOtherLiveIncarnation,
IdentityLeaseCondition::HeldByExpiredIncarnation,
];
let resources = [
IdentityResourceCondition::Unavailable,
IdentityResourceCondition::Missing,
IdentityResourceCondition::Matching,
IdentityResourceCondition::Divergent,
IdentityResourceCondition::Malformed,
];
let sessions = [
IdentitySessionCondition::Unavailable,
IdentitySessionCondition::Missing,
IdentitySessionCondition::Matching,
IdentitySessionCondition::RecoverableDivergence,
IdentitySessionCondition::AmbiguousDivergence,
IdentitySessionCondition::Malformed,
IdentitySessionCondition::IrrecoverablyCorrupt,
];
let mut classified = 0usize;
for intent in intents {
for lease in leases {
for session in sessions {
for runtime in resources {
for member in resources {
for wiring in resources {
let mut facts = matching_facts();
facts.intent = intent;
facts.lease = lease;
facts.session = session;
facts.runtime = runtime;
facts.member = member;
facts.wiring = wiring;
let expected = classify_identity_reconciliation(facts);
assert_eq!(
generated_transition_decision(&mut authority, facts),
expected,
);
classified += 1;
}
}
}
}
}
}
assert_eq!(classified, 31_500);
assert_eq!(authority.state(), &initial_state);
}
#[test]
fn classifier_is_total_over_operation_garbage_cross_shapes() {
let mut authority = crate::machines::mob_machine::MobMachineAuthority::new();
let initial_state = authority.state().clone();
let receipts = [
IdentityReceiptCondition::NotRequired,
IdentityReceiptCondition::Unavailable,
IdentityReceiptCondition::Missing,
IdentityReceiptCondition::Matching,
IdentityReceiptCondition::Conflicting,
IdentityReceiptCondition::Malformed,
];
let trusts = [
IdentityExternalTrustCondition::NotRequired,
IdentityExternalTrustCondition::Unavailable,
IdentityExternalTrustCondition::Matching,
IdentityExternalTrustCondition::Absent,
IdentityExternalTrustCondition::Contradictory,
IdentityExternalTrustCondition::Indeterminate,
IdentityExternalTrustCondition::Malformed,
];
let ceremonies = [
IdentityExternalCeremonyCondition::NotRequired,
IdentityExternalCeremonyCondition::FreshAvailable,
IdentityExternalCeremonyCondition::TemporarilyUnavailable,
IdentityExternalCeremonyCondition::AwaitFresh,
IdentityExternalCeremonyCondition::SpentOrUnknown,
];
let deliveries = [
IdentityInitialDeliveryCondition::NotRequired,
IdentityInitialDeliveryCondition::Unavailable,
IdentityInitialDeliveryCondition::ProvenAbsent,
IdentityInitialDeliveryCondition::AcceptedPendingExact,
IdentityInitialDeliveryCondition::CommittedExact,
IdentityInitialDeliveryCondition::ContentOnlyMatch,
IdentityInitialDeliveryCondition::OperationCollision,
IdentityInitialDeliveryCondition::Contradictory,
IdentityInitialDeliveryCondition::Indeterminate,
IdentityInitialDeliveryCondition::Malformed,
];
let mut classified = 0usize;
for external_required in [false, true] {
for delivery_required in [false, true] {
for external_receipt in receipts {
for external_trust in trusts {
for ceremony in ceremonies {
for delivery_receipt in receipts {
for initial_delivery in deliveries {
let mut facts = matching_facts();
facts.external_binding_required = external_required;
facts.initial_delivery_required = delivery_required;
facts.external_binding_receipt = external_receipt;
facts.external_trust = external_trust;
facts.external_ceremony = ceremony;
facts.initial_delivery_receipt = delivery_receipt;
facts.initial_delivery = initial_delivery;
let expected = classify_identity_reconciliation(facts);
assert_eq!(
generated_transition_decision(&mut authority, facts),
expected,
);
classified += 1;
}
}
}
}
}
}
}
assert_eq!(classified, 50_400);
assert_eq!(authority.state(), &initial_state);
}
#[test]
fn classifier_orders_one_obligation_per_pass() {
let mut facts = matching_facts();
facts.intent = IdentityAuthorityCondition::PresentCreateIfAbsent;
facts.session_creation_receipt = IdentityReceiptCondition::Missing;
facts.wiring = IdentityResourceCondition::Divergent;
facts.member = IdentityResourceCondition::Missing;
facts.runtime = IdentityResourceCondition::Divergent;
facts.session = IdentitySessionCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureSessionAuthority
);
facts.session = IdentitySessionCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::SealSessionCreationConsumed
);
facts.session_creation_receipt = IdentityReceiptCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureRuntimeRegistration
);
facts.runtime = IdentityResourceCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureMemberMaterialization
);
facts.member = IdentityResourceCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
}
#[test]
fn wiring_drift_and_cleanup_share_one_reconciliation_obligation() {
let mut facts = matching_facts();
facts.wiring = IdentityResourceCondition::Divergent;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.wiring = IdentityResourceCondition::Matching;
facts.session = IdentitySessionCondition::Malformed;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.intent = IdentityAuthorityCondition::Absent;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
}
#[test]
fn divergent_member_is_repair_blocked_before_create_only_actuation() {
let mut facts = matching_facts();
facts.member = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureMemberMaterialization
);
facts.member = IdentityResourceCondition::Divergent;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn malformed_evidence_beats_unrelated_unavailability() {
let mut facts = matching_facts();
facts.runtime = IdentityResourceCondition::Malformed;
facts.wiring = IdentityResourceCondition::Unavailable;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn unrelated_later_observations_do_not_block_session_or_runtime_obligations() {
let mut facts = matching_facts();
facts.intent = IdentityAuthorityCondition::PresentCreateIfAbsent;
facts.session_creation_receipt = IdentityReceiptCondition::Missing;
facts.session = IdentitySessionCondition::Missing;
facts.runtime = IdentityResourceCondition::Malformed;
facts.member = IdentityResourceCondition::Unavailable;
facts.wiring = IdentityResourceCondition::Unavailable;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureSessionAuthority
);
facts.session = IdentitySessionCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::SealSessionCreationConsumed
);
facts.session_creation_receipt = IdentityReceiptCondition::Matching;
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureRuntimeRegistration
);
}
#[test]
fn absent_cleanup_drains_known_residue_before_ambiguous_session_blocks() {
let mut facts = matching_facts();
facts.intent = IdentityAuthorityCondition::Absent;
facts.retirement_receipt = IdentityReceiptCondition::Missing;
facts.session = IdentitySessionCondition::Malformed;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.wiring = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireMemberMaterialization
);
facts.member = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireRuntimeRegistration
);
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn present_malformed_session_drains_derived_residue_before_blocking() {
let mut facts = matching_facts();
facts.session = IdentitySessionCondition::Malformed;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.wiring = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireMemberMaterialization
);
facts.member = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireRuntimeRegistration
);
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn require_existing_never_fabricates_missing_or_ambiguous_history() {
for session in [
IdentitySessionCondition::Missing,
IdentitySessionCondition::AmbiguousDivergence,
] {
let mut facts = matching_facts();
facts.session = session;
facts.wiring = IdentityResourceCondition::Missing;
facts.member = IdentityResourceCondition::Missing;
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
}
#[test]
fn create_if_absent_is_consumed_after_the_first_matching_store_authority() {
let mut facts = matching_facts();
facts.intent = IdentityAuthorityCondition::PresentCreateIfAbsent;
facts.session_creation_receipt = IdentityReceiptCondition::Missing;
facts.session = IdentitySessionCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureSessionAuthority
);
facts.session = IdentitySessionCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::SealSessionCreationConsumed
);
facts.session_creation_receipt = IdentityReceiptCondition::Matching;
facts.session = IdentitySessionCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn absent_retires_every_resource_then_seals_proof() {
let mut facts = matching_facts();
facts.intent = IdentityAuthorityCondition::Absent;
facts.retirement_receipt = IdentityReceiptCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.wiring = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireMemberMaterialization
);
facts.member = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireRuntimeRegistration
);
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReleaseSessionAuthority
);
facts.session = IdentitySessionCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::SealRetirementProven
);
facts.retirement_receipt = IdentityReceiptCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::Tombstoned
);
}
#[test]
fn external_binding_never_replays_without_fresh_ceremony() {
let mut facts = matching_facts();
facts.external_binding_required = true;
facts.external_binding_receipt = IdentityReceiptCondition::Matching;
facts.external_trust = IdentityExternalTrustCondition::Absent;
facts.external_ceremony = IdentityExternalCeremonyCondition::AwaitFresh;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::AwaitExternalBindingCeremony
);
facts.external_ceremony = IdentityExternalCeremonyCondition::FreshAvailable;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureExternalBinding
);
facts.external_ceremony = IdentityExternalCeremonyCondition::SpentOrUnknown;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RepairBlocked
);
}
#[test]
fn initial_delivery_uses_exact_runtime_input_identity() {
let mut facts = matching_facts();
facts.initial_delivery_required = true;
facts.initial_delivery_receipt = IdentityReceiptCondition::Missing;
facts.initial_delivery = IdentityInitialDeliveryCondition::ProvenAbsent;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureInitialDeliveryReceipt
);
facts.initial_delivery_receipt = IdentityReceiptCondition::Matching;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::EnsureInitialDelivery
);
facts.initial_delivery = IdentityInitialDeliveryCondition::AcceptedPendingExact;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::AwaitInitialDelivery
);
facts.initial_delivery = IdentityInitialDeliveryCondition::CommittedExact;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::Converged
);
}
#[test]
fn expired_or_missing_lease_is_reclaimable_by_construction() {
let mut facts = matching_facts();
facts.lease = IdentityLeaseCondition::HeldByExpiredIncarnation;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::AcquireLease
);
facts.lease = IdentityLeaseCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::AcquireLease
);
}
#[test]
fn corrupt_transcript_is_preserved_after_derived_residue_is_retired() {
let mut facts = matching_facts();
facts.session = IdentitySessionCondition::IrrecoverablyCorrupt;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::ReconcileWiring
);
facts.wiring = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireMemberMaterialization
);
facts.member = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::RetireRuntimeRegistration
);
facts.runtime = IdentityResourceCondition::Missing;
assert_eq!(
classify_identity_reconciliation(facts),
IdentityReconcileDecision::Quarantined
);
}
#[test]
fn external_desired_binding_cannot_carry_bootstrap_authority() {
let identity = serde_json::json!({
"kind": "ed25519_public_key",
"public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
});
let with_token = serde_json::json!({
"execution": "external",
"address": "tcp://127.0.0.1:4242",
"identity": identity.clone(),
"bootstrap_token": "must-not-decode"
});
assert!(serde_json::from_value::<DesiredExecution>(with_token).is_err());
let query_secret = serde_json::json!({
"execution": "external",
"address": "tcp://127.0.0.1:4242?mob_supervisor_bootstrap_token=secret",
"identity": identity.clone()
});
assert!(serde_json::from_value::<DesiredExecution>(query_secret).is_err());
let valid = serde_json::json!({
"execution": "external",
"address": "tcp://127.0.0.1:4242",
"identity": identity
});
assert!(serde_json::from_value::<DesiredExecution>(valid).is_ok());
}
#[test]
fn local_callback_tool_contract_is_strict_and_canonical() {
let alpha = DesiredLocalCallbackTool::new(
"alpha",
"Alpha callback",
serde_json::json!({"type": "object"}),
)
.unwrap();
let beta = DesiredLocalCallbackTool::new(
"beta",
"Beta callback",
serde_json::json!({"type": "object"}),
)
.unwrap();
validate_required_local_callback_tools(&[alpha.clone(), beta.clone()], true).unwrap();
assert!(validate_required_local_callback_tools(&[beta, alpha.clone()], true).is_err());
assert!(validate_required_local_callback_tools(&[alpha.clone(), alpha], true).is_err());
assert!(
DesiredLocalCallbackTool::new(
"invalid",
"Invalid callback",
serde_json::json!({"type": 7}),
)
.is_err()
);
assert!(
serde_json::from_value::<DesiredLocalCallbackTool>(serde_json::json!({
"name": "extra",
"description": "Extra callback",
"input_schema": {},
"handler": "must-not-enter-durable-material"
}))
.is_err()
);
}
#[test]
fn authority_digest_seals_tombstone_and_cleanup_targets() {
let identity = AgentIdentity::from("parent-1");
let intent = IdentityIntent::Absent { identity };
let mut record = IdentityIntentRecord {
schema_version: IDENTITY_INTENT_SCHEMA_VERSION,
mob_id: MobId::from("homecore"),
intent_revision: 7,
declaration_scope: None,
declaration_revision: None,
tombstone_generation: Some(3),
initial_delivery_generation_highwater: 0,
retirement_plan: IdentityRetirementPlan::NoKnownRealization,
intent_digest: intent.digest().unwrap(),
authority_digest: String::new(),
intent,
};
record.authority_digest = record.canonical_authority_digest().unwrap();
record.validate().unwrap();
let mut legacy_shape = serde_json::to_value(&record).unwrap();
legacy_shape.as_object_mut().unwrap().remove("mob_id");
assert!(
serde_json::from_value::<IdentityIntentRecord>(legacy_shape).is_err(),
"mob authority scope is required and must never default during decode"
);
let mut transplanted = record.clone();
transplanted.mob_id = MobId::from("other-mob");
assert!(matches!(
transplanted.validate(),
Err(IdentityIntentError::DigestMismatch)
));
record.tombstone_generation = Some(4);
assert!(matches!(
record.validate(),
Err(IdentityIntentError::DigestMismatch)
));
}
#[test]
fn initial_delivery_identity_is_content_sealed() {
let mut delivery =
DesiredInitialDelivery::new(1, InputId::new(), ContentInput::from("hello once"))
.unwrap();
delivery.validate().unwrap();
delivery.message = ContentInput::from("different message");
assert!(matches!(
delivery.validate(),
Err(IdentityIntentError::InvalidInitialDelivery)
));
}
#[test]
fn permit_requires_scope_incarnation_and_unexpired_claim() {
let permit = IdentityActuationPermit {
mob_id: MobId::from("homecore"),
identity: AgentIdentity::from("parent-1"),
target: IdentityActuatorTarget::Runtime,
intent_revision: 4,
intent_digest: format!("sha256:{}", "1".repeat(64)),
intent_authority_digest: format!("sha256:{}", "2".repeat(64)),
lease_epoch: 9,
lease_holder_id: "controller-a".to_string(),
lease_incarnation_id: "incarnation-a".to_string(),
lease_expires_at_ms: 200,
target_observation: IdentityTargetObservationVersion::Absent {
absence_version: "absence:17".to_string(),
},
};
permit.validate_for_write(199).unwrap();
assert!(matches!(
permit.validate_for_write(200),
Err(IdentityIntentError::ExpiredActuationPermit)
));
let mut receipt_permit = permit.clone();
receipt_permit.target = IdentityActuatorTarget::InitialDeliveryReceipt;
assert!(matches!(
receipt_permit.validate_for_write(199),
Err(IdentityIntentError::InvalidActuationPermit)
));
receipt_permit.target_observation = IdentityTargetObservationVersion::InsertIfAbsent;
receipt_permit.validate_for_write(199).unwrap();
let mut resource_with_receipt_cas = permit;
resource_with_receipt_cas.target_observation =
IdentityTargetObservationVersion::InsertIfAbsent;
assert!(matches!(
resource_with_receipt_cas.validate_for_write(199),
Err(IdentityIntentError::InvalidActuationPermit)
));
}
#[test]
fn bounded_lease_lifetime_is_enforced() {
let valid = IdentityLeaseRecord {
schema_version: IDENTITY_LEASE_SCHEMA_VERSION,
epoch_highwater: 1,
active: Some(IdentityLeaseClaim {
holder_id: "controller".to_string(),
incarnation_id: "process-a".to_string(),
epoch: 1,
renewed_at_ms: 10,
expires_at_ms: 10 + IDENTITY_LEASE_MAX_TTL_MS,
}),
};
valid.validate().unwrap();
let mut invalid = valid;
invalid.active.as_mut().unwrap().expires_at_ms += 1;
assert!(matches!(
invalid.validate(),
Err(IdentityIntentError::InvalidLeaseLifetime)
));
}
}