use crate::dagster::{self, DagsterBackend};
use crate::error;
use crate::models::*;
use crate::service::ProcessesService;
use actix_files::NamedFile;
use actix_web::{
http::header::ContentEncoding, http::StatusCode, web, Either, HttpRequest, HttpResponse,
};
use bbox_core::service::ServiceEndpoints;
use log::{info, warn};
use serde_json::json;
async fn process_list(_req: HttpRequest) -> HttpResponse {
let backend = DagsterBackend::new();
let jobs = backend.process_list().await.unwrap_or_else(|e| {
warn!("Dagster backend error: {e}");
Vec::new()
});
let processes = jobs
.iter()
.map(|job| {
let mut process = ProcessSummary::new(job.name.clone(), "1.0.0".to_string());
process.description.clone_from(&job.description);
process
})
.collect::<Vec<_>>();
let resp = ProcessList {
processes,
links: Vec::new(),
};
HttpResponse::Ok().json(resp)
}
async fn get_process_description(process_id: web::Path<String>) -> HttpResponse {
let backend = DagsterBackend::new();
match backend.get_process_description(&process_id).await {
Ok(descr) => HttpResponse::Ok().json(descr), Err(error::Error::NotFound(type_)) => HttpResponse::NotFound().json(Exception::new(type_)),
Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
}
}
async fn execute(
process_id: web::Path<String>,
parameters: web::Json<dagster::Execute>,
req: HttpRequest,
) -> JobResultResponse {
info!("Execute `{process_id}` with parameters `{parameters:?}`");
let backend = DagsterBackend::new();
let prefer_async = req
.headers()
.get("Prefer")
.and_then(|headerval| {
headerval
.to_str()
.ok()
.map(|headerstr| headerstr.contains("respond-async"))
})
.unwrap_or(false);
if prefer_async {
let resp = match backend.execute(&process_id, ¶meters).await {
Ok(status) => HttpResponse::build(StatusCode::CREATED).json(status),
Err(error::Error::NotFound(type_)) => {
HttpResponse::NotFound().json(Exception::new(type_))
}
Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
};
Either::Left(resp)
} else {
let job_result = backend.execute_sync(&process_id, ¶meters).await;
job_result_response(job_result)
}
}
async fn get_jobs() -> HttpResponse {
let backend = DagsterBackend::new();
match backend.get_jobs().await {
Ok(jobs) => HttpResponse::Ok().json(jobs), Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
}
}
async fn get_status(job_id: web::Path<String>) -> HttpResponse {
let backend = DagsterBackend::new();
match backend.get_status(&job_id).await {
Ok(status) => HttpResponse::Ok().json(status),
Err(error::Error::NotFound(type_)) => HttpResponse::NotFound().json(Exception::new(type_)),
Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
}
}
async fn dismiss(job_id: web::Path<String>) -> HttpResponse {
HttpResponse::InternalServerError().json(job_id.to_string())
}
pub enum JobResult {
FilePath(String),
Json(serde_json::Value),
}
type JobResultResponse = Either<HttpResponse, std::result::Result<NamedFile, std::io::Error>>;
async fn get_result(job_id: web::Path<String>) -> JobResultResponse {
let backend = DagsterBackend::new();
let job_result = backend.get_result(&job_id).await;
job_result_response(job_result)
}
fn job_result_response(job_result: crate::error::Result<JobResult>) -> JobResultResponse {
match job_result {
Ok(result) => match result {
JobResult::Json(json) => {
Either::Left(HttpResponse::Ok().json(json!({ "result": {"value": json }})))
}
JobResult::FilePath(path) => {
info!("get_result from {path}");
let file = NamedFile::open(path)
.unwrap()
.set_content_encoding(ContentEncoding::Identity);
Either::Right(Ok(file))
}
},
Err(error::Error::NotFound(type_)) => {
Either::Left(HttpResponse::NotFound().json(Exception::new(type_)))
}
Err(e) => Either::Left(HttpResponse::InternalServerError().json(Exception::from(e))),
}
}
impl ServiceEndpoints for ProcessesService {
fn register_endpoints(&self, cfg: &mut web::ServiceConfig) {
if self.backend.is_none() {
return;
}
cfg.service(web::resource("/processes").route(web::get().to(process_list)))
.service(
web::resource("/processes/{processID}")
.route(web::get().to(get_process_description)),
)
.service(
web::resource("/processes/{processID}/execution").route(web::post().to(execute)),
)
.service(web::resource("/jobs").route(web::get().to(get_jobs)))
.service(web::resource("/jobs/{jobId}").route(web::get().to(get_status)))
.service(web::resource("/jobs/{jobId}").route(web::delete().to(dismiss)))
.service(web::resource("/jobs/{jobId}/results").route(web::get().to(get_result)));
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::ProcessesServiceCfg;
use actix_web::{body, dev::Service, http, test, App, Error};
#[actix_web::test]
#[ignore]
async fn test_process_list() -> Result<(), Error> {
if !ProcessesServiceCfg::from_config().has_backend() {
return Ok(());
}
let app = test::init_service(
App::new().service(web::resource("/processes").route(web::get().to(process_list))),
)
.await;
let req = test::TestRequest::get().uri("/processes").to_request();
let resp = app.call(req).await.unwrap();
assert_eq!(resp.status(), http::StatusCode::OK);
let response_body = body::to_bytes(resp.into_body()).await?;
println!("{response_body:?}");
assert!(response_body.starts_with(b"{\"processes\":["));
Ok(())
}
}