bbox_processes_server/
endpoints.rs1use crate::dagster::{self, DagsterBackend};
4use crate::error;
5use crate::models::*;
6use crate::service::ProcessesService;
7use actix_files::NamedFile;
8use actix_web::{
9 http::header::ContentEncoding, http::StatusCode, web, Either, HttpRequest, HttpResponse,
10};
11use bbox_core::service::ServiceEndpoints;
12use log::{info, warn};
13use serde_json::json;
14
15async fn process_list(_req: HttpRequest) -> HttpResponse {
17 let backend = DagsterBackend::new();
18 let jobs = backend.process_list().await.unwrap_or_else(|e| {
19 warn!("Dagster backend error: {e}");
20 Vec::new()
21 });
22 let processes = jobs
23 .iter()
24 .map(|job| {
25 let mut process = ProcessSummary::new(job.name.clone(), "1.0.0".to_string());
26 process.description.clone_from(&job.description);
27 process
28 })
29 .collect::<Vec<_>>();
30 let resp = ProcessList {
89 processes,
90 links: Vec::new(),
91 };
92
93 HttpResponse::Ok().json(resp)
94}
95
96async fn get_process_description(process_id: web::Path<String>) -> HttpResponse {
98 let backend = DagsterBackend::new();
99 match backend.get_process_description(&process_id).await {
100 Ok(descr) => HttpResponse::Ok().json(descr), Err(error::Error::NotFound(type_)) => HttpResponse::NotFound().json(Exception::new(type_)),
102 Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
103 }
104}
105
106async fn execute(
108 process_id: web::Path<String>,
109 parameters: web::Json<dagster::Execute>,
110 req: HttpRequest,
111) -> JobResultResponse {
112 info!("Execute `{process_id}` with parameters `{parameters:?}`");
113 let backend = DagsterBackend::new();
114 let prefer_async = req
115 .headers()
116 .get("Prefer")
117 .and_then(|headerval| {
118 headerval
119 .to_str()
120 .ok()
121 .map(|headerstr| headerstr.contains("respond-async"))
122 })
123 .unwrap_or(false);
124 if prefer_async {
126 let resp = match backend.execute(&process_id, ¶meters).await {
127 Ok(status) => HttpResponse::build(StatusCode::CREATED).json(status),
132 Err(error::Error::NotFound(type_)) => {
133 HttpResponse::NotFound().json(Exception::new(type_))
134 }
135 Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
136 };
137 Either::Left(resp)
138 } else {
139 let job_result = backend.execute_sync(&process_id, ¶meters).await;
140 job_result_response(job_result)
142 }
143}
144
145async fn get_jobs() -> HttpResponse {
147 let backend = DagsterBackend::new();
148 match backend.get_jobs().await {
149 Ok(jobs) => HttpResponse::Ok().json(jobs), Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
151 }
152}
153
154async fn get_status(job_id: web::Path<String>) -> HttpResponse {
156 let backend = DagsterBackend::new();
157 match backend.get_status(&job_id).await {
158 Ok(status) => HttpResponse::Ok().json(status),
159 Err(error::Error::NotFound(type_)) => HttpResponse::NotFound().json(Exception::new(type_)),
160 Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
161 }
162}
163
164async fn dismiss(job_id: web::Path<String>) -> HttpResponse {
166 HttpResponse::InternalServerError().json(job_id.to_string())
167}
168
169pub enum JobResult {
170 FilePath(String),
171 Json(serde_json::Value),
172}
173
174type JobResultResponse = Either<HttpResponse, std::result::Result<NamedFile, std::io::Error>>;
175
176async fn get_result(job_id: web::Path<String>) -> JobResultResponse {
178 let backend = DagsterBackend::new();
179 let job_result = backend.get_result(&job_id).await;
180 job_result_response(job_result)
181}
182
183fn job_result_response(job_result: crate::error::Result<JobResult>) -> JobResultResponse {
184 match job_result {
185 Ok(result) => match result {
186 JobResult::Json(json) => {
187 Either::Left(HttpResponse::Ok().json(json!({ "result": {"value": json }})))
189 }
190 JobResult::FilePath(path) => {
191 info!("get_result from {path}");
192 let file = NamedFile::open(path)
200 .unwrap()
201 .set_content_encoding(ContentEncoding::Identity);
202 Either::Right(Ok(file))
203 }
204 },
205 Err(error::Error::NotFound(type_)) => {
206 Either::Left(HttpResponse::NotFound().json(Exception::new(type_)))
207 }
208 Err(e) => Either::Left(HttpResponse::InternalServerError().json(Exception::from(e))),
209 }
210}
211
212impl ServiceEndpoints for ProcessesService {
213 fn register_endpoints(&self, cfg: &mut web::ServiceConfig) {
214 if self.backend.is_none() {
215 return;
216 }
217 cfg.service(web::resource("/processes").route(web::get().to(process_list)))
218 .service(
219 web::resource("/processes/{processID}")
220 .route(web::get().to(get_process_description)),
221 )
222 .service(
223 web::resource("/processes/{processID}/execution").route(web::post().to(execute)),
224 )
225 .service(web::resource("/jobs").route(web::get().to(get_jobs)))
226 .service(web::resource("/jobs/{jobId}").route(web::get().to(get_status)))
227 .service(web::resource("/jobs/{jobId}").route(web::delete().to(dismiss)))
228 .service(web::resource("/jobs/{jobId}/results").route(web::get().to(get_result)));
229 }
230}
231
232#[cfg(test)]
233mod tests {
234 use super::*;
235 use crate::config::ProcessesServiceCfg;
236 use actix_web::{body, dev::Service, http, test, App, Error};
237
238 #[actix_web::test]
239 #[ignore]
240 async fn test_process_list() -> Result<(), Error> {
241 if !ProcessesServiceCfg::from_config().has_backend() {
242 return Ok(());
243 }
244 let app = test::init_service(
245 App::new().service(web::resource("/processes").route(web::get().to(process_list))),
246 )
247 .await;
248
249 let req = test::TestRequest::get().uri("/processes").to_request();
250 let resp = app.call(req).await.unwrap();
251
252 assert_eq!(resp.status(), http::StatusCode::OK);
253
254 let response_body = body::to_bytes(resp.into_body()).await?;
255 println!("{response_body:?}");
256 assert!(response_body.starts_with(b"{\"processes\":["));
257
258 Ok(())
259 }
260}