Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 0 additions & 19 deletions core/lomas_core/models/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,22 +148,3 @@ def get_lomas_logger(name: str, level: str = "NOTSET") -> logging.Logger:
logging.getLogger(name).setLevel(level)

return logging.getLogger(name)


# Exceptions
# -----------------------------------------------------------------------------


class ExceptionType(StrEnum):
"""Lomas server exception types.

To be used as discriminator when parsing corresponding models
"""

INVALID_QUERY = "InvalidQueryException"
USER_NOT_FOUND = "UserNotFoundException"
DATASET_NOT_FOUND = "DatasetNotFoundException"
JOB_NOT_FOUND = "JobNotFoundException"
EXTERNAL_LIBRARY = "ExternalLibraryException"
UNAUTHORIZED_ACCESS = "UnauthorizedAccessException"
INTERNAL_SERVER = "InternalServerException"
29 changes: 28 additions & 1 deletion server/lomas_server/admin_database/local_database.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
import sqlite3
from collections.abc import Generator
from contextlib import AbstractContextManager, closing, contextmanager, nullcontext
from datetime import timedelta
from pathlib import Path
from tempfile import SpooledTemporaryFile
from typing import Any, BinaryIO, override
Expand Down Expand Up @@ -232,11 +233,37 @@ def get_job(self, uid: UUID, current_conn: sqlite3.Connection | None = None) ->

return Job.model_validate_json(row[0])

@db_span("db.expire_jobs", table="admin-db")
def expire_jobs(self, delay: timedelta = timedelta(seconds=2)) -> list[UUID]:
ADMINDB_QUERY_COUNTER.add(1, {"operation": "exipre_jobs"})

with _sqlite_connection(self._db_path) as conn:
rows = conn.execute(
"""
UPDATE jobs
SET
status = ?
WHERE
status = ?
AND
(unixepoch('now') - started_at) > ?
RETURNING uid;
""",
(str(JobStatus.PENDING), str(JobStatus.IN_PROGRESS), int(delay.total_seconds())),
).fetchall()

for row in rows:
logger.debug(f"expiring Job {row[0]}")

return [UUID(row[0]) for row in rows]

@override
@db_span("db.get_job_pending", table="admin-db")
def get_job_pending(self) -> Job | None:
ADMINDB_QUERY_COUNTER.add(1, {"operation": "get_job_pending"})

self.expire_jobs()

with _sqlite_connection(self._db_path) as conn:
row = conn.execute(
"SELECT job_json FROM jobs WHERE status = ? ORDER BY started_at LIMIT 1",
Expand All @@ -257,7 +284,7 @@ def put_job(self, job: Job) -> None:
conn.execute(
"INSERT INTO jobs "
"(uid, user_name, dataset_name, status, started_at, job_json) "
"VALUES (?, ?, ?, ?, 'now', ?)",
"VALUES (?, ?, ?, ?, unixepoch('now'), ?)",
(
str(job.uid),
job.requested_by,
Expand Down
1 change: 0 additions & 1 deletion server/lomas_server/routes/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,6 @@ def handle_query_to_job(

new_task = Job(requested_by=user.name, dataset_name=dataset_name, query=query)

# app.state.jobs[str(new_task.uid)] = new_task
admin_database.put_job(new_task)

return new_task
Expand Down
Loading