use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1::{
GetProjectionReceiptReply, GetProjectionReceiptRequest, PublishProjectionManifestReply,
PublishProjectionManifestRequest, RecordProjectionManifestReply,
RecordProjectionManifestRequest, RegisterProjectionPublisherReply,
RegisterProjectionPublisherRequest, ResolveProjectionManifestReply,
ResolveProjectionManifestRequest, StateProjectionCatalogService,
StateProjectionCatalogServiceExt,
};
use polyc_state::{
context::CallContext,
error::StateError,
id::{CommandId, OwnerId},
projection::{
self, ProjectionCatalog, ProjectionGeneration, ProjectionHead, ProjectionKey,
ProjectionManifest, PublishManifest, PublisherFence, PublisherId, ReadProjectionReceipt,
RecordManifest, RegisterPublisher, ResolveManifest,
},
revision::JournalSource,
};
use crate::{
admission::{
AudienceBinding, PeerIdentity, check_audience_binding, check_call_context_version,
check_not_draining, check_transport_deadline, state_audience,
},
error::to_connect_error as state_to_connect_error,
projection::error::to_connect_error,
trace::adopt_caller_trace,
wire::{DeclaredCall, Kernel, declared_call},
};
pub struct ProjectionCatalogSvc {
catalog: Arc<dyn ProjectionCatalog>,
draining: Arc<AtomicBool>,
binding: Arc<AudienceBinding>,
}
impl ProjectionCatalogSvc {
#[must_use]
pub const fn new(
catalog: Arc<dyn ProjectionCatalog>,
draining: Arc<AtomicBool>,
binding: Arc<AudienceBinding>,
) -> Self {
Self {
catalog,
draining,
binding,
}
}
#[must_use]
pub fn register_on(self, router: Router) -> Router {
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<polyc_proto::proto::polychrome::state::v1::CallContext>>,
method: &'static str,
) -> Result<(CallContext, tracing::Span), ConnectError> {
check_not_draining(self.draining.load(Ordering::Relaxed))?;
let declared: DeclaredCall =
declared_call(context).map_err(|error| state_to_connect_error(&error))?;
let family = projection::family();
check_call_context_version(declared.version)
.map_err(|error| state_to_connect_error(&error))?;
check_audience_binding(
&Self::peer(ctx),
&declared.audience,
&state_audience(),
&self.binding,
&family,
)
.map_err(|error| state_to_connect_error(&error))?;
check_transport_deadline(ctx.time_remaining(), &family)
.map_err(|error| state_to_connect_error(&error))?;
Ok((
declared.origin_relative_context(),
adopt_caller_trace(ctx.headers(), method),
))
}
}
fn required<T: Default, P: buffa::ProtoBox<T>>(
field: &str,
reason: &str,
value: buffa::MessageField<T, P>,
) -> Result<T, ConnectError> {
value.into_option().ok_or_else(|| {
to_connect_error(
&StateError::Malformed {
field: field.to_owned(),
reason: reason.to_owned(),
}
.into(),
)
})
}
#[allow(refining_impl_trait)]
impl StateProjectionCatalogService for ProjectionCatalogSvc {
async fn register_publisher(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, RegisterProjectionPublisherRequest>,
) -> ServiceResult<RegisterProjectionPublisherReply> {
let RegisterProjectionPublisherRequest {
context,
metadata,
key,
source,
publisher,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
let (context, span) = self.admit(&ctx, context, "RegisterProjectionPublisher")?;
let _entered = span.enter();
let (fence, receipt) = self
.catalog
.register(
RegisterPublisher::new(
Kernel::try_from(required(
"metadata",
"a registration carries metadata",
metadata,
)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionKey>::try_from(required(
"key",
"a registration names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<JournalSource>::try_from(required(
"source",
"a registration names one source",
source,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
PublisherId::new(publisher),
OwnerId::new(caller),
),
&context,
)
.map_err(|error| to_connect_error(&error))?;
Response::ok(RegisterProjectionPublisherReply {
fence: buffa::MessageField::some(Kernel(&fence).into()),
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
async fn record_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, RecordProjectionManifestRequest>,
) -> ServiceResult<RecordProjectionManifestReply> {
let RecordProjectionManifestRequest {
context,
metadata,
manifest,
caller,
signed_artifact,
__buffa_unknown_fields: _,
} = request.to_owned_message();
let (context, span) = self.admit(&ctx, context, "RecordProjectionManifest")?;
let _entered = span.enter();
let receipt = self
.catalog
.record(
RecordManifest::new(
Kernel::try_from(required("metadata", "a record carries metadata", metadata)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionManifest>::try_from(required(
"manifest",
"a record carries one manifest",
manifest,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
signed_artifact,
OwnerId::new(caller),
),
&context,
)
.map_err(|error| to_connect_error(&error))?;
Response::ok(RecordProjectionManifestReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
async fn publish_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, PublishProjectionManifestRequest>,
) -> ServiceResult<PublishProjectionManifestReply> {
let PublishProjectionManifestRequest {
context,
metadata,
key,
source,
expected,
publish,
publisher,
fence,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
let (context, span) = self.admit(&ctx, context, "PublishProjectionManifest")?;
let _entered = span.enter();
let receipt = self
.catalog
.publish(
PublishManifest::new(
Kernel::try_from(required(
"metadata",
"a publication carries metadata",
metadata,
)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionKey>::try_from(required(
"key",
"a publication names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<JournalSource>::try_from(required(
"source",
"a publication names one source",
source,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
ProjectionGeneration::new(expected),
ProjectionGeneration::new(publish),
PublisherId::new(publisher),
Kernel::<PublisherFence>::try_from(required(
"fence",
"a publication carries authority",
fence,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
OwnerId::new(caller),
),
&context,
)
.map_err(|error| to_connect_error(&error))?;
Response::ok(PublishProjectionManifestReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
async fn resolve_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, ResolveProjectionManifestRequest>,
) -> ServiceResult<ResolveProjectionManifestReply> {
let ResolveProjectionManifestRequest {
context,
key,
source,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
let (context, span) = self.admit(&ctx, context, "ResolveProjectionManifest")?;
let _entered = span.enter();
let resolution = self
.catalog
.resolve(
ResolveManifest::new(
Kernel::<ProjectionKey>::try_from(required(
"key",
"a resolve names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<JournalSource>::try_from(required(
"source",
"a resolve names one source",
source,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
OwnerId::new(caller),
),
&context,
)
.map_err(|error| to_connect_error(&error))?;
let (head, pending, highest_recorded) = resolution.into_parts();
let (manifest, superseded_generation, superseded_source) = match &head {
ProjectionHead::Absent => (
buffa::MessageField::default(),
None,
buffa::MessageField::default(),
),
ProjectionHead::Current(manifest) => (
buffa::MessageField::some(Kernel(manifest.as_ref()).into()),
None,
buffa::MessageField::default(),
),
ProjectionHead::Superseded { generation, source } => (
buffa::MessageField::default(),
Some(generation.get()),
buffa::MessageField::some(Kernel(source.as_ref()).into()),
),
};
Response::ok(ResolveProjectionManifestReply {
manifest,
superseded_generation,
superseded_source,
pending_manifest: pending
.as_deref()
.map_or_else(buffa::MessageField::default, |manifest| {
buffa::MessageField::some(Kernel(manifest).into())
}),
highest_recorded_generation: highest_recorded.get(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
async fn get_receipt(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, GetProjectionReceiptRequest>,
) -> ServiceResult<GetProjectionReceiptReply> {
let GetProjectionReceiptRequest {
context,
key,
source,
command_id,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
let (context, span) = self.admit(&ctx, context, "GetProjectionReceipt")?;
let _entered = span.enter();
let receipt = self
.catalog
.receipt(
ReadProjectionReceipt::new(
Kernel::<ProjectionKey>::try_from(required(
"key",
"a receipt read names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<JournalSource>::try_from(required(
"source",
"a receipt read names one source",
source,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
CommandId::new(command_id),
OwnerId::new(caller),
),
&context,
)
.map_err(|error| to_connect_error(&error))?;
Response::ok(GetProjectionReceiptReply {
receipt: receipt
.as_ref()
.map_or_else(buffa::MessageField::default, |receipt| {
buffa::MessageField::some(Kernel(receipt).into())
}),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}