File size: 8,908 Bytes
0bd4ab4
 
 
 
 
 
 
 
 
2593ebc
0bd4ab4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
dce8249
0bd4ab4
 
 
 
 
 
 
 
dce8249
18482ee
803fd56
 
 
0bd4ab4
 
dce8249
 
 
 
 
d77f57b
b017d78
dce8249
d77f57b
 
 
 
dce8249
 
 
d77f57b
dce8249
d77f57b
 
dce8249
 
d77f57b
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
449681d
 
d77f57b
449681d
d77f57b
 
 
 
08f3ab0
d77f57b
 
 
 
 
08f3ab0
d77f57b
 
 
 
 
 
449681d
 
d77f57b
 
 
 
9035689
 
 
 
 
 
 
 
 
 
dce8249
d77f57b
 
9035689
d77f57b
08f3ab0
9035689
d77f57b
dce8249
 
d77f57b
 
 
 
dce8249
 
 
 
 
 
 
 
 
 
 
 
 
 
616ce9e
dce8249
 
616ce9e
 
 
dce8249
 
 
 
d77f57b
 
dce8249
 
 
2839032
766e6f5
2839032
 
766e6f5
2839032
a4aa7df
 
766e6f5
a4aa7df
 
766e6f5
 
 
 
 
 
 
 
2839032
 
 
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
190
191
192
193
194
195
196
197
198
199
"""
Service pour exécution des analyses
"""
from sqlalchemy.orm import Session
from uuid import UUID
from app.db.repositories.analysis_repo import AnalysisRepository
from app.db.models.analysis_run import AnalysisRun
from datetime import datetime
import json
from app.core.logging import logger

