use std::future::Future;
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
context::CallContext,
error::StateError,
id::{CommandId, NamespaceId},
page::PageCompleteness,
persona::{
LinkCodeDigest, ParticipationScope, PersonaChangeRead, PersonaDirectory, PersonaRead,
PersonaWrite,
},
revision::{JournalPosition, PartitionIncarnation},
versioned::{
VersionedRead, VersionedTransact,
authority::{PersonaId, ScopeId},
},
};
use crate::{
admission::{
AudienceBinding, PeerIdentity, check_audience_binding, check_call_context_version,
check_not_draining, check_transport_deadline, slot_wait, state_audience,
},
blocking,
error::to_connect_error,
trace::adopt_caller_trace,
wire::{DeclaredCall, Kernel, declared_call, fixed_bytes},
};
use super::wire::{
attempt_key_from_wire, attempts_to_wire, challenge_to_wire, change_to_wire, command_from_wire,
identity_from_wire, link_to_wire, merge_id_ticket_to_wire, merge_request_to_wire,
participation_to_wire, profile_to_wire, visibility_to_wire,
};
pub trait PersonaAuthority: PersonaRead + PersonaWrite + PersonaChangeRead {
fn namespace(&self) -> &NamespaceId;
}
impl<V> PersonaAuthority for PersonaDirectory<V>
where
V: VersionedRead + VersionedTransact,
{
fn namespace(&self) -> &NamespaceId {
self.namespace()
}
}
pub struct PersonaSvc {
authority: Arc<dyn PersonaAuthority>,
draining: Arc<AtomicBool>,
binding: Arc<AudienceBinding>,
}
impl PersonaSvc {
#[must_use]
pub const fn new(
authority: Arc<dyn PersonaAuthority>,
draining: Arc<AtomicBool>,
binding: Arc<AudienceBinding>,
) -> Self {
Self {
authority,
draining,
binding,
}
}
#[must_use]
pub fn register_on(self, router: Router) -> Router {
use pb::StatePersonaServiceExt as _;
Arc::new(self).register(router)
}
fn peer(ctx: &RequestContext) -> PeerIdentity {
PeerIdentity::from_verified_leaf(
ctx.peer_certs().and_then(<[_]>::first).map(|leaf| &**leaf),
)
}
fn admit(
&self,
ctx: &RequestContext,
context: impl Into<Option<pb::CallContext>>,
method: &'static str,
) -> Result<(CallContext, tracing::Span, std::time::Duration), ConnectError> {
check_not_draining(self.draining.load(Ordering::Relaxed))?;
let declared: DeclaredCall =
declared_call(context).map_err(|error| to_connect_error(&error))?;
let family = polyc_state::versioned::family();
check_call_context_version(declared.version).map_err(|error| to_connect_error(&error))?;
check_audience_binding(
&Self::peer(ctx),
&declared.audience,
&state_audience(),
&self.binding,
&family,
)
.map_err(|error| to_connect_error(&error))?;
check_transport_deadline(ctx.time_remaining(), &family)
.map_err(|error| to_connect_error(&error))?;
Ok((
declared.origin_relative_context(),
adopt_caller_trace(ctx.headers(), method),
slot_wait(ctx.time_remaining(), &declared),
))
}
fn namespace(&self, value: &str) -> Result<(), ConnectError> {
if value == self.authority.namespace().as_str() {
Ok(())
} else {
Err(to_connect_error(&StateError::Denied {
family: polyc_state::versioned::family(),
}))
}
}
}
#[allow(refining_impl_trait)]
impl pb::StatePersonaService for PersonaSvc {
fn transact(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::TransactPersonaRequest>,
) -> impl Future<Output = ServiceResult<pb::TransactPersonaReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, message.context, "TransactPersona")?;
let metadata = message.metadata.into_option().ok_or_else(|| {
to_connect_error(&StateError::Malformed {
field: "metadata".into(),
reason: "a persona command carries metadata".into(),
})
})?;
self.namespace(&metadata.namespace)?;
let command = command_from_wire(metadata, message.mutations)
.map_err(|error| to_connect_error(&error))?;
let authority = Arc::clone(&self.authority);
let receipt =
blocking::mutation(span, polyc_state::versioned::family(), wait, move || {
authority.transact(command, &context)
})
.await?;
Response::ok(pb::TransactPersonaReply {
receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(&receipt))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn resolve_identity(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::ResolvePersonaIdentityRequest>,
) -> impl Future<Output = ServiceResult<pb::ResolvePersonaIdentityReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "ResolvePersonaIdentity")?;
self.namespace(&message.namespace)?;
let identity = message.identity.into_option().ok_or_else(|| {
to_connect_error(&StateError::Malformed {
field: "identity".into(),
reason: "a resolution names an external identity".into(),
})
})?;
let identity = identity_from_wire(identity);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.resolve(&identity, &context)
})
.await?;
Response::ok(pb::ResolvePersonaIdentityReply {
link: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(link_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_profile(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaProfileRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaProfileReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, message.context, "GetPersonaProfile")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.profile(&persona, &context)
})
.await?;
Response::ok(pb::GetPersonaProfileReply {
profile: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(profile_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_snapshot(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaSnapshotRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaSnapshotReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, message.context, "GetPersonaSnapshot")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.snapshot(&persona, &context)
})
.await?;
Response::ok(pb::GetPersonaSnapshotReply {
profile: fact
.value()
.profile
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(profile_to_wire(value))
}),
identity_links: fact
.value()
.identity_links
.iter()
.map(link_to_wire)
.collect(),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_participation(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaParticipationRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaParticipationReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaParticipation")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let scope = ScopeId::new(message.scope);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.participation(&persona, &scope, &context)
})
.await?;
Response::ok(pb::GetPersonaParticipationReply {
participation: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(participation_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_search_visibility(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaSearchVisibilityRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaSearchVisibilityReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaSearchVisibility")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let scope = ScopeId::new(message.scope);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.search_visibility(&persona, &scope, &context)
})
.await?;
Response::ok(pb::GetPersonaSearchVisibilityReply {
visibility: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(visibility_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn list_participations(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::ListPersonaParticipationsRequest>,
) -> impl Future<Output = ServiceResult<pb::ListPersonaParticipationsReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "ListPersonaParticipations")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let cap = usize::try_from(message.cap).unwrap_or(usize::MAX);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.participations(&persona, cap, &context)
})
.await?;
Response::ok(pb::ListPersonaParticipationsReply {
participations: fact.value().iter().map(participation_to_wire).collect(),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn list_personas(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::ListPersonasRequest>,
) -> impl Future<Output = ServiceResult<pb::ListPersonasReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, message.context, "ListPersonas")?;
self.namespace(&message.namespace)?;
let cap = usize::try_from(message.cap).unwrap_or(usize::MAX);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.personas(cap, &context)
})
.await?;
Response::ok(pb::ListPersonasReply {
personas: fact
.value()
.iter()
.map(|persona| persona.as_str().to_owned())
.collect(),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn resolve_participation_scope(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::ResolvePersonaParticipationScopeRequest>,
) -> impl Future<Output = ServiceResult<pb::ResolvePersonaParticipationScopeReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "ResolveParticipationScope")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let cap = usize::try_from(message.cap).unwrap_or(usize::MAX);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.participation_scope(&persona, cap, &context)
})
.await?;
let (visible_scopes, refused_count) = match fact.value() {
ParticipationScope::Resolved(scopes) => (
scopes
.iter()
.map(|scope| scope.as_str().to_owned())
.collect(),
None,
),
ParticipationScope::RefusedOverCap { count } => {
(Vec::new(), Some(u64::try_from(*count).unwrap_or(u64::MAX)))
}
};
Response::ok(pb::ResolvePersonaParticipationScopeReply {
visible_scopes,
refused_count,
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_link_challenge(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaLinkChallengeRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaLinkChallengeReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaLinkChallenge")?;
self.namespace(&message.namespace)?;
let digest = LinkCodeDigest::new(message.digest);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.link_challenge(&digest, &context)
})
.await?;
Response::ok(pb::GetPersonaLinkChallengeReply {
challenge: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(challenge_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_link_attempts(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaLinkAttemptsRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaLinkAttemptsReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaLinkAttempts")?;
self.namespace(&message.namespace)?;
let key = attempt_key_from_wire(message.key.into_option().ok_or_else(|| {
to_connect_error(&StateError::Malformed {
field: "key".into(),
reason: "a link-attempt read names its counter".into(),
})
})?)
.map_err(|error| to_connect_error(&error))?;
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.link_attempts(&key, &context)
})
.await?;
Response::ok(pb::GetPersonaLinkAttemptsReply {
attempts: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(attempts_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_pending_invite(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaPendingInviteRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaPendingInviteReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaPendingInvite")?;
self.namespace(&message.namespace)?;
let identity = message.identity.into_option().ok_or_else(|| {
to_connect_error(&StateError::Malformed {
field: "identity".into(),
reason: "a pending-invite read names its target identity".into(),
})
})?;
let identity = identity_from_wire(identity);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.pending_invite(&identity, &context)
})
.await?;
Response::ok(pb::GetPersonaPendingInviteReply {
digest: fact.value().as_ref().map(|value| value.as_str().to_owned()),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_merge_request(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaMergeRequestRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaMergeRequestReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaMergeRequest")?;
self.namespace(&message.namespace)?;
let digest = LinkCodeDigest::new(message.digest);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.merge_request(&digest, &context)
})
.await?;
Response::ok(pb::GetPersonaMergeRequestReply {
request: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(merge_request_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_outgoing_merge(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaOutgoingMergeRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaOutgoingMergeReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaOutgoingMerge")?;
self.namespace(&message.namespace)?;
let initiator = PersonaId::new(message.initiator);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.outgoing_merge(&initiator, &context)
})
.await?;
Response::ok(pb::GetPersonaOutgoingMergeReply {
digest: fact.value().as_ref().map(|value| value.as_str().to_owned()),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_merge_id_ticket(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaMergeIdTicketRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaMergeIdTicketReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaMergeIdTicket")?;
self.namespace(&message.namespace)?;
let digest = LinkCodeDigest::new(message.digest);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.merge_id_ticket(&digest, &context)
})
.await?;
Response::ok(pb::GetPersonaMergeIdTicketReply {
ticket: fact
.value()
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(merge_id_ticket_to_wire(value))
}),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_outgoing_merge_id(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaOutgoingMergeIdRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaOutgoingMergeIdReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) =
self.admit(&ctx, message.context, "GetPersonaOutgoingMergeId")?;
self.namespace(&message.namespace)?;
let persona = PersonaId::new(message.persona);
let authority = Arc::clone(&self.authority);
let fact = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.outgoing_merge_id(&persona, &context)
})
.await?;
Response::ok(pb::GetPersonaOutgoingMergeIdReply {
digest: fact.value().as_ref().map(|value| value.as_str().to_owned()),
snapshot_revision: fact.snapshot_revision().get(),
entry_revision: fact
.entry_revision()
.map(polyc_state::revision::Revision::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn list_changes(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::ListStatePersonaChangesRequest>,
) -> impl Future<Output = ServiceResult<pb::ListStatePersonaChangesReply>> {
let message = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, message.context, "ListPersonaChanges")?;
self.namespace(&message.namespace)?;
let after = message.after.map(JournalPosition::new);
let expected_lineage = message
.expected_lineage
.as_deref()
.map(|bytes| {
fixed_bytes::<32>("expected_lineage", bytes)
.map(PartitionIncarnation::from_bytes)
})
.transpose()
.map_err(|error| to_connect_error(&error))?;
let limit = message.limit;
let authority = Arc::clone(&self.authority);
let page = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.changes(after, expected_lineage, limit, &context)
})
.await?;
Response::ok(pb::ListStatePersonaChangesReply {
changes: page.changes().iter().map(change_to_wire).collect(),
next_after: page.next().map(|cursor| cursor.position().get()),
completeness: buffa::EnumValue::Known(match page.completeness() {
PageCompleteness::Complete => pb::PageCompleteness::Complete,
PageCompleteness::Truncated => pb::PageCompleteness::Truncated,
}),
frontier: page.frontier().get(),
lineage: page.lineage().as_bytes().to_vec(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_receipt(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, pb::GetPersonaReceiptRequest>,
) -> impl Future<Output = ServiceResult<pb::GetPersonaReceiptReply>> {
let message = request.to_owned_message();
async move {
let (_context, span, wait) = self.admit(&ctx, message.context, "GetPersonaReceipt")?;
self.namespace(&message.namespace)?;
let namespace = NamespaceId::new(message.namespace);
let command_id = CommandId::new(message.command_id);
let authority = Arc::clone(&self.authority);
let receipt = blocking::read(span, polyc_state::versioned::family(), wait, move || {
authority.committed_receipt(&namespace, &command_id)
})
.await?;
Response::ok(pb::GetPersonaReceiptReply {
receipt: receipt
.as_ref()
.map_or_else(buffa::MessageField::none, |value| {
buffa::MessageField::some(pb::Receipt::from(Kernel(value)))
}),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
}