185 lines
6.9 KiB
Python
185 lines
6.9 KiB
Python
"""
|
|
Scheduler module for FactsDB service
|
|
Handles automated fact extraction jobs
|
|
"""
|
|
|
|
import os
|
|
import time
|
|
import threading
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from apscheduler.triggers.interval import IntervalTrigger
|
|
from typing import Dict, Any, List
|
|
from .file_processor import FileProcessor
|
|
from .ai_processor import AIProcessor
|
|
from .database import DatabaseManager
|
|
from .config import Config, AIEndpointConfig
|
|
from .monitoring import increment_fact_extraction, increment_file_processing, increment_error
|
|
|
|
class FactExtractionJob:
|
|
"""Represents a single fact extraction job"""
|
|
|
|
def __init__(self, config: Config, db_manager: DatabaseManager,
|
|
file_processor: FileProcessor, ai_processor: AIProcessor):
|
|
self.config = config
|
|
self.db_manager = db_manager
|
|
self.file_processor = file_processor
|
|
self.ai_processor = ai_processor
|
|
self.running = False
|
|
|
|
def execute(self, directory_path: str, table_name: str, prompt: str, model: str):
|
|
"""Execute fact extraction on a directory"""
|
|
try:
|
|
print(f"Starting fact extraction for directory: {directory_path}")
|
|
increment_fact_extraction()
|
|
|
|
# Get all files in directory
|
|
if not os.path.exists(directory_path):
|
|
raise Exception(f"Directory does not exist: {directory_path}")
|
|
|
|
files = []
|
|
for root, dirs, filenames in os.walk(directory_path):
|
|
for filename in filenames:
|
|
file_path = os.path.join(root, filename)
|
|
if self.file_processor.is_supported_file_type(file_path):
|
|
files.append(file_path)
|
|
|
|
print(f"Found {len(files)} files to process")
|
|
|
|
# Process each file
|
|
processed_count = 0
|
|
for file_path in files:
|
|
# Check if file has already been processed
|
|
if self.db_manager.is_file_processed(file_path):
|
|
print(f"Skipping already processed file: {file_path}")
|
|
continue
|
|
|
|
try:
|
|
# Extract text from file
|
|
print(f"Processing file: {file_path}")
|
|
increment_file_processing()
|
|
text_content = self.file_processor.extract_text_from_file(file_path)
|
|
|
|
# Extract facts using AI
|
|
fact_data = self.ai_processor.extract_facts_from_text(text_content, prompt, model)
|
|
|
|
# Store fact in database
|
|
self.db_manager.store_fact(fact_data, table_name, file_path)
|
|
|
|
# Mark file as processed
|
|
self.db_manager.mark_file_processed(file_path, table_name)
|
|
|
|
processed_count += 1
|
|
print(f"Successfully processed: {file_path}")
|
|
|
|
except Exception as e:
|
|
print(f"Error processing file {file_path}: {str(e)}")
|
|
increment_error()
|
|
continue
|
|
|
|
print(f"Fact extraction completed. Processed {processed_count} files.")
|
|
|
|
except Exception as e:
|
|
print(f"Fact extraction job failed: {str(e)}")
|
|
increment_error()
|
|
raise
|
|
|
|
class FactExtractionScheduler:
|
|
"""Manages scheduled fact extraction jobs"""
|
|
|
|
def __init__(self, config: Config):
|
|
self.config = config
|
|
self.scheduler = BackgroundScheduler()
|
|
self.jobs = {}
|
|
self.is_running = False
|
|
self.file_processor = FileProcessor()
|
|
self.ai_processor = AIProcessor(config.ai_endpoint)
|
|
# Fix: Add error handling for database initialization in scheduler
|
|
try:
|
|
self.db_manager = DatabaseManager(config.database)
|
|
except Exception as e:
|
|
print(f"Warning: Scheduler database initialization failed: {e}")
|
|
# Create a minimal database manager or handle gracefully
|
|
# For now, we'll skip database initialization for scheduler status
|
|
self.db_manager = None
|
|
|
|
def start(self):
|
|
"""Start the scheduler"""
|
|
if not self.is_running:
|
|
# Add the default job to run every 10 minutes
|
|
self.scheduler.add_job(
|
|
func=self._run_all_jobs,
|
|
trigger=IntervalTrigger(minutes=10),
|
|
id='fact_extraction_job',
|
|
name='Fact Extraction Job'
|
|
)
|
|
|
|
self.scheduler.start()
|
|
self.is_running = True
|
|
print("Fact extraction scheduler started")
|
|
|
|
def stop(self):
|
|
"""Stop the scheduler"""
|
|
if self.is_running:
|
|
self.scheduler.shutdown()
|
|
self.is_running = False
|
|
print("Fact extraction scheduler stopped")
|
|
|
|
def _run_all_jobs(self):
|
|
"""Run all scheduled jobs"""
|
|
print("Running scheduled fact extraction jobs...")
|
|
# In a real implementation, this would iterate through configured jobs
|
|
# For now, we'll just run a basic check
|
|
pass
|
|
|
|
def add_job(self, directory_path: str, table_name: str, prompt: str = "",
|
|
model: str = "gpt-oss", interval_minutes: int = 10):
|
|
"""Add a new scheduled job"""
|
|
job_id = f"job_{len(self.jobs) + 1}"
|
|
|
|
# Create job function
|
|
def job_function():
|
|
job = FactExtractionJob(self.config, self.db_manager, self.file_processor, self.ai_processor)
|
|
job.execute(directory_path, table_name, prompt, model)
|
|
|
|
# Add to scheduler
|
|
self.scheduler.add_job(
|
|
func=job_function,
|
|
trigger=IntervalTrigger(minutes=interval_minutes),
|
|
id=job_id,
|
|
name=f"Fact Extraction: {directory_path}"
|
|
)
|
|
|
|
self.jobs[job_id] = {
|
|
'directory_path': directory_path,
|
|
'table_name': table_name,
|
|
'prompt': prompt,
|
|
'model': model,
|
|
'interval_minutes': interval_minutes
|
|
}
|
|
|
|
print(f"Added scheduled job for: {directory_path}")
|
|
|
|
def remove_job(self, job_id: str):
|
|
"""Remove a scheduled job"""
|
|
if job_id in self.jobs:
|
|
self.scheduler.remove_job(job_id)
|
|
del self.jobs[job_id]
|
|
print(f"Removed job: {job_id}")
|
|
|
|
def get_jobs(self) -> List[Dict[str, Any]]:
|
|
"""Get list of all scheduled jobs"""
|
|
return list(self.jobs.values())
|
|
|
|
def is_job_running(self, job_id: str) -> bool:
|
|
"""Check if a job is running"""
|
|
return job_id in self.jobs
|
|
|
|
# Global scheduler instance
|
|
_scheduler = None
|
|
|
|
def get_scheduler(config: Config) -> FactExtractionScheduler:
|
|
"""Get the global scheduler instance"""
|
|
global _scheduler
|
|
if _scheduler is None:
|
|
_scheduler = FactExtractionScheduler(config)
|
|
return _scheduler |