polyc-state-connect 2026.10.2

State plane transport adapter: capability-specific Connect clients and server-trait glue mapping the generated wire types onto the polyc-state kernel — typed outcomes, per-call admission, and the conformance surface the authenticated shell proves itself against (docs/proposals/separated-planes.md).
//! Server adapter for the authoritative persona-memory journal.

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::{
    error::{OutageReach, StateError},
    persona_memory::journal::{
        MemoryJournalPartition, MemoryReplayRange, PersonaMemoryHistory,
        PersonaMemoryJournalAuthority,
    },
    revision::PartitionIncarnation,
};

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::{
    append_from_wire, destroy_from_wire, history_to_wire, lineage_page_to_wire, migrate_from_wire,
    partition_page_to_wire, replay_to_wire, rewrite_from_wire,
};

/// State's authoritative persona-memory journal service.
pub struct PersonaMemoryJournalSvc {
    authority: Arc<dyn PersonaMemoryJournalAuthority>,
    history: Arc<dyn PersonaMemoryHistory>,
    draining: Arc<AtomicBool>,
    binding: Arc<AudienceBinding>,
}

impl PersonaMemoryJournalSvc {
    /// Builds the service.
    #[must_use]
    pub const fn new(
        authority: Arc<dyn PersonaMemoryJournalAuthority>,
        history: Arc<dyn PersonaMemoryHistory>,
        draining: Arc<AtomicBool>,
        binding: Arc<AudienceBinding>,
    ) -> Self {
        Self {
            authority,
            history,
            draining,
            binding,
        }
    }

