Spaces:
Sleeping
Sleeping
File size: 6,166 Bytes
e1f18d9 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 | """
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]
|