durable-actors 0.3.1

Standalone regional durable-actors control plane, host, and durability runtime
Documentation
use std::sync::Arc;

use tonic::{Request, Response, Status};
use tracing::warn;

use super::proto::{
    ActivateActorReply, ActivateActorRequest, Empty, HostInvokeActorRequest,
    HostSocketEventRequest, InvokeActorReply, PublishSocketEffectsRequest,
    actor_host_service_server::{ActorHostService, ActorHostServiceServer},
};
use crate::{
    actor::{
        ActorExecutionResult, ActorInvocation, ActorInvocationFailure,
        MAX_ACTOR_EXECUTOR_MESSAGE_BYTES,
    },
    control_plane::{ActorJwtVerifier, ActorPrincipal},
    host::{ActorHost, HostId},
};

pub(crate) struct ActorHostGrpcService {
    host_id: HostId,
    session_id: String,
    host: Arc<ActorHost>,
    invocation_auth: ActorJwtVerifier,
    sockets: Arc<crate::host::sockets::HostSockets>,
}

impl ActorHostGrpcService {
    pub(crate) fn new(
        host: Arc<ActorHost>,
        session_id: String,
        auth: ActorJwtVerifier,
        sockets: Arc<crate::host::sockets::HostSockets>,
    ) -> Self {
        Self {
            host_id: host.id().clone(),
            session_id,
            host,
            invocation_auth: auth,
            sockets,
        }
    }

    pub(crate) fn into_service(self) -> ActorHostServiceServer<Self> {
        ActorHostServiceServer::new(self)
            .max_decoding_message_size(MAX_ACTOR_EXECUTOR_MESSAGE_BYTES)
            .max_encoding_message_size(MAX_ACTOR_EXECUTOR_MESSAGE_BYTES)
    }
}

#[tonic::async_trait]
impl ActorHostService for ActorHostGrpcService {
    async fn activate(
        &self,
        request: Request<ActivateActorRequest>,
    ) -> Result<Response<ActivateActorReply>, Status> {
        let principal = self.invocation_auth.authenticate(&request).await?;
        validate_activation(&principal, &self.host_id, &self.session_id)?;
        let actor: crate::actor::ActorKey = request
            .into_inner()
            .actor
            .ok_or_else(|| Status::invalid_argument("actor is required"))?
            .into();
        let activation = self
            .host
            .activate_actor(actor)
            .await
            .map_err(|error| Status::unavailable(format!("{error:#}")))?;
        Ok(Response::new(ActivateActorReply {
            owner_epoch: activation.owner_epoch,
        }))
    }

    async fn invoke(
        &self,
        request: Request<HostInvokeActorRequest>,
    ) -> Result<Response<InvokeActorReply>, Status> {
        let invocation = self.authorize_invocation(request).await?;
        self.invoke_authorized(invocation).await
    }

    async fn publish_socket_effects(
        &self,
        request: Request<PublishSocketEffectsRequest>,
    ) -> Result<Response<Empty>, Status> {
        let principal = self.invocation_auth.authenticate(&request).await?;
        if principal.session_id != self.session_id {
            return Err(Status::permission_denied(
                "actor credential belongs to another host session",
            ));
        }
        let request = request.into_inner();
        let actor: crate::actor::ActorKey = request
            .actor
            .ok_or_else(|| Status::invalid_argument("actor is required"))?
            .into();
        actor
            .validate()
            .map_err(|error| Status::invalid_argument(error.to_string()))?;
        validate_host_request(&principal, &self.host_id, &actor, request.owner_epoch)?;
        let effects =
            serde_json::from_slice::<Vec<crate::actor::ActorSocketEffect>>(&request.effects_json)
                .map_err(|_| Status::invalid_argument("invalid socket effects JSON"))?;
        crate::actor::validate_socket_effects(&effects)
            .map_err(|error| Status::invalid_argument(error.to_string()))?;
        self.sockets
            .publish_authorized(&actor, &self.host_id, request.owner_epoch, effects)
            .await
            .map_err(|error| {
                Status::unavailable(format!("socket effects could not be delivered: {error:#}"))
            })?;
        Ok(Response::new(Empty {}))
    }

