Coverage for pygeoapi/process/manager/base.py: 67%
157 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 08:15 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 08:15 +0000
1# =================================================================
2#
3# Authors: Tom Kralidis <tomkralidis@gmail.com>
4# Ricardo Garcia Silva <ricardo.garcia.silva@geobeyond.it>
5# Francesco Martinelli <francesco.martinelli@ingv.it>
6#
7# Copyright (c) 2024 Tom Kralidis
8# (c) 2023 Ricardo Garcia Silva
9# (c) 2026 Francesco Martinelli
10#
11# Permission is hereby granted, free of charge, to any person
12# obtaining a copy of this software and associated documentation
13# files (the "Software"), to deal in the Software without
14# restriction, including without limitation the rights to use,
15# copy, modify, merge, publish, distribute, sublicense, and/or sell
16# copies of the Software, and to permit persons to whom the
17# Software is furnished to do so, subject to the following
18# conditions:
19#
20# The above copyright notice and this permission notice shall be
21# included in all copies or substantial portions of the Software.
22#
23# THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,
24# EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES
25# OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND
26# NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT
27# HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY,
28# WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING
29# FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR
30# OTHER DEALINGS IN THE SOFTWARE.
31#
32# =================================================================
34import collections
35import json
36import logging
37from multiprocessing import dummy
38from pathlib import Path
39from typing import Any, Dict, Tuple, Optional, OrderedDict
40import uuid
42import requests
44from pygeoapi.plugin import load_plugin
45from pygeoapi.process.base import (
46 BaseProcessor,
47 JobNotFoundError,
48 JobResultNotFoundError,
49 ProcessorExecuteError,
50 UnknownProcessError,
51)
52from pygeoapi.util import (
53 get_current_datetime,
54 is_request_allowed,
55 JobStatus,
56 ProcessExecutionMode,
57 RequestedProcessExecutionMode,
58 RequestedResponse,
59 Subscriber
60)
62LOGGER = logging.getLogger(__name__)
65class BaseManager:
66 """generic Manager ABC"""
67 processes: OrderedDict[str, Dict]
69 def __init__(self, manager_def: dict):
70 """
71 Initialize object
73 :param manager_def: manager definition
75 :returns: `pygeoapi.process.manager.base.BaseManager`
76 """
78 self.name = manager_def['name']
79 self.is_async = False
80 self.supports_subscribing = False
81 self.connection = manager_def.get('connection')
82 self.output_dir = manager_def.get('output_dir')
84 if self.output_dir is not None: 84 ↛ 92line 84 didn't jump to line 92 because the condition on line 84 was always true
85 self.output_dir = Path(self.output_dir)
87 # Note: There are two different things named OrderedDict here - one
88 # is coming from typing.OrderedDict (type annotation), the other is
89 # coming from collections.OrderedDict (actual type we want to use here)
90 # - this will not be needed anymore when pygeoapi moves to requiring
91 # Python 3.9 as the minimum supported Python version
92 self.processes = collections.OrderedDict()
93 for id_, process_conf in manager_def.get('processes', {}).items():
94 self.processes[id_] = dict(process_conf)
96 def get_processor(self, process_id: str) -> BaseProcessor:
97 """Instantiate a processor.
99 :param process_id: Identifier of the process
101 :raises UnknownProcessError: if the processor cannot be created
102 :returns: instance of the processor
103 """
105 try:
106 process_conf = self.processes[process_id]
107 except KeyError as err:
108 raise UnknownProcessError('Invalid process identifier') from err
109 else:
110 pp = load_plugin('process', process_conf['processor'])
111 pp.allow_internal_requests = process_conf.get(
112 'allow_internal_requests', False)
114 return pp
116 def get_jobs(self,
117 status: JobStatus = None,
118 limit: Optional[int] = None,
119 offset: Optional[int] = None
120 ) -> dict:
121 """
122 Get process jobs, optionally filtered by status
124 :param status: job status (accepted, running, successful,
125 failed, results) (default is all)
126 :param limit: number of jobs to return
127 :param offset: pagination offset
129 :returns: dict of list of jobs (identifier, status, process identifier)
130 and numberMatched
131 """
133 raise NotImplementedError()
135 def add_job(self, job_metadata: dict) -> str:
136 """
137 Add a job
139 :param job_metadata: `dict` of job metadata
141 :returns: `str` added job identifier
142 """
144 raise NotImplementedError()
146 def update_job(self, job_id: str, update_dict: dict) -> bool:
147 """
148 Updates a job
150 :param job_id: job identifier
151 :param update_dict: `dict` of property updates
153 :returns: `bool` of status result
154 """
156 raise NotImplementedError()
158 def get_job(self, job_id: str) -> dict:
159 """
160 Get a job (!)
162 :param job_id: job identifier
164 :raises JobNotFoundError: if the job_id does not correspond to a
165 known job
166 :returns: `dict` of job result
167 """
169 raise JobNotFoundError()
171 def get_job_result(self, job_id: str) -> Tuple[str, Any]:
172 """
173 Returns the actual output from a completed process
175 :param job_id: job identifier
177 :raises JobNotFoundError: if the job_id does not correspond to a
178 known job
179 :raises JobResultNotFoundError: if the job-related result cannot
180 be returned
181 :returns: `tuple` of mimetype and raw output
182 """
184 raise JobResultNotFoundError()
186 def delete_job(self, job_id: str) -> bool:
187 """
188 Deletes a job and associated results/outputs
190 :param job_id: job identifier
192 :raises JobNotFoundError: if the job_id does not correspond to a
193 known job
194 :returns: `bool` of status result
195 """
197 raise JobNotFoundError()
199 def _execute_handler_async(self, p: BaseProcessor, job_id: str,
200 data_dict: dict,
201 requested_outputs: Optional[dict] = None,
202 subscriber: Optional[Subscriber] = None,
203 requested_response: Optional[RequestedResponse] = RequestedResponse.raw.value # noqa
204 ) -> Tuple[str, None, JobStatus]:
205 """
206 This private execution handler executes a process in a background
207 thread using `multiprocessing.dummy`
209 https://docs.python.org/3/library/multiprocessing.html#module-multiprocessing.dummy # noqa
211 :param p: `pygeoapi.process` object
212 :param job_id: job identifier
213 :param data_dict: `dict` of data parameters
214 :param requested_outputs: `dict` optionally specifying the subset of
215 required outputs - defaults to all outputs.
216 The value of any key may be an object and
217 include the property `transmissionMode`
218 (defaults to `value`)
219 Note: 'optional' is for backward
220 compatibility.
221 :param subscriber: optional `Subscriber` specifying callback URLs
222 :param requested_response: `RequestedResponse` optionally specifying
223 raw or document (default is `raw`)
225 :returns: tuple of None (i.e. initial response payload)
226 and JobStatus.accepted (i.e. initial job status)
227 """
229 args = (p, job_id, data_dict, requested_outputs, subscriber,
230 requested_response)
232 _process = dummy.Process(target=self._execute_handler_sync, args=args)
233 _process.start()
235 return 'application/json', None, JobStatus.accepted
237 def _execute_handler_sync(self, p: BaseProcessor, job_id: str,
238 data_dict: dict,
239 requested_outputs: Optional[dict] = None,
240 subscriber: Optional[Subscriber] = None,
241 requested_response: Optional[RequestedResponse] = RequestedResponse.raw.value # noqa
242 ) -> Tuple[str, Any, JobStatus]:
243 """
244 Synchronous execution handler
246 If the manager has defined `output_dir`, then the result
247 will be written to disk
248 output store. There is no clean-up of old process outputs.
250 :param p: `pygeoapi.process` object
251 :param job_id: job identifier
252 :param data_dict: `dict` of data parameters
253 :param requested_outputs: `dict` optionally specifying the subset of
254 required outputs - defaults to all outputs.
255 The value of any key may be an object and
256 include the property `transmissionMode`
257 (defaults to `value`)
258 Note: 'optional' is for backward
259 compatibility.
260 :param subscriber: optional `Subscriber` specifying callback URLs
261 :param requested_response: `RequestedResponse` optionally specifying
262 raw or document (default is `raw`)
264 :returns: tuple of MIME type, response payload and status
265 """
267 extra_execute_parameters = {}
269 # only pass requested_outputs if supported,
270 # otherwise this breaks existing processes
271 if p.supports_outputs: 271 ↛ 274line 271 didn't jump to line 274 because the condition on line 271 was always true
272 extra_execute_parameters['outputs'] = requested_outputs
274 self._send_in_progress_notification(subscriber)
276 try:
277 if self.output_dir is not None: 277 ↛ 281line 277 didn't jump to line 281 because the condition on line 277 was always true
278 filename = f"{p.metadata['id']}-{job_id}"
279 job_filename = self.output_dir / filename
280 else:
281 job_filename = None
283 current_status = JobStatus.running
284 jfmt, outputs = p.execute(data_dict, **extra_execute_parameters)
286 if isinstance(outputs, bytes):
287 outputs = outputs.decode('utf-8')
289 if requested_response == RequestedResponse.document.value:
290 outputs = {
291 'outputs': [outputs]
292 }
294 self.update_job(job_id, {
295 'updated': get_current_datetime(),
296 'status': current_status.value,
297 'message': 'Writing job output',
298 'progress': 95
299 })
301 if self.output_dir is not None:
302 LOGGER.debug(f'writing output to {job_filename}')
303 if isinstance(outputs, (dict, list)):
304 mode = 'w'
305 data = json.dumps(outputs, sort_keys=True, indent=4)
306 encoding = 'utf-8'
307 elif isinstance(outputs, bytes):
308 mode = 'wb'
309 data = outputs
310 encoding = None
311 elif isinstance(outputs, str):
312 mode = 'w'
313 data = outputs
314 encoding = None
315 with job_filename.open(mode=mode, encoding=encoding) as fh:
316 fh.write(data)
318 current_status = JobStatus.successful
320 job_update_metadata = {
321 'finished': get_current_datetime(),
322 'updated': get_current_datetime(),
323 'status': current_status.value,
324 'location': str(job_filename),
325 'mimetype': jfmt,
326 'message': 'Job complete',
327 'progress': 100
328 }
330 self.update_job(job_id, job_update_metadata)
331 self._send_success_notification(subscriber, outputs=outputs)
333 except Exception as err:
334 # TODO assess correct exception type and description to help users
335 # NOTE, the /results endpoint should return the error HTTP status
336 # for jobs that failed, the specification says that failing jobs
337 # must still be able to be retrieved with their error message
338 # intact, and the correct HTTP error status at the /results
339 # endpoint, even if the /result endpoint correctly returns the
340 # failure information (i.e. what one might assume is a 200
341 # response).
343 current_status = JobStatus.failed
344 code = 'InvalidParameterValue'
345 outputs = {
346 'type': code,
347 'code': code,
348 'description': f'Error executing process: {err}'
349 }
350 LOGGER.exception(err)
351 job_metadata = {
352 'finished': get_current_datetime(),
353 'updated': get_current_datetime(),
354 'status': current_status.value,
355 'location': None,
356 'mimetype': 'application/octet-stream',
357 'message': f'{code}: {outputs["description"]}'
358 }
360 jfmt = 'application/json'
362 self.update_job(job_id, job_metadata)
363 self._send_failed_notification(subscriber)
365 return jfmt, outputs, current_status
367 def execute_process(
368 self,
369 process_id: str,
370 data_dict: dict,
371 execution_mode: Optional[RequestedProcessExecutionMode] = None,
372 requested_outputs: Optional[dict] = None,
373 subscriber: Optional[Subscriber] = None,
374 requested_response: Optional[RequestedResponse] = RequestedResponse.raw.value # noqa
375 ) -> Tuple[str, Any, JobStatus, Optional[Dict[str, str]]]:
376 """
377 Default process execution handler
379 :param process_id: process identifier
380 :param data_dict: `dict` of data parameters
381 :param execution_mode: `str` optionally specifying sync or async
382 processing.
383 :param requested_outputs: `dict` optionally specifying the subset of
384 required outputs - defaults to all outputs.
385 The value of any key may be an object and
386 include the property `transmissionMode`
387 (default is `value`)
388 Note: 'optional' is for backward
389 compatibility.
390 :param subscriber: `Subscriber` optionally specifying callback urls
391 :param requested_response: `RequestedResponse` optionally specifying
392 raw or document (default is `raw`)
395 :raises UnknownProcessError: if the input process_id does not
396 correspond to a known process
397 :returns: tuple of job_id, MIME type, response payload, status and
398 optionally additional HTTP headers to include in the final
399 response
400 """
402 job_id = str(uuid.uuid1())
403 self.processor = self.get_processor(process_id)
404 self.processor.set_job_id(job_id)
405 extra_execute_handler_parameters = {
406 'requested_response': requested_response
407 }
409 job_control_options = self.processor.metadata.get(
410 'jobControlOptions', [])
412 if execution_mode == RequestedProcessExecutionMode.respond_async:
413 # client wants async - do we support it?
414 process_supports_async = (
415 ProcessExecutionMode.async_execute.value in job_control_options
416 )
417 if self.is_async and process_supports_async: 417 ↛ 425line 417 didn't jump to line 425 because the condition on line 417 was always true
418 LOGGER.debug('Asynchronous execution')
419 handler = self._execute_handler_async
420 response_headers = {
421 'Preference-Applied': (
422 RequestedProcessExecutionMode.respond_async.value)
423 }
424 else:
425 LOGGER.debug('Synchronous execution')
426 handler = self._execute_handler_sync
427 response_headers = {
428 'Preference-Applied': (
429 RequestedProcessExecutionMode.wait.value)
430 }
431 else: # client has no preference or clients wants sync
432 # do we support sync?
433 process_supports_sync = (
434 ProcessExecutionMode.sync_execute.value in job_control_options
435 )
436 if not process_supports_sync: 436 ↛ 437line 436 didn't jump to line 437 because the condition on line 436 was never true
437 LOGGER.debug('Asynchronous execution')
438 handler = self._execute_handler_async
439 response_headers = {
440 'Preference-Applied': (
441 RequestedProcessExecutionMode.respond_async.value)
442 }
443 else:
444 # according to OAPI - Processes spec we ought to
445 # respond with sync
446 LOGGER.debug('Synchronous execution')
447 handler = self._execute_handler_sync
448 if execution_mode == RequestedProcessExecutionMode.wait: 448 ↛ 449line 448 didn't jump to line 449 because the condition on line 448 was never true
449 response_headers = None
450 else:
451 response_headers = {
452 'Preference-Applied': (
453 RequestedProcessExecutionMode.wait.value)
454 }
456 # Add Job before returning any response.
457 current_status = JobStatus.accepted
458 job_metadata = {
459 'type': 'process',
460 'identifier': job_id,
461 'process_id': process_id,
462 'created': get_current_datetime(),
463 'started': get_current_datetime(),
464 'updated': get_current_datetime(),
465 'finished': None,
466 'status': current_status.value,
467 'location': None,
468 'mimetype': 'application/octet-stream',
469 'message': 'Job accepted and ready for execution',
470 'progress': 5
471 }
472 self.add_job(job_metadata)
474 # only pass subscriber if supported, otherwise this breaks
475 # existing managers
476 if self.supports_subscribing: 476 ↛ 481line 476 didn't jump to line 481 because the condition on line 476 was always true
477 extra_execute_handler_parameters['subscriber'] = subscriber
479 # TODO: handler's response could also be allowed to include more HTTP
480 # headers
481 mime_type, outputs, status = handler(
482 self.processor,
483 job_id,
484 data_dict,
485 requested_outputs,
486 **extra_execute_handler_parameters)
488 return job_id, mime_type, outputs, status, response_headers
490 def _send_in_progress_notification(self, subscriber: Optional[Subscriber]):
491 if subscriber and subscriber.in_progress_uri:
492 self.__do_subscriber_request(subscriber.in_progress_uri)
494 def _send_success_notification(
495 self, subscriber: Optional[Subscriber], outputs: Any
496 ):
497 if subscriber and subscriber.success_uri:
498 self.__do_subscriber_request(subscriber.success_uri, outputs)
500 def _send_failed_notification(self, subscriber: Optional[Subscriber]):
501 if subscriber and subscriber.failed_uri:
502 self.__do_subscriber_request(subscriber.failed_uri)
504 def __do_subscriber_request(self, url: str, data: dict = {}) -> None:
505 """
506 Helper function to execute a subscriber URL via HTTP POST
508 :param url: `str` of URL
509 :param data: `dict` of request payload
511 :returns: `None`
512 """
514 if not is_request_allowed(url, self.processor.allow_internal_requests):
515 msg = 'URL not allowed'
516 LOGGER.error(f'{msg}: {url}')
517 raise ProcessorExecuteError(msg)
519 response = requests.post(url, json=data)
520 LOGGER.debug(
521 f'Response: {response.status_code}'
522 )
524 def __repr__(self):
525 return f'<BaseManager> {self.name}'
528def get_manager(config: Dict) -> BaseManager:
529 """Instantiate process manager from the supplied configuration.
531 :param config: pygeoapi configuration
533 :returns: The pygeoapi process manager object
534 """
535 manager_conf = config.get('server', {}).get(
536 'manager',
537 {
538 'name': 'Dummy',
539 'connection': None,
540 'output_dir': None
541 }
542 )
543 processes_conf = {}
544 for id_, resource_conf in config.get('resources', {}).items():
545 if resource_conf.get('type') == 'process':
546 processes_conf[id_] = resource_conf
547 manager_conf['processes'] = processes_conf
548 if manager_conf.get('name') == 'Dummy': 548 ↛ 549line 548 didn't jump to line 549 because the condition on line 548 was never true
549 LOGGER.info('Starting dummy manager')
550 return load_plugin('process_manager', manager_conf)