Spaces:
Sleeping
Sleeping
| """ | |
| 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 | |
| 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] | |