use crate::api::ApiError::Internal;
use crate::batch_messages::{
ActivityBatchRequestMessageType, ActivityBatchResponseMessageType, HasPEPBatchInfo,
PEPBatchMessageType, Record,
};
use crate::data_source::DataSource;
use crate::messages::{
ActivityRequestMessageType, ActivityResponseMessageType, HasPEPParticipantInfo, PEPMessageType,
};
use crate::pseudonym_service_factory::{create_pseudonym_service_pool, OauthAuth};
use crate::pseudonym_service_pool::PseudonymServicePool;
use actix_web::web::ServiceConfig;
use actix_web::{error::ErrorUnauthorized, middleware::Logger, web, HttpResponse, ResponseError};
use actix_web_httpauth::middleware::HttpAuthentication;
use libpep::data::json::EncryptedPEPJSONValue;
use log::error;
use oauth_token_service::TokenServiceConfig;
use paas_client::prelude::PAASConfig;
use paas_client::pseudonym_service::PseudonymServiceError;
use std::fmt::{Display, Formatter};
use std::future::Future;
use std::sync::Arc;
#[derive(Debug)]
pub enum ApiError {
NotFound,
Internal(String),
BadRequest(String),
}
impl Display for ApiError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
ApiError::NotFound => write!(f, "Not found"),
ApiError::Internal(msg) => write!(f, "Internal error: {}", msg),
ApiError::BadRequest(msg) => write!(f, "Bad request: {}", msg),
}
}
}
impl ResponseError for ApiError {
fn error_response(&self) -> HttpResponse {
match self {
ApiError::NotFound => HttpResponse::NotFound().json("Not found"),
ApiError::Internal(msg) => HttpResponse::InternalServerError().json(msg),
ApiError::BadRequest(msg) => HttpResponse::BadRequest().json(msg),
}
}
}
#[derive(Clone)]
pub struct DataSourceAPI {
connector: Arc<dyn DataSource>,
pseudonym_service: PseudonymServicePool,
nest_auth_token: String,
}
impl DataSourceAPI {
pub async fn init(
connector: impl DataSource,
nest_auth_token: String,
ps_config: PAASConfig,
oauth_config: TokenServiceConfig,
pool_size: usize,
) -> Self {
let oauth_auth = OauthAuth::new(Some(oauth_config))
.await
.expect("Failed to create OAuth authentication");
let pseudonym_service =
create_pseudonym_service_pool(oauth_auth, ps_config, pool_size).await;
Self {
connector: Arc::new(connector),
pseudonym_service,
nest_auth_token,
}
}
pub fn configure_service(&self, cfg: &mut ServiceConfig) {
let auth_token = self.nest_auth_token.clone();
let connector = self.connector.clone();
let auth = HttpAuthentication::bearer(move |req, credentials| {
let auth_token = auth_token.clone();
async move {
if credentials.token() == auth_token {
Ok(req)
} else {
Err((ErrorUnauthorized("Invalid token"), req))
}
}
});
cfg.app_data(web::Data::new(connector)).service(
web::scope("/api")
.app_data(web::Data::new(self.pseudonym_service.clone()))
.wrap(auth)
.route(
"/activities/create",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.create_participant_activity(r), req, ps
)
.await
}
}),
)
.route(
"/activities/create_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.create_participant_activity_batch(r), req, ps
)
.await
}
}),
)
.route(
"/activities/update",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.update_participant_activity(r), req, ps
)
.await
}
}),
)
.route(
"/activities/update_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.update_participant_activity_batch(r), req, ps
)
.await
}
}),
)
.route(
"/activities/delete",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.delete_participant_activity(r), req, ps
)
.await
}
}),
)
.route(
"/activities/delete_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.delete_participant_activity_batch(r), req, ps
)
.await
}
}),
)
.route(
"/activities/details",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.participant_activity_details(r), req, ps
)
.await
}
}),
)
.route(
"/activities/details_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.participant_activity_details_batch(r), req, ps
)
.await
}
}),
)
.route(
"/activities/progress",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.participant_activity_progress(r), req, ps
)
.await
}
}),
)
.route(
"/activities/progress_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.participant_activity_progress_batch(r),
req,
ps,
)
.await
}
}),
)
.route(
"/activities/result",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep::<
ActivityRequestMessageType,
ActivityResponseMessageType,
_,
_,
>(
|r| connector.participant_activity_result(r), req, ps
)
.await
}
}), )
.route(
"/activities/result_batch",
web::post().to(|connector: web::Data<Arc<dyn DataSource>>, req, ps| {
let connector = connector.clone();
async move {
handle_pep_batch::<
ActivityBatchRequestMessageType,
ActivityBatchResponseMessageType,
EncryptedPEPJSONValue,
Record,
_,
_,
>(
|r| connector.participant_activity_result_batch(r), req, ps
)
.await
}
}),
),
);
}
pub async fn run_standalone(self, host: &str, port: u16) -> std::io::Result<()> {
use actix_web::{App, HttpServer};
HttpServer::new(move || {
App::new()
.wrap(Logger::default())
.configure(|cfg| self.configure_service(cfg))
})
.bind((host, port))?
.run()
.await
}
}
async fn handle_pep<Req, Res, F, Fut>(
handle: F,
request: web::Json<Req::PEPMessage>,
pseudonym_service: web::Data<PseudonymServicePool>,
) -> Result<HttpResponse, ApiError>
where
Req: PEPMessageType, Res: PEPMessageType, F: FnOnce(Req::Message) -> Fut, Fut: Future<Output = Result<Option<Res::Message>, ApiError>>, {
let mut rng = rand::rng();
let pep_request = request.into_inner();
let domain_to = pep_request.domain_to().clone();
let request = {
let mut ps = pseudonym_service.acquire().await;
Req::unpack(pep_request, &mut ps).await?
};
let response: Option<Res::Message> = handle(request).await?;
if let Some(response) = response {
let mut ps = pseudonym_service.acquire().await;
let pep_response = Res::pack(response, domain_to, &mut ps, &mut rng).await?;
Ok(HttpResponse::Ok().json(pep_response))
} else {
Ok(HttpResponse::Ok().finish())
}
}
async fn handle_pep_batch<Req, Res, PepT, T, F, Fut>(
handle: F,
request: web::Json<Req::PEPBatchMessage>,
pseudonym_service: web::Data<PseudonymServicePool>,
) -> Result<HttpResponse, ApiError>
where
Req: PEPBatchMessageType<PepT, T>, Res: PEPBatchMessageType<PepT, T>, F: FnOnce(Req::BatchMessage) -> Fut, Fut: Future<Output = Result<Option<Res::BatchMessage>, ApiError>>, {
let pep_request = request.into_inner();
let domain_from = pep_request.domain().clone();
let request = {
let mut ps = pseudonym_service.acquire().await;
Req::unpack(pep_request, &mut ps).await?
};
let response: Option<Res::BatchMessage> = handle(request).await?;
if let Some(response) = response {
let mut ps = pseudonym_service.acquire().await;
let pep_response = Res::pack(response, domain_from, &mut ps)?;
Ok(HttpResponse::Ok().json(pep_response))
} else {
Ok(HttpResponse::Ok().finish())
}
}
pub(crate) trait PseudonymServiceErrorHandler<T> {
fn handle_pseudonym_error(self) -> Result<T, ApiError>;
}
impl<T> PseudonymServiceErrorHandler<T> for Result<T, PseudonymServiceError> {
fn handle_pseudonym_error(self) -> Result<T, ApiError> {
self.map_err(|e| {
error!("Pseudonym service error: {e:?}");
Internal("Pseudonym service operation failed".to_string())
})
}
}