use std::{
future::Future,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
};
use connectrpc::{ConnectError, RequestContext, Response, Router, ServiceRequest, ServiceResult};
use polyc_proto::proto::polychrome::state::v1::{
BeginCollectionReply, BeginCollectionRequest, CompleteCollectionReply,
CompleteCollectionRequest, GetCollectionStateReply, GetCollectionStateRequest,
GetManifestByGenerationReply, GetManifestByGenerationRequest, GetProjectionReceiptReply,
GetProjectionReceiptRequest, OpenCollectionState, PublishProjectionManifestReply,
PublishProjectionManifestRequest, RecordProjectionManifestReply,
RecordProjectionManifestRequest, RegisterProjectionPublisherReply,
RegisterProjectionPublisherRequest, ResolveProjectionManifestReply,
ResolveProjectionManifestRequest, RetireGenerationsReply, RetireGenerationsRequest,
StateProjectionCatalogService, StateProjectionCatalogServiceExt,
};
use polyc_state::{
context::CallContext,
error::StateError,
feed::ProjectionSource,
id::{CommandId, OwnerId},
projection::{
self, ProjectionCatalog, ProjectionGeneration, ProjectionHead, ProjectionKey,
ProjectionManifest, PublishManifest, PublisherFence, PublisherId, ReadCollectionState,
ReadManifestByGeneration, ReadProjectionReceipt, RecordCollectionIntent,
RecordCollectionReceipt, RecordManifest, RegisterPublisher, ResolveManifest,
RetireGenerations, artifact::ExactObjectRef,
},
};
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 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, std::time::Duration), 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),
slot_wait(ctx.time_remaining(), &declared),
))
}
}
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 {
fn register_publisher(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, RegisterProjectionPublisherRequest>,
) -> impl Future<Output = ServiceResult<RegisterProjectionPublisherReply>> {
let RegisterProjectionPublisherRequest {
context,
metadata,
key,
source_authority,
publisher,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "RegisterProjectionPublisher")?;
let command = 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::<ProjectionSource>::try_from(required(
"source_authority",
"a registration names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
PublisherId::new(publisher),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let (fence, receipt) =
blocking::mutation(span, projection::family(), wait, move || {
catalog.register(command, &context)
})
.await?;
Response::ok(RegisterProjectionPublisherReply {
fence: buffa::MessageField::some(Kernel(&fence).into()),
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn record_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, RecordProjectionManifestRequest>,
) -> impl Future<Output = ServiceResult<RecordProjectionManifestReply>> {
let RecordProjectionManifestRequest {
context,
metadata,
manifest,
caller,
signed_artifact,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "RecordProjectionManifest")?;
let command = 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),
);
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::mutation(span, projection::family(), wait, move || {
catalog.record(command, &context)
})
.await?;
Response::ok(RecordProjectionManifestReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn publish_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, PublishProjectionManifestRequest>,
) -> impl Future<Output = ServiceResult<PublishProjectionManifestReply>> {
let PublishProjectionManifestRequest {
context,
metadata,
key,
source_authority,
expected,
publish,
publisher,
fence,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "PublishProjectionManifest")?;
let command = 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::<ProjectionSource>::try_from(required(
"source_authority",
"a publication names one source",
source_authority,
)?)
.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),
);
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::mutation(span, projection::family(), wait, move || {
catalog.publish(command, &context)
})
.await?;
Response::ok(PublishProjectionManifestReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn resolve_manifest(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, ResolveProjectionManifestRequest>,
) -> impl Future<Output = ServiceResult<ResolveProjectionManifestReply>> {
let ResolveProjectionManifestRequest {
context,
key,
source_authority,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "ResolveProjectionManifest")?;
let query = 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::<ProjectionSource>::try_from(required(
"source_authority",
"a resolve names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let resolution = blocking::read(span, projection::family(), wait, move || {
catalog.resolve(query, &context)
})
.await?;
let (head, pending, highest_recorded) = resolution.into_parts();
let (manifest, superseded_generation, superseded_source_authority) = 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_authority,
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(),
})
}
}
fn get_receipt(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, GetProjectionReceiptRequest>,
) -> impl Future<Output = ServiceResult<GetProjectionReceiptReply>> {
let GetProjectionReceiptRequest {
context,
key,
source_authority,
command_id,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "GetProjectionReceipt")?;
let query = 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::<ProjectionSource>::try_from(required(
"source_authority",
"a receipt read names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
CommandId::new(command_id),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::read(span, projection::family(), wait, move || {
catalog.receipt(query, &context)
})
.await?;
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(),
})
}
}
fn retire(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, RetireGenerationsRequest>,
) -> impl Future<Output = ServiceResult<RetireGenerationsReply>> {
let RetireGenerationsRequest {
context,
metadata,
key,
source_authority,
expected,
floor,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "RetireGenerations")?;
let command = RetireGenerations::new(
Kernel::try_from(required("metadata", "a retire carries metadata", metadata)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionKey>::try_from(required("key", "a retire names one key", key)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionSource>::try_from(required(
"source_authority",
"a retire names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
ProjectionGeneration::new(expected),
ProjectionGeneration::new(floor),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::mutation(span, projection::family(), wait, move || {
catalog.retire(command, &context)
})
.await?;
Response::ok(RetireGenerationsReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn begin_collection(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, BeginCollectionRequest>,
) -> impl Future<Output = ServiceResult<BeginCollectionReply>> {
let BeginCollectionRequest {
context,
metadata,
key,
source_authority,
generation,
objects,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "BeginCollection")?;
let objects: Vec<ExactObjectRef> = objects
.into_iter()
.map(|object| {
Kernel::<ExactObjectRef>::try_from(object)
.map(Kernel::into_inner)
.map_err(|error: StateError| to_connect_error(&error.into()))
})
.collect::<Result<_, _>>()?;
let command = RecordCollectionIntent::try_new(
Kernel::try_from(required(
"metadata",
"a begin-collection carries metadata",
metadata,
)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionKey>::try_from(required(
"key",
"a begin-collection names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionSource>::try_from(required(
"source_authority",
"a begin-collection names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
ProjectionGeneration::new(generation),
objects,
OwnerId::new(caller),
)
.map_err(|error: StateError| to_connect_error(&error.into()))?;
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::mutation(span, projection::family(), wait, move || {
catalog.begin_collection(command, &context)
})
.await?;
Response::ok(BeginCollectionReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn complete_collection(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, CompleteCollectionRequest>,
) -> impl Future<Output = ServiceResult<CompleteCollectionReply>> {
let CompleteCollectionRequest {
context,
metadata,
key,
source_authority,
generation,
intent_command_id,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "CompleteCollection")?;
let command = RecordCollectionReceipt::new(
Kernel::try_from(required(
"metadata",
"a complete-collection carries metadata",
metadata,
)?)
.map_err(|error: StateError| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionKey>::try_from(required(
"key",
"a complete-collection names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionSource>::try_from(required(
"source_authority",
"a complete-collection names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
ProjectionGeneration::new(generation),
CommandId::new(intent_command_id),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let receipt = blocking::mutation(span, projection::family(), wait, move || {
catalog.complete_collection(command, &context)
})
.await?;
Response::ok(CompleteCollectionReply {
receipt: buffa::MessageField::some(Kernel(&receipt).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_collection_state(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, GetCollectionStateRequest>,
) -> impl Future<Output = ServiceResult<GetCollectionStateReply>> {
let GetCollectionStateRequest {
context,
key,
source_authority,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "GetCollectionState")?;
let query = ReadCollectionState::new(
Kernel::<ProjectionKey>::try_from(required(
"key",
"a collection-state read names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionSource>::try_from(required(
"source_authority",
"a collection-state read names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let state = blocking::read(span, projection::family(), wait, move || {
catalog.collection_state(query, &context)
})
.await?;
let open = state.open().map(|open| OpenCollectionState {
generation: open.generation().get(),
objects: open
.objects()
.iter()
.map(|object| Kernel(object).into())
.collect(),
intent_command_id: open.intent().as_str().to_owned(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
});
Response::ok(GetCollectionStateReply {
retired_floor: state.retired_floor().get(),
retired_at_nanos: state
.retired_at()
.map(polyc_state::deadline::MonotonicInstant::as_nanos),
open: open.map_or_else(buffa::MessageField::default, buffa::MessageField::some),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
fn get_manifest_by_generation(
&self,
ctx: RequestContext,
request: ServiceRequest<'_, GetManifestByGenerationRequest>,
) -> impl Future<Output = ServiceResult<GetManifestByGenerationReply>> {
let GetManifestByGenerationRequest {
context,
key,
source_authority,
generation,
caller,
__buffa_unknown_fields: _,
} = request.to_owned_message();
async move {
let (context, span, wait) = self.admit(&ctx, context, "GetManifestByGeneration")?;
let query = ReadManifestByGeneration::new(
Kernel::<ProjectionKey>::try_from(required(
"key",
"a manifest read names one key",
key,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
Kernel::<ProjectionSource>::try_from(required(
"source_authority",
"a manifest read names one source",
source_authority,
)?)
.map_err(|error| to_connect_error(&error.into()))?
.into_inner(),
ProjectionGeneration::new(generation),
OwnerId::new(caller),
);
let catalog = Arc::clone(&self.catalog);
let manifest = blocking::read(span, projection::family(), wait, move || {
catalog.manifest_by_generation(query, &context)
})
.await?;
Response::ok(GetManifestByGenerationReply {
manifest: buffa::MessageField::some(Kernel(&manifest).into()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
}
}