diff --git a/swagger_server/controllers/jobs_controller.py b/swagger_server/controllers/jobs_controller.py index 278d3a1..02f5b59 100644 --- a/swagger_server/controllers/jobs_controller.py +++ b/swagger_server/controllers/jobs_controller.py @@ -11,19 +11,21 @@ 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) from threading import Thread -from swagger_server.utils import cl_utils +from swagger_server.utils import cl_utils, db_utils from swagger_server.utils import txt_utils from werkzeug.datastructures import FileStorage import threading import time -CLASSLA_CONCURANCE_LIMIT = 3 -DOC2TEXT_CONCURANCE_LIMIT = 3 -ATEAPI_CONCURANCE_LIMIT = 2 +CLASSLA_CONCURANCE_LIMIT = 4 +DOC2TEXT_CONCURANCE_LIMIT = 4 +ATEAPI_CONCURANCE_LIMIT = 4 +IZLUSCI_PO_ISKANJU_CONCURANCE_LIMIT = 4 classla_sem = threading.Semaphore(CLASSLA_CONCURANCE_LIMIT) doc2text_sem = threading.Semaphore(DOC2TEXT_CONCURANCE_LIMIT) ateapi_sem = threading.Semaphore(ATEAPI_CONCURANCE_LIMIT) +izluscipoiskanju_sem = threading.Semaphore(IZLUSCI_PO_ISKANJU_CONCURANCE_LIMIT) def delete_job(job_id): # noqa: E501 @@ -75,10 +77,28 @@ def clear_up_unfinished_jobs(): async def try_do_jobs(): - with cf.ThreadPoolExecutor(max_workers=3) as ex: + with cf.ThreadPoolExecutor(max_workers=4) as ex: ex.submit(try_do_jobs_classla) ex.submit(try_do_jobs_doc2text) ex.submit(try_do_jobs_ateapi) + ex.submit(try_do_jobs_izluscipoiskanju) + + +### Job looping +def try_do_jobs_izluscipoiskanju(): + while True: + try: + if izluscipoiskanju_sem._value > 0: + unfinished_jobs = Job.select() \ + .where(Job.finished_on.is_null(), Job.started_on.is_null(), Job.job_type == 5) \ + .limit(izluscipoiskanju_sem._value) + with cf.ThreadPoolExecutor(max_workers=IZLUSCI_PO_ISKANJU_CONCURANCE_LIMIT) as ex: + [ex.submit(execute_izluscipoiskanju_job, job) for job in unfinished_jobs] + except Exception as e: + print(f"Exception in try_do_jobs_ateapi") + traceback.print_exc() + finally: + time.sleep(3) ### Job looping @@ -220,6 +240,21 @@ def execute_ateapi_job(job: Job): ateapi_sem.release() +def execute_izluscipoiskanju_job(job: Job): + try: + izluscipoiskanju_sem.acquire() + job.started_on = datetime.datetime.utcnow() + job.save() + info = json.loads(job.job_input) + terKand = db_utils.vrni_oss_terminoloske_kandidate(info['leta'], info['vrste'], info['kljucnebesede'], + info['udk']) + job.job_output = terKand + job.finished_on = datetime.datetime.utcnow() + job.save() + finally: + izluscipoiskanju_sem.release() + + clear_up_unfinished_jobs() loop = asyncio.get_event_loop() diff --git a/swagger_server/controllers/oss_controller.py b/swagger_server/controllers/oss_controller.py index 02d7194..ffc02f3 100644 --- a/swagger_server/controllers/oss_controller.py +++ b/swagger_server/controllers/oss_controller.py @@ -1,3 +1,8 @@ +import json + +import connexion + +from swagger_server.requets_db.models.vrsta import JobManager from swagger_server.utils import db_utils from swagger_server import util from flask import send_file @@ -43,8 +48,8 @@ def get_extracted_words(leta=None, vrste=None, kljucnebesede=None, udk=None): # :rtype: List[TerminoloskiKandidat] """ - files = db_utils.vrni_oss_terminoloske_kandidate(leta, vrste, kljucnebesede, udk) - return files, 200 + terKand = db_utils.vrni_oss_terminoloske_kandidate(leta, vrste, kljucnebesede, udk) + return terKand, 200 def get_extracted_words_async(leta=None, vrste=None, kljucnebesede=None, udk=None): # noqa: E501 @@ -63,7 +68,12 @@ def get_extracted_words_async(leta=None, vrste=None, kljucnebesede=None, udk=Non :rtype: str """ - return 'do some magic!' + job, is_old_job = JobManager.create_job(5, json.dumps( + {'leta': leta, 'vrste': vrste, 'kljucnebesede': kljucnebesede, 'udk': udk})) + if job is None: + return "Something went wrong", 500 + ret = {'check_job_url': f'{connexion.request.url_root}/job/{job.id}'} + return ret, 200 def get_files(leta=None, vrste=None, kljucnebesede=None, udk=None): # noqa: E501 diff --git a/swagger_server/requets_db/models/vrsta.py b/swagger_server/requets_db/models/vrsta.py index 01c4558..dc56453 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]) # TODO: After pushing this, comment it and push again +# db.drop_tables([Job]) If pushing this uncommented, comment it and push again db.create_tables([Job]) @@ -55,12 +55,13 @@ class JobManager: :possibilities: # 1 = pretvori datoteko v besedilo, 2 = oznaci besedilo, 12 = oboje # 3 = pretvori dat v besedilo OCR, 2 = oznaci besedilo, 32 = oboje - # 4 = izlusci async + # 4 = izlusci async glede na vhodne conlluje + # 5 = izlusci po iskanju async :return: Job object, Did already exist boolean """ try: - if job_type in [2, 4]: + if job_type in [2, 4, 5]: job, is_new = Job.get_or_create(job_type=job_type, job_input=job_input, input_size=len(job_input)) elif job_type in [1, 3, 12, 32]: tmp_file = ""