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

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# ================================================================= 

33 

34import collections 

35import json 

36import logging 

37from multiprocessing import dummy 

38from pathlib import Path 

39from typing import Any, Dict, Tuple, Optional, OrderedDict 

40import uuid 

41 

42import requests 

43 

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) 

61 

62LOGGER = logging.getLogger(__name__) 

63 

64 

65class BaseManager: 

66 """generic Manager ABC""" 

67 processes: OrderedDict[str, Dict] 

68 

69 def __init__(self, manager_def: dict): 

70 """ 

71 Initialize object 

72 

73 :param manager_def: manager definition 

74 

75 :returns: `pygeoapi.process.manager.base.BaseManager` 

76 """ 

77 

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') 

83 

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) 

86 

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) 

95 

96 def get_processor(self, process_id: str) -> BaseProcessor: 

97 """Instantiate a processor. 

98 

99 :param process_id: Identifier of the process 

100 

101 :raises UnknownProcessError: if the processor cannot be created 

102 :returns: instance of the processor 

103 """ 

104 

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) 

113 

114 return pp 

115 

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 

123 

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 

128 

129 :returns: dict of list of jobs (identifier, status, process identifier) 

130 and numberMatched 

131 """ 

132 

133 raise NotImplementedError() 

134 

135 def add_job(self, job_metadata: dict) -> str: 

136 """ 

137 Add a job 

138 

139 :param job_metadata: `dict` of job metadata 

140 

141 :returns: `str` added job identifier 

142 """ 

143 

144 raise NotImplementedError() 

145 

146 def update_job(self, job_id: str, update_dict: dict) -> bool: 

147 """ 

148 Updates a job 

149 

150 :param job_id: job identifier 

151 :param update_dict: `dict` of property updates 

152 

153 :returns: `bool` of status result 

154 """ 

155 

156 raise NotImplementedError() 

157 

158 def get_job(self, job_id: str) -> dict: 

159 """ 

160 Get a job (!) 

161 

162 :param job_id: job identifier 

163 

164 :raises JobNotFoundError: if the job_id does not correspond to a 

165 known job 

166 :returns: `dict` of job result 

167 """ 

168 

169 raise JobNotFoundError() 

170 

171 def get_job_result(self, job_id: str) -> Tuple[str, Any]: 

172 """ 

173 Returns the actual output from a completed process 

174 

175 :param job_id: job identifier 

176 

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 """ 

183 

184 raise JobResultNotFoundError() 

185 

186 def delete_job(self, job_id: str) -> bool: 

187 """ 

188 Deletes a job and associated results/outputs 

189 

190 :param job_id: job identifier 

191 

192 :raises JobNotFoundError: if the job_id does not correspond to a 

193 known job 

194 :returns: `bool` of status result 

195 """ 

196 

197 raise JobNotFoundError() 

198 

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` 

208 

209 https://docs.python.org/3/library/multiprocessing.html#module-multiprocessing.dummy # noqa 

210 

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`) 

224 

225 :returns: tuple of None (i.e. initial response payload) 

226 and JobStatus.accepted (i.e. initial job status) 

227 """ 

228 

229 args = (p, job_id, data_dict, requested_outputs, subscriber, 

230 requested_response) 

231 

232 _process = dummy.Process(target=self._execute_handler_sync, args=args) 

233 _process.start() 

234 

235 return 'application/json', None, JobStatus.accepted 

236 

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 

245 

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. 

249 

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`) 

263 

264 :returns: tuple of MIME type, response payload and status 

265 """ 

266 

267 extra_execute_parameters = {} 

268 

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 

273 

274 self._send_in_progress_notification(subscriber) 

275 

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 

282 

283 current_status = JobStatus.running 

284 jfmt, outputs = p.execute(data_dict, **extra_execute_parameters) 

285 

286 if isinstance(outputs, bytes): 

287 outputs = outputs.decode('utf-8') 

288 

289 if requested_response == RequestedResponse.document.value: 

290 outputs = { 

291 'outputs': [outputs] 

292 } 

293 

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 }) 

300 

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) 

317 

318 current_status = JobStatus.successful 

319 

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 } 

329 

330 self.update_job(job_id, job_update_metadata) 

331 self._send_success_notification(subscriber, outputs=outputs) 

332 

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). 

342 

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 } 

359 

360 jfmt = 'application/json' 

361 

362 self.update_job(job_id, job_metadata) 

363 self._send_failed_notification(subscriber) 

364 

365 return jfmt, outputs, current_status 

366 

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 

378 

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`) 

393 

394 

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 """ 

401 

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 } 

408 

409 job_control_options = self.processor.metadata.get( 

410 'jobControlOptions', []) 

411 

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 } 

455 

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) 

473 

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 

478 

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) 

487 

488 return job_id, mime_type, outputs, status, response_headers 

489 

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) 

493 

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) 

499 

500 def _send_failed_notification(self, subscriber: Optional[Subscriber]): 

501 if subscriber and subscriber.failed_uri: 

502 self.__do_subscriber_request(subscriber.failed_uri) 

503 

504 def __do_subscriber_request(self, url: str, data: dict = {}) -> None: 

505 """ 

506 Helper function to execute a subscriber URL via HTTP POST 

507 

508 :param url: `str` of URL 

509 :param data: `dict` of request payload 

510 

511 :returns: `None` 

512 """ 

513 

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) 

518 

519 response = requests.post(url, json=data) 

520 LOGGER.debug( 

521 f'Response: {response.status_code}' 

522 ) 

523 

524 def __repr__(self): 

525 return f'<BaseManager> {self.name}' 

526 

527 

528def get_manager(config: Dict) -> BaseManager: 

529 """Instantiate process manager from the supplied configuration. 

530 

531 :param config: pygeoapi configuration 

532 

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)