""" Queue service for managing duplication jobs and other background tasks. """ from datetime import datetime from typing import Optional from uuid import uuid4 from motor.motor_asyncio import AsyncIOMotorCollection from config.database import kittykat_agent_queue_collection from logger import logger class QueueService: """Service for managing queue items.""" def __init__(self, collection: AsyncIOMotorCollection = kittykat_agent_queue_collection): self.collection = collection @staticmethod async def delete_queue_items(user_id: str, queue_item_ids: list[str]): """Delete queue items by IDs from user's queue array.""" from config.database import kittykat_agent_queue_collection await kittykat_agent_queue_collection.update_one( {"_id": user_id}, {"$pull": {"queue": {"id": {"$in": queue_item_ids}}}} ) async def create_duplication_queue_item( self, title: str, description: str, brand_id: str, user_id: str, source_campaign_id: Optional[str] = None, target_campaign_id: Optional[str] = None, source_subfolder_id: Optional[str] = None, target_subfolder_id: Optional[str] = None, total_items: int = 0 ) -> dict: """ Create a queue item for duplication job. Args: title: Queue item title description: Queue item description brand_id: Brand ID user_id: User performing duplication source_campaign_id: Source campaign ID target_campaign_id: Target campaign ID source_subfolder_id: Source subfolder ID target_subfolder_id: Target subfolder ID total_items: Total items to duplicate Returns: Created queue item document """ queue_item_id = str(uuid4()) queue_item = { "id": queue_item_id, "title": title, "description": description, "type": "duplication", "status": "processing", "created_at": datetime.now(), "metadata": { "brand_id": brand_id, "source_campaign_id": source_campaign_id, "target_campaign_id": target_campaign_id, "source_subfolder_id": source_subfolder_id, "target_subfolder_id": target_subfolder_id, "total_items": total_items, "processed_items": 0, "failed_items": 0 } } # Add to user's queue array await self.collection.update_one( {"_id": user_id}, { "$push": {"queue": queue_item}, "$setOnInsert": {"_id": user_id} }, upsert=True ) logger.info(f"Created duplication queue item: {queue_item_id} for user: {user_id}") return queue_item async def update_queue_item( self, user_id: str, queue_id: str, status: Optional[str] = None, processed_items: Optional[int] = None, failed_items: Optional[int] = None, error_message: Optional[str] = None ) -> bool: """ Update queue item status and metadata. Args: user_id: User ID queue_id: Queue item ID status: New status processed_items: Number of items processed failed_items: Number of items failed error_message: Error message if failed Returns: True if updated successfully """ update_fields = {"queue.$.updated_at": datetime.now()} if status: update_fields["queue.$.status"] = status if processed_items is not None: update_fields["queue.$.metadata.processed_items"] = processed_items if failed_items is not None: update_fields["queue.$.metadata.failed_items"] = failed_items if error_message is not None: update_fields["queue.$.metadata.error_message"] = error_message result = await self.collection.update_one( {"_id": user_id, "queue.id": queue_id}, {"$set": update_fields} ) return result.modified_count > 0 async def get_queue_item(self, user_id: str, queue_id: str) -> Optional[dict]: """Get queue item by ID from user's queue.""" user_doc = await self.collection.find_one( {"_id": user_id, "queue.id": queue_id}, {"queue.$": 1} ) if user_doc and "queue" in user_doc and len(user_doc["queue"]) > 0: return user_doc["queue"][0] return None async def get_user_queue_items( self, user_id: str, queue_type: Optional[str] = None, status: Optional[str] = None, skip: int = 0, limit: int = 50 ) -> list[dict]: """ Get queue items for a user. Args: user_id: User ID queue_type: Filter by queue type (optional) status: Filter by status (optional) skip: Number of items to skip limit: Maximum items to return Returns: List of queue items """ user_doc = await self.collection.find_one({"_id": user_id}) if not user_doc or "queue" not in user_doc: return [] queue_items = user_doc["queue"] # Apply filters if queue_type: queue_items = [item for item in queue_items if item.get("type") == queue_type] if status: queue_items = [item for item in queue_items if item.get("status") == status] # Sort by created_at descending queue_items.sort(key=lambda x: x.get("created_at", datetime.min), reverse=True) # Apply pagination return queue_items[skip:skip + limit]