    /// Registers the capability-specific routes.
    #[must_use]
    pub fn register_on(self, router: Router) -> Router {
        use pb::StatePersonaMemoryJournalServiceExt 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<
        (
            polyc_state::context::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::persona_memory::journal::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, namespace: &str) -> Result<(), ConnectError> {
        if namespace == self.authority.namespace().as_str() {
            Ok(())
        } else {
            Err(to_connect_error(&StateError::Denied {
                family: polyc_state::persona_memory::journal::family(),
            }))
        }
    }
}

#[allow(refining_impl_trait)]
impl pb::StatePersonaMemoryJournalService for PersonaMemoryJournalSvc {
    async fn append(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::AppendPersonaMemoryRecordsRequest>,
    ) -> ServiceResult<pb::AppendPersonaMemoryRecordsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "AppendPersonaMemoryRecords")?;
        let command = append_from_wire(required(message.command, "append")?)
            .map_err(|error| to_connect_error(&error))?;
        self.namespace(command.write().metadata().scope().namespace().as_str())?;
        let authority = Arc::clone(&self.authority);
        let outcome = blocking::mutation(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.append(command, &context),
        )
        .await?;
        Response::ok(pb::AppendPersonaMemoryRecordsReply {
            receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(outcome.receipt()))),
            positions: outcome.positions().to_vec(),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn rewrite(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::RewritePersonaMemoryRecordsRequest>,
    ) -> ServiceResult<pb::RewritePersonaMemoryRecordsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "RewritePersonaMemoryRecords")?;
        let command = rewrite_from_wire(required(message.command, "rewrite")?)
            .map_err(|error| to_connect_error(&error))?;
        self.namespace(command.write().metadata().scope().namespace().as_str())?;
        let authority = Arc::clone(&self.authority);
        let outcome = blocking::mutation(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.rewrite(command, &context),
        )
        .await?;
        Response::ok(pb::RewritePersonaMemoryRecordsReply {
            receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(outcome.receipt()))),
            dropped: outcome.dropped() as u64,
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn migrate(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::MigratePersonaMemoryRecordsRequest>,
    ) -> ServiceResult<pb::MigratePersonaMemoryRecordsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "MigratePersonaMemoryRecords")?;
        let command = migrate_from_wire(required(message.command, "migration")?)
            .map_err(|error| to_connect_error(&error))?;
        self.namespace(command.write().metadata().scope().namespace().as_str())?;
        let authority = Arc::clone(&self.authority);
        let outcome = blocking::mutation(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.migrate(command, &context),
        )
        .await?;
        Response::ok(pb::MigratePersonaMemoryRecordsReply {
            receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(outcome.receipt()))),
            copied: outcome.copied() as u64,
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn destroy(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::DestroyPersonaMemoryRecordsRequest>,
    ) -> ServiceResult<pb::DestroyPersonaMemoryRecordsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "DestroyPersonaMemoryRecords")?;
        let command = destroy_from_wire(required(message.command, "destruction")?)
            .map_err(|error| to_connect_error(&error))?;
        self.namespace(command.write().metadata().scope().namespace().as_str())?;
        let authority = Arc::clone(&self.authority);
        let receipt = blocking::mutation(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.destroy(command, &context),
        )
        .await?;
        Response::ok(pb::DestroyPersonaMemoryRecordsReply {
            receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(&receipt))),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        })
    }

    async fn replay(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::ReplayPersonaMemoryRecordsRequest>,
    ) -> ServiceResult<pb::ReplayPersonaMemoryRecordsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "ReplayPersonaMemoryRecords")?;
        self.namespace(&message.namespace)?;
        let partition = MemoryJournalPartition::parse(message.partition)
            .map_err(|error| to_connect_error(&error))?;
        let range = MemoryReplayRange::new(message.start, message.max_records, message.max_bytes);
        let authority = Arc::clone(&self.authority);
        let page = blocking::read(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.replay(&partition, range, &context),
        )
        .await?;
        if page.signed_head().is_none() {
            return Err(to_connect_error(&StateError::Unavailable {
                family: polyc_state::persona_memory::journal::family(),
                reach: OutageReach::PossiblyApplied,
            }));
        }
        Response::ok(replay_to_wire(&page))
    }

    async fn list(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::ListPersonaMemoryPartitionsRequest>,
    ) -> ServiceResult<pb::ListPersonaMemoryPartitionsReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "ListPersonaMemoryPartitions")?;
        self.namespace(&message.namespace)?;
        let authority = Arc::clone(&self.authority);
        let page = blocking::read(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || authority.list(message.after.as_deref(), message.limit, &context),
        )
        .await?;
        Response::ok(partition_page_to_wire(&page))
    }

    async fn read_history(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::ReadPersonaMemoryHistoryRequest>,
    ) -> ServiceResult<pb::ReadPersonaMemoryHistoryReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "ReadPersonaMemoryHistory")?;
        self.namespace(&message.namespace)?;
        let partition = MemoryJournalPartition::parse(message.partition)
            .map_err(|error| to_connect_error(&error))?;
        let expected_incarnation = message
            .expected_incarnation
            .map(|bytes| {
                fixed_bytes::<{ PartitionIncarnation::LEN }>("expected_incarnation", &bytes)
                    .map(PartitionIncarnation::from_bytes)
            })
            .transpose()
            .map_err(|error: StateError| to_connect_error(&error))?;
        let range = MemoryReplayRange::new(message.start, message.max_records, message.max_bytes);
        let history = Arc::clone(&self.history);
        let page = blocking::read(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || history.read_history(&partition, expected_incarnation, range, &context),
        )
        .await?;
        Response::ok(history_to_wire(&page))
    }

    async fn list_lineage(
        &self,
        ctx: RequestContext,
        request: ServiceRequest<'_, pb::ListPersonaMemoryLineageRequest>,
    ) -> ServiceResult<pb::ListPersonaMemoryLineageReply> {
        let message = request.to_owned_message();
        let (context, span, wait) =
            self.admit(&ctx, message.context, "ListPersonaMemoryLineage")?;
        self.namespace(&message.namespace)?;
        let history = Arc::clone(&self.history);
        let page = blocking::read(
            span,
            polyc_state::persona_memory::journal::family(),
            wait,
            move || history.list_lineage(message.after.as_deref(), message.limit, &context),
        )
        .await?;
        Response::ok(lineage_page_to_wire(&page))
    }
}

// Named over `P: buffa::ProtoBox<T>` rather than `impl Into<Option<T>>`: see
// `crate::wire::required`'s doc comment.
fn required<T: Default, P: buffa::ProtoBox<T>>(
    field: buffa::MessageField<T, P>,
    name: &'static str,
) -> Result<T, ConnectError> {
    field.into_option().ok_or_else(|| {
        to_connect_error(&StateError::Malformed {
            field: name.into(),
            reason: format!("a request carries its {name}"),
        })
    })
}