Skip to content

Commit

Permalink
Dh-4844/adding the timeout to the question endpoint (#208)
Browse files Browse the repository at this point in the history
* Dh-4844/adding the timeout to the question endpoint

* DH-4844/reformat with black
  • Loading branch information
MohammadrezaPourreza authored Oct 13, 2023
1 parent 51727bf commit 2e9d382
Show file tree
Hide file tree
Showing 3 changed files with 42 additions and 0 deletions.
6 changes: 6 additions & 0 deletions dataherald/api/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,12 @@ def scan_db(
def answer_question(self, question_request: QuestionRequest) -> Response:
pass

@abstractmethod
def answer_question_with_timeout(
self, question_request: QuestionRequest
) -> Response:
pass

@abstractmethod
def get_questions(self, db_connection_id: str | None = None) -> list[Question]:
pass
Expand Down
33 changes: 33 additions & 0 deletions dataherald/api/fastapi.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import json
import logging
import os
import threading
import time
from typing import List

Expand Down Expand Up @@ -163,6 +165,37 @@ def answer_question(self, question_request: QuestionRequest) -> Response:
response_repository = ResponseRepository(self.storage)
return response_repository.insert(generated_answer)

@override
def answer_question_with_timeout(
self, question_request: QuestionRequest
) -> Response:
result = None
exception = None
user_question = Question(
question=question_request.question,
db_connection_id=question_request.db_connection_id,
)
stop_event = threading.Event()

def run_and_catch_exceptions():
nonlocal result, exception
if not stop_event.is_set():
result = self.answer_question(question_request)

thread = threading.Thread(target=run_and_catch_exceptions)
thread.start()
thread.join(timeout=int(os.getenv("DH_ENGINE_TIMEOUT")))
if thread.is_alive():
stop_event.set()
return JSONResponse(
status_code=400,
content={
"question_id": user_question.id,
"error_message": "Timeout Error",
},
)
return result

@override
def create_database_connection(
self, database_connection_request: DatabaseConnectionRequest
Expand Down
3 changes: 3 additions & 0 deletions dataherald/server/fastapi/__init__.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import os
from typing import Any, List

import fastapi
Expand Down Expand Up @@ -216,6 +217,8 @@ def scan_db(
return self._api.scan_db(scanner_request, background_tasks)

def answer_question(self, question_request: QuestionRequest) -> Response:
if os.getenv("DH_ENGINE_TIMEOUT", None):
return self._api.answer_question_with_timeout(question_request)
return self._api.answer_question(question_request)

def get_questions(self, db_connection_id: str | None = None) -> list[Question]:
Expand Down

0 comments on commit 2e9d382

Please sign in to comment.