annator-atom / backend /integrations /atom_projects_memory_pipeline.py
techprotrade's picture
Full stack ATOM backend + AIMONEYFLOW clients (port 7860) (part 4)
f0ba3c6 verified
Raw
History Blame Contribute Delete
4.5 kB
"""
ATOM Projects Data Memory Pipeline
Background ingestion for Asana and Jira data into LanceDB for AI/RAG.
"""
import asyncio
from datetime import datetime
import json
import logging
import os
from typing import Any, Dict, List, Optional
from core.websockets import manager
from integrations.atom_communication_ingestion_pipeline import (
CommunicationData,
LanceDBMemoryManager,
get_memory_manager,
)
from integrations.jira_service import get_jira_service
# from integrations.asana_service import asana_service # Import when ready
logger = logging.getLogger(__name__)
class ProjectsMemoryPipeline:
"""
Ingests Project data (Tasks, Issues) into the shared LanceDB memory.
"""
def __init__(self, workspace_id: Optional[str] = None):
self.memory_manager = get_memory_manager(workspace_id)
async def run_pipeline(self):
"""Main entry point for scheduled ingestion"""
logger.info("Starting Projects Memory Pipeline...")
await self._ingest_jira()
# await self._ingest_asana()
logger.info("Projects Memory Pipeline Completed.")
# Broadcast Status Update
try:
await manager.broadcast_event("communication_stats", "status_update", {
"pipeline": "projects",
"status": "completed",
"timestamp": datetime.now().isoformat()
})
except Exception as e:
logger.error(f"Failed to broadcast projects status: {e}")
async def _ingest_jira(self):
"""Fetch recent Jira issues and ingest"""
try:
logger.info("Fetching Jira Issues for Memory Ingestion...")
jira = get_jira_service()
# Check connection
if not jira.test_connection().get("authenticated"):
logger.warning("Skipping Jira ingestion: Not Authenticated")
return
# Fetch recent updated issues
jql = "order by updated DESC"
results = jira.search_issues(jql=jql, max_results=50)
issues = results.get("issues", [])
count = 0
for issue in issues:
if self._ingest_task("jira", issue):
count += 1
logger.info(f"Successfully ingested {count} Jira issues into memory.")
except Exception as e:
logger.error(f"Jira Ingestion Failed: {e}")
def _ingest_task(self, source: str, task_data: Dict[str, Any]) -> bool:
"""Map task/issue to CommunicationData structure and ingest"""
try:
# Mapping logic specific to Jira
fields = task_data.get("fields", {})
summary = fields.get("summary", "Untitled Task")
description = fields.get("description") or ""
status = fields.get("status", {}).get("name", "Unknown")
content = f"Task: {summary}\nStatus: {status}\nDescription: {description}\nSource: {source.title()}"
data = CommunicationData(
id=f"{source}_task_{task_data.get('key')}", # Use Key for Jira
app_type=f"{source}_task",
timestamp=datetime.now(),
direction="inbound",
sender=fields.get("creator", {}).get("displayName", "system"),
recipient="atom",
subject=f"Task Update: {summary} ({task_data.get('key')})",
content=content,
attachments=[],
metadata={
"task_id": task_data.get('id'),
"key": task_data.get('key'),
"status": status,
"priority": fields.get("priority", {}).get("name"),
"raw_data": json.dumps(task_data) # Be careful with size
},
status="active", # Active in memory
priority="normal",
tags=["project", "task", source],
vector_embedding=None
)
return self.memory_manager.ingest_communication(data)
except Exception as e:
logger.error(f"Error mapping/ingesting task {task_data.get('key')}: {e}")
return False
# Global Instance
projects_pipeline = ProjectsMemoryPipeline()
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
asyncio.run(projects_pipeline.run_pipeline())