use std::path::Path;
use std::sync::Arc;
use async_trait::async_trait;
use thiserror::Error;
use crate::config::ClassifierProfile;
use crate::domain::{
ArchivedSource, ClassificationAttemptRecord, ClassificationBatchSummary,
ClassificationCacheKey, ClassificationFailure, ClassificationOutcome, ClassificationRecord,
DocumentClassification, DocumentIntelligence, ExtractedDocument, IngestionRun,
IntelligenceCacheKey, ModelAvailability, ResearchDraft, ResearchNote, SourceDigest,
StagedSource, StoredNote,
};
#[async_trait]
pub trait DocumentTools: Send + Sync {
async fn pdf_text(&self, source: &Path) -> Result<String, ExtractionError>;
async fn pdf_ocr(&self, source: &Path) -> Result<String, ExtractionError>;
async fn image_text(&self, source: &Path) -> Result<String, ExtractionError>;
async fn html_text(&self, source: &Path) -> Result<String, ExtractionError>;
}
#[async_trait]
pub trait DocumentExtractor: Send + Sync {
async fn extract(&self, source: &StagedSource) -> Result<ExtractedDocument, ExtractionError>;
}
#[async_trait]
pub trait DocumentClassifier: Send + Sync {
fn cache_keys(&self, extracted_characters: usize) -> Vec<ClassificationCacheKey>;
async fn classify(
&self,
document: &ExtractedDocument,
) -> Result<ClassificationOutcome, ClassificationFailure>;
}
#[async_trait]
pub trait ProfileClassifier: Send + Sync {
async fn classify_with_profile(
&self,
document: &ExtractedDocument,
profile: &ClassifierProfile,
) -> Result<DocumentClassification, ClassificationFailure>;
}
#[async_trait]
pub trait ModelProbe: Send + Sync {
async fn probe(
&self,
profile: &ClassifierProfile,
) -> Result<ModelAvailability, ClassificationFailure>;
}
#[async_trait]
pub trait ResearchAnalyzer: Send + Sync {
async fn analyze(&self, document: &ExtractedDocument) -> Result<ResearchDraft, AnalysisError>;
}
#[async_trait]
pub trait DocumentIntelligenceAnalyzer: Send + Sync {
async fn analyze_intelligence(
&self,
document: &ExtractedDocument,
classification: &DocumentClassification,
) -> Result<DocumentIntelligence, AnalysisError>;
}
pub trait DocumentIntelligenceStore: Send + Sync {
fn load_intelligence(
&self,
digest: &SourceDigest,
key: &IntelligenceCacheKey,
) -> Result<Option<DocumentIntelligence>, IntelligenceStoreError>;
fn save_intelligence(
&self,
key: &IntelligenceCacheKey,
intelligence: &DocumentIntelligence,
) -> Result<(), IntelligenceStoreError>;
}
pub trait VaultStore: Send + Sync {
fn find_by_digest(&self, digest: &SourceDigest) -> Result<Option<StoredNote>, VaultError>;
fn create_note(&self, note: &ResearchNote) -> Result<StoredNote, VaultError>;
fn archive_source(&self, source: &StagedSource) -> Result<ArchivedSource, VaultError>;
}
pub trait IngestionStore: Send + Sync {
fn load(&self, digest: &SourceDigest) -> Result<Option<IngestionRun>, IngestionStoreError>;
fn save(&self, run: &IngestionRun) -> Result<(), IngestionStoreError>;
}
pub trait ClassificationStore: Send + Sync {
fn load(
&self,
digest: &SourceDigest,
key: &ClassificationCacheKey,
) -> Result<Option<ClassificationRecord>, ClassificationStoreError>;
fn save(
&self,
key: &ClassificationCacheKey,
record: &ClassificationRecord,
) -> Result<(), ClassificationStoreError>;
}
pub trait ClassificationRunStore: Send + Sync {
fn begin_batch(&self, batch_id: &str, started_at: &str)
-> Result<(), ClassificationStoreError>;
fn save_attempt(
&self,
attempt: &ClassificationAttemptRecord,
) -> Result<(), ClassificationStoreError>;
fn latest_attempts(
&self,
batch_id: &str,
) -> Result<Vec<ClassificationAttemptRecord>, ClassificationStoreError>;
fn finish_batch(
&self,
summary: &ClassificationBatchSummary,
) -> Result<(), ClassificationStoreError>;
}
impl<T> ClassificationStore for Arc<T>
where
T: ClassificationStore + ?Sized,
{
fn load(
&self,
digest: &SourceDigest,
key: &ClassificationCacheKey,
) -> Result<Option<ClassificationRecord>, ClassificationStoreError> {
(**self).load(digest, key)
}
fn save(
&self,
key: &ClassificationCacheKey,
record: &ClassificationRecord,
) -> Result<(), ClassificationStoreError> {
(**self).save(key, record)
}
}
impl<T> ClassificationRunStore for Arc<T>
where
T: ClassificationRunStore + ?Sized,
{
fn begin_batch(
&self,
batch_id: &str,
started_at: &str,
) -> Result<(), ClassificationStoreError> {
(**self).begin_batch(batch_id, started_at)
}
fn save_attempt(
&self,
attempt: &ClassificationAttemptRecord,
) -> Result<(), ClassificationStoreError> {
(**self).save_attempt(attempt)
}
fn latest_attempts(
&self,
batch_id: &str,
) -> Result<Vec<ClassificationAttemptRecord>, ClassificationStoreError> {
(**self).latest_attempts(batch_id)
}
fn finish_batch(
&self,
summary: &ClassificationBatchSummary,
) -> Result<(), ClassificationStoreError> {
(**self).finish_batch(summary)
}
}
impl<T> DocumentIntelligenceStore for Arc<T>
where
T: DocumentIntelligenceStore + ?Sized,
{
fn load_intelligence(
&self,
digest: &SourceDigest,
key: &IntelligenceCacheKey,
) -> Result<Option<DocumentIntelligence>, IntelligenceStoreError> {
(**self).load_intelligence(digest, key)
}
fn save_intelligence(
&self,
key: &IntelligenceCacheKey,
intelligence: &DocumentIntelligence,
) -> Result<(), IntelligenceStoreError> {
(**self).save_intelligence(key, intelligence)
}
}
#[async_trait]
pub trait VaultIndexer: Send + Sync {
async fn index(&self) -> Result<(), IndexError>;
}
#[derive(Debug, Error, Clone, PartialEq, Eq)]
#[error("ingestion store failed: {0}")]
pub struct IngestionStoreError(pub String);
#[derive(Debug, Error, Clone, PartialEq, Eq)]
#[error("classification store failed: {0}")]
pub struct ClassificationStoreError(pub String);
#[derive(Debug, Error, Clone, PartialEq, Eq)]
#[error("document intelligence store failed: {0}")]
pub struct IntelligenceStoreError(pub String);
#[derive(Debug, Error, Clone, PartialEq, Eq)]
#[error("vault indexing failed: {0}")]
pub struct IndexError(pub String);
#[derive(Debug, Error)]
pub enum VaultError {
#[error("vault note already exists: {0}")]
AlreadyExists(String),
#[error("vault filesystem operation failed: {0}")]
Filesystem(#[source] std::io::Error),
#[error("vault note rendering failed: {0}")]
Rendering(String),
#[error("source digest changed before archival")]
DigestMismatch,
#[error("unsafe vault path: {0}")]
UnsafePath(String),
}
#[derive(Debug, Error, Clone, PartialEq, Eq)]
pub enum AnalysisError {
#[error("local BAML analysis failed: {0}")]
Model(String),
#[error("invalid generated analysis output: {0}")]
Invalid(String),
}
#[derive(Debug, Error, Clone, PartialEq, Eq)]
pub enum ExtractionError {
#[error("document tool failed: {0}")]
Tool(String),
#[error("document contained no usable text")]
Empty,
#[error("invalid extracted document: {0}")]
Invalid(String),
#[error("staged source changed during extraction")]
SourceChanged,
}