feat(faces): indexing pipeline + DI wiring
Phase 2 increment 5: - FaceIndexingService: a FileLifecycleHook that, on image upload, detects + embeds faces in a background task and stores them. Dedup-aware (clones an identical blob's faces instead of re-running inference), reindexes on overwrite, and relies on the DB cascade for deletes. Completely inert when no model is ready. - DI: registers the hook in the FileLifecycleService chain and exposes PeopleService in AppState — both gated on OXICLOUD_ENABLE_FACES, both using the default no-op analyzer until the operator wires a real ONNX model. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JW6ghFMDtnRYuYNzZhb47M
This commit is contained in:
@@ -17,6 +17,7 @@ use crate::application::services::folder_service::FolderService;
|
||||
use crate::application::services::i18n_application_service::I18nApplicationService;
|
||||
use crate::application::services::nextcloud_file_id_service::NextcloudFileIdService;
|
||||
use crate::application::services::nextcloud_login_flow_service::NextcloudLoginFlowService;
|
||||
use crate::application::services::people_service::PeopleService;
|
||||
use crate::application::services::places_service::PlacesService;
|
||||
use crate::application::services::recent_service::RecentService;
|
||||
use crate::application::services::search_service::SearchService;
|
||||
@@ -360,6 +361,9 @@ impl AppServiceFactory {
|
||||
fls = fls.with_hook(audio.clone());
|
||||
}
|
||||
fls = fls.with_hook(media_metadata_service.clone());
|
||||
if self.config.features.enable_faces {
|
||||
fls = fls.with_hook(self.create_face_indexing_service(db_pool));
|
||||
}
|
||||
let file_lifecycle = Arc::new(fls);
|
||||
|
||||
Ok(CoreServices {
|
||||
@@ -810,6 +814,34 @@ impl AppServiceFactory {
|
||||
service
|
||||
}
|
||||
|
||||
/// Creates the face-indexing lifecycle hook (People feature). Uses the
|
||||
/// default no-op analyzer until the operator wires a real ONNX model.
|
||||
pub fn create_face_indexing_service(
|
||||
&self,
|
||||
db_pool: &Arc<PgPool>,
|
||||
) -> Arc<crate::infrastructure::services::face_indexing_service::FaceIndexingService> {
|
||||
let blob_root = self.storage_path.join(".blobs");
|
||||
let analyzer: Arc<dyn crate::application::ports::face_ports::FaceAnalyzerPort> =
|
||||
Arc::new(crate::infrastructure::services::noop_face_analyzer::NoopFaceAnalyzer);
|
||||
Arc::new(
|
||||
crate::infrastructure::services::face_indexing_service::FaceIndexingService::new(
|
||||
db_pool.clone(),
|
||||
blob_root,
|
||||
analyzer,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
/// Creates the People (faces) read/clustering service.
|
||||
pub fn create_people_service(&self, db_pool: &Arc<PgPool>) -> Arc<PeopleService> {
|
||||
let repo = Arc::new(
|
||||
crate::infrastructure::repositories::pg::FacePgRepository::new(db_pool.clone()),
|
||||
);
|
||||
let service = Arc::new(PeopleService::new(repo));
|
||||
tracing::info!("People service initialized");
|
||||
service
|
||||
}
|
||||
|
||||
/// Preloads translations for every locale in the registry. Build
|
||||
/// the registry at startup via `LocaleRegistry::discover` and pass
|
||||
/// the resulting list here.
|
||||
@@ -1018,6 +1050,7 @@ impl AppServiceFactory {
|
||||
let favorites_service: Option<Arc<FavoritesService>>;
|
||||
let recent_service: Option<Arc<RecentService>>;
|
||||
let places_service: Option<Arc<PlacesService>>;
|
||||
let people_service: Option<Arc<PeopleService>>;
|
||||
let storage_usage_service: Option<Arc<StorageUsageService>>;
|
||||
let mut auth_services: Option<crate::common::di::AuthServices> = None;
|
||||
let mut nextcloud_services: Option<NextcloudServices> = None;
|
||||
@@ -1046,6 +1079,12 @@ impl AppServiceFactory {
|
||||
None
|
||||
};
|
||||
|
||||
people_service = if core.config.features.enable_faces {
|
||||
Some(self.create_people_service(&pool))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
storage_usage_service = Some(storage_usage.clone());
|
||||
|
||||
self.start_tree_etag_flush_job(&maintenance_pool);
|
||||
@@ -1273,6 +1312,7 @@ impl AppServiceFactory {
|
||||
favorites_service,
|
||||
recent_service,
|
||||
places_service,
|
||||
people_service,
|
||||
storage_usage_service,
|
||||
calendar_service: None,
|
||||
contact_service: None,
|
||||
@@ -1720,6 +1760,7 @@ pub struct AppState {
|
||||
pub favorites_service: Option<Arc<FavoritesService>>,
|
||||
pub recent_service: Option<Arc<RecentService>>,
|
||||
pub places_service: Option<Arc<PlacesService>>,
|
||||
pub people_service: Option<Arc<PeopleService>>,
|
||||
pub storage_usage_service: Option<Arc<StorageUsageService>>,
|
||||
pub calendar_service: Option<Arc<CalendarService>>,
|
||||
pub contact_service: Option<Arc<ContactStorageAdapter>>,
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
//! Face indexing as a `FileLifecycleHook`.
|
||||
//!
|
||||
//! On image upload it detects + embeds faces (off the request path, in a
|
||||
//! background task) and stores them. It mirrors `MediaMetadataService`: reads
|
||||
//! the blob from the local `.blobs` tree, is dedup-aware (identical uploads
|
||||
//! clone an existing file's faces instead of re-running inference), and is
|
||||
//! completely inert when no model is configured (`FaceAnalyzerPort::is_ready()
|
||||
//! == false`) — so the feature compiles and runs with the default no-op
|
||||
//! analyzer until the operator wires a real ONNX model.
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
|
||||
use chrono::Utc;
|
||||
use sqlx::PgPool;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::application::ports::face_ports::{FaceAnalyzerPort, FaceRepository};
|
||||
use crate::application::ports::file_lifecycle::FileLifecycleHook;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::domain::entities::face::Face;
|
||||
use crate::infrastructure::repositories::pg::FacePgRepository;
|
||||
|
||||
/// Minimum detector confidence for a face to be stored.
|
||||
const MIN_DET_SCORE: f32 = 0.6;
|
||||
|
||||
fn is_image(content_type: &str) -> bool {
|
||||
content_type.starts_with("image/")
|
||||
}
|
||||
|
||||
pub struct FaceIndexingService {
|
||||
pool: Arc<PgPool>,
|
||||
repo: Arc<FacePgRepository>,
|
||||
analyzer: Arc<dyn FaceAnalyzerPort>,
|
||||
blob_root: PathBuf,
|
||||
}
|
||||
|
||||
impl FaceIndexingService {
|
||||
pub fn new(pool: Arc<PgPool>, blob_root: PathBuf, analyzer: Arc<dyn FaceAnalyzerPort>) -> Self {
|
||||
let repo = Arc::new(FacePgRepository::new(pool.clone()));
|
||||
Self {
|
||||
pool,
|
||||
repo,
|
||||
analyzer,
|
||||
blob_root,
|
||||
}
|
||||
}
|
||||
|
||||
/// Local path of a blob: `.blobs/{prefix}/{hash}.blob`.
|
||||
fn blob_path(&self, hash: &str) -> PathBuf {
|
||||
let prefix = if hash.len() >= 2 { &hash[0..2] } else { hash };
|
||||
self.blob_root.join(prefix).join(format!("{hash}.blob"))
|
||||
}
|
||||
|
||||
/// Spawn a background indexing task. `reuse_dedup` clones faces from an
|
||||
/// existing file with the same blob hash instead of re-running inference;
|
||||
/// `delete_first` clears prior faces (used on overwrite).
|
||||
fn spawn_index(&self, file_id: Uuid, blob_hash: String, reuse_dedup: bool, delete_first: bool) {
|
||||
let pool = self.pool.clone();
|
||||
let repo = self.repo.clone();
|
||||
let analyzer = self.analyzer.clone();
|
||||
let blob_path = self.blob_path(&blob_hash);
|
||||
tokio::spawn(async move {
|
||||
if delete_first {
|
||||
let _ = repo.delete_faces_for_file(file_id).await;
|
||||
}
|
||||
if let Err(e) = index_file(
|
||||
&pool,
|
||||
&repo,
|
||||
analyzer.as_ref(),
|
||||
file_id,
|
||||
&blob_path,
|
||||
&blob_hash,
|
||||
reuse_dedup,
|
||||
)
|
||||
.await
|
||||
{
|
||||
tracing::warn!(target: "oxicloud::faces", "face indexing failed for {file_id}: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
impl FileLifecycleHook for FaceIndexingService {
|
||||
fn on_file_created(
|
||||
&self,
|
||||
file_id: &str,
|
||||
blob_hash: &str,
|
||||
content_type: &str,
|
||||
is_new_blob: bool,
|
||||
) {
|
||||
if !is_image(content_type) || !self.analyzer.is_ready() {
|
||||
return;
|
||||
}
|
||||
if let Ok(fid) = file_id.parse::<Uuid>() {
|
||||
// Dedup hit (blob already existed) → clone an existing file's faces.
|
||||
self.spawn_index(fid, blob_hash.to_string(), !is_new_blob, false);
|
||||
}
|
||||
}
|
||||
|
||||
fn on_file_copied(
|
||||
&self,
|
||||
file_id: &str,
|
||||
blob_hash: &str,
|
||||
content_type: &str,
|
||||
_source_file_id: &str,
|
||||
) {
|
||||
if !is_image(content_type) || !self.analyzer.is_ready() {
|
||||
return;
|
||||
}
|
||||
if let Ok(fid) = file_id.parse::<Uuid>() {
|
||||
self.spawn_index(fid, blob_hash.to_string(), true, false);
|
||||
}
|
||||
}
|
||||
|
||||
fn on_file_updated(&self, file_id: &str, blob_hash: &str, content_type: &str) {
|
||||
if !is_image(content_type) || !self.analyzer.is_ready() {
|
||||
return;
|
||||
}
|
||||
if let Ok(fid) = file_id.parse::<Uuid>() {
|
||||
self.spawn_index(fid, blob_hash.to_string(), false, true);
|
||||
}
|
||||
}
|
||||
|
||||
fn on_file_deleted(&self, _file_id: &str) {
|
||||
// faces.faces.file_id has ON DELETE CASCADE — the DB cleans up.
|
||||
}
|
||||
}
|
||||
|
||||
async fn lookup_user(pool: &PgPool, file_id: Uuid) -> Result<Uuid, DomainError> {
|
||||
let row: (Uuid,) = sqlx::query_as("SELECT user_id FROM storage.files WHERE id = $1")
|
||||
.bind(file_id)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("Faces", format!("lookup user: {e}")))?;
|
||||
Ok(row.0)
|
||||
}
|
||||
|
||||
async fn index_file(
|
||||
pool: &PgPool,
|
||||
repo: &FacePgRepository,
|
||||
analyzer: &dyn FaceAnalyzerPort,
|
||||
file_id: Uuid,
|
||||
blob_path: &Path,
|
||||
blob_hash: &str,
|
||||
reuse_dedup: bool,
|
||||
) -> Result<(), DomainError> {
|
||||
let user_id = lookup_user(pool, file_id).await?;
|
||||
|
||||
// Dedup-aware fast path: reuse faces already computed for an identical blob.
|
||||
if reuse_dedup {
|
||||
let peers = repo.faces_for_blob(user_id, blob_hash).await?;
|
||||
let cloned: Vec<Face> = peers
|
||||
.into_iter()
|
||||
.filter(|f| f.file_id != file_id)
|
||||
.map(|f| Face {
|
||||
id: Uuid::new_v4(),
|
||||
file_id,
|
||||
..f
|
||||
})
|
||||
.collect();
|
||||
if !cloned.is_empty() {
|
||||
repo.save_faces(&cloned).await?;
|
||||
return Ok(());
|
||||
}
|
||||
// No peer found — fall through and analyze.
|
||||
}
|
||||
|
||||
let bytes = tokio::fs::read(blob_path)
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("Faces", format!("read blob: {e}")))?;
|
||||
let detected = analyzer.analyze(&bytes).await?;
|
||||
|
||||
let faces: Vec<Face> = detected
|
||||
.into_iter()
|
||||
.filter(|d| d.det_score >= MIN_DET_SCORE)
|
||||
.map(|d| Face {
|
||||
id: Uuid::new_v4(),
|
||||
file_id,
|
||||
user_id,
|
||||
person_id: None,
|
||||
bbox: d.bbox,
|
||||
det_score: d.det_score,
|
||||
quality: d.quality,
|
||||
embedding: d.embedding,
|
||||
blob_hash: Some(blob_hash.to_string()),
|
||||
created_at: Utc::now(),
|
||||
})
|
||||
.collect();
|
||||
repo.save_faces(&faces).await
|
||||
}
|
||||
@@ -6,6 +6,7 @@ pub mod compression_service;
|
||||
pub mod dedup_service;
|
||||
pub mod encrypted_blob_backend;
|
||||
pub mod exif_service;
|
||||
pub mod face_indexing_service;
|
||||
pub mod file_content_cache;
|
||||
pub mod file_system_i18n_service;
|
||||
pub mod image_transcode_service;
|
||||
|
||||
Reference in New Issue
Block a user