use std::collections::{HashMap, HashSet};
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use rayon::iter::{IntoParallelIterator, IntoParallelRefIterator, ParallelIterator};
use scryer_db::{DependencyPackage, Project, ProjectDependency, ScryerDb};
use crate::dependency::{CargoDependencyProvider, DependencyProvider, DependencySourceType};
use crate::freshness::{Freshness, FreshnessCache, MAX_REFRESH_FILES, RefreshReport};
use crate::hasher::{ChangeDetector, FileChange};
use crate::index_state::{IndexState, IndexStates};
use crate::ingest::{
BatchIngestionActor, DEFAULT_WRITER_BATCH_CAP, ExistingFileStatus, IngestionMessage,
};
use crate::parsers::{ParseMode, parse_file};
use crate::payload::{ParsedFilePayload, RawTypeContract};
use crate::resolve::CrossFileResolver;
use crate::scanner::WorkspaceScanner;
use crate::search::{SearchIndexCache, SearchQuery, SearchResults};
use crate::stats::{DependencyState, IndexStats, ParseFailure, ParseFailures, SkippedFiles};
use crate::watcher::WorkspaceWatcher;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct IndexReport {
pub scanned_files: usize,
pub added_files: usize,
pub modified_files: usize,
pub unchanged_files: usize,
pub deleted_files: usize,
pub total_scopes: usize,
pub total_symbols: usize,
pub total_references: usize,
pub total_edges: usize,
pub unresolved_references: usize,
pub dependencies_indexed: usize,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct IndexOptions {
pub writer_batch_cap: usize,
pub file_order_seed: Option<u64>,
pub dependencies: bool,
}
impl Default for IndexOptions {
fn default() -> Self {
Self {
writer_batch_cap: DEFAULT_WRITER_BATCH_CAP,
file_order_seed: None,
dependencies: true,
}
}
}
fn root_package_name(root: &Path) -> Option<String> {
let manifest = fs::read_to_string(root.join("Cargo.toml")).ok()?;
let mut in_package = false;
for line in manifest.lines() {
let line = line.trim();
if line.starts_with('[') {
in_package = line == "[package]";
} else if in_package
&& let Some(rest) = line.strip_prefix("name")
&& let Some(value) = rest.trim_start().strip_prefix('=')
{
let name = value.trim().trim_matches(|c| c == '"' || c == '\'');
return (!name.is_empty()).then(|| name.replace('-', "_"));
}
}
None
}
fn dependency_state(linked: usize, mut missing: Vec<String>) -> DependencyState {
if missing.is_empty() {
return DependencyState::Indexed { packages: linked };
}
missing.sort();
missing.dedup();
DependencyState::Partial {
packages: linked.saturating_sub(missing.len()),
missing,
}
}
fn seeded_order_key(seed: u64, path: &Path) -> [u8; 8] {
let mut hasher = blake3::Hasher::new();
hasher.update(&seed.to_le_bytes());
hasher.update(path.to_string_lossy().as_bytes());
let mut key = [0u8; 8];
key.copy_from_slice(&hasher.finalize().as_bytes()[..8]);
key
}
#[derive(Clone)]
pub struct EngineService {
db: ScryerDb,
rayon_pool: Arc<rayon::ThreadPool>,
index_locks: Arc<Mutex<HashMap<u64, Arc<tokio::sync::Mutex<()>>>>>,
search: Arc<SearchIndexCache>,
states: Arc<IndexStates>,
freshness: Arc<FreshnessCache>,
options: IndexOptions,
}
impl EngineService {
pub fn new(db: ScryerDb) -> Self {
let pool = rayon::ThreadPoolBuilder::new()
.build()
.expect("Failed to initialize Rayon thread pool");
Self {
db,
rayon_pool: Arc::new(pool),
index_locks: Arc::default(),
search: Arc::default(),
states: Arc::default(),
freshness: Arc::default(),
options: IndexOptions::default(),
}
}
pub fn with_thread_count(db: ScryerDb, thread_count: usize) -> Self {
let pool = rayon::ThreadPoolBuilder::new()
.num_threads(thread_count)
.build()
.expect("Failed to initialize custom Rayon thread pool");
Self {
db,
rayon_pool: Arc::new(pool),
index_locks: Arc::default(),
search: Arc::default(),
states: Arc::default(),
freshness: Arc::default(),
options: IndexOptions::default(),
}
}
pub fn with_index_options(mut self, options: IndexOptions) -> Self {
self.options = options;
self
}
pub fn db(&self) -> &ScryerDb {
&self.db
}
pub fn index_state(&self, project_id: u64) -> IndexState {
self.states.get(project_id)
}
pub fn mark_building(&self, project_id: u64) {
self.states.mark_building(project_id);
}
pub async fn wait_until_ready(&self, project_id: u64, limit: Duration) -> IndexState {
self.states.wait_until_ready(project_id, limit).await
}
pub fn states(&self) -> &Arc<IndexStates> {
&self.states
}
pub async fn wait_watcher_idle(&self, limit: Duration) {
self.states.wait_watcher_idle(limit).await;
}
async fn index_guard(&self, project_id: u64) -> tokio::sync::OwnedMutexGuard<()> {
let lock = {
let mut locks = self.index_locks.lock().unwrap_or_else(|e| e.into_inner());
Arc::clone(locks.entry(project_id).or_default())
};
lock.lock_owned().await
}
async fn ensure_project_exists(&self, project_id: u64) -> anyhow::Result<()> {
let mut guard = self.db.lock().await;
let project = Project::filter(Project::fields().id().eq(project_id))
.first()
.exec(&mut *guard)
.await?;
anyhow::ensure!(project.is_some(), "Project {project_id} is not registered");
Ok(())
}
async fn touch_project(&self, project_id: u64) -> anyhow::Result<()> {
self.search.bump(project_id);
let stmt = format!(
"UPDATE project SET updated_at = '{}' WHERE id = {project_id};",
scryer_db::time::now_rfc3339()
);
let mut guard = self.db.lock().await;
toasty::sql::statement(&stmt).exec(&mut *guard).await?;
Ok(())
}
pub async fn purge_project(&self, project_id: u64) -> anyhow::Result<()> {
let _guard = self.index_guard(project_id).await;
{
let mut db_guard = self.db.lock().await;
let mut tx = db_guard.transaction().await?;
for table in scryer_db::PROJECT_OWNED_TABLES {
let stmt = format!("DELETE FROM {table} WHERE project_id = {project_id};");
toasty::sql::statement(&stmt).exec(&mut tx).await?;
}
tx.commit().await?;
}
self.search.invalidate(project_id).await;
Ok(())
}
async fn store_index_stats(&self, project_id: u64, stats: &IndexStats) -> anyhow::Result<()> {
let json = serde_json::to_string(stats)?;
let mut db_guard = self.db.lock().await;
let mut tx = db_guard.transaction().await?;
toasty::sql::statement(format!(
"DELETE FROM index_stats WHERE project_id = {project_id};"
))
.exec(&mut tx)
.await?;
scryer_db::IndexStats::create()
.project_id(project_id)
.updated_at(scryer_db::time::now_rfc3339())
.stats_json(json)
.exec(&mut tx)
.await?;
tx.commit().await?;
Ok(())
}
pub async fn index_stats(&self, project_id: u64) -> anyhow::Result<Option<IndexStats>> {
crate::health::load_index_stats(&self.db, project_id).await
}
pub async fn index_project(&self, project_id: u64, root: &Path) -> anyhow::Result<IndexReport> {
let _guard = self.index_guard(project_id).await;
let _ready = self.states.ready_on_drop(project_id);
self.ensure_project_exists(project_id).await?;
if self.index_stats(project_id).await?.is_none() {
self.states.mark_building(project_id);
}
let canonical_root = dunce::canonicalize(root).unwrap_or_else(|_| root.to_path_buf());
let total_timer = std::time::Instant::now();
let phase_timer = std::time::Instant::now();
let scanner = WorkspaceScanner::new(&canonical_root);
let scan = scanner.scan_with_skips();
let too_large = scan.too_large.len();
let scanned = scan.files;
let scan_ms = phase_timer.elapsed().as_millis();
let phase_timer = std::time::Instant::now();
let change_report =
ChangeDetector::detect_changes(&self.db, project_id, &canonical_root, &scanned).await?;
let detect_ms = phase_timer.elapsed().as_millis();
let mut report = IndexReport {
scanned_files: scanned.len(),
added_files: change_report.added_count,
modified_files: change_report.modified_count,
unchanged_files: change_report.unchanged_count,
deleted_files: change_report.deleted_count,
total_scopes: 0,
total_symbols: 0,
total_references: 0,
total_edges: 0,
unresolved_references: 0,
dependencies_indexed: 0,
};
let skipped = SkippedFiles {
too_large,
unreadable: change_report.unreadable.len(),
};
if change_report.is_empty() {
let phase_timer = std::time::Instant::now();
let mut dependencies = None;
if self.options.dependencies && canonical_root.join("Cargo.toml").exists() {
match self
.index_dependencies_with_state(project_id, &canonical_root)
.await
{
Ok((deps_indexed, state)) => {
report.dependencies_indexed = deps_indexed;
dependencies = Some(state);
}
Err(err) => {
tracing::warn!(
"Failed to index dependencies for project {project_id}: {err}"
);
dependencies = Some(DependencyState::DiscoveryFailed {
error: format!("{err:#}"),
});
}
}
}
let mut stats = self.index_stats(project_id).await?.unwrap_or_default();
stats.skipped = skipped;
stats.parse_failures = ParseFailures::default();
if let Some(state) = dependencies {
stats.dependencies = state;
}
self.store_index_stats(project_id, &stats).await?;
tracing::info!(
"index_project timings (no changes) project={project_id}: scan={scan_ms}ms detect={detect_ms}ms deps={}ms total={}ms",
phase_timer.elapsed().as_millis(),
total_timer.elapsed().as_millis()
);
return Ok(report);
}
let ingestion_handle = BatchIngestionActor::spawn(
self.db.clone(),
project_id,
4096,
self.options.writer_batch_cap,
);
let sender = ingestion_handle.sender();
let mut files_to_parse = Vec::new();
let mut file_statuses = HashMap::new();
let mut deleted_paths = Vec::new();
let mut unchanged_paths = Vec::new();
for change in change_report.changes {
match change {
FileChange::Added {
rel_path,
content_hash,
content,
} => {
file_statuses.insert(rel_path.clone(), ExistingFileStatus::New);
files_to_parse.push((rel_path, content_hash, content));
}
FileChange::Modified {
file_id,
rel_path,
content_hash,
content,
} => {
file_statuses.insert(rel_path.clone(), ExistingFileStatus::Existing(file_id));
files_to_parse.push((rel_path, content_hash, content));
}
FileChange::Deleted { rel_path, .. } => {
deleted_paths.push(rel_path.clone());
let _ = sender.send(IngestionMessage::DeleteFile(rel_path)).await;
}
FileChange::Unchanged { rel_path, .. } => unchanged_paths.push(rel_path),
}
}
let pending: std::collections::BTreeSet<String> = files_to_parse
.iter()
.map(|(p, _, _)| p.to_string_lossy().to_string())
.chain(
deleted_paths
.iter()
.map(|p| p.to_string_lossy().to_string()),
)
.collect();
self.states.set_total(project_id, files_to_parse.len() * 2);
self.states.mark_stale(project_id, pending);
match self.options.file_order_seed {
Some(seed) => {
files_to_parse.sort_by_cached_key(|(path, _, _)| seeded_order_key(seed, path))
}
None => files_to_parse.sort_by(|a, b| a.0.cmp(&b.0)),
}
let phase_timer = std::time::Instant::now();
let pool = Arc::clone(&self.rayon_pool);
let files_to_parse = Arc::new(files_to_parse);
let files_for_parse = Arc::clone(&files_to_parse);
let parse_states = Arc::clone(&self.states);
let parsed_files: Vec<anyhow::Result<ParsedFilePayload>> =
tokio::task::spawn_blocking(move || {
pool.install(|| {
files_for_parse
.par_iter()
.map(|(rel_path, hash, content)| {
let parsed = parse_file(rel_path, content, hash, ParseMode::Workspace);
parse_states.add_progress(project_id, 1);
parsed
})
.collect()
})
})
.await?;
let mut payloads = Vec::new();
let mut parse_failures = ParseFailures::default();
for (res, (path, _, _)) in parsed_files.into_iter().zip(files_to_parse.iter()) {
match res {
Ok(payload) => payloads.push(payload),
Err(err) => {
tracing::warn!("Skipping file that failed to parse: {err:#}");
parse_failures.count += 1;
let failure = ParseFailure {
path: path.to_string_lossy().to_string(),
error: format!("{err:#}"),
};
if parse_failures
.first
.as_ref()
.is_none_or(|first| failure.path < first.path)
{
parse_failures.first = Some(failure);
}
}
}
}
let files_to_parse_paths: Vec<String> = files_to_parse
.iter()
.map(|(p, _, _)| p.to_string_lossy().to_string())
.collect();
let syntax_error_paths: Vec<String> = payloads
.iter()
.filter(|p| p.syntax_errors)
.map(|p| p.relative_path.to_string_lossy().to_string())
.collect();
let parsed_paths: HashSet<&Path> =
payloads.iter().map(|p| p.relative_path.as_path()).collect();
file_statuses.retain(|path, _| parsed_paths.contains(path.as_path()));
let parse_ms = phase_timer.elapsed().as_millis();
let root_crate = root_package_name(&canonical_root);
let mut file_sources = Vec::new();
let mut all_symbols = HashMap::new();
let mut all_scopes = HashMap::new();
let content_map: HashMap<&Path, &[u8]> = files_to_parse
.iter()
.map(|(p, _, c)| (p.as_path(), c.as_slice()))
.collect();
for p in &payloads {
if let Some(bytes) = content_map.get(p.relative_path.as_path()) {
if let Ok(src) = std::str::from_utf8(bytes) {
file_sources.push((p.relative_path.clone(), src.to_string()));
}
} else {
let full_path = canonical_root.join(&p.relative_path);
if let Ok(src) = fs::read_to_string(&full_path) {
file_sources.push((p.relative_path.clone(), src));
}
}
all_symbols.insert(p.relative_path.clone(), p.symbols.clone());
all_scopes.insert(p.relative_path.clone(), p.scopes.clone());
}
let phase_timer = std::time::Instant::now();
let extracted_map = CrossFileResolver::resolve_project(
&file_sources,
&all_symbols,
&all_scopes,
root_crate.as_deref(),
)?;
let resolve_ms = phase_timer.elapsed().as_millis();
let phase_timer = std::time::Instant::now();
let mut total_scopes = 0;
let mut total_symbols = 0;
for mut payload in payloads {
total_scopes += payload.scopes.len();
total_symbols += payload.symbols.len();
if let Some(extracted) = extracted_map.get(&payload.relative_path) {
payload.references = extracted.references.clone();
payload.edges = extracted.edges.clone();
}
let status = file_statuses
.remove(&payload.relative_path)
.unwrap_or(ExistingFileStatus::Unknown);
let written_path = payload.relative_path.to_string_lossy().to_string();
sender
.send(IngestionMessage::UpsertFile {
payload,
status,
dependency_package_id: None,
})
.await
.map_err(|e| anyhow::anyhow!("Writer channel closed: {e}"))?;
self.states.add_progress(project_id, 1);
self.states.file_done(project_id, &written_path);
}
report.total_scopes = total_scopes;
report.total_symbols = total_symbols;
drop(sender);
let outcome = ingestion_handle.finish().await?;
let write_ms = phase_timer.elapsed().as_millis();
let phase_timer = std::time::Instant::now();
let mut dependencies = DependencyState::NotApplicable;
if self.options.dependencies && canonical_root.join("Cargo.toml").exists() {
match self
.index_dependencies_with_state(project_id, &canonical_root)
.await
{
Ok((deps_indexed, state)) => {
report.dependencies_indexed = deps_indexed;
dependencies = state;
}
Err(err) => {
tracing::warn!("Failed to index dependencies for project {project_id}: {err}");
dependencies = DependencyState::DiscoveryFailed {
error: format!("{err:#}"),
};
}
}
}
let deps_ms = phase_timer.elapsed().as_millis();
let phase_timer = std::time::Instant::now();
let (first_pass_references, first_pass_edges) =
(outcome.inserted_references, outcome.inserted_edges);
let link = BatchIngestionActor::new(self.db.clone(), project_id)
.link_deferred(outcome)
.await?;
report.total_references = first_pass_references + link.linked_references;
report.total_edges = first_pass_edges + link.linked_edges;
report.unresolved_references = link.unresolved_references;
let link_ms = phase_timer.elapsed().as_millis();
let touched: HashSet<String> = files_to_parse_paths
.into_iter()
.chain(
deleted_paths
.iter()
.map(|p| p.to_string_lossy().to_string()),
)
.collect();
let mut syntax_error_files: Vec<String> = self
.index_stats(project_id)
.await?
.map(|s| s.syntax_error_files)
.unwrap_or_default()
.into_iter()
.filter(|p| !touched.contains(p))
.collect();
syntax_error_files.extend(syntax_error_paths);
syntax_error_files.sort();
syntax_error_files.dedup();
self.store_index_stats(
project_id,
&IndexStats {
unresolved_references: link.unresolved_references,
parse_failures,
syntax_error_files,
last_run_files: report.added_files + report.modified_files,
skipped,
dependencies,
},
)
.await?;
self.touch_project(project_id).await?;
tracing::info!(
"index_project timings project={project_id}: scan={scan_ms}ms detect={detect_ms}ms parse={parse_ms}ms resolve={resolve_ms}ms write={write_ms}ms deps={deps_ms}ms final_link={link_ms}ms total={}ms",
total_timer.elapsed().as_millis()
);
Ok(report)
}
pub async fn index_file(
&self,
project_id: u64,
root: &Path,
rel_path: &Path,
) -> anyhow::Result<()> {
let full_path = root.join(rel_path);
if fs::metadata(&full_path).is_ok_and(|m| m.len() > crate::scanner::MAX_SOURCE_FILE_BYTES) {
tracing::warn!(
"Not indexing {}: over the file size cap",
full_path.display()
);
return self.remove_file(project_id, rel_path).await;
}
let _guard = self.index_guard(project_id).await;
self.ensure_project_exists(project_id).await?;
let content = fs::read(&full_path)?;
let content_hash = crate::hasher::hash_bytes(&content);
let mut parsed = parse_file(rel_path, &content, &content_hash, ParseMode::Workspace)?;
let syntax_errors = parsed.syntax_errors;
if let Ok(src) = std::str::from_utf8(&content) {
let mut all_symbols = HashMap::new();
let mut all_scopes = HashMap::new();
all_symbols.insert(rel_path.to_path_buf(), parsed.symbols.clone());
all_scopes.insert(rel_path.to_path_buf(), parsed.scopes.clone());
if let Ok(extracted) = CrossFileResolver::resolve_file_references(
rel_path,
src,
&all_symbols,
&all_scopes,
root_package_name(root).as_deref(),
) {
parsed.references = extracted.references;
parsed.edges = extracted.edges;
}
}
let actor = BatchIngestionActor::new(self.db.clone(), project_id);
let outcome = actor.ingest_payload(parsed).await?;
actor.link_deferred(outcome).await?;
self.record_syntax_errors(project_id, rel_path, syntax_errors)
.await?;
self.touch_project(project_id).await?;
Ok(())
}
async fn record_syntax_errors(
&self,
project_id: u64,
rel_path: &Path,
has_errors: bool,
) -> anyhow::Result<()> {
let Some(mut stats) = self.index_stats(project_id).await? else {
return Ok(());
};
let path = rel_path.to_string_lossy().to_string();
if stats.syntax_error_files.contains(&path) == has_errors {
return Ok(());
}
if has_errors {
stats.syntax_error_files.push(path);
stats.syntax_error_files.sort();
} else {
stats.syntax_error_files.retain(|p| *p != path);
}
self.store_index_stats(project_id, &stats).await
}
pub async fn refresh_files(
&self,
project_id: u64,
root: &Path,
rel_paths: &[PathBuf],
) -> anyhow::Result<RefreshReport> {
let mut report = RefreshReport::default();
if self.states.get(project_id).is_indexing() {
return Ok(report);
}
let scanner = WorkspaceScanner::new(root);
let mut ignore: Option<crate::scanner::IgnoreFilter> = None;
for (n, rel) in rel_paths.iter().enumerate() {
if n >= MAX_REFRESH_FILES {
report.skipped.extend(
rel_paths[n..]
.iter()
.map(|p| p.to_string_lossy().to_string()),
);
break;
}
let abs = root.join(rel);
let name = rel.to_string_lossy().to_string();
let indexed_hash = {
let mut guard = self.db.lock().await;
scryer_db::SourceFile::filter(
scryer_db::SourceFile::fields()
.project_id()
.eq(project_id)
.and(scryer_db::SourceFile::fields().path().eq(name.clone())),
)
.first()
.exec(&mut *guard)
.await?
.map(|f| f.content_hash)
};
match indexed_hash {
Some(hash) => match self.freshness.check(&abs, &hash) {
(Freshness::Fresh, _) => {}
(Freshness::Changed, signature) => {
self.index_file(project_id, root, rel).await?;
self.freshness.record(&abs, signature);
report.refreshed.push(name);
}
(Freshness::Missing, _) => {
self.remove_file(project_id, rel).await?;
self.freshness.forget(&abs);
report.removed.push(name);
}
},
None => {
if abs.is_file()
&& scanner.is_supported(rel)
&& !ignore
.get_or_insert_with(|| crate::scanner::IgnoreFilter::new(root))
.is_ignored(&abs)
{
let signature = std::fs::metadata(&abs)
.ok()
.map(|m| (m.len(), m.modified().ok()));
self.index_file(project_id, root, rel).await?;
self.freshness.record(&abs, signature);
report.refreshed.push(name);
}
}
}
}
Ok(report)
}
pub async fn remove_file(&self, project_id: u64, rel_path: &Path) -> anyhow::Result<()> {
let _guard = self.index_guard(project_id).await;
let actor = BatchIngestionActor::new(self.db.clone(), project_id);
let outcome = actor.delete_file(rel_path).await?;
actor.link_deferred(outcome).await?;
self.record_syntax_errors(project_id, rel_path, false)
.await?;
self.touch_project(project_id).await?;
Ok(())
}
pub async fn index_dependencies(&self, project_id: u64, root: &Path) -> anyhow::Result<usize> {
Ok(self
.index_dependencies_with_state(project_id, root)
.await?
.0)
}
pub async fn index_dependencies_with_state(
&self,
project_id: u64,
root: &Path,
) -> anyhow::Result<(usize, DependencyState)> {
let manifest_path = if root.is_file() {
root.to_path_buf()
} else {
root.join("Cargo.toml")
};
if !manifest_path.exists() {
return Ok((0, DependencyState::NotApplicable));
}
let deps_timer = std::time::Instant::now();
let provider = CargoDependencyProvider::new();
let discovered = match provider.discover_dependencies(root).await {
Ok(deps) => deps,
Err(err) => {
tracing::warn!(
"Failed to discover cargo dependencies for {}: {err}",
root.display()
);
return Ok((
0,
DependencyState::DiscoveryFailed {
error: format!("{err:#}"),
},
));
}
};
let discover_ms = deps_timer.elapsed().as_millis();
let direct_deps: Vec<_> = discovered.into_iter().filter(|d| d.is_direct).collect();
if direct_deps.is_empty() {
if self
.prune_dependency_links(project_id, &HashSet::new())
.await?
{
self.search.bump(project_id);
}
return Ok((0, DependencyState::Indexed { packages: 0 }));
}
let mut to_index = Vec::new();
let mut links_changed = false;
let mut linked_packages = HashSet::new();
let mut missing_sources: Vec<String> = Vec::new();
{
let mut db_guard = self.db.lock().await;
for dep in direct_deps {
let pkg_hash = dep
.package_hash
.clone()
.unwrap_or_else(|| format!("{}-{}", dep.name, dep.version));
let existing_pkg = DependencyPackage::filter(
DependencyPackage::fields()
.name()
.eq(&dep.name)
.and(DependencyPackage::fields().version().eq(&dep.version))
.and(DependencyPackage::fields().package_hash().eq(&pkg_hash)),
)
.first()
.exec(&mut *db_guard)
.await?;
let (package_id, needs_index, retry) = match existing_pkg {
Some(pkg) => (pkg.id, pkg.indexed_at.is_empty(), pkg.indexed_at.is_empty()),
None => {
let source_type_str = match &dep.source_type {
DependencySourceType::CratesIo => "cratesio",
DependencySourceType::Git => "git",
DependencySourceType::Path => "path",
DependencySourceType::Unknown(s) => s.as_str(),
};
let new_pkg = DependencyPackage::create()
.name(dep.name.clone())
.version(dep.version.clone())
.source_type(source_type_str.to_string())
.package_hash(pkg_hash)
.root_path(dep.root_path.to_string_lossy().to_string())
.manifest_path(dep.manifest_path.to_string_lossy().to_string())
.indexed_at(String::new())
.exec(&mut *db_guard)
.await?;
(new_pkg.id, true, false)
}
};
let existing_link = ProjectDependency::filter(
ProjectDependency::fields().project_id().eq(project_id).and(
ProjectDependency::fields()
.dependency_package_id()
.eq(package_id),
),
)
.first()
.exec(&mut *db_guard)
.await?;
linked_packages.insert(package_id);
if existing_link.is_none() {
links_changed = true;
ProjectDependency::create()
.project_id(project_id)
.dependency_package_id(package_id)
.is_direct(dep.is_direct)
.features(dep.features.join(","))
.exec(&mut *db_guard)
.await?;
}
let source_dir = provider
.resolve_source_directory(&dep)
.unwrap_or_else(|_| dep.root_path.clone());
if !source_dir.exists() {
missing_sources.push(dep.name.clone());
}
if needs_index {
to_index.push((dep, package_id, retry, source_dir));
}
}
}
links_changed |= self
.prune_dependency_links(project_id, &linked_packages)
.await?;
let register_ms = deps_timer.elapsed().as_millis() - discover_ms;
if links_changed {
self.search.bump(project_id);
}
if to_index.is_empty() {
return Ok((
0,
dependency_state(linked_packages.len(), std::mem::take(&mut missing_sources)),
));
}
let ingestion_handle =
BatchIngestionActor::spawn(self.db.clone(), 0, 2000, self.options.writer_batch_cap);
let sender = ingestion_handle.sender();
let concurrency = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(4)
.clamp(2, 8);
let semaphore = Arc::new(tokio::sync::Semaphore::new(concurrency));
let mut join_set = tokio::task::JoinSet::new();
for (dep, package_id, retry, source_dir) in to_index {
let permit = Arc::clone(&semaphore)
.acquire_owned()
.await
.map_err(|e| anyhow::anyhow!("Semaphore closed: {e}"))?;
let sender = sender.clone();
let rayon_pool = Arc::clone(&self.rayon_pool);
join_set.spawn(async move {
let _permit = permit;
let target_dir = if source_dir.join("src").is_dir() {
source_dir.join("src")
} else {
source_dir.clone()
};
let mut files_to_parse = Vec::new();
for entry in ignore::WalkBuilder::new(&target_dir)
.hidden(false)
.parents(false)
.build()
.flatten()
{
let path = entry.path();
if path.is_file()
&& path.extension().and_then(|s| s.to_str()) == Some("rs")
&& let Ok(rel_path) = path.strip_prefix(&source_dir)
&& let Ok(content) = fs::read(path)
{
let hash = crate::hasher::hash_bytes(&content);
files_to_parse.push((rel_path.to_path_buf(), hash, content));
}
}
if !files_to_parse.is_empty() {
let crate_name = dep.name.clone();
let parsed_files: Vec<anyhow::Result<ParsedFilePayload>> =
tokio::task::spawn_blocking(move || {
rayon_pool.install(|| {
files_to_parse
.into_par_iter()
.map(|(rel_path, hash, content)| {
crate::parsers::rust::RustAstParser::parse(
&rel_path,
&content,
&hash,
ParseMode::DependencyPublicSurface {
crate_name: crate_name.clone(),
},
)
})
.collect()
})
})
.await
.map_err(|e| anyhow::anyhow!("Join error in parse task: {e}"))?;
let status = || {
if retry {
ExistingFileStatus::Unknown
} else {
ExistingFileStatus::New
}
};
for payload in parsed_files {
let payload = match payload {
Ok(p) => p,
Err(e) => {
tracing::warn!("Skipping unparsable file in {}: {e}", dep.name);
continue;
}
};
sender
.send(IngestionMessage::UpsertFile {
payload,
status: status(),
dependency_package_id: Some(package_id),
})
.await
.map_err(|e| anyhow::anyhow!("Writer channel closed: {e}"))?;
}
}
Ok::<u64, anyhow::Error>(package_id)
});
}
let mut completed = Vec::new();
while let Some(res) = join_set.join_next().await {
match res {
Ok(Ok(package_id)) => completed.push(package_id),
Ok(Err(e)) => tracing::warn!("Dependency indexing failed: {e}"),
Err(e) => tracing::warn!("Dependency indexing task failed: {e}"),
}
}
let crates_ms = deps_timer.elapsed().as_millis() - discover_ms - register_ms;
drop(sender);
let outcome = ingestion_handle.finish().await?;
BatchIngestionActor::new(self.db.clone(), 0)
.link_deferred(outcome)
.await?;
self.search.bump(0);
let indexed_count = completed.len();
{
let mut db_guard = self.db.lock().await;
for package_id in completed {
if let Some(mut pkg) =
DependencyPackage::filter(DependencyPackage::fields().id().eq(package_id))
.first()
.exec(&mut *db_guard)
.await?
{
pkg.update()
.indexed_at(scryer_db::time::now_rfc3339())
.exec(&mut *db_guard)
.await?;
}
}
}
tracing::info!(
"index_dependencies timings: discover={discover_ms}ms register={register_ms}ms walk_parse_send={crates_ms}ms drain_writer={}ms crates={indexed_count}",
deps_timer.elapsed().as_millis() - discover_ms - register_ms - crates_ms
);
Ok((
indexed_count,
dependency_state(linked_packages.len(), missing_sources),
))
}
async fn prune_dependency_links(
&self,
project_id: u64,
keep: &HashSet<u64>,
) -> anyhow::Result<bool> {
let mut db_guard = self.db.lock().await;
let links =
ProjectDependency::filter(ProjectDependency::fields().project_id().eq(project_id))
.exec(&mut *db_guard)
.await?;
let mut removed = false;
for link in links
.into_iter()
.filter(|l| !keep.contains(&l.dependency_package_id))
{
let stmt = format!(
"DELETE FROM project_dependency WHERE project_id = {project_id} AND id = {};",
link.id
);
toasty::sql::statement(&stmt).exec(&mut *db_guard).await?;
removed = true;
}
Ok(removed)
}
pub async fn get_type_contract(
&self,
project_id: u64,
root: &Path,
type_name: &str,
file_path_hint: Option<&str>,
member_query: Option<&str>,
) -> anyhow::Result<RawTypeContract> {
let mut guard = self.db.lock().await;
let symbols =
scryer_db::Symbol::filter(scryer_db::Symbol::fields().project_id().eq(project_id))
.exec(&mut *guard)
.await?;
let matching_type = symbols.iter().find(|s| {
if let Some(hint) = file_path_hint {
let _ = hint;
}
(s.name == type_name || s.qualified_name == type_name)
&& matches!(
s.kind.as_str(),
"struct" | "enum" | "trait" | "type" | "class" | "interface"
)
});
let (type_sym, is_external, dep_root) = if let Some(sym) = matching_type {
(sym.clone(), false, None)
} else {
let dep_links = scryer_db::ProjectDependency::filter(
scryer_db::ProjectDependency::fields()
.project_id()
.eq(project_id),
)
.exec(&mut *guard)
.await?;
let dep_pkg_ids: Vec<u64> = dep_links
.into_iter()
.map(|l| l.dependency_package_id)
.collect();
let dep_symbols =
scryer_db::Symbol::filter(scryer_db::Symbol::fields().project_id().eq(0))
.exec(&mut *guard)
.await?;
let target_name = type_name.rsplit("::").next().unwrap_or(type_name);
let found_dep_sym = dep_symbols.into_iter().find(|s| {
let matches_name = if type_name.contains("::") {
s.qualified_name == type_name
|| s.qualified_name.ends_with(&format!("::{type_name}"))
} else {
s.name == target_name
};
matches_name
&& matches!(
s.kind.as_str(),
"struct" | "enum" | "trait" | "type" | "class" | "interface"
)
&& s.dependency_package_id
.is_some_and(|id| dep_pkg_ids.contains(&id))
});
if let Some(sym) = found_dep_sym {
let pkg_id = sym.dependency_package_id.unwrap();
let pkg = scryer_db::DependencyPackage::filter(
scryer_db::DependencyPackage::fields().id().eq(pkg_id),
)
.first()
.exec(&mut *guard)
.await?;
(sym, true, pkg.map(|p| PathBuf::from(p.root_path)))
} else {
return Ok(RawTypeContract {
found: false,
type_name: type_name.to_string(),
kind: None,
signature: None,
docstring: None,
file_path: None,
members: Vec::new(),
});
}
};
let file_proj_id = if is_external { 0 } else { project_id };
let file = scryer_db::SourceFile::filter(
scryer_db::SourceFile::fields()
.project_id()
.eq(file_proj_id)
.and(scryer_db::SourceFile::fields().id().eq(type_sym.file_id)),
)
.first()
.exec(&mut *guard)
.await?;
let file_rel_path = file.map(|f| f.path);
let mut members = Vec::new();
let mut final_file_path = None;
if let Some(rel) = &file_rel_path {
let full_disk_path = if is_external && let Some(dep_dir) = &dep_root {
dep_dir.join(rel)
} else {
root.join(rel)
};
final_file_path = Some(full_disk_path.to_string_lossy().to_string());
let extension = full_disk_path
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_ascii_lowercase();
if let Ok(source) = fs::read_to_string(&full_disk_path) {
members = match extension.as_str() {
"py" => Vec::new(),
"ts" | "js" => crate::parsers::typescript::extract_type_members(
&source,
&type_sym.name,
false,
),
"tsx" | "jsx" => crate::parsers::typescript::extract_type_members(
&source,
&type_sym.name,
true,
),
_ => crate::parsers::rust::extract_type_members(&source, &type_sym.name),
};
}
let dotted_names = matches!(extension.as_str(), "py" | "ts" | "tsx" | "js" | "jsx");
let pool_symbols = if is_external {
scryer_db::Symbol::filter(
scryer_db::Symbol::fields().project_id().eq(0).and(
scryer_db::Symbol::fields()
.dependency_package_id()
.eq(type_sym.dependency_package_id),
),
)
.exec(&mut *guard)
.await?
} else {
symbols
};
let crate_root = type_sym.qualified_name.split("::").next();
for s in &pool_symbols {
if dotted_names {
let in_class = s
.qualified_name
.rsplit_once('.')
.is_some_and(|(parent, _)| parent == type_sym.qualified_name);
let entry = match s.kind.as_str() {
"fn" | "method" if in_class => format!("method: {}", s.signature),
"var" | "const" if in_class => format!("field: {}", s.signature),
_ => continue,
};
if !members.contains(&entry) {
members.push(entry);
}
continue;
}
let Some((parent, _)) = s.qualified_name.rsplit_once("::") else {
continue;
};
let owned_by_type = parent == type_sym.qualified_name
|| (parent.rsplit("::").next() == Some(type_sym.name.as_str())
&& s.qualified_name.split("::").next() == crate_root);
if s.kind == "fn" && owned_by_type {
let entry = format!("method: {}", s.signature);
if !members.contains(&entry) {
members.push(entry);
}
}
}
}
if let Some(mq) = member_query {
let q_lower = mq.to_lowercase();
members.retain(|m| m.to_lowercase().contains(&q_lower));
}
Ok(RawTypeContract {
found: true,
type_name: type_sym.name,
kind: Some(type_sym.kind),
signature: Some(type_sym.signature),
docstring: type_sym.docstring,
file_path: final_file_path,
members,
})
}
pub async fn search_symbols(
&self,
project_id: u64,
query: SearchQuery,
) -> anyhow::Result<SearchResults> {
crate::search::symbols::search_symbols(
&self.db,
&self.rayon_pool,
&self.search,
project_id,
query,
)
.await
}
pub fn watcher(&self) -> anyhow::Result<WorkspaceWatcher> {
WorkspaceWatcher::start(self.clone(), Duration::from_millis(200))
}
pub fn watch_project(&self, project_id: u64, root: &Path) -> anyhow::Result<WorkspaceWatcher> {
let watcher = self.watcher()?;
watcher.add_project(project_id, root)?;
Ok(watcher)
}
}