Keerthana-1117's picture
feat: implement product extraction functionality with Pub/Sub integration and update queue models
5610532
Raw
History Blame Contribute Delete
2.64 kB
from datetime import datetime
from config.database import kittykat_agent_queue_collection
from kittykat_agent.queue.models import QueueStatus, QueueItem
class QueueService:
@staticmethod
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")
@staticmethod
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")
@staticmethod
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
)
@staticmethod
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")