Spaces:
Running
Running
feat: implement product extraction functionality with Pub/Sub integration and update queue models
5610532 | from datetime import datetime | |
| from config.database import kittykat_agent_queue_collection | |
| from kittykat_agent.queue.models import QueueStatus, QueueItem | |
| class QueueService: | |
| async def update_queue_item(user_id: str, queue_item_id: str, status: QueueStatus): | |
| collection = kittykat_agent_queue_collection | |
| # Update the status of the queue item | |
| result = await collection.update_one( | |
| {"_id": user_id, "queue.id": queue_item_id}, | |
| {"$set": {"queue.$.status": status.value, | |
| "queue.$.updated_at": datetime.now()}}, | |
| ) | |
| if result.modified_count == 0: | |
| raise Exception("Failed to update queue item status") | |
| async def update_queue_item_metadata(user_id: str, queue_item_id: str, metadata: dict): | |
| """Update metadata fields for a queue item.""" | |
| collection = kittykat_agent_queue_collection | |
| # Build update dict with metadata fields | |
| update_fields = { | |
| f"queue.$.metadata.{key}": value | |
| for key, value in metadata.items() | |
| } | |
| update_fields["queue.$.updated_at"] = datetime.now() | |
| result = await collection.update_one( | |
| {"_id": user_id, "queue.id": queue_item_id}, | |
| {"$set": update_fields} | |
| ) | |
| if result.modified_count == 0: | |
| raise Exception("Failed to update queue item metadata") | |
| async def create_queue_item(user_id: str, queue_item_id: str, queue_item: QueueItem): | |
| collection = kittykat_agent_queue_collection | |
| item = { | |
| "id": queue_item_id, | |
| "title": queue_item.title, | |
| "description": queue_item.description, | |
| "type": queue_item.type, | |
| "status": queue_item.status.value, | |
| "created_at": datetime.now(), | |
| } | |
| if queue_item.metadata: | |
| item["metadata"] = queue_item.metadata | |
| # Create a new queue item | |
| await collection.update_one( | |
| {"_id": user_id}, | |
| { | |
| "$push": { | |
| "queue": item | |
| } | |
| }, | |
| upsert=True | |
| ) | |
| async def delete_queue_items(user_id: str, queue_item_ids: list[str]): | |
| collection = kittykat_agent_queue_collection | |
| # Delete specified queue items | |
| result = await collection.update_one( | |
| {"_id": user_id}, | |
| {"$pull": {"queue": {"id": {"$in": queue_item_ids}}}} | |
| ) | |
| if result.modified_count == 0: | |
| raise Exception("Failed to delete queue items") | |