use std::sync::Arc;
use pensieve_artifact_graph::content::ArtifactContentIndexer;
use pensieve_artifact_graph::ArtifactGraphWriter;
use pensieve_catalog::PostgresCatalog;
use pensieve_core::catalog::Catalog;
use pensieve_core::segment_format::SegmentFormat;
use pensieve_core::tenant::TenantId;
use object_store::ObjectStore;
pub async fn sync_artifact_nodes(
catalog: Arc<dyn Catalog>,
format: Arc<dyn SegmentFormat>,
tenant: TenantId,
) -> anyhow::Result<usize> {
let Some(pg) = catalog.as_ref_any().downcast_ref::<PostgresCatalog>() else {
return Ok(0);
};
let records = pg.list_live_artifacts(tenant).await?;
let writer = ArtifactGraphWriter::new(catalog.clone(), format);
writer.sync(&records).await
}
pub async fn sync_artifact_content(
catalog: Arc<dyn Catalog>,
format: Arc<dyn SegmentFormat>,
store: Arc<dyn ObjectStore>,
tenant: TenantId,
) -> anyhow::Result<usize> {
let Some(pg) = catalog.as_ref_any().downcast_ref::<PostgresCatalog>() else {
return Ok(0);
};
let embed = match pensieve_memory::shared_embedding().await {
Ok(e) => e,
Err(e) => {
tracing::warn!(error = %e, "no embedding backend; skipping artifact content indexing");
return Ok(0);
}
};
let records = pg.list_live_artifacts(tenant).await?;
let indexer = ArtifactContentIndexer::new(catalog.clone(), format, store, embed);
indexer.index(&records).await
}
pub async fn sync_artifact_nodes_all_tenants(
catalog: Arc<dyn Catalog>,
format: Arc<dyn SegmentFormat>,
store: Option<Arc<dyn ObjectStore>>,
) -> anyhow::Result<usize> {
let Some(pg) = catalog.as_ref_any().downcast_ref::<PostgresCatalog>() else {
return Ok(0);
};
let tenants = pg.list_artifact_tenants().await?;
let mut total = 0;
for tenant in tenants {
match sync_artifact_nodes(catalog.clone(), format.clone(), tenant).await {
Ok(n) => total += n,
Err(e) => tracing::warn!(error = %e, ?tenant, "artifact-graph sync failed for tenant"),
}
if let Some(store) = &store {
match sync_artifact_content(catalog.clone(), format.clone(), store.clone(), tenant)
.await
{
Ok(n) if n > 0 => {
tracing::info!(chunks = n, ?tenant, "indexed artifact content")
}
Ok(_) => {}
Err(e) => {
tracing::warn!(error = %e, ?tenant, "artifact content indexing failed for tenant")
}
}
}
}
Ok(total)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use pensieve_format_tlm::TelemetryFormat;
use object_store::memory::InMemory;
#[tokio::test]
async fn non_postgres_catalog_is_noop() {
let catalog: Arc<dyn Catalog> =
Arc::new(pensieve_catalog_sqlite::SqliteCatalog::connect_in_memory().await.unwrap());
let store = Arc::new(InMemory::new());
let format: Arc<dyn SegmentFormat> =
Arc::new(TelemetryFormat::new(store, "pensieve-test"));
let n = sync_artifact_nodes(catalog, format, pensieve_core::tenant::DEFAULT_TENANT)
.await
.unwrap();
assert_eq!(n, 0);
}
#[tokio::test]
async fn all_tenants_non_postgres_is_noop() {
let catalog: Arc<dyn Catalog> =
Arc::new(pensieve_catalog_sqlite::SqliteCatalog::connect_in_memory().await.unwrap());
let store = Arc::new(InMemory::new());
let format: Arc<dyn SegmentFormat> =
Arc::new(TelemetryFormat::new(store, "pensieve-test"));
let n = sync_artifact_nodes_all_tenants(catalog, format, None)
.await
.unwrap();
assert_eq!(n, 0);
}
}