Skip to content

Commit 7839018

Browse files
author
Ilyas Gasanov
committed
[DOP-19790] Implement scheduler
1 parent 374a130 commit 7839018

4 files changed

Lines changed: 34 additions & 23 deletions

File tree

syncmaster/scheduler/__main__.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,15 +2,17 @@
22
# SPDX-License-Identifier: Apache-2.0
33
import asyncio
44

5+
from syncmaster.config import Settings
56
from syncmaster.scheduler.transfer_fetcher import TransferFetcher
67
from syncmaster.scheduler.transfer_job_manager import TransferJobManager
78

89
TIMEOUT = 180 # seconds
910

1011

1112
async def main():
12-
transfer_fetcher = TransferFetcher()
13-
transfer_job_manager = TransferJobManager()
13+
settings = Settings()
14+
transfer_fetcher = TransferFetcher(settings)
15+
transfer_job_manager = TransferJobManager(settings)
1416
transfer_job_manager.scheduler.start()
1517

1618
while True:

syncmaster/scheduler/transfer_fetcher.py

Lines changed: 8 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,28 +1,24 @@
11
# SPDX-FileCopyrightText: 2023-2024 MTS PJSC
22
# SPDX-License-Identifier: Apache-2.0
33
from sqlalchemy import select
4-
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
5-
from sqlalchemy.orm import sessionmaker
64

75
from syncmaster.config import Settings
86
from syncmaster.db.models import Transfer
9-
10-
DATABASE_URL = Settings().build_db_connection_uri()
11-
engine = create_async_engine(DATABASE_URL)
12-
AsyncSessionLocal = sessionmaker(bind=engine, class_=AsyncSession, expire_on_commit=False)
7+
from syncmaster.scheduler.utils import get_async_session
138

149

1510
class TransferFetcher:
16-
def __init__(self):
11+
def __init__(self, settings: Settings):
12+
self.settings = settings
1713
self.last_updated_at = None
1814

1915
async def fetch_updated_jobs(self):
2016

21-
async with AsyncSessionLocal() as session:
22-
if self.last_updated_at is None:
23-
query = select(Transfer).filter(Transfer.is_scheduled)
24-
else:
25-
query = select(Transfer).filter(Transfer.updated_at > self.last_updated_at)
17+
async with get_async_session(self.settings) as session:
18+
query = select(Transfer)
19+
if self.last_updated_at is not None:
20+
query = query.filter(Transfer.updated_at > self.last_updated_at)
21+
2622
result = await session.execute(query)
2723
transfers = result.scalars().all()
2824

syncmaster/scheduler/transfer_job_manager.py

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -10,15 +10,16 @@
1010
from syncmaster.config import Settings
1111
from syncmaster.db.models import RunType, Status, Transfer
1212
from syncmaster.exceptions.run import CannotConnectToTaskQueueError
13-
from syncmaster.scheduler.transfer_fetcher import AsyncSessionLocal
13+
from syncmaster.scheduler.transfer_fetcher import get_async_session
1414
from syncmaster.schemas.v1.connections.connection import ReadAuthDataSchema
1515
from syncmaster.worker.config import celery
1616

1717

1818
class TransferJobManager:
19-
def __init__(self):
20-
self.scheduler = AsyncIOScheduler(timezone=Settings().TZ)
21-
self.scheduler.add_jobstore("sqlalchemy", url=Settings().build_db_connection_uri(driver="psycopg2"))
19+
def __init__(self, settings: Settings):
20+
self.scheduler = AsyncIOScheduler(timezone=settings.TZ)
21+
self.scheduler.add_jobstore("sqlalchemy", url=settings.build_db_connection_uri(driver="psycopg2"))
22+
self.settings = settings
2223

2324
def update_jobs(self, transfers: list[Transfer]) -> None:
2425
for transfer in transfers:
@@ -34,21 +35,21 @@ def update_jobs(self, transfers: list[Transfer]) -> None:
3435
self.scheduler.modify_job(
3536
job_id=job_id,
3637
trigger=CronTrigger.from_crontab(transfer.schedule),
37-
args=(transfer.id,),
38+
args=(transfer.id, self.settings),
3839
)
3940
else:
4041
self.scheduler.add_job(
4142
func=TransferJobManager.send_job_to_celery,
4243
id=job_id,
4344
trigger=CronTrigger.from_crontab(transfer.schedule),
44-
args=(transfer.id,),
45+
args=(transfer.id, self.settings),
4546
)
4647

4748
@staticmethod
48-
async def send_job_to_celery(transfer_id: int) -> None:
49+
async def send_job_to_celery(transfer_id: int, settings: Settings) -> None:
4950

50-
async with AsyncSessionLocal() as session:
51-
unit_of_work = UnitOfWork(session=session, settings=Settings())
51+
async with get_async_session(settings) as session:
52+
unit_of_work = UnitOfWork(session=session, settings=settings)
5253

5354
transfer = await unit_of_work.transfer.read_by_id(transfer_id)
5455
credentials_source = await unit_of_work.credentials.read(transfer.source_connection_id)

syncmaster/scheduler/utils.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
# SPDX-FileCopyrightText: 2023-2024 MTS PJSC
2+
# SPDX-License-Identifier: Apache-2.0
3+
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
4+
5+
from syncmaster.config import Settings
6+
from syncmaster.db.factory import create_session_factory
7+
8+
9+
def get_async_session(settings: Settings) -> AsyncSession:
10+
engine = create_async_engine(settings.build_db_connection_uri())
11+
session_factory = create_session_factory(engine=engine)
12+
return session_factory()

0 commit comments

Comments
 (0)