    async fn handle_socket(
        &self,
        request: Request<HostSocketEventRequest>,
    ) -> Result<Response<InvokeActorReply>, Status> {
        let principal = self.invocation_auth.authenticate(&request).await?;
        if principal.session_id != self.session_id {
            return Err(Status::permission_denied(
                "actor credential belongs to another host session",
            ));
        }
        let request = request.into_inner();
        let owner_epoch = request.owner_epoch;

        let invocation: crate::actor::ActorSocketInvocation = request
            .try_into()
            .map_err(|error| Status::invalid_argument(format!("{error:#}")))?;
        validate_host_request(&principal, &self.host_id, &invocation.actor, owner_epoch)?;
        let result = self
            .host
            .handle_socket_event(invocation, owner_epoch)
            .await
            .map_err(|error| {
                Status::unavailable(format!("actor socket event failed: {error:#}"))
            })?;
        Ok(Response::new(InvokeActorReply::from(result)))
    }
}

impl ActorHostGrpcService {
    async fn authorize_invocation(
        &self,
        request: Request<HostInvokeActorRequest>,
    ) -> Result<AuthorizedHostInvocation, Status> {
        let principal = self.invocation_auth.authenticate(&request).await?;
        if principal.session_id != self.session_id {
            return Err(Status::permission_denied(
                "actor credential belongs to another host session",
            ));
        }
        let request = request.into_inner();
        let invocation: ActorInvocation = request
            .invocation
            .ok_or_else(|| Status::invalid_argument("actor invocation is required"))?
            .try_into()
            .map_err(|error| Status::invalid_argument(format!("{error:#}")))?;
        validate_host_request(
            &principal,
            &self.host_id,
            &invocation.actor,
            request.owner_epoch,
        )?;
        Ok(AuthorizedHostInvocation {
            invocation,
            owner_epoch: request.owner_epoch,
        })
    }

    async fn invoke_authorized(
        &self,
        request: AuthorizedHostInvocation,
    ) -> Result<Response<InvokeActorReply>, Status> {
        let request_id = request.invocation.request_id.clone();
        let result = match self
            .host
            .invoke_actor(request.invocation, request.owner_epoch)
            .await
        {
            Ok(result) => result,
            Err(error) => {
                warn!(request_id, error = %format!("{error:#}"), "actor invocation failed before execution");
                ActorExecutionResult::Failed {
                    failure: ActorInvocationFailure {
                        code: "unavailable".into(),
                        message: "actor could not start because its state was unavailable".into(),
                    },
                }
            }
        };
        Ok(Response::new(InvokeActorReply::from(result)))
    }
}

fn validate_activation(
    principal: &ActorPrincipal,
    host: &HostId,
    session: &str,
) -> Result<(), Status> {
    if principal.host_id != *host
        || principal.session_id != session
        || principal.invocation.is_some()
    {
        return Err(Status::permission_denied(
            "host activation authority is required",
        ));
    }
    Ok(())
}

fn validate_host_request(
    principal: &ActorPrincipal,
    host_id: &HostId,
    actor: &crate::actor::ActorKey,
    owner_epoch: u64,
) -> Result<(), Status> {
    if owner_epoch == 0 {
        return Err(Status::invalid_argument(
            "actor ownership capability is incomplete",
        ));
    }
    if principal.host_id != *host_id {
        return Err(Status::permission_denied(
            "actor invocation credential is not for this host",
        ));
    }
    if let Some(capability) = &principal.invocation
        && (capability.actor != *actor
            || capability.host_id != *host_id
            || capability.owner_epoch != owner_epoch)
    {
        return Err(Status::permission_denied(
            "actor invocation does not match its direct capability",
        ));
    }
    Ok(())
}

struct AuthorizedHostInvocation {
    invocation: ActorInvocation,
    owner_epoch: u64,
}

#[cfg(test)]
#[path = "../../tests/unit/grpc/service.rs"]
mod tests;