fix: fix services accessig directly to localstorage

Prevent services accessing directly to localstorage and prefer using an astraction layer
to expose full blob. The abstraction layer (dedup services) will cover backend storage
election (local, s3, ...), encryption, etc

This change permit audio_metadata_service, media_metadaa_service, face_indexing_service to handle
blobs without worring of the backend.

note: prefered way to handle blob is the streamed way. Some services may not have this possibility
This commit is contained in:
Edouard Vanbelle
2026-08-02 19:20:23 +02:00
parent 74866bea81
commit 7663f803d3
11 changed files with 621 additions and 78 deletions
+8
View File
@@ -14,3 +14,11 @@ Non-obvious rules that trip up new code. Terse on purpose.
- Any new endpoint that mints or consumes credentials/tokens must consult one of the `is_*_login_allowed()` helpers, not the raw allowlist.
- Any new "policy-disabled" refusal must emit an `audit`-target line before returning — matches `auth.login_rejected`, `magic_link.redemption_rejected` conventions.
## Storage backend access
- **Read blob content through `Arc<DedupService>`.** It's the ONE canonical read abstraction — CDC-manifest-aware (`file.blob_hash` may reference a chunk manifest, not a blob), backend-agnostic (Local/S3/Azure), wrapper-transparent (encryption/retry/cache). Never take `Arc<dyn BlobStorageBackend>` directly in a service that reads content; you'll silently break on any file ≥ 64 KiB (`CDC_MIN_CHUNK`). Follow `thumbnail_service`, `audio_metadata_service`, `media_metadata_service`, `face_indexing_service`, `search_index::content_index_worker` as reference impls.
- **Reads use `DedupService` methods**: `dedup.read_blob_bytes(hash)` for byte-slice analyzers (ONNX, EXIF, ID3-via-Reader), `dedup.stream_blob_to_tempfile(hash, &temp_dir, ".ext")` for crates that only accept `&Path` (mp3_duration, ffprobe, `nom-exif` video), `dedup.read_blob_stream(hash)` for streaming to a downstream `Stream` consumer.
- **Never hand-craft blob paths.** No `blob_root: PathBuf` fields, no `<storage>/.blobs/<xx>/<hash>.blob` constructions. `BlobStorageBackend::local_blob_path` returns `None` under `EncryptedBlobBackend`; do not rely on it. The three services that did this pre-2026-08 (audio/media/face) are the anti-pattern — see memory `project_services_bypassing_blob_backend`.
- **Persistent state = backend**, not `<storage_path>/*` sidecars. Local sidecars (`.thumbnails/`, `.transcoded/`, `.blob-cache/`, `.search-index/`, `.plugin-logs/`, `.uploads/`) are only for caches (regenerable) or truly-temp scratch (deleted on drop). Anything a user would notice losing → blob backend. Tier-2 migration plan: `docs/plan/derived-blobs.md`.
- **Temp files use `OXICLOUD_TEMP_DIR`** via the shared config path (`AppConfig::temp_dir`) — not raw `std::env::temp_dir()`. Ops point it at real disk on RAM-constrained Linux deployments (default `/tmp` = tmpfs = RAM).
+18
View File
@@ -2252,6 +2252,19 @@ pub struct AppConfig {
pub storage_path: PathBuf,
/// Static files directory path
pub static_path: PathBuf,
/// Directory for tier-1 temporary data — pure scratch, safe to
/// lose at reboot. Backend `stream_to_tempfile` writes here so
/// extractors that require a `&Path` (id3, mp3_duration,
/// ffprobe, nom-exif video) can operate on a local file
/// without the service ever seeing a raw blob path.
///
/// Env: `OXICLOUD_TEMP_DIR`. Default: `std::env::temp_dir()`
/// (respects `$TMPDIR`). On Linux this is typically `/tmp`,
/// often mounted as tmpfs (RAM-backed); production
/// deployments concerned about physical RAM under
/// high concurrency should point this at a disk-backed
/// directory (e.g. `/var/lib/oxicloud/tmp`).
pub temp_dir: PathBuf,
/// Server port
pub server_port: u16,
/// Server host
@@ -2356,6 +2369,7 @@ impl Default for AppConfig {
Self {
storage_path: PathBuf::from("./storage"),
static_path: PathBuf::from("./static"),
temp_dir: env::temp_dir(),
server_port: 8086,
server_host: "127.0.0.1".to_string(),
cache: CacheConfig::default(),
@@ -2394,6 +2408,10 @@ impl AppConfig {
config.static_path = PathBuf::from(static_path);
}
if let Ok(temp_dir) = env::var("OXICLOUD_TEMP_DIR") {
config.temp_dir = PathBuf::from(temp_dir);
}
if let Ok(server_port) = env::var("OXICLOUD_SERVER_PORT")
&& let Ok(port) = server_port.parse::<u16>()
{
+14 -9
View File
@@ -468,11 +468,11 @@ impl AppServiceFactory {
);
// Audio metadata service — created here so it can be wired into file_lifecycle.
let audio_metadata_service = self.create_audio_metadata_service(db_pool);
let audio_metadata_service = self.create_audio_metadata_service(db_pool, &dedup_service);
// Image/video capture-metadata service — extracts EXIF/container capture
// dates so the Photos timeline groups by real capture time, not upload time.
let media_metadata_service = self.create_media_metadata_service(db_pool);
let media_metadata_service = self.create_media_metadata_service(db_pool, &dedup_service);
// ThumbnailRefreshHook: handles FileLifecycleHook events (create/update/delete).
// Implemented on ThumbnailRefreshHook (not ThumbnailService) to avoid circular Arc:
@@ -544,7 +544,7 @@ impl AppServiceFactory {
}
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));
fls = fls.with_hook(self.create_face_indexing_service(db_pool, &dedup_service));
}
let file_lifecycle = Arc::new(fls);
@@ -950,15 +950,16 @@ impl AppServiceFactory {
pub fn create_audio_metadata_service(
&self,
db_pool: &Arc<PgPool>,
dedup: &Arc<crate::infrastructure::services::dedup_service::DedupService>,
) -> Option<Arc<AudioMetadataService>> {
if !self.config.features.enable_music {
tracing::info!("Audio metadata service is disabled (music feature disabled)");
return None;
}
let blob_root = self.storage_path.join(".blobs");
Some(Arc::new(AudioMetadataService::new(
db_pool.clone(),
blob_root,
dedup.clone(),
self.config.temp_dir.clone(),
)))
}
@@ -967,9 +968,13 @@ impl AppServiceFactory {
pub fn create_media_metadata_service(
&self,
db_pool: &Arc<PgPool>,
dedup: &Arc<crate::infrastructure::services::dedup_service::DedupService>,
) -> Arc<MediaMetadataService> {
let blob_root = self.storage_path.join(".blobs");
Arc::new(MediaMetadataService::new(db_pool.clone(), blob_root))
Arc::new(MediaMetadataService::new(
db_pool.clone(),
dedup.clone(),
self.config.temp_dir.clone(),
))
}
/// Creates the trash service
@@ -1137,13 +1142,13 @@ impl AppServiceFactory {
pub fn create_face_indexing_service(
&self,
db_pool: &Arc<PgPool>,
dedup: &Arc<crate::infrastructure::services::dedup_service::DedupService>,
) -> Arc<crate::infrastructure::services::face_indexing_service::FaceIndexingService> {
let blob_root = self.storage_path.join(".blobs");
let analyzer = self.build_face_analyzer();
Arc::new(
crate::infrastructure::services::face_indexing_service::FaceIndexingService::new(
db_pool.clone(),
blob_root,
dedup.clone(),
analyzer,
),
)
@@ -8,6 +8,7 @@ use uuid::Uuid;
use crate::application::ports::file_lifecycle::FileLifecycleHook;
use crate::common::errors::DomainError;
use crate::infrastructure::services::dedup_service::DedupService;
#[derive(Debug, FromRow)]
pub struct AudioFileRow {
@@ -17,22 +18,32 @@ pub struct AudioFileRow {
pub struct AudioMetadataService {
pool: Arc<PgPool>,
blob_root: PathBuf,
/// CDC-aware blob reader. Same abstraction `thumbnail_service` uses —
/// hides both the chunk-manifest concatenation and the underlying
/// `BlobStorageBackend` wrapper stack.
dedup: Arc<DedupService>,
/// Tier-1 scratch directory for `stream_blob_to_tempfile`. Pulled
/// from `AppConfig::temp_dir` (env `OXICLOUD_TEMP_DIR`) at DI time.
temp_dir: PathBuf,
}
impl AudioMetadataService {
pub fn new(pool: Arc<PgPool>, blob_root: PathBuf) -> Self {
Self { pool, blob_root }
pub fn new(pool: Arc<PgPool>, dedup: Arc<DedupService>, temp_dir: PathBuf) -> Self {
Self {
pool,
dedup,
temp_dir,
}
}
pub fn is_audio_file(mime_type: &str) -> bool {
mime_type.starts_with("audio/")
}
pub fn spawn_extraction_background(service: Arc<Self>, file_id: Uuid, file_path: PathBuf) {
pub fn spawn_extraction_background(service: Arc<Self>, file_id: Uuid, blob_hash: String) {
tokio::spawn(async move {
tracing::info!("🎵 Extracting audio metadata for: {}", file_id);
if let Err(e) = service.extract_and_save(&file_id, &file_path).await {
if let Err(e) = service.extract_and_save(&file_id, &blob_hash).await {
tracing::warn!("Failed to extract audio metadata: {}", e);
}
});
@@ -41,22 +52,17 @@ impl AudioMetadataService {
pub fn spawn_extraction_with_delete_background(
service: Arc<Self>,
file_id: Uuid,
file_path: PathBuf,
blob_hash: String,
) {
tokio::spawn(async move {
tracing::info!("🎵 Updating audio metadata for: {}", file_id);
let _ = service.delete_metadata(&file_id).await;
if let Err(e) = service.extract_and_save(&file_id, &file_path).await {
if let Err(e) = service.extract_and_save(&file_id, &blob_hash).await {
tracing::warn!("Failed to update audio metadata: {}", e);
}
});
}
fn blob_path(&self, hash: &str) -> PathBuf {
let prefix = &hash[0..2];
self.blob_root.join(prefix).join(format!("{}.blob", hash))
}
/// Extract ID3 tag and MP3 duration from a file.
///
/// All I/O is synchronous (id3 + mp3_duration crates), so this MUST
@@ -104,15 +110,26 @@ impl AudioMetadataService {
pub async fn extract_and_save(
&self,
file_id: &Uuid,
file_path: &Path,
blob_hash: &str,
) -> Result<(), DomainError> {
info!(
"AudioMetadataService: blob_root={:?}, file_id={}, file_path={:?}",
self.blob_root, file_id, file_path,
"AudioMetadataService: extracting file_id={}, blob_hash={}",
file_id, blob_hash,
);
// Stream the blob (CDC-aware — chunks concatenated on the fly for
// chunked files; wrapper stack handles encryption + retry + cache)
// to a tempfile in the configured tier-1 temp dir, then hand its
// `.path()` to the id3 + mp3_duration crates which only expose
// `from_path` APIs. Peak process-heap = one chunk (~1 MiB)
// regardless of MP3 size. Guard drops → tempfile auto-removed.
let named = self
.dedup
.stream_blob_to_tempfile(blob_hash, &self.temp_dir, ".mp3")
.await?;
let path = named.path().to_path_buf();
// ── Sync I/O on the blocking thread pool (never stalls Tokio workers) ──
let path = file_path.to_path_buf();
let metadata = tokio::task::spawn_blocking(move || Self::extract_metadata_blocking(&path))
.await
.map_err(|e| {
@@ -121,6 +138,10 @@ impl AudioMetadataService {
format!("spawn_blocking join error: {e}"),
)
})?;
// Explicitly hold `named` alive until after the extraction — the
// `spawn_blocking` closure only borrows the raw `path`, so the
// guard must not drop while the extractor is running.
drop(named);
let Some(m) = metadata else {
return Ok(());
@@ -208,8 +229,10 @@ impl AudioMetadataService {
let audio_file = row.map_err(|e| {
DomainError::database_error(format!("Failed to fetch audio file row: {}", e))
})?;
let file_path = self.blob_path(&audio_file.blob_hash);
match self.extract_and_save(&audio_file.file_id, &file_path).await {
match self
.extract_and_save(&audio_file.file_id, &audio_file.blob_hash)
.await
{
Ok(()) => processed += 1,
Err(e) => {
warn!(
@@ -313,8 +336,7 @@ impl AudioMetadataService {
}
Ok(_) => {
// No existing metadata found — original not yet processed; fall back.
let file_path = service.blob_path(&blob_hash);
if let Err(e) = service.extract_and_save(&new_file_id, &file_path).await {
if let Err(e) = service.extract_and_save(&new_file_id, &blob_hash).await {
warn!(
"Failed to extract audio metadata for {}: {}",
new_file_id, e
@@ -368,10 +390,11 @@ impl FileLifecycleHook for AudioMetadataService {
};
let service = Arc::new(Self {
pool: self.pool.clone(),
blob_root: self.blob_root.clone(),
dedup: self.dedup.clone(),
temp_dir: self.temp_dir.clone(),
});
if is_new_blob {
Self::spawn_extraction_background(service, uuid, self.blob_path(blob_hash));
Self::spawn_extraction_background(service, uuid, blob_hash.to_string());
} else {
Self::clone_or_extract_background(service, uuid, blob_hash.to_string());
}
@@ -400,7 +423,8 @@ impl FileLifecycleHook for AudioMetadataService {
};
let service = Arc::new(Self {
pool: self.pool.clone(),
blob_root: self.blob_root.clone(),
dedup: self.dedup.clone(),
temp_dir: self.temp_dir.clone(),
});
Self::clone_from_source_background(service, uuid, source_uuid, blob_hash.to_string());
}
@@ -415,9 +439,10 @@ impl FileLifecycleHook for AudioMetadataService {
};
let service = Arc::new(Self {
pool: self.pool.clone(),
blob_root: self.blob_root.clone(),
dedup: self.dedup.clone(),
temp_dir: self.temp_dir.clone(),
});
Self::spawn_extraction_with_delete_background(service, uuid, self.blob_path(blob_hash));
Self::spawn_extraction_with_delete_background(service, uuid, blob_hash.to_string());
}
fn on_file_deleted(&self, _file_id: &str) {
@@ -2056,6 +2056,56 @@ impl DedupService {
Ok(Bytes::from(data))
}
/// Stream a blob to a temp file for extractors that only accept a
/// filesystem `Path` (id3, mp3_duration, ffprobe, nom-exif video).
/// CDC-aware — reads through [`Self::read_blob_stream`] so a chunked
/// file's chunks are concatenated on the fly. Peak process-heap =
/// one chunk (~1 MiB) regardless of blob size.
///
/// `temp_dir` is the destination directory (typically
/// `AppConfig::temp_dir`, from env `OXICLOUD_TEMP_DIR`). `suffix`
/// is appended to the tempfile name (e.g. `".mp3"`, `".jpg"`) so
/// content-sniffing extractors get a hint. The returned
/// `NamedTempFile` auto-removes on drop; callers pass `.path()`
/// to the extractor, then let the guard fall out of scope.
pub async fn stream_blob_to_tempfile(
&self,
hash: &str,
temp_dir: &std::path::Path,
suffix: &str,
) -> Result<tempfile::NamedTempFile, DomainError> {
use tokio::io::AsyncWriteExt;
let named = tempfile::Builder::new()
.prefix("oxi-blob-")
.suffix(suffix)
.tempfile_in(temp_dir)
.map_err(|e| {
DomainError::internal_error("Dedup", format!("mktemp in {:?}: {e}", temp_dir))
})?;
// Re-open with tokio's async File so we can await writes.
let path = named.path().to_path_buf();
let mut file = tokio::fs::OpenOptions::new()
.write(true)
.truncate(true)
.open(&path)
.await
.map_err(|e| DomainError::internal_error("Dedup", format!("reopen temp: {e}")))?;
// CDC-aware: manifest lookup + chunk concat OR legacy backend passthrough.
let mut stream = self.read_blob_stream(hash).await?;
while let Some(chunk) = stream.next().await {
let bytes = chunk
.map_err(|e| DomainError::internal_error("Dedup", format!("stream chunk: {e}")))?;
file.write_all(&bytes)
.await
.map_err(|e| DomainError::internal_error("Dedup", format!("temp write: {e}")))?;
}
file.flush()
.await
.map_err(|e| DomainError::internal_error("Dedup", format!("temp flush: {e}")))?;
drop(file);
Ok(named)
}
/// Stream a byte range — CDC-aware with legacy fallback.
///
/// For CDC files: calculates which chunks overlap the requested range,
@@ -3971,6 +4021,83 @@ mod rechunk_integration_tests {
cleanup(&pool, &hash, &files).await;
}
// ─── stream_blob_to_tempfile — CDC-aware read to a filesystem path ───
//
// Regression tests for the fix landed on `fix/services-use-blob-abstraction`:
// audio_metadata_service, media_metadata_service, and face_indexing_service
// all read blob content via DedupService (`read_blob_bytes` /
// `stream_blob_to_tempfile`), NOT the raw `BlobStorageBackend`. If someone
// reverts a service to `backend.get_blob_stream(hash)`, this test fails
// because `hash` is a chunk-manifest hash — the physical backend has no
// blob at that key. Bug returns silently otherwise; these tests catch it.
/// Local backend: seed a > 64 KiB blob, rechunk to CDC, then call
/// `stream_blob_to_tempfile` and verify the tempfile contents match
/// the original. Proves the CDC chunk-concat path works.
#[tokio::test]
async fn stream_blob_to_tempfile_reads_cdc_chunked_local() {
let pool = test_pool().await;
let dir = TempDir::new().unwrap();
let svc = local_svc(&pool, &dir).await;
// 200 KiB → forced multi-chunk after rechunk_legacy_blobs.
let data = content(200 * 1024, 33);
let (hash, files) = seed_legacy(&svc, &pool, &dir, &data, 1, "cdc-local", None).await;
svc.rechunk_legacy_blobs().await.expect("sweep");
// Sanity: rechunk actually produced a manifest (i.e. we're on the
// CDC path, not the legacy-fallback branch of read_blob_stream).
assert!(
manifest(&pool, &hash).await.is_some(),
"expected a CDC manifest after rechunk (test wouldn't cover the bug otherwise)"
);
// New method — the entry point audio/media services use.
let temp_dir = TempDir::new().unwrap();
let named = svc
.stream_blob_to_tempfile(&hash, temp_dir.path(), ".bin")
.await
.expect("stream_blob_to_tempfile must succeed on CDC-chunked blob");
let round_tripped = tokio::fs::read(named.path()).await.expect("read tempfile");
assert_eq!(round_tripped, data, "tempfile content must match original");
cleanup(&pool, &hash, &files).await;
}
/// Encrypted backend variant — proves the wrapper stack (decryption
/// on read) is honoured. Same regression class: if a service reads
/// raw ciphertext instead of going through DedupService, this fails.
#[tokio::test]
async fn stream_blob_to_tempfile_reads_cdc_chunked_encrypted() {
let pool = test_pool().await;
let dir = TempDir::new().unwrap();
let svc = encrypted_svc(&pool, &dir).await;
let data = content(150 * 1024, 77);
let (hash, files) = seed_legacy(&svc, &pool, &dir, &data, 1, "cdc-enc", None).await;
svc.rechunk_legacy_blobs().await.expect("sweep");
assert!(
manifest(&pool, &hash).await.is_some(),
"expected a CDC manifest after rechunk"
);
let temp_dir = TempDir::new().unwrap();
let named = svc
.stream_blob_to_tempfile(&hash, temp_dir.path(), ".bin")
.await
.expect("stream_blob_to_tempfile must succeed on encrypted CDC blob");
let round_tripped = tokio::fs::read(named.path()).await.expect("read tempfile");
assert_eq!(
round_tripped, data,
"tempfile content must match original plaintext (wrapper stack must decrypt transparently)"
);
cleanup(&pool, &hash, &files).await;
}
}
// ─────────────────────────────────────────────────────────────────────────────
@@ -1,14 +1,15 @@
//! 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.
//! background task) and stores them. Mirrors `ThumbnailService`: reads the
//! blob through `DedupService` (CDC-manifest lookup, wrapper-stack
//! delegation, encryption transparency — the service sees none of that),
//! 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;
@@ -20,6 +21,7 @@ use crate::application::ports::file_lifecycle::FileLifecycleHook;
use crate::common::errors::DomainError;
use crate::domain::entities::face::Face;
use crate::infrastructure::repositories::pg::FacePgRepository;
use crate::infrastructure::services::dedup_service::DedupService;
/// Minimum detector confidence for a face to be stored.
const MIN_DET_SCORE: f32 = 0.6;
@@ -48,7 +50,10 @@ pub struct FaceIndexingService {
pool: Arc<PgPool>,
repo: Arc<FacePgRepository>,
analyzer: Arc<dyn FaceAnalyzerPort>,
blob_root: PathBuf,
/// CDC-aware blob reader. Same abstraction `thumbnail_service` uses —
/// hides both the chunk-manifest concatenation and the underlying
/// `BlobStorageBackend` wrapper stack.
dedup: Arc<DedupService>,
/// Bounds concurrent indexing tasks. The lifecycle hooks spawn one
/// task per uploaded/copied image with no ceiling, so a bulk upload
/// used to fan out N simultaneous full-image reads + decodes +
@@ -60,23 +65,21 @@ pub struct FaceIndexingService {
}
impl FaceIndexingService {
pub fn new(pool: Arc<PgPool>, blob_root: PathBuf, analyzer: Arc<dyn FaceAnalyzerPort>) -> Self {
pub fn new(
pool: Arc<PgPool>,
dedup: Arc<DedupService>,
analyzer: Arc<dyn FaceAnalyzerPort>,
) -> Self {
let repo = Arc::new(FacePgRepository::new(pool.clone()));
Self {
pool,
repo,
analyzer,
blob_root,
dedup,
index_semaphore: Arc::new(tokio::sync::Semaphore::new(max_concurrent_index())),
}
}
/// 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).
@@ -84,7 +87,7 @@ impl FaceIndexingService {
let pool = self.pool.clone();
let repo = self.repo.clone();
let analyzer = self.analyzer.clone();
let blob_path = self.blob_path(&blob_hash);
let dedup = self.dedup.clone();
let semaphore = self.index_semaphore.clone();
tokio::spawn(async move {
// Queue behind the concurrency budget BEFORE touching the
@@ -102,7 +105,7 @@ impl FaceIndexingService {
&repo,
analyzer.as_ref(),
file_id,
&blob_path,
&dedup,
&blob_hash,
reuse_dedup,
)
@@ -161,7 +164,12 @@ impl FileLifecycleHook for FaceIndexingService {
}
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")
// Post-D7: `storage.files.user_id` was dropped in
// migrations/20260904000000_drop_files_folders_user_id.sql —
// provenance moved to `created_by` / `updated_by`. For the
// faces.user_id anchor, the file's original creator is what we
// want (matches the pre-D7 semantic of the dropped column).
let row: (Uuid,) = sqlx::query_as("SELECT created_by FROM storage.files WHERE id = $1")
.bind(file_id)
.fetch_one(pool)
.await
@@ -174,7 +182,7 @@ async fn index_file(
repo: &FacePgRepository,
analyzer: &dyn FaceAnalyzerPort,
file_id: Uuid,
blob_path: &Path,
dedup: &Arc<DedupService>,
blob_hash: &str,
reuse_dedup: bool,
) -> Result<(), DomainError> {
@@ -199,9 +207,12 @@ async fn index_file(
// 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}")))?;
// CDC-aware, backend-agnostic read: `DedupService` concatenates chunks
// for CDC files, delegates straight through for legacy whole-file
// blobs, and inherits the backend wrapper stack (encryption, retry,
// cache) transparently. Peak process-heap = image size, already
// bounded by `index_semaphore` above.
let bytes = dedup.read_blob_bytes(blob_hash).await?;
let detected = analyzer.analyze(&bytes).await?;
let faces: Vec<Face> = detected
@@ -32,6 +32,7 @@ use uuid::Uuid;
use crate::application::ports::file_lifecycle::FileLifecycleHook;
use crate::common::errors::DomainError;
use crate::infrastructure::repositories::pg::file_metadata_repository::FileMetadataRepository;
use crate::infrastructure::services::dedup_service::DedupService;
use crate::infrastructure::services::exif_service::{ExifMetadata, ExifService};
#[derive(Debug, FromRow)]
@@ -50,12 +51,22 @@ pub struct MetadataExtractionResult {
pub struct MediaMetadataService {
pool: Arc<PgPool>,
blob_root: PathBuf,
/// CDC-aware blob reader. Same abstraction `thumbnail_service` uses —
/// hides both the chunk-manifest concatenation and the underlying
/// `BlobStorageBackend` wrapper stack.
dedup: Arc<DedupService>,
/// Tier-1 scratch directory for `stream_blob_to_tempfile`. Pulled
/// from `AppConfig::temp_dir` (env `OXICLOUD_TEMP_DIR`) at DI time.
temp_dir: PathBuf,
}
impl MediaMetadataService {
pub fn new(pool: Arc<PgPool>, blob_root: PathBuf) -> Self {
Self { pool, blob_root }
pub fn new(pool: Arc<PgPool>, dedup: Arc<DedupService>, temp_dir: PathBuf) -> Self {
Self {
pool,
dedup,
temp_dir,
}
}
pub fn is_image_file(mime_type: &str) -> bool {
@@ -71,15 +82,11 @@ impl MediaMetadataService {
Self::is_image_file(mime_type) || Self::is_video_file(mime_type)
}
fn blob_path(&self, hash: &str) -> PathBuf {
let prefix = &hash[0..2];
self.blob_root.join(prefix).join(format!("{}.blob", hash))
}
fn arc(&self) -> Arc<Self> {
Arc::new(Self {
pool: self.pool.clone(),
blob_root: self.blob_root.clone(),
dedup: self.dedup.clone(),
temp_dir: self.temp_dir.clone(),
})
}
@@ -130,13 +137,32 @@ impl MediaMetadataService {
/// Extract metadata for one file and persist it (no-op when nothing useful
/// could be extracted).
///
/// Streams the blob (CDC-aware — chunks concatenated on the fly for
/// chunked files; wrapper stack handles encryption + retry + cache)
/// to a tempfile in the configured tier-1 temp dir, then hands its
/// `.path()` to the sync extractors (kamadak-exif via
/// `std::fs::read(path)` for images, `nom-exif` video track reader
/// for videos — both are path-based). Peak process-heap = one chunk
/// (~1 MiB) regardless of media size.
pub async fn extract_and_save(
&self,
file_id: &Uuid,
file_path: &Path,
blob_hash: &str,
mime_type: &str,
) -> Result<(), DomainError> {
let path = file_path.to_path_buf();
// File-extension hint for the tempfile suffix — helps
// `nom-exif`'s content sniffer land the right parser branch.
let suffix = if Self::is_video_file(mime_type) {
".mp4"
} else {
".jpg"
};
let named = self
.dedup
.stream_blob_to_tempfile(blob_hash, &self.temp_dir, suffix)
.await?;
let path = named.path().to_path_buf();
let mime = mime_type.to_string();
let meta = tokio::task::spawn_blocking(move || Self::extract_blocking(&path, &mime))
.await
@@ -146,6 +172,9 @@ impl MediaMetadataService {
format!("spawn_blocking join error: {e}"),
)
})?;
// Keep the tempfile alive until after extraction — the
// spawn_blocking closure only borrowed the raw path.
drop(named);
let Some(meta) = meta else {
return Ok(());
@@ -175,13 +204,13 @@ impl MediaMetadataService {
pub fn spawn_extraction_background(
service: Arc<Self>,
file_id: Uuid,
file_path: PathBuf,
blob_hash: String,
mime_type: String,
) {
tokio::spawn(async move {
tracing::info!("📷 Extracting capture metadata for: {}", file_id);
if let Err(e) = service
.extract_and_save(&file_id, &file_path, &mime_type)
.extract_and_save(&file_id, &blob_hash, &mime_type)
.await
{
tracing::warn!("Failed to extract capture metadata: {}", e);
@@ -192,14 +221,14 @@ impl MediaMetadataService {
pub fn spawn_extraction_with_delete_background(
service: Arc<Self>,
file_id: Uuid,
file_path: PathBuf,
blob_hash: String,
mime_type: String,
) {
tokio::spawn(async move {
tracing::info!("📷 Updating capture metadata for: {}", file_id);
let _ = service.delete_metadata(&file_id).await;
if let Err(e) = service
.extract_and_save(&file_id, &file_path, &mime_type)
.extract_and_save(&file_id, &blob_hash, &mime_type)
.await
{
tracing::warn!("Failed to update capture metadata: {}", e);
@@ -287,9 +316,8 @@ impl MediaMetadataService {
info!("Cloned capture metadata for file {}", new_file_id);
}
Ok(_) => {
let file_path = service.blob_path(&blob_hash);
if let Err(e) = service
.extract_and_save(&new_file_id, &file_path, &mime_type)
.extract_and_save(&new_file_id, &blob_hash, &mime_type)
.await
{
warn!(
@@ -335,9 +363,8 @@ impl MediaMetadataService {
let media = row.map_err(|e| {
DomainError::database_error(format!("Failed to fetch media file row: {}", e))
})?;
let file_path = self.blob_path(&media.blob_hash);
match self
.extract_and_save(&media.file_id, &file_path, &media.mime_type)
.extract_and_save(&media.file_id, &media.blob_hash, &media.mime_type)
.await
{
Ok(()) => processed += 1,
@@ -530,7 +557,7 @@ impl FileLifecycleHook for MediaMetadataService {
Self::spawn_extraction_background(
service,
uuid,
self.blob_path(blob_hash),
blob_hash.to_string(),
content_type.to_string(),
);
} else {
@@ -584,7 +611,7 @@ impl FileLifecycleHook for MediaMetadataService {
Self::spawn_extraction_with_delete_background(
self.arc(),
uuid,
self.blob_path(blob_hash),
blob_hash.to_string(),
content_type.to_string(),
);
}