Merge pull request #485 from AtalayaLabs/claude/zealous-faraday-58s1at
feat: Photos evolution — Places (map) & People (faces) + gallery polish
This commit is contained in:
@@ -0,0 +1,322 @@
|
||||
//! PostgreSQL repository for the People (faces) feature.
|
||||
//!
|
||||
//! Embeddings are stored as `BYTEA` (512 × little-endian `f32`); there is no
|
||||
//! pgvector dependency. Similarity search / clustering is done in-app over the
|
||||
//! decoded vectors (see `PeopleService`).
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use chrono::{DateTime, Utc};
|
||||
use sqlx::PgPool;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::application::ports::face_ports::FaceRepository;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::domain::entities::face::{BoundingBox, Face, Person};
|
||||
|
||||
/// Row shape for `faces.faces` selects (avoids `clippy::type_complexity`).
|
||||
type FaceRow = (
|
||||
Uuid, // id
|
||||
Uuid, // file_id
|
||||
Uuid, // user_id
|
||||
Option<Uuid>, // person_id
|
||||
Vec<f32>, // bbox (REAL[])
|
||||
f32, // det_score
|
||||
Option<f32>, // quality
|
||||
Vec<u8>, // embedding (BYTEA)
|
||||
Option<String>, // blob_hash
|
||||
DateTime<Utc>, // created_at
|
||||
);
|
||||
|
||||
type PersonRow = (
|
||||
Uuid, // id
|
||||
Uuid, // user_id
|
||||
Option<String>, // display_name
|
||||
Option<Uuid>, // cover_face_id
|
||||
bool, // is_hidden
|
||||
DateTime<Utc>, // created_at
|
||||
);
|
||||
|
||||
fn embedding_to_bytes(e: &[f32]) -> Vec<u8> {
|
||||
let mut out = Vec::with_capacity(e.len() * 4);
|
||||
for v in e {
|
||||
out.extend_from_slice(&v.to_le_bytes());
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
fn bytes_to_embedding(b: &[u8]) -> Vec<f32> {
|
||||
b.chunks_exact(4)
|
||||
.map(|c| f32::from_le_bytes([c[0], c[1], c[2], c[3]]))
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn row_to_face(r: FaceRow) -> Face {
|
||||
let (
|
||||
id,
|
||||
file_id,
|
||||
user_id,
|
||||
person_id,
|
||||
bbox,
|
||||
det_score,
|
||||
quality,
|
||||
embedding,
|
||||
blob_hash,
|
||||
created_at,
|
||||
) = r;
|
||||
Face {
|
||||
id,
|
||||
file_id,
|
||||
user_id,
|
||||
person_id,
|
||||
bbox: BoundingBox::from_slice(&bbox),
|
||||
det_score,
|
||||
quality,
|
||||
embedding: bytes_to_embedding(&embedding),
|
||||
blob_hash,
|
||||
created_at,
|
||||
}
|
||||
}
|
||||
|
||||
fn row_to_person(r: PersonRow) -> Person {
|
||||
let (id, user_id, display_name, cover_face_id, is_hidden, created_at) = r;
|
||||
Person {
|
||||
id,
|
||||
user_id,
|
||||
display_name,
|
||||
cover_face_id,
|
||||
is_hidden,
|
||||
created_at,
|
||||
}
|
||||
}
|
||||
|
||||
fn db_err(ctx: &'static str, e: sqlx::Error) -> DomainError {
|
||||
DomainError::internal_error("FacePg", format!("{ctx}: {e}"))
|
||||
}
|
||||
|
||||
const FACE_COLS: &str =
|
||||
"id, file_id, user_id, person_id, bbox, det_score, quality, embedding, blob_hash, created_at";
|
||||
const PERSON_COLS: &str = "id, user_id, display_name, cover_face_id, is_hidden, created_at";
|
||||
|
||||
pub struct FacePgRepository {
|
||||
pool: Arc<PgPool>,
|
||||
}
|
||||
|
||||
impl FacePgRepository {
|
||||
pub fn new(pool: Arc<PgPool>) -> Self {
|
||||
Self { pool }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl FaceRepository for FacePgRepository {
|
||||
async fn save_faces(&self, faces: &[Face]) -> Result<(), DomainError> {
|
||||
if faces.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
let mut tx = self.pool.begin().await.map_err(|e| db_err("begin", e))?;
|
||||
for f in faces {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO faces.faces
|
||||
(id, file_id, user_id, person_id, bbox, det_score, quality, embedding, blob_hash)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind(f.id)
|
||||
.bind(f.file_id)
|
||||
.bind(f.user_id)
|
||||
.bind(f.person_id)
|
||||
.bind(f.bbox.to_array())
|
||||
.bind(f.det_score)
|
||||
.bind(f.quality)
|
||||
.bind(embedding_to_bytes(&f.embedding))
|
||||
.bind(f.blob_hash.as_deref())
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| db_err("save_faces", e))?;
|
||||
}
|
||||
tx.commit().await.map_err(|e| db_err("commit", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn faces_for_file(&self, file_id: Uuid) -> Result<Vec<Face>, DomainError> {
|
||||
let sql = format!("SELECT {FACE_COLS} FROM faces.faces WHERE file_id = $1");
|
||||
let rows: Vec<FaceRow> = sqlx::query_as(&sql)
|
||||
.bind(file_id)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("faces_for_file", e))?;
|
||||
Ok(rows.into_iter().map(row_to_face).collect())
|
||||
}
|
||||
|
||||
async fn delete_faces_for_file(&self, file_id: Uuid) -> Result<(), DomainError> {
|
||||
sqlx::query("DELETE FROM faces.faces WHERE file_id = $1")
|
||||
.bind(file_id)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("delete_faces_for_file", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn faces_for_user(&self, user_id: Uuid) -> Result<Vec<Face>, DomainError> {
|
||||
let sql = format!("SELECT {FACE_COLS} FROM faces.faces WHERE user_id = $1");
|
||||
let rows: Vec<FaceRow> = sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("faces_for_user", e))?;
|
||||
Ok(rows.into_iter().map(row_to_face).collect())
|
||||
}
|
||||
|
||||
async fn faces_for_blob(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
blob_hash: &str,
|
||||
) -> Result<Vec<Face>, DomainError> {
|
||||
let sql =
|
||||
format!("SELECT {FACE_COLS} FROM faces.faces WHERE user_id = $1 AND blob_hash = $2");
|
||||
let rows: Vec<FaceRow> = sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(blob_hash)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("faces_for_blob", e))?;
|
||||
Ok(rows.into_iter().map(row_to_face).collect())
|
||||
}
|
||||
|
||||
async fn assign_person(
|
||||
&self,
|
||||
face_id: Uuid,
|
||||
person_id: Option<Uuid>,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query("UPDATE faces.faces SET person_id = $2 WHERE id = $1")
|
||||
.bind(face_id)
|
||||
.bind(person_id)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("assign_person", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn create_person(&self, person: &Person) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO faces.persons (id, user_id, display_name, cover_face_id, is_hidden)
|
||||
VALUES ($1, $2, $3, $4, $5)
|
||||
"#,
|
||||
)
|
||||
.bind(person.id)
|
||||
.bind(person.user_id)
|
||||
.bind(person.display_name.as_deref())
|
||||
.bind(person.cover_face_id)
|
||||
.bind(person.is_hidden)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("create_person", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn persons_for_user(&self, user_id: Uuid) -> Result<Vec<Person>, DomainError> {
|
||||
let sql = format!(
|
||||
"SELECT {PERSON_COLS} FROM faces.persons WHERE user_id = $1 ORDER BY created_at"
|
||||
);
|
||||
let rows: Vec<PersonRow> = sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("persons_for_user", e))?;
|
||||
Ok(rows.into_iter().map(row_to_person).collect())
|
||||
}
|
||||
|
||||
async fn rename_person(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
person_id: Uuid,
|
||||
name: Option<String>,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
"UPDATE faces.persons SET display_name = $3, updated_at = now() WHERE id = $2 AND user_id = $1",
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(person_id)
|
||||
.bind(name)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("rename_person", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn set_person_cover(
|
||||
&self,
|
||||
person_id: Uuid,
|
||||
cover_face_id: Uuid,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
"UPDATE faces.persons SET cover_face_id = $2, updated_at = now() WHERE id = $1",
|
||||
)
|
||||
.bind(person_id)
|
||||
.bind(cover_face_id)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("set_person_cover", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn set_person_hidden(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
person_id: Uuid,
|
||||
hidden: bool,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
"UPDATE faces.persons SET is_hidden = $3, updated_at = now() WHERE id = $2 AND user_id = $1",
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(person_id)
|
||||
.bind(hidden)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("set_person_hidden", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn files_for_person(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
person_id: Uuid,
|
||||
) -> Result<Vec<Uuid>, DomainError> {
|
||||
let rows: Vec<(Uuid,)> = sqlx::query_as(
|
||||
r#"
|
||||
SELECT file_id
|
||||
FROM faces.faces
|
||||
WHERE user_id = $1 AND person_id = $2
|
||||
GROUP BY file_id
|
||||
ORDER BY max(created_at) DESC
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(person_id)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("files_for_person", e))?;
|
||||
Ok(rows.into_iter().map(|(id,)| id).collect())
|
||||
}
|
||||
|
||||
async fn delete_all_for_user(&self, user_id: Uuid) -> Result<(), DomainError> {
|
||||
let mut tx = self.pool.begin().await.map_err(|e| db_err("begin", e))?;
|
||||
sqlx::query("DELETE FROM faces.faces WHERE user_id = $1")
|
||||
.bind(user_id)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| db_err("delete_all_faces", e))?;
|
||||
sqlx::query("DELETE FROM faces.persons WHERE user_id = $1")
|
||||
.bind(user_id)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| db_err("delete_all_persons", e))?;
|
||||
tx.commit().await.map_err(|e| db_err("commit", e))?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -20,6 +20,8 @@ type MediaFileRow = (
|
||||
String, // blob_hash
|
||||
Option<Uuid>, // user_id
|
||||
i64, // sort_date
|
||||
Option<i32>, // width
|
||||
Option<i32>, // height
|
||||
);
|
||||
|
||||
use bytes::Bytes;
|
||||
@@ -30,6 +32,7 @@ use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::application::dtos::geo_dto::{GeoBounds, GeoCluster};
|
||||
use crate::application::dtos::search_dto::SearchCriteriaDto;
|
||||
use crate::application::ports::storage_ports::FileReadPort;
|
||||
use crate::common::errors::DomainError;
|
||||
@@ -416,7 +419,7 @@ impl FileBlobReadRepository {
|
||||
owner_id: Uuid,
|
||||
before: Option<i64>,
|
||||
limit: i64,
|
||||
) -> Result<(Vec<File>, Vec<i64>), DomainError> {
|
||||
) -> Result<(Vec<File>, Vec<i64>, Vec<(Option<i32>, Option<i32>)>), DomainError> {
|
||||
let rows: Vec<MediaFileRow> = sqlx::query_as(
|
||||
r#"
|
||||
SELECT fi.id::text, fi.name, fi.folder_id::text, fo.path,
|
||||
@@ -425,9 +428,11 @@ impl FileBlobReadRepository {
|
||||
EXTRACT(EPOCH FROM fi.updated_at)::bigint,
|
||||
fi.blob_hash,
|
||||
fi.user_id,
|
||||
EXTRACT(EPOCH FROM fi.media_sort_date)::bigint AS sort_date
|
||||
EXTRACT(EPOCH FROM fi.media_sort_date)::bigint AS sort_date,
|
||||
fm.width, fm.height
|
||||
FROM storage.files fi
|
||||
LEFT JOIN storage.folders fo ON fo.id = fi.folder_id
|
||||
LEFT JOIN storage.file_metadata fm ON fm.file_id = fi.id
|
||||
WHERE fi.user_id = $1
|
||||
AND NOT fi.is_trashed
|
||||
AND (fi.mime_type LIKE 'image/%' OR fi.mime_type LIKE 'video/%')
|
||||
@@ -446,15 +451,67 @@ impl FileBlobReadRepository {
|
||||
|
||||
let mut files = Vec::with_capacity(rows.len());
|
||||
let mut sort_dates = Vec::with_capacity(rows.len());
|
||||
let mut dims = Vec::with_capacity(rows.len());
|
||||
|
||||
for (id, name, fid, fpath, size, mime, ca, ma, blob_hash, uid, sd) in rows {
|
||||
for (id, name, fid, fpath, size, mime, ca, ma, blob_hash, uid, sd, w, h) in rows {
|
||||
files.push(Self::row_to_file(
|
||||
id, name, fid, fpath, size, mime, ca, ma, blob_hash, uid,
|
||||
)?);
|
||||
sort_dates.push(sd);
|
||||
dims.push((w, h));
|
||||
}
|
||||
|
||||
Ok((files, sort_dates))
|
||||
Ok((files, sort_dates, dims))
|
||||
}
|
||||
|
||||
/// Aggregate the caller's geotagged photos into grid cells of side `cell`
|
||||
/// (degrees) within `bounds`. Plain SQL (no PostGIS), scoped to `user_id`.
|
||||
/// Returns one cluster per non-empty cell with its centroid, photo count
|
||||
/// and a representative photo id (for the cluster thumbnail).
|
||||
pub async fn list_geo_clusters(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
bounds: GeoBounds,
|
||||
cell: f64,
|
||||
) -> Result<Vec<GeoCluster>, DomainError> {
|
||||
let rows: Vec<(i64, f64, f64, String)> = sqlx::query_as(
|
||||
r#"
|
||||
SELECT count(*) AS n,
|
||||
avg(fm.longitude) AS clng,
|
||||
avg(fm.latitude) AS clat,
|
||||
min(fm.file_id::text) AS sample_id
|
||||
FROM storage.file_metadata fm
|
||||
JOIN storage.files fi ON fi.id = fm.file_id
|
||||
WHERE fi.user_id = $1
|
||||
AND NOT fi.is_trashed
|
||||
AND fm.latitude IS NOT NULL
|
||||
AND fm.longitude IS NOT NULL
|
||||
AND fm.longitude BETWEEN $2 AND $3
|
||||
AND fm.latitude BETWEEN $4 AND $5
|
||||
GROUP BY round(fm.longitude / $6), round(fm.latitude / $6)
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(bounds.west)
|
||||
.bind(bounds.east)
|
||||
.bind(bounds.south)
|
||||
.bind(bounds.north)
|
||||
.bind(cell)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FileBlobRead", format!("list_geo_clusters: {e}"))
|
||||
})?;
|
||||
|
||||
Ok(rows
|
||||
.into_iter()
|
||||
.map(|(n, clng, clat, sample_id)| GeoCluster {
|
||||
lng: clng,
|
||||
lat: clat,
|
||||
count: n,
|
||||
sample_file_id: sample_id,
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ mod contact_group_pg_repository;
|
||||
mod contact_persistence_dto;
|
||||
mod contact_pg_repository;
|
||||
mod device_code_pg_repository;
|
||||
mod face_pg_repository;
|
||||
mod favorites_pg_repository;
|
||||
pub mod file_metadata_repository;
|
||||
mod magic_link_token_pg_repository;
|
||||
@@ -33,6 +34,7 @@ pub use contact_group_pg_repository::ContactGroupPgRepository;
|
||||
pub use contact_persistence_dto::*;
|
||||
pub use contact_pg_repository::ContactPgRepository;
|
||||
pub use device_code_pg_repository::DeviceCodePgRepository;
|
||||
pub use face_pg_repository::FacePgRepository;
|
||||
pub use favorites_pg_repository::FavoritesPgRepository;
|
||||
pub use file_blob_read_repository::FileBlobReadRepository;
|
||||
pub use file_blob_write_repository::FileBlobWriteRepository;
|
||||
|
||||
Reference in New Issue
Block a user