durable-actors 0.7.10

Standalone regional durable-actors control plane, host, and durability runtime
use axum::{
    Json,
    extract::{Path, State, rejection::JsonRejection},
    http::{HeaderMap, StatusCode, header},
    response::{IntoResponse, Response},
};
use serde::Deserialize;
use serde_json::{Value, json};

use super::public_api::{
    ActorPath, ActorTargetReply, ApiError, FindActorRequest, PublicApiState, resolve_actor_target,
};
use crate::actor::ActorInvocation;

pub(super) async fn invoke(
    State(state): State<PublicApiState>,
    Path(path): Path<ActorPath>,
    headers: HeaderMap,
    body: Result<Json<InvokeRequest>, JsonRejection>,
) -> Result<Response, ApiError> {
    let actor = path.into_actor();
    let Json(request) = body.map_err(ApiError::json)?;
    if let Some(epoch) = request.owner_epoch {
        return invoke_cached(&state, actor, &headers, request, epoch).await;
    }
    state
        .admin
        .authorize_discovery(
            headers
                .get(header::AUTHORIZATION)
                .and_then(|value| value.to_str().ok())
                .unwrap_or(""),
            &actor.project_id,
        )
        .map_err(|_| ApiError::unauthorized("actor invocation credential was rejected"))?;
    let invocation = ActorInvocation {
        actor,
        request_id: request.request_id,
        method: request.method,
        args: request.args,
    };
    invocation.validate().map_err(ApiError::bad_request)?;
    let (target, outcome) =
        resolve_and_dispatch(&state, &headers, request.home_region, &invocation).await?;
    Ok((
        [(header::CACHE_CONTROL, "no-store")],
        Json(json!({"target": target, "outcome": outcome})),
    )
        .into_response())
}

async fn resolve_and_dispatch(
    state: &PublicApiState,
    headers: &HeaderMap,
    home_region: Option<String>,
    invocation: &ActorInvocation,
) -> Result<(ActorTargetReply, Value), ApiError> {
    let deadline = tokio::time::Instant::now() + super::CONTROL_PLANE_REQUEST_TIMEOUT;
    loop {
        let target = resolve_actor_target(
            state,
            &invocation.actor,
            headers,
            Ok(Json(FindActorRequest {
                home_region: home_region.clone(),
            })),
        )
        .await?;
        let outcome = dispatch(
            &state.hosts,
            &target.backend_route,
            &target.token,
            target.owner_epoch,
            invocation,
        )
        .await?;
        if !matches!(
            outcome["type"].as_str(),
            Some("not_executed" | "unauthenticated")
        ) || tokio::time::Instant::now() >= deadline
        {
            return Ok((target, outcome));
        }
        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
    }
}

async fn invoke_cached(
    state: &PublicApiState,
    actor: crate::actor::ActorKey,
    headers: &HeaderMap,
    request: InvokeRequest,
    epoch: u64,
) -> Result<Response, ApiError> {
    let authorization = headers
        .get(header::AUTHORIZATION)
        .and_then(|value| value.to_str().ok())
        .unwrap_or("");
    let gateway = state
        .invocations
        .gateway
        .as_ref()
        .ok_or_else(|| ApiError::unauthorized("gateway is not configured"))?;
    let route = gateway
        .invocation_route(&actor, epoch, authorization)
        .map_err(|_| ApiError::unauthorized("actor invocation credential was rejected"))?;
    let token = authorization
        .strip_prefix("Bearer ")
        .ok_or_else(|| ApiError::unauthorized("bearer token required"))?;
    let invocation = ActorInvocation {
        actor,
        request_id: request.request_id,
        method: request.method,
        args: request.args,
    };
    invocation.validate().map_err(ApiError::bad_request)?;
    let outcome = dispatch(&state.hosts, &route, token, epoch, &invocation).await?;
    Ok(([(header::CACHE_CONTROL, "no-store")], Json(outcome)).into_response())
}

