111 lines
2.9 KiB
Python
111 lines
2.9 KiB
Python
"""
|
|
Task scheduler for periodic AMO CRM data refresh jobs.
|
|
|
|
This module sets up scheduled tasks using TaskIQ-FastStream integration
|
|
for periodic data synchronization and maintenance operations.
|
|
"""
|
|
|
|
import logging
|
|
from taskiq_faststream import StreamScheduler
|
|
from taskiq.schedule_sources import LabelScheduleSource
|
|
from workers.broker import broker
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Schedule periodic data refresh jobs
|
|
@broker.task(
|
|
message={"entity_type": "deals"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */6 * * *"}] # Every 6 hours
|
|
)
|
|
async def scheduled_deals_refresh():
|
|
"""Scheduled refresh for deals data."""
|
|
logger.info("Scheduled deals refresh triggered")
|
|
|
|
|
|
@broker.task(
|
|
message={"entity_type": "contacts"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */4 * * *"}] # Every 4 hours
|
|
)
|
|
async def scheduled_contacts_refresh():
|
|
"""Scheduled refresh for contacts data."""
|
|
logger.info("Scheduled contacts refresh triggered")
|
|
|
|
|
|
@broker.task(
|
|
message={"entity_type": "companies"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */8 * * *"}] # Every 8 hours
|
|
)
|
|
async def scheduled_companies_refresh():
|
|
"""Scheduled refresh for companies data."""
|
|
logger.info("Scheduled companies refresh triggered")
|
|
|
|
|
|
@broker.task(
|
|
message={"entity_type": "users"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */12 * * *"}] # Every 12 hours
|
|
)
|
|
async def scheduled_users_refresh():
|
|
"""Scheduled refresh for users data."""
|
|
logger.info("Scheduled users refresh triggered")
|
|
|
|
|
|
@broker.task(
|
|
message={"entity_type": "pipelines"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */24 * * *"}] # Daily
|
|
)
|
|
async def scheduled_pipelines_refresh():
|
|
"""Scheduled refresh for pipelines data."""
|
|
logger.info("Scheduled pipelines refresh triggered")
|
|
|
|
|
|
@broker.task(
|
|
message={"entity_type": "events"},
|
|
channel="refresh-jobs",
|
|
schedule=[{"cron": "0 */2 * * *"}] # Every 2 hours
|
|
)
|
|
async def scheduled_events_refresh():
|
|
"""Scheduled refresh for events data."""
|
|
logger.info("Scheduled events refresh triggered")
|
|
|
|
|
|
# Initialize scheduler
|
|
scheduler = StreamScheduler(
|
|
broker=broker,
|
|
sources=[LabelScheduleSource(broker)]
|
|
)
|
|
|
|
|
|
async def start_scheduler():
|
|
"""Start the task scheduler."""
|
|
logger.info("Starting FastStream task scheduler")
|
|
await scheduler.startup()
|
|
|
|
|
|
async def stop_scheduler():
|
|
"""Stop the task scheduler."""
|
|
logger.info("Stopping FastStream task scheduler")
|
|
await scheduler.shutdown()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
# This allows running the scheduler directly
|
|
import asyncio
|
|
|
|
async def main():
|
|
await start_scheduler()
|
|
try:
|
|
# Keep the scheduler running
|
|
while True:
|
|
await asyncio.sleep(60)
|
|
except KeyboardInterrupt:
|
|
logger.info("Scheduler interrupted by user")
|
|
finally:
|
|
await stop_scheduler()
|
|
|
|
asyncio.run(main())
|