From e95d534e785861f3cf48b6ccb57f207f89ece7e1 Mon Sep 17 00:00:00 2001 From: Kikimanox Date: Mon, 24 Oct 2022 16:54:58 +0200 Subject: [PATCH] Added deleting finished/not started jobs. Modified code to work with new ATEapi docker --- rsdo5.json | 2 +- .../controllers/extract_controller.py | 48 +++++++++++------- swagger_server/controllers/jobs_controller.py | 50 ++++++++++++++++--- swagger_server/models/job_response.py | 2 +- swagger_server/requets_db/models/vrsta.py | 2 +- swagger_server/swagger/swagger.yaml | 3 +- 6 files changed, 79 insertions(+), 28 deletions(-) diff --git a/rsdo5.json b/rsdo5.json index 3fefbba..464fd83 100644 --- a/rsdo5.json +++ b/rsdo5.json @@ -1083,7 +1083,7 @@ "properties": { "job_status": { "type": "string", - "enum": ["waiting in que", "currently processing", "finished processing"] + "enum": ["waiting in que", "currently processing", "finished processing (OK)", "finished processing (ERROR)"] }, "finished_on": { "type": "string", diff --git a/swagger_server/controllers/extract_controller.py b/swagger_server/controllers/extract_controller.py index a526727..b7f6b11 100644 --- a/swagger_server/controllers/extract_controller.py +++ b/swagger_server/controllers/extract_controller.py @@ -1,6 +1,6 @@ import codecs import os - +from flask import Response import connexion import json from pathlib import Path @@ -29,28 +29,40 @@ def do_izlusci(conllus, prepovedane_besede): ('file', ('temp_1.conllu', fp, 'application/octet-stream')) ] res = requests.post(ATEapi_endpoint, files=files) - data = json.loads(res.text) + try: + data = json.loads(res.text) + except: + data = "ATEapi error" + except Exception as e: + return Response(f"Exception in izlusci ({str(e)})", 500) finally: fp.close() os.remove(tmp_file_path) - ret = {'terminoloski_kandidati': [ - { - 'POSoznake': tk['msd'], - 'kandidat': tk['terms'], # more to bit lemma al terms? - 'kanonicnaoblika': tk['canonical'], - 'ranking': tk['ranking'], - 'podporneutezi': [ - 0.0, # ???????? - 0.0 # ?????? - ], - 'pogostostpojavljanja': [0, 0] # ??????? - } - for tk in data if tk['terms'] not in prepovedane_besede - ]} - return ret, 200 + if res.status_code != 200: + return Response(str(data), 400) + # return str(data), 400 + try: + ret = {'terminoloski_kandidati': [ + { + 'POSoznake': tk['term_example_msd'], + 'kandidat': tk['lemma'], # more to bit lemma al terms? + 'kanonicnaoblika': tk['canonical'], + 'ranking': tk['ranking'], + 'podporneutezi': [ + 0.0, # ???????? + 0.0 # ?????? + ], + # 'pogostostpojavljanja': [0, 0] # ??????? + 'pogostostpojavljanja': tk['frequency'] # ??????? + } + for tk in data if tk['lemma'] not in prepovedane_besede + ]} + except: + return Response(data, 400) + return Response(ret, 200) except Exception as e: - return str(e), 500 + return Response(str(e), 500) def get_candidates_async(body): # noqa: E501 diff --git a/swagger_server/controllers/jobs_controller.py b/swagger_server/controllers/jobs_controller.py index 02f5b59..54e466f 100644 --- a/swagger_server/controllers/jobs_controller.py +++ b/swagger_server/controllers/jobs_controller.py @@ -6,7 +6,7 @@ import traceback import peewee import asyncio import concurrent.futures as cf - +from flask import Response from swagger_server.controllers.extract_controller import do_izlusci from swagger_server.models.job_response import JobResponse # noqa: E501 from swagger_server.requets_db.models.vrsta import (Job) @@ -27,6 +27,8 @@ doc2text_sem = threading.Semaphore(DOC2TEXT_CONCURANCE_LIMIT) ateapi_sem = threading.Semaphore(ATEAPI_CONCURANCE_LIMIT) izluscipoiskanju_sem = threading.Semaphore(IZLUSCI_PO_ISKANJU_CONCURANCE_LIMIT) +running_threads = {} + def delete_job(job_id): # noqa: E501 """Izbriše job @@ -38,7 +40,14 @@ def delete_job(job_id): # noqa: E501 :rtype: str """ - return 'Endpoint currently disabled' + try: + job = Job.get_by_id(job_id) + if job.started_on is not None and job.finished_on is None: + return Response("Cancelling ongoing jobs currently not implemented.", 400) + job.delete_instance() + return f"Job with the ID {job_id} was removed." + except peewee.DoesNotExist: + return Response("Job with this ID does not exist", 404) def get_job_status(job_id): # noqa: E501 @@ -58,8 +67,13 @@ def get_job_status(job_id): # noqa: E501 if job.started_on is not None and job.finished_on is None: return JobResponse(job_status="currently processing", created_on=job.created_on, started_on=job.started_on), 200 - if job.started_on is not None and job.finished_on is not None: - return JobResponse(job_status="finished processing", created_on=job.created_on, started_on=job.started_on, + if job.started_on is not None and job.finished_on is not None and not job.job_output.startswith("ERROR -"): + return JobResponse(job_status="finished processing (OK)", created_on=job.created_on, + started_on=job.started_on, + finished_on=job.finished_on, job_result=job.job_output), 200 + if job.started_on is not None and job.finished_on is not None and job.job_output.startswith("ERROR -"): + return JobResponse(job_status="finished processing (ERROR)", created_on=job.created_on, + started_on=job.started_on, finished_on=job.finished_on, job_result=job.job_output), 200 except peewee.DoesNotExist: return "Job with this ID does not exist", 404 @@ -110,7 +124,18 @@ def try_do_jobs_ateapi(): .where(Job.finished_on.is_null(), Job.started_on.is_null(), Job.job_type == 4) \ .limit(ateapi_sem._value) with cf.ThreadPoolExecutor(max_workers=ATEAPI_CONCURANCE_LIMIT) as ex: - [ex.submit(execute_ateapi_job, job) for job in unfinished_jobs] + # [ex.submit(execute_ateapi_job, job) for job in unfinished_jobs] + with cf.ThreadPoolExecutor(max_workers=ATEAPI_CONCURANCE_LIMIT) as ex: + for job in unfinished_jobs: + ex.submit(execute_ateapi_job, job) + # sub = ex.submit(execute_ateapi_job, job) + # running_threads[job.id] = sub + # time.sleep(1) + # running_threads[job.id] + # preveri ce je done: ```running_threads[job.id].done()``` (vrne true false) + + # was_canceled = running_threads[job.id].cancel() + # Todo: Mogoce kaksna druga opcija? Ampak verjetno ne, ne vidim (še?) kak prekicat ONGOING job except Exception as e: print(f"Exception in try_do_jobs_ateapi") traceback.print_exc() @@ -232,7 +257,20 @@ def execute_ateapi_job(job: Job): job.started_on = datetime.datetime.utcnow() job.save() info = json.loads(job.job_input) - ret_json, _ = do_izlusci(info['conllus'], info['prepovedane_besede']) + _res = do_izlusci(info['conllus'], info['prepovedane_besede']) + try: + if type(_res.response) is dict: + ret_json = str(_res.response) + else: + try: + ret_json = _res.response[0].decode('utf-8') + except: + ret_json = "Unknown exception." + + if _res.status_code != 200: + ret_json = f'ERROR - {ret_json}' + except: + ret_json = "ERROR - Unknown exception." job.job_output = ret_json job.finished_on = datetime.datetime.utcnow() job.save() diff --git a/swagger_server/models/job_response.py b/swagger_server/models/job_response.py index 2a0bfc0..5ea0c50 100644 --- a/swagger_server/models/job_response.py +++ b/swagger_server/models/job_response.py @@ -78,7 +78,7 @@ class JobResponse(Model): :param job_status: The job_status of this JobResponse. :type job_status: str """ - allowed_values = ["waiting in que", "currently processing", "finished processing"] # noqa: E501 + allowed_values = ["waiting in que", "currently processing", "finished processing (OK)", "finished processing (ERROR)"] # noqa: E501 if job_status not in allowed_values: raise ValueError( "Invalid value for `job_status` ({0}), must be one of {1}" diff --git a/swagger_server/requets_db/models/vrsta.py b/swagger_server/requets_db/models/vrsta.py index dc56453..438efa5 100644 --- a/swagger_server/requets_db/models/vrsta.py +++ b/swagger_server/requets_db/models/vrsta.py @@ -43,7 +43,7 @@ class Job(BaseModel): input_file = TextField(index=True, null=True) -# db.drop_tables([Job]) If pushing this uncommented, comment it and push again +# db.drop_tables([Job]) # If pushing this uncommented, comment it and push again db.create_tables([Job]) diff --git a/swagger_server/swagger/swagger.yaml b/swagger_server/swagger/swagger.yaml index 0c003e4..b4a8379 100644 --- a/swagger_server/swagger/swagger.yaml +++ b/swagger_server/swagger/swagger.yaml @@ -741,7 +741,8 @@ components: enum: - waiting in que - currently processing - - finished processing + - finished processing (OK) + - finished processing (ERROR) finished_on: type: string format: date-time