From 42076386f3e24820a15b0aa8756952d28498a9c6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kristjan=20=C5=BDagar?= Date: Fri, 12 Aug 2022 11:39:24 +0200 Subject: [PATCH] Dodana logika vrste, ustvarjanje job-ov, urejen in popravljen swagger.yaml, poizvedba koncanih job-ov, protibitev statusa ce je koncano ali ne. --- .gitignore | 5 +- rsdo5.json | 62 ++++++++ swagger_server/controllers/jobs_controller.py | 105 +++++++++++++ .../controllers/marktext_controller.py | 10 +- swagger_server/db_utils.py | 19 ++- swagger_server/models/__init__.py | 1 + swagger_server/models/job_response.py | 142 ++++++++++++++++++ swagger_server/requets_db/models/vrsta.py | 56 +++++++ swagger_server/swagger/swagger.yaml | 52 +++++++ swagger_server/test/test_jobs_controller.py | 32 ++++ 10 files changed, 476 insertions(+), 8 deletions(-) create mode 100644 swagger_server/controllers/jobs_controller.py create mode 100644 swagger_server/models/job_response.py create mode 100644 swagger_server/requets_db/models/vrsta.py create mode 100644 swagger_server/test/test_jobs_controller.py diff --git a/.gitignore b/.gitignore index 95f80af..bdee654 100644 --- a/.gitignore +++ b/.gitignore @@ -67,4 +67,7 @@ target/ .idea # Related to Development folder -mnt/ \ No newline at end of file +mnt/ + +gen/ +swagger_server/requets_db/dbs/* \ No newline at end of file diff --git a/rsdo5.json b/rsdo5.json index 0ed681d..935b497 100644 --- a/rsdo5.json +++ b/rsdo5.json @@ -11,6 +11,46 @@ } ], "paths": { + "/job/{job_id}": { + "get": { + "tags": [ + "jobs" + ], + "summary": "Vrne status", + "operationId": "getJobStatus", + "parameters": [ + { + "name": "job_id", + "in": "path", + "required": true, + "schema": { + "type": "integer", + "format": "int64" + } + }, + { + "name": "show_estimated_completion", + "in": "query", + "schema": { + "type": "boolean" + }, + "description": "Calculate estimate time remaining based on various factors (could be inaccurate)" + } + ], + "responses": { + "200": { + "description": "OK", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/JobResponse" + } + } + } + } + } + } + }, "/oznaciBesedilo": { "post": { "tags": [ @@ -667,6 +707,28 @@ } } } + }, + "JobResponse": { + "type": "object", + "properties": { + "finished_job": { + "type": "boolean" + }, + "completed_at": { + "type": "string", + "format": "date-time" + }, + "estimated_completion": { + "type": "string", + "format": "date-time" + }, + "job_result": { + "type": "string" + } + }, + "required": [ + "finished_job" + ] } } } diff --git a/swagger_server/controllers/jobs_controller.py b/swagger_server/controllers/jobs_controller.py new file mode 100644 index 0000000..833a8a2 --- /dev/null +++ b/swagger_server/controllers/jobs_controller.py @@ -0,0 +1,105 @@ +import datetime +import random + +import connexion +import peewee +import six +import asyncio + +from swagger_server.models.job_response import JobResponse # noqa: E501 +from swagger_server import util +from swagger_server.requets_db.models.vrsta import (Job, JobManager) +from threading import Semaphore, Thread +from swagger_server.classla import cl_utils + +CLASSLA_CONCURANCE_LIMIT = 4 +classla_sem = asyncio.Semaphore(CLASSLA_CONCURANCE_LIMIT) + + +def get_job_status(job_id, show_estimated_completion=None): # noqa: E501 + """Vrne status + + # noqa: E501 + + :param job_id: + :type job_id: int + :param show_estimated_completion: Calculate estimate time remaining based on various factors (could be inaccurate) + :type show_estimated_completion: bool + + :rtype: JobResponse + """ + try: + job = Job.get_by_id(job_id) + if not job.finished_on: + est_com = None + # todo: if estimate completion: calculate it and set it to est_com + return JobResponse(finished_job=False, estimated_completion=est_com), 200 + return JobResponse(finished_job=True, completed_at=job.finished_on, job_result=job.job_output), 200 + except peewee.DoesNotExist as e: + return "Job with this ID does not exist", 404 + + +def clear_up_unfinished_jobs(): + """ + In case server crashed while jobs were in queue... + """ + Job.update(started_on=None).where(Job.started_on != None, Job.finished_on == None).execute() + + +### Job looping +async def try_do_jobs(): + await asyncio.sleep(15) # wait for tokenizers to load for classla ... + while True: + try: + if cl_utils.nlp_loaded: + # print(cl_utils.raw_text_to_conllu("Danes je lep soncen dan. Res je!")) + if classla_sem._value > 0: + # classla + unfinished_jobs = Job.select() \ + .where(Job.finished_on == None, Job.started_on == None, Job.job_type == 1) \ + .limit(classla_sem._value) + tasks = [ + asyncio.ensure_future(execute_classla_job(job)) + for job in unfinished_jobs + ] + asyncio.get_event_loop().create_task(prep_classla_jobs(tasks)) + # await asyncio.gather(*tasks) + + + + else: + pass + except Exception as e: + print(f"Exception in ... {e}") + finally: + await asyncio.sleep(3) + + +async def prep_classla_jobs(tasks): + await asyncio.gather(*tasks) + + +async def execute_classla_job(job: Job): + async with classla_sem: + job.started_on = datetime.datetime.utcnow() + job.save() + conllu, _ = cl_utils.raw_text_to_conllu(job.job_input) + job.job_output = conllu + job.finished_on = datetime.datetime.utcnow() + job.save() + + +clear_up_unfinished_jobs() +loop = asyncio.get_event_loop() +def loop_in_thread(loop): + asyncio.set_event_loop(loop) + loop.run_until_complete(try_do_jobs()) + + +t = Thread(target=loop_in_thread, args=(loop,)) +t.start() + + +# clear_up_unfinished_jobs() +# loop = asyncio.get_event_loop() +# loop.run_until_complete(try_do_jobs()) this version seems more at home, but it blocks the thread, fix that? diff --git a/swagger_server/controllers/marktext_controller.py b/swagger_server/controllers/marktext_controller.py index c3da9d4..2eb0cef 100644 --- a/swagger_server/controllers/marktext_controller.py +++ b/swagger_server/controllers/marktext_controller.py @@ -4,6 +4,7 @@ import six from swagger_server.models.oznaci_besedilo_body import OznaciBesediloBody # noqa: E501 from swagger_server import util from swagger_server.classla import cl_utils +from swagger_server.requets_db.models.vrsta import (Job, JobManager) def get_text(body): # noqa: E501 @@ -18,6 +19,11 @@ def get_text(body): # noqa: E501 """ if connexion.request.is_json: body = OznaciBesediloBody.from_dict(connexion.request.get_json()) # noqa: E501 - conllu = cl_utils.raw_text_to_conllu(body.besedilo) - return conllu + # conllu = cl_utils.raw_text_to_conllu(body.besedilo) + # return conllu + job, is_old_job = JobManager.create_job(1, body.besedilo) + if job is None: + return "Something went wrong", 500 + ret = {'check_job_url': f'{connexion.request.url_root}?job_id={job.id}'} + return ret # Todo: Update swagger to the newest response template later diff --git a/swagger_server/db_utils.py b/swagger_server/db_utils.py index f0b8725..877a9e0 100644 --- a/swagger_server/db_utils.py +++ b/swagger_server/db_utils.py @@ -4,11 +4,20 @@ import sys # todo: DO NOT PUSH THIS TO GIT database_info = { - 'database': 'TEMPORARYVALJUSTFORTHECOMMITPURPOSE', - 'host': 'TEMPORARYVALJUSTFORTHECOMMITPURPOSE', - 'port': 'TEMPORARYVALJUSTFORTHECOMMITPURPOSE', - 'user': 'TEMPORARYVALJUSTFORTHECOMMITPURPOSE', - 'password': 'TEMPORARYVALJUSTFORTHECOMMITPURPOSE' + 'database': 'conllus_150k', + 'host': '164.8.252.72', + 'port': 3306, + 'user': 'kiki', + 'password': 'kiki123' +} + +# PUT THIS IN INSTEAD WHEN COMMITING, THIS IS ONLY TEMPORARY UNTIL AN .env FILE IS ADDED +database_info_tmp = { + 'database': '...', + 'host': '...', + 'port': 0000, + 'user': '...', + 'password': '...' } # Connect to MariaDB Platform diff --git a/swagger_server/models/__init__.py b/swagger_server/models/__init__.py index a9718ec..49ce18d 100644 --- a/swagger_server/models/__init__.py +++ b/swagger_server/models/__init__.py @@ -8,3 +8,4 @@ from swagger_server.models.datoteka_v_besedilo_ocr_body import DatotekaVBesedilo from swagger_server.models.izlusci_body import IzlusciBody from swagger_server.models.oznaci_besedilo_body import OznaciBesediloBody from swagger_server.models.terminoloski_kandidat import TerminoloskiKandidat +from swagger_server.models.job_response import JobResponse diff --git a/swagger_server/models/job_response.py b/swagger_server/models/job_response.py new file mode 100644 index 0000000..c4507a3 --- /dev/null +++ b/swagger_server/models/job_response.py @@ -0,0 +1,142 @@ +# coding: utf-8 + +from __future__ import absolute_import +from datetime import date, datetime # noqa: F401 + +from typing import List, Dict # noqa: F401 + +from swagger_server.models.base_model_ import Model +from swagger_server import util + + +class JobResponse(Model): + """NOTE: This class is auto generated by the swagger code generator program. + + Do not edit the class manually. + """ + def __init__(self, finished_job: bool=None, completed_at: datetime=None, estimated_completion: datetime=None, job_result: str=None): # noqa: E501 + """JobResponse - a model defined in Swagger + + :param finished_job: The finished_job of this JobResponse. # noqa: E501 + :type finished_job: bool + :param completed_at: The completed_at of this JobResponse. # noqa: E501 + :type completed_at: datetime + :param estimated_completion: The estimated_completion of this JobResponse. # noqa: E501 + :type estimated_completion: datetime + :param job_result: The job_result of this JobResponse. # noqa: E501 + :type job_result: str + """ + self.swagger_types = { + 'finished_job': bool, + 'completed_at': datetime, + 'estimated_completion': datetime, + 'job_result': str + } + + self.attribute_map = { + 'finished_job': 'finished_job', + 'completed_at': 'completed_at', + 'estimated_completion': 'estimated_completion', + 'job_result': 'job_result' + } + self._finished_job = finished_job + self._completed_at = completed_at + self._estimated_completion = estimated_completion + self._job_result = job_result + + @classmethod + def from_dict(cls, dikt) -> 'JobResponse': + """Returns the dict as a model + + :param dikt: A dict. + :type: dict + :return: The JobResponse of this JobResponse. # noqa: E501 + :rtype: JobResponse + """ + return util.deserialize_model(dikt, cls) + + @property + def finished_job(self) -> bool: + """Gets the finished_job of this JobResponse. + + + :return: The finished_job of this JobResponse. + :rtype: bool + """ + return self._finished_job + + @finished_job.setter + def finished_job(self, finished_job: bool): + """Sets the finished_job of this JobResponse. + + + :param finished_job: The finished_job of this JobResponse. + :type finished_job: bool + """ + if finished_job is None: + raise ValueError("Invalid value for `finished_job`, must not be `None`") # noqa: E501 + + self._finished_job = finished_job + + @property + def completed_at(self) -> datetime: + """Gets the completed_at of this JobResponse. + + + :return: The completed_at of this JobResponse. + :rtype: datetime + """ + return self._completed_at + + @completed_at.setter + def completed_at(self, completed_at: datetime): + """Sets the completed_at of this JobResponse. + + + :param completed_at: The completed_at of this JobResponse. + :type completed_at: datetime + """ + + self._completed_at = completed_at + + @property + def estimated_completion(self) -> datetime: + """Gets the estimated_completion of this JobResponse. + + + :return: The estimated_completion of this JobResponse. + :rtype: datetime + """ + return self._estimated_completion + + @estimated_completion.setter + def estimated_completion(self, estimated_completion: datetime): + """Sets the estimated_completion of this JobResponse. + + + :param estimated_completion: The estimated_completion of this JobResponse. + :type estimated_completion: datetime + """ + + self._estimated_completion = estimated_completion + + @property + def job_result(self) -> str: + """Gets the job_result of this JobResponse. + + + :return: The job_result of this JobResponse. + :rtype: str + """ + return self._job_result + + @job_result.setter + def job_result(self, job_result: str): + """Sets the job_result of this JobResponse. + + + :param job_result: The job_result of this JobResponse. + :type job_result: str + """ + + self._job_result = job_result diff --git a/swagger_server/requets_db/models/vrsta.py b/swagger_server/requets_db/models/vrsta.py new file mode 100644 index 0000000..33c447e --- /dev/null +++ b/swagger_server/requets_db/models/vrsta.py @@ -0,0 +1,56 @@ +from peewee import * +from datetime import datetime +import os + +if not os.path.exists('dbs'): + os.mkdir('dbs') + +DB = 'requets_db/dbs/jobs.db' +db = SqliteDatabase(DB, pragmas={ + # 'journal_mode': 'wal', + 'cache_size': -1 * 128 * 1024, # 128MB + 'foreign_keys': 1}) + + +class BaseModel(Model): + class Meta: + database = db + + +class Job(BaseModel): + id = AutoField() + job_type = IntegerField() # 1 = oznaci besedilo + job_input = TextField(index=True) + job_output = TextField(null=True) + created_on = DateTimeField(default=datetime.utcnow) + finished_on = DateTimeField(null=True) + input_size = IntegerField() + started_on = DateTimeField(null=True) + + +# db.drop_tables([Job]) +db.create_tables([Job]) + + +class JobManager: + @staticmethod + def create_job(job_type, job_input) -> Tuple(Job, bool): + """ + :return: Job object, Did already exist boolean + """ + try: + job, is_new = Job.get_or_create(job_type=job_type, job_input=job_input, input_size=len(job_input)) + return job, not is_new + except Exception as e: + print(f'Exception at creating a job: {e}') + return None, False + + # try: + # job = Job.get_or_none(Job.job_input == job_input, Job.job_type == job_type) + # if job: + # return job, True + # job = Job.create(job_type=job_type, job_input=job_input, input_size=len(job_input)) + # return job, False + # except Exception as e: + # print(f'Exception at creating a job: {e}') + # return None, False diff --git a/swagger_server/swagger/swagger.yaml b/swagger_server/swagger/swagger.yaml index abcf60b..be9680a 100644 --- a/swagger_server/swagger/swagger.yaml +++ b/swagger_server/swagger/swagger.yaml @@ -6,6 +6,38 @@ servers: - url: http://localhost:8089 description: Generated server url paths: + /job/{job_id}: + get: + tags: + - jobs + summary: Vrne status + operationId: get_job_status + parameters: + - name: job_id + in: path + required: true + style: simple + explode: false + schema: + type: integer + format: int64 + - name: show_estimated_completion + in: query + description: Calculate estimate time remaining based on various factors (could + be inaccurate) + required: false + style: form + explode: true + schema: + type: boolean + responses: + "200": + description: OK + content: + application/json: + schema: + $ref: '#/components/schemas/JobResponse' + x-openapi-router-controller: swagger_server.controllers.jobs_controller /oznaciBesedilo: post: tags: @@ -485,6 +517,26 @@ components: kanonicnaoblika: kanonicnaoblika nosilnautez: 0.8008282 kandidat: kandidat + JobResponse: + required: + - finished_job + type: object + properties: + finished_job: + type: boolean + completed_at: + type: string + format: date-time + estimated_completion: + type: string + format: date-time + job_result: + type: string + example: + completed_at: 2000-01-23T04:56:07.000+00:00 + estimated_completion: 2000-01-23T04:56:07.000+00:00 + job_result: job_result + finished_job: true oznaciBesedilo_body: type: object properties: diff --git a/swagger_server/test/test_jobs_controller.py b/swagger_server/test/test_jobs_controller.py new file mode 100644 index 0000000..942f5ce --- /dev/null +++ b/swagger_server/test/test_jobs_controller.py @@ -0,0 +1,32 @@ +# coding: utf-8 + +from __future__ import absolute_import + +from flask import json +from six import BytesIO + +from swagger_server.models.job_response import JobResponse # noqa: E501 +from swagger_server.test import BaseTestCase + + +class TestJobsController(BaseTestCase): + """JobsController integration test stubs""" + + def test_get_job_status(self): + """Test case for get_job_status + + Vrne status + """ + query_string = [('job_id', 789), + ('show_estimated_completion', True)] + response = self.client.open( + '/job', + method='GET', + query_string=query_string) + self.assert200(response, + 'Response body is : ' + response.data.decode('utf-8')) + + +if __name__ == '__main__': + import unittest + unittest.main()