Skip to main content

bbox_processes_server/
endpoints.rs

1//! Endpoints according to <https://ogcapi.ogc.org/processes/> API
2
3use 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
15/// retrieve the list of available processes
16async 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    /* Example:
31    {
32      "processes": [
33        {
34          "id": "EchoProcess",
35          "title": "EchoProcess",
36          "version": "1.0.0",
37          "jobControlOptions": [
38            "async-execute",
39            "sync-execute"
40          ],
41          "outputTransmission": [
42            "value",
43            "reference"
44          ],
45          "additionalParameters": {
46            "title": "string",
47            "role": "string",
48            "href": "string",
49            "parameters": [
50              {
51                "name": "string",
52                "value": [
53                  "string",
54                  0,
55                  0,
56                  [
57                    null
58                  ],
59                  {}
60                ]
61              }
62            ]
63          },
64          "links": [
65            {
66              "href": "https://processing.example.org/oapi-p/processes/EchoProcess",
67              "type": "application/json",
68              "rel": "self",
69              "title": "process description"
70            }
71          ]
72        }
73      ],
74      "links": [
75        {
76          "href": "https://processing.example.org/oapi-p/processes?f=json",
77          "rel": "self",
78          "type": "application/json"
79        },
80        {
81          "href": "https://processing.example.org/oapi-p/processes?f=html",
82          "rel": "alternate",
83          "type": "text/html"
84        }
85      ]
86    }
87    */
88    let resp = ProcessList {
89        processes,
90        links: Vec::new(),
91    };
92
93    HttpResponse::Ok().json(resp)
94}
95
96/// retrieve a process description
97async 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), // TODO: type ProcessDescription
101        Err(error::Error::NotFound(type_)) => HttpResponse::NotFound().json(Exception::new(type_)),
102        Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
103    }
104}
105
106/// execute a process
107async 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    // TODO: support sync/async-only processes
125    if prefer_async {
126        let resp = match backend.execute(&process_id, &parameters).await {
127            /* responses:
128                    200:
129                      $ref: 'http://schemas.opengis.net/ogcapi/processes/part1/1.0/openapi/responses/ExecuteSync.yaml'
130            */
131            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, &parameters).await;
140        // TODO: respect parameters.response != "raw"
141        job_result_response(job_result)
142    }
143}
144
145/// retrieve the list of jobs
146async fn get_jobs() -> HttpResponse {
147    let backend = DagsterBackend::new();
148    match backend.get_jobs().await {
149        Ok(jobs) => HttpResponse::Ok().json(jobs), // TODO: type JobList
150        Err(e) => HttpResponse::InternalServerError().json(Exception::from(e)),
151    }
152}
153
154/// retrieve the status of a job
155async 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
164/// cancel a job execution, remove a finished job
165async 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
176/// retrieve the result(s) of a job
177async 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                // qualifiedInputValue
188                Either::Left(HttpResponse::Ok().json(json!({ "result": {"value": json }})))
189            }
190            JobResult::FilePath(path) => {
191                info!("get_result from {path}");
192                // Prevent file compression for now.
193                // With compression enabled, files are returned compressed with the following headers:
194                // * content-encoding: gzip
195                // * vary: accept-encoding
196                // * content-type: application/pdf
197                // * content-disposition: attachment; filename="12575280.pdf"
198                // This seems correct, but clients don't decompress the attached file!?
199                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}