File size: 7,967 Bytes
f8ee2b7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
"""
Worker process for executing crawl jobs from the queue.
"""

import asyncio
import logging
from typing import Optional, Dict
from pathlib import Path

from layer_crawler_etl.queue_workers.job_queue import JobQueue, Job, JobStatus, QueueBackend
from layer_crawler_etl.layer0_source_registry.source_registry import SourceRegistry, Source
from layer_crawler_etl.layer1_crawlers import (
    CodeCrawler, DependencyCrawler, LicenseCrawler, 
    SecurityCrawler, TestBuildCrawler, BrowserRuntimeCrawler
)
from layer_crawler_etl.layer2_etl import Extractor, Transformer, Loader
from layer_crawler_etl.layer3_scoring import Scorer
from layer_crawler_etl.layer4_action import ActionEngine

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)


class CrawlerWorker:
    """Worker that processes crawl jobs from the queue."""
    
    CRAWLER_MAP = {
        "code": CodeCrawler,
        "dependency": DependencyCrawler,
        "license": LicenseCrawler,
        "security": SecurityCrawler,
        "test_build": TestBuildCrawler,
        "browser_runtime": BrowserRuntimeCrawler
    }
    
    def __init__(
        self,
        job_queue: JobQueue,
        source_registry: SourceRegistry,
        storage_path: Optional[Path] = None
    ):
        self.job_queue = job_queue
        self.source_registry = source_registry
        self.storage_path = storage_path or Path("layer_crawler_etl/storage")
        
        # Initialize ETL components
        self.extractor = Extractor(storage_path)
        self.transformer = Transformer()
        self.loader = Loader(storage_path / "normalized")
        self.scorer = Scorer()
        self.action_engine = ActionEngine()
        
        # Initialize crawlers
        self.crawlers = {
            crawler_type: crawler_class(storage_path)
            for crawler_type, crawler_class in self.CRAWLER_MAP.items()
        }
    
    async def process_job(self, job: Job) -> bool:
        """Process a single crawl job."""
        logger.info(f"Processing job {job.job_id} for source {job.source_id}")
        
        try:
            # Get source from registry
            source = self.source_registry.get_source(job.source_id)
            if not source:
                raise ValueError(f"Source {job.source_id} not found in registry")
            
            # Get appropriate crawler
            crawler_type = job.crawler_type
            if crawler_type not in self.crawlers:
                raise ValueError(f"Unknown crawler type: {crawler_type}")
            
            crawler = self.crawlers[crawler_type]
            
            # Validate source
            if not crawler.validate_source(source):
                raise ValueError(f"Source not valid for crawler {crawler_type}")
            
            # Execute crawl
            logger.info(f"Crawling source {source.source_id} with {crawler_type}")
            crawl_result = await crawler.crawl(source)
            
            # Extract
            logger.info(f"Extracting data from crawl result")
            extracted = self.extractor.extract_from_crawl_result(crawl_result)
            
            # Transform
            logger.info(f"Transforming extracted data")
            normalized = self.transformer.transform(extracted)
            
            # Load
            logger.info(f"Loading normalized record")
            load_result = self.loader.load(normalized)
            
            # Score
            logger.info(f"Scoring normalized record")
            score_result = self.scorer.score(normalized)
            
            # Generate actions
            logger.info(f"Generating actions")
            action_result = self.action_engine.generate_actions(score_result, normalized.data)
            
            # Update job with result
            result = {
                "crawl_success": crawl_result.success,
                "load_success": load_result.records_loaded > 0,
                "evidence_score": score_result.evidence_score,
                "prod_score": score_result.prod_score,
                "overall_score": score_result.overall_score,
                "risk_level": score_result.risk_level.value,
                "actions_count": len(action_result.actions_generated),
                "blocked": len(action_result.blocked_claims) > 0
            }
            
            await self.job_queue.update_job_status(
                job,
                JobStatus.COMPLETED,
                result=result
            )
            
            # Update source status in registry
            if crawl_result.success:
                from layer_crawler_etl.layer0_source_registry.source_registry import SourceStatus
                self.source_registry.update_source_status(source.source_id, SourceStatus.COMPLETED)
            else:
                from layer_crawler_etl.layer0_source_registry.source_registry import SourceStatus
                self.source_registry.update_source_status(source.source_id, SourceStatus.FAILED)
            
            logger.info(f"Job {job.job_id} completed successfully")
            return True
            
        except Exception as e:
            logger.error(f"Job {job.job_id} failed: {str(e)}")
            
            # Check if we should retry
            job.retry_count += 1
            if job.retry_count < job.max_retries:
                logger.info(f"Retrying job {job.job_id} (attempt {job.retry_count}/{job.max_retries})")
                await self.job_queue.update_job_status(job, JobStatus.RETRY, error=str(e))
                return False
            else:
                logger.error(f"Job {job.job_id} failed after {job.max_retries} retries")
                await self.job_queue.update_job_status(job, JobStatus.FAILED, error=str(e))
                
                # Update source status
                from layer_crawler_etl.layer0_source_registry.source_registry import SourceStatus
                self.source_registry.update_source_status(job.source_id, SourceStatus.FAILED)
                
                return False
    
    async def run(self, poll_interval: float = 1.0):
        """Run worker continuously, polling for jobs."""
        logger.info("Starting crawler worker")
        
        while True:
            try:
                job = await self.job_queue.get_next_job()
                
                if job:
                    await self.process_job(job)
                else:
                    # No jobs available, wait
                    await asyncio.sleep(poll_interval)
                    
            except Exception as e:
                logger.error(f"Worker error: {str(e)}")
                await asyncio.sleep(poll_interval)
    
    async def run_batch(self, max_jobs: int = 10) -> Dict:
        """Run worker for a batch of jobs."""
        logger.info(f"Starting batch processing (max {max_jobs} jobs)")
        
        processed = 0
        succeeded = 0
        failed = 0
        
        while processed < max_jobs:
            job = await self.job_queue.get_next_job()
            
            if not job:
                logger.info("No more jobs in queue")
                break
            
            success = await self.process_job(job)
            processed += 1
            
            if success:
                succeeded += 1
            else:
                failed += 1
        
        logger.info(f"Batch processing complete: {processed} jobs, {succeeded} succeeded, {failed} failed")
        
        return {
            "processed": processed,
            "succeeded": succeeded,
            "failed": failed
        }


async def main():
    """Main entry point for worker."""
    # Initialize components
    source_registry = SourceRegistry()
    job_queue = JobQueue(backend=QueueBackend.MEMORY)
    worker = CrawlerWorker(job_queue, source_registry)
    
    # Run worker
    await worker.run()


if __name__ == "__main__":
    asyncio.run(main())