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]