class AnalysisService:
    """Service pour orchestration des analyses K2 Think"""
    
    def __init__(self, db: Session):
        self.analysis_repo = AnalysisRepository(db)
    
    def create_analysis_run(
        self,
        project_id: UUID,
        model_used: str
    ):
        """Crée une nouvelle run d' d'analyse"""
        analysis = AnalysisRun(
            project_id=project_id,
            model_used=model_used,
            status="PENDING",
            started_at=datetime.utcnow()
        )
        self.analysis_repo.create(analysis)
        return analysis
    
    def complete_analysis(
        self,
        analysis_id: UUID,
        result: dict = None,
        status: str = "COMPLETED"
    ):
        """Marque une analyse comme complétée"""
        analysis = self.analysis_repo.get_by_id(analysis_id)
        
        if analysis:
            analysis.status = status
            analysis.completed_at = datetime.utcnow()
            if result:
                analysis.result_data = result # SQLAlchemy handles JSON serialization
            self.analysis_repo.db.add(analysis)
            self.analysis_repo.db.commit()
            self.analysis_repo.db.refresh(analysis)
        
        return analysis

    async def process_analysis(self, analysis_id: str, request):
        """Tâche de fond pour traiter l'analyse"""
        from app.services.k2_think_engine import K2ThinkEngine
        from app.db.repositories.paper_repo import PaperRepository
        from app.db.models.research_paper import ResearchPaper
        from app.models.schemas import AnalysisRequest as K2AnalysisRequest, ScientificDocument
        from app.db.models.reasoning_trace import ReasoningTrace
        from app.services.arxiv_service import ArXivService
        from app.services.doaj_service import DOAJService
        from app.services.pubmed_service import PubMedService
        from app.services.openalex_service import OpenAlexService
        
        db = self.analysis_repo.db
        try:
            # 1. Fetch papers from DB or Remote
            paper_repo = PaperRepository(db)
            analysis_run = self.analysis_repo.get_by_id(UUID(analysis_id))
            project_id = analysis_run.project_id
            
            docs = []
            for pid in request.paper_ids:
                # Try UUID search first
                paper = None
                try:
                    paper_uuid = UUID(pid)
                    paper = paper_repo.get_by_id(paper_uuid)
                except (ValueError, AttributeError):
                    # Try Remote ID search in DB
                    paper = db.query(ResearchPaper).filter(
                        ResearchPaper.remote_id == pid,
                        ResearchPaper.project_id == project_id
                    ).first()
                
                if not paper:
                    # FETCH FROM REMOTE AND SAVE
                    logger.info(f"Paper {pid} not in DB. Fetching metadata on-demand...")
                    remote_data = None
                    try:
                        if "doaj_" in pid:
                            svc = DOAJService(); remote_data = svc.fetch_papers(pid.replace("doaj_", ""), 1)
                        elif "pubmed_" in pid:
                            svc = PubMedService(); remote_data = svc.fetch_papers(pid.replace("pubmed_", ""), 1)
                        elif "openalex_" in pid:
                            svc = OpenAlexService(); remote_data = svc.fetch_papers(pid.replace("openalex_", ""), 1)
                        elif "arxiv_" in pid:
                            svc = ArXivService(); remote_data = svc.fetch_papers(pid.replace("arxiv_", ""), 1)
                        else:
                            # Fallback to ArXiv if no prefix (for old saved papers without prefix)
                            svc = ArXivService(); remote_data = svc.fetch_papers(pid, 1)
                        
                        if remote_data:
                            p_info = remote_data[0]
                            snippet = p_info.get("content", p_info.get("summary", ""))[:10000]
                            paper = ResearchPaper(
                                project_id=project_id,
                                remote_id=pid,
                                title=p_info.get("title", "Unknown"),
                                authors=", ".join(p_info.get("authors", [])),
                                summary=snippet,
                                publication_year=datetime.now().year # Fallback
                            )
                            db.add(paper)
                            db.commit()
                            db.refresh(paper)
                            logger.info(f"Saved remote paper {pid} to project {project_id}")
                        else:
                            logger.warning(f"No remote data found for {pid}")
                    except Exception as fetch_err:
                        logger.error(f"Failed to fetch/save remote paper {pid}: {fetch_err}")

                if paper:
                    # Convert authors string to list
                    author_list = [a.strip() for a in paper.authors.split(",")] if paper.authors else ["Unknown"]
                    
                    # Determine document type based on remote_id prefix
                    dtype = "custom"
                    if paper.remote_id:
                        if "pubmed_" in paper.remote_id: dtype = "pubmed"
                        elif "arxiv_" in paper.remote_id or "." in paper.remote_id: dtype = "arxiv"
                        else: dtype = "arxiv" # default for most research platforms

                    docs.append(ScientificDocument(
                        id=str(paper.id),
                        title=paper.title,
                        authors=author_list,
                        abstract=paper.summary or "",
                        content=(paper.summary or "")[:10000], 
                        document_type=dtype,
                        url=paper.pdf_path or ""
                    ))

            if not docs:
                logger.error(f"No documents found for analysis {analysis_id}")
                raise ValueError("Zero documents identified for analysis. Aborting.")

            # 2. Prepare K2 Request
            k2_req = K2AnalysisRequest(
                documents=docs,
                user_id=request.user_id,
                user_profile=request.user_profile,
                reasoning_depth=request.reasoning_depth,
                ethics_rigor=request.ethics_rigor,
                info_density=request.info_density
            )

            # 3. Call K2 Engine
            engine = K2ThinkEngine()
            result = await engine.process_analysis_request(k2_req)

            # 4. Save Reasoning Trace (Detailed audit log)
            trace = ReasoningTrace(
                analysis_id=UUID(analysis_id),
                step_number=1,
                reasoning=result.reasoning_summary or "Analysis complete",
                source_chunks=result.reasoning_trace # JSON field
            )
            db.add(trace)
            
            # 5. Complete Analysis
            # result is a Pydantic model, convert to dict for storage
            self.complete_analysis(UUID(analysis_id), result=result.model_dump())
            db.commit()

        except Exception as e:
            db.rollback() # CLEAN TRANSACTION
            logger.error(f"FATAL ERROR in process_analysis for {analysis_id}: {str(e)}")
            import traceback
            logger.error(traceback.format_exc())
            
            try:
                error_data = {
                    "status": "FAILED",
                    "reasoning_summary": f"Erreur technique lors de l'analyse : {str(e)}\n\nTrace: {traceback.format_exc()[:500]}...",
                    "confidence_overall": 0
                }
                # Use a safe UUID conversion
                try:
                    target_uuid = UUID(analysis_id) if isinstance(analysis_id, str) else analysis_id
                    self.complete_analysis(target_uuid, status="FAILED", result=error_data)
                    db.commit()
                    logger.info(f"Marked analysis {analysis_id} as FAILED in DB")
                except Exception as uuid_err:
                    logger.error(f"Could not convert {analysis_id} to UUID or save failure: {uuid_err}")
            except Exception as final_err:
                logger.error(f"Failed to even mark analysis as FAILED: {final_err}")
                db.rollback()