pub(super) async fn publish(
    State(state): State<PublicApiState>,
    Path(path): Path<ActorPath>,
    headers: HeaderMap,
    body: Result<Json<PublishRequest>, JsonRejection>,
) -> Result<Response, ApiError> {
    let actor = path.into_actor();
    let authorization = headers
        .get(header::AUTHORIZATION)
        .and_then(|value| value.to_str().ok())
        .unwrap_or("");
    let gateway = state
        .invocations
        .gateway
        .as_ref()
        .ok_or_else(|| ApiError::unauthorized("gateway is not configured"))?;
    if authorization.is_empty() {
        return Err(ApiError::unauthorized("bearer token required"));
    }
    let Json(request) = body.map_err(ApiError::json)?;
    let route = gateway
        .invocation_route(&actor, request.owner_epoch, authorization)
        .map_err(|_| ApiError::unauthorized("actor invocation credential was rejected"))?;
    crate::actor::validate_socket_effects(&request.effects).map_err(ApiError::bad_request)?;
    let url = format!(
        "{}/v1/projects/{}/actors/{}/{}/socket-effects",
        route.trim_end_matches('/'),
        actor.project_id,
        actor.actor_name,
        actor.actor_id
    );
    let response = state
        .hosts
        .post(url)
        .header(header::AUTHORIZATION, authorization)
        .json(&request)
        .send()
        .await
        .map_err(|_| outcome_unknown())?;
    let status = response.status();
    let body = response.bytes().await.map_err(|_| outcome_unknown())?;
    Ok((status, [(header::CACHE_CONTROL, "no-store")], body).into_response())
}

#[derive(serde::Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct PublishRequest {
    owner_epoch: u64,
    effects: Vec<crate::actor::ActorSocketEffect>,
}

async fn dispatch(
    client: &reqwest::Client,
    route: &str,
    token: &str,
    owner_epoch: u64,
    invocation: &ActorInvocation,
) -> Result<Value, ApiError> {
    let actor = &invocation.actor;
    let url = format!(
        "{}/v1/projects/{}/actors/{}/{}/invoke",
        route.trim_end_matches('/'),
        actor.project_id,
        actor.actor_name,
        actor.actor_id
    );
    let response = client
        .post(url)
        .bearer_auth(token)
        .json(&json!({
            "requestId":invocation.request_id, "ownerEpoch":owner_epoch,
            "method":invocation.method, "args":invocation.args,
        }))
        .send()
        .await;
    let response = match response {
        Ok(response) => response,
        Err(error) if error.is_connect() => {
            return Ok(json!({"type":"not_executed", "reason":"upstream_not_reached"}));
        }
        Err(_) => return Err(outcome_unknown()),
    };
    if response.status() == StatusCode::UNAUTHORIZED {
        return Ok(json!({"type":"unauthenticated"}));
    }
    if !response.status().is_success() {
        return Err(outcome_unknown());
    }
    let reply = read_reply(response).await?;
    match reply.get("type").and_then(Value::as_str) {
        Some("not_executed")
            if matches!(
                reply["reason"].as_str(),
                Some("stale_owner" | "host_unavailable" | "upstream_not_reached")
            ) =>
        {
            Ok(reply)
        }
        Some("completed") if reply.get("result").is_some() => Ok(reply),
        Some("failed")
            if reply["code"].as_str().is_some_and(|code| !code.is_empty())
                && reply["message"].is_string() =>
        {
            Ok(reply)
        }
        _ => Err(outcome_unknown()),
    }
}

async fn read_reply(response: reqwest::Response) -> Result<Value, ApiError> {
    response.json().await.map_err(|_| outcome_unknown())
}

fn outcome_unknown() -> ApiError {
    ApiError::new(
        StatusCode::BAD_GATEWAY,
        "outcome_unknown",
        "actor invocation was dispatched but its outcome could not be confirmed",
    )
}

#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct InvokeRequest {
    request_id: String,
    owner_epoch: Option<u64>,
    method: String,
    args: Vec<Value>,
    home_region: Option<String>,
}