kittykat-backend / src /kittykat_agent /queue /duplication_service.py
arunj-0's picture
KKP-738 | feat: implement campaign and subfolder duplication functionality with queue management
e1f18d9
Raw
History Blame Contribute Delete
6.17 kB
"""
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]