amo-server/servers/job_server.py
Maxim Snesarev 33d6bb7ebd Refactor AMO CRM Data Collection Service to use PostgreSQL
- Updated database configuration to switch from SQLite to PostgreSQL, including changes to alembic.ini, Docker Compose, and environment settings.
- Refactored application code to utilize PostgreSQL database adapters, ensuring compatibility with the new database structure.
- Enhanced API routes and data handling to support the new database, including adjustments in data models and query logic.
- Introduced new job processing mechanisms for full synchronization of AMO CRM entities, leveraging FastStream for background tasks.
- Improved logging and error handling across the application to facilitate better monitoring and debugging.
- Removed obsolete SQLite adapter files and migrations, streamlining the project structure for PostgreSQL integration.
2025-11-05 00:38:37 +03:00

302 lines
9.7 KiB
Python

"""
Job server for managing background job publishing and status tracking.
This module provides the JobServer class for queuing various types of jobs
using FastStream Redis broker and tracking their status in the database.
"""
import uuid
import logging
import time
from datetime import datetime
from typing import Dict, Any, Optional, List
from faststream.redis import RedisBroker
from sqlalchemy import desc
from adapters.postgres.database import SessionLocal
from adapters.postgres.models import ExportJob
logger = logging.getLogger(__name__)
class JobServer:
"""Server for managing background job operations."""
def __init__(self, broker: RedisBroker):
"""
Initialize JobServer with Redis broker.
Args:
broker: FastStream Redis broker instance
"""
self.broker = broker
async def queue_export_job(self, configuration_id: int) -> str:
"""
Queue an export job for processing.
Args:
configuration_id: ID of the export configuration to process
Returns:
job_id: Unique identifier for the queued job
"""
job_id = str(uuid.uuid4())
job_data = {
"job_id": job_id,
"job_type": "export",
"configuration_id": configuration_id,
"created_at": datetime.utcnow().isoformat(),
"status": "queued"
}
try:
# Create job record in database
await self._create_job_record(job_data)
# Publish to export-jobs channel
await self.broker.publish(job_data, "export-jobs")
logger.info(f"Export job {job_id} queued successfully")
return job_id
except Exception as e:
logger.error(f"Failed to queue export job: {str(e)}")
await self._update_job_status(job_id, "failed", error_message=str(e))
raise
async def queue_sync_job(
self,
entity_type: str,
limit: Optional[int] = None,
page: Optional[int] = None,
**kwargs
) -> str:
"""
Queue a data synchronization job.
Args:
entity_type: Type of entity to sync (deals, contacts, companies, etc.)
limit: Maximum number of records to sync
page: Page number for pagination
**kwargs: Additional parameters for the sync job
Returns:
job_id: Unique identifier for the queued job
"""
job_id = str(uuid.uuid4())
sync_data = {
"job_id": job_id,
"job_type": "sync",
"entity_type": entity_type,
"parameters": {
"limit": limit,
"page": page,
**kwargs
},
"created_at": datetime.utcnow().isoformat(),
"status": "queued"
}
try:
# Create job record in database
await self._create_job_record(sync_data)
# Publish to sync-jobs channel
await self.broker.publish(sync_data, "sync-jobs")
logger.info(f"Sync job {job_id} for {entity_type} queued successfully")
return job_id
except Exception as e:
logger.error(f"Failed to queue sync job: {str(e)}")
await self._update_job_status(job_id, "failed", error_message=str(e))
raise
async def schedule_refresh_job(self, entity_type: str) -> None:
"""
Schedule periodic data refresh for an entity type.
Args:
entity_type: Type of entity to refresh
"""
try:
# Publish to refresh-jobs channel
await self.broker.publish(entity_type, "refresh-jobs")
logger.info(f"Refresh job for {entity_type} scheduled successfully")
except Exception as e:
logger.error(f"Failed to schedule refresh job for {entity_type}: {str(e)}")
raise
async def get_job_status(self, job_id: str) -> Optional[Dict[str, Any]]:
"""
Get the status of a specific job.
Args:
job_id: Unique identifier of the job (UUID string)
Returns:
Job status information or None if not found
"""
try:
db = SessionLocal()
try:
# Query by the job_id UUID stored in records_processed field temporarily
# Note: In production, you might want a dedicated job_uuid column
job = db.query(ExportJob).filter(
ExportJob.id == int(job_id) if job_id.isdigit() else None
).first()
if not job:
return None
return {
"job_id": job.id,
"configuration_id": job.configuration_id,
"status": job.status,
"records_processed": job.records_processed,
"total_records": job.total_records,
"created_at": job.created_at,
"started_at": job.started_at,
"completed_at": job.completed_at,
"error_message": job.error_message
}
finally:
db.close()
except Exception as e:
logger.error(f"Failed to get job status for {job_id}: {str(e)}")
return None
async def list_jobs(
self,
status: Optional[str] = None,
job_type: Optional[str] = None,
limit: int = 20,
offset: int = 0
) -> Dict[str, Any]:
"""
List jobs with optional filtering.
Args:
status: Filter by job status (pending, running, completed, failed)
job_type: Filter by job type (export, sync) - currently only export supported
limit: Maximum number of jobs to return
offset: Number of jobs to skip
Returns:
Dictionary containing jobs list and pagination info
"""
try:
db = SessionLocal()
try:
# Build query with filters
query = db.query(ExportJob)
if status:
query = query.filter(ExportJob.status == status)
# Get total count
total = query.count()
# Get paginated results
jobs = query.order_by(desc(ExportJob.created_at)).offset(offset).limit(limit).all()
job_list = []
for job in jobs:
job_list.append({
"job_id": job.id,
"configuration_id": job.configuration_id,
"status": job.status,
"records_processed": job.records_processed,
"total_records": job.total_records,
"created_at": job.created_at,
"started_at": job.started_at,
"completed_at": job.completed_at,
"error_message": job.error_message
})
return {
"jobs": job_list,
"total": total,
"limit": limit,
"offset": offset
}
finally:
db.close()
except Exception as e:
logger.error(f"Failed to list jobs: {str(e)}")
return {"jobs": [], "total": 0, "limit": limit, "offset": offset}
async def _create_job_record(self, job_data: Dict[str, Any]) -> None:
"""
Create a job record in the database.
Note: For export jobs, the record is created in the export router.
For sync jobs, we could create a similar table or just track via logs.
Args:
job_data: Job information to store
"""
# For sync jobs, we're not creating database records yet
# They're tracked through logs and Redis
# Export jobs are created in the export router before queuing
logger.debug(f"Job queued: {job_data.get('job_id')}")
async def _update_job_status(
self,
job_id: str,
status: str,
error_message: Optional[str] = None,
**kwargs
) -> None:
"""
Update job status in the database.
Args:
job_id: Unique identifier of the job
status: New status (pending, running, completed, failed)
error_message: Error message if status is failed
**kwargs: Additional fields to update
"""
try:
db = SessionLocal()
try:
# Try to find job by ID
job = None
if job_id.isdigit():
job = db.query(ExportJob).filter(ExportJob.id == int(job_id)).first()
if not job:
logger.warning(f"Job {job_id} not found for status update")
return
# Update status
job.status = status
if error_message:
job.error_message = error_message
# Update timestamps based on status
if status == "running" and not job.started_at:
job.started_at = int(time.time())
elif status in ["completed", "failed"] and not job.completed_at:
job.completed_at = int(time.time())
# Update any additional fields
for key, value in kwargs.items():
if hasattr(job, key):
setattr(job, key, value)
db.commit()
logger.info(f"Updated job {job_id} status to {status}")
finally:
db.close()
except Exception as e:
logger.error(f"Failed to update job status for {job_id}: {str(e)}")
raise