use std::sync::Arc;
use khive_gate::{check_with_mailbox_policy, mailbox_read_owner};
use khive_storage::note::{Note, NoteMailboxScope};
use serde_json::Value;
use sha2::{Digest, Sha256};
use crate::{
engine_config::ActorConfig, ActorRef, Gate, GateDecision, GateError, GateRef, GateRequest,
KhiveRuntime, MailboxPolicyError, MailboxReadGate, NamespaceToken, RuntimeError, RuntimeResult,
};
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct MailboxView {
pub actor_id: String,
pub delegated: bool,
}
impl MailboxView {
fn legacy_local(&self, token: &NamespaceToken) -> bool {
!self.delegated && token.actor().is_anonymous() && token.actor().id == "local"
}
pub fn note_scope(&self, token: &NamespaceToken) -> NoteMailboxScope {
NoteMailboxScope {
actor_id: self.actor_id.clone(),
legacy_local: self.legacy_local(token),
}
}
pub fn permits_message_note(&self, token: &NamespaceToken, note: &Note) -> bool {
permits_message_note(&self.actor_id, self.legacy_local(token), note)
}
pub fn scope_permits_message_note(scope: &NoteMailboxScope, note: &Note) -> bool {
permits_message_note(&scope.actor_id, scope.legacy_local, note)
}
}
fn permits_message_note(actor_id: &str, legacy_local: bool, note: &Note) -> bool {
if note.kind != "message" {
return true;
}
let properties = note.properties.as_ref();
let text = |key: &str| -> Option<&str> {
properties
.and_then(|value| value.get(key))
.and_then(Value::as_str)
};
match text("direction") {
Some("inbound") => {
text("to_actor") == Some(actor_id)
|| (legacy_local
&& properties
.and_then(|p| p.get("to_actor"))
.is_none_or(Value::is_null))
}
Some("outbound") => {
text("from_actor") == Some(actor_id)
|| (legacy_local
&& properties
.and_then(|p| p.get("from_actor"))
.is_none_or(Value::is_null))
}
None => {
legacy_local
&& properties
.and_then(|p| p.get("direction"))
.is_none_or(Value::is_null)
&& text("from_actor").is_none()
&& text("to_actor").is_none()
}
_ => false,
}
}
pub(crate) fn validate_mailbox_request(req: &GateRequest) -> RuntimeResult<()> {
mailbox_read_owner(req)
.map(|_| ())
.map_err(|error| RuntimeError::InvalidInput(error.to_string()))
}
impl KhiveRuntime {
pub fn authorize_mailbox_view(
&self,
token: &NamespaceToken,
verb: &str,
selector: Option<&str>,
args: &Value,
) -> RuntimeResult<MailboxView> {
let selector_field = match verb {
"comm.inbox" | "comm.thread" => Some("mailbox_actor"),
"comm.probe" => Some("actor"),
"list" | "search" | "get" | "context" | "neighbors" => None,
_ => {
return Err(RuntimeError::InvalidInput(
"mailbox views are supported only by comm.inbox, comm.thread, comm.probe, list, search, get, context and neighbors"
.into(),
));
}
};
let req = GateRequest::new(
token.actor().clone(),
token.gate_namespace().clone(),
verb,
args.clone(),
);
validate_mailbox_request(&req)?;
if let Some(selector_field) = selector_field {
if args.get(selector_field).and_then(Value::as_str) != selector {
return Err(RuntimeError::InvalidInput(format!(
"mailbox selector must match the original {selector_field} argument"
)));
}
} else if selector.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"{verb} has no cross-actor mailbox selector"
)));
}
match check_with_mailbox_policy(self.config().gate.as_ref(), &req) {
Ok(GateDecision::Allow { .. }) => {
let owner = mailbox_read_owner(&req)
.map_err(|error| RuntimeError::InvalidInput(error.to_string()))?;
Ok(MailboxView {
delegated: owner.is_some(),
actor_id: owner.map_or_else(|| token.actor().id.clone(), |owner| owner.id),
})
}
Ok(GateDecision::Deny { reason }) => Err(RuntimeError::permission_denied(verb, reason)),
Err(error) => Err(RuntimeError::GateUnavailable {
verb: verb.to_string(),
reason: error.wire_reason().to_string(),
}),
}
}
}
impl ActorConfig {
pub(crate) fn mailbox_gate(
&self,
inner: GateRef,
) -> Result<Option<MailboxReadGate>, MailboxPolicyError> {
if self.mailbox_readers.is_empty() {
return Ok(None);
}
if self.mailbox_readers.len() > 256 {
return Err(MailboxPolicyError::TooManyReaders);
}
if self
.mailbox_readers
.iter()
.any(|id| !crate::is_valid_mailbox_actor_label(id))
{
return Err(MailboxPolicyError::InvalidReader);
}
let owner = self.id.as_deref().ok_or(MailboxPolicyError::InvalidOwner)?;
crate::Namespace::parse(owner).map_err(|_| MailboxPolicyError::InvalidOwner)?;
let readers = self
.mailbox_readers
.iter()
.map(|id| ActorRef::new("actor", id.clone()))
.collect();
MailboxReadGate::new(inner, ActorRef::new("actor", owner), readers).map(Some)
}
}
pub(crate) fn configured_mailbox_gate(actor: &ActorConfig, inner: GateRef) -> GateRef {
match actor.mailbox_gate(inner.clone()) {
Ok(Some(gate)) => Arc::new(gate),
Ok(None) => inner,
Err(_) => {
let mut hash = Sha256::new();
hash.update(b"khive.invalid-mailbox-read-policy.v1\0");
for field in [inner.configuration_fingerprint(), actor.id.as_deref()] {
match field {
Some(value) => {
hash.update([1]);
hash.update((value.len() as u64).to_be_bytes());
hash.update(value.as_bytes());
}
None => hash.update([0]),
}
}
hash.update((actor.mailbox_readers.len() as u64).to_be_bytes());
for reader in &actor.mailbox_readers {
hash.update((reader.len() as u64).to_be_bytes());
hash.update(reader.as_bytes());
}
Arc::new(InvalidMailboxConfigGate {
fingerprint: format!("sha256:{:x}", hash.finalize()),
})
}
}
}
#[derive(Debug)]
struct InvalidMailboxConfigGate {
fingerprint: String,
}
impl Gate for InvalidMailboxConfigGate {
fn check(&self, _req: &GateRequest) -> Result<GateDecision, GateError> {
Err(GateError::Policy(
"invalid mailbox reader configuration".into(),
))
}
fn configuration_fingerprint(&self) -> Option<&str> {
Some(&self.fingerprint)
}
}
#[cfg(test)]
mod tests;