Sync face data as sealed shards, so a second device does not re-index
Indexing 23,500 images is about two hours of CPU, and the result is byte-identical on every device: the same model over the same proxy produces the same embedding. Paying for it once per account rather than once per device is the point. Shards rather than the catalog snapshot, because the snapshot goes up whole on every sync and a fully indexed library carries roughly 30 MB of embeddings. That is exactly the cost the thumbnail store's 25 MB cap exists to bound, so face shards use the same cap -- imported from dr_thumbs rather than restated, since the number is a statement about sync cost and the two must not drift apart. The split follows the one already there: bulk immutable data in sealed shards, small mutable data in the catalog snapshot. Faces, landmarks, embeddings and run markers shard; people, names and assignments ride the catalog and merge by uuid. Keyed on oc:fileid throughout, never on image_id, because a row id means nothing on another device. The run marker travels with the faces it describes. Without it a receiving device cannot tell an image with no faces from one never examined, and would re-detect every landscape it had just adopted. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -12,6 +12,10 @@ dr-types.workspace = true
|
||||
# features are off, so this brings in no ONNX runtime and no weights: only the
|
||||
# model-free half compiles here.
|
||||
dr-face.workspace = true
|
||||
# For `SHARD_MAX_BYTES` alone. The face shards are capped at the same 25 MB the
|
||||
# thumbnail shards are, and sharing the constant is what keeps them from
|
||||
# drifting apart — the cap is a statement about sync cost, not about thumbnails.
|
||||
dr-thumbs.workspace = true
|
||||
# The `Storage` trait, and nothing else from it. A scan has to read a real
|
||||
# directory, and this is how `core/` reaches the platform without a
|
||||
# `#[cfg(target_os)]` of its own (ARCH §4.1: calls go downward).
|
||||
|
||||
@@ -0,0 +1,958 @@
|
||||
//! TRACES: FR-CULL-8 | FR-NC-7 | NFR-RES-4
|
||||
//! Face data as sealed shards, so a second device does not re-index the library.
|
||||
//!
|
||||
//! Indexing a 23,500-image library is on the order of two hours of CPU
|
||||
//! (docs/faces.md §12.2). It is also **byte-identical on every device**: the
|
||||
//! same model over the same proxy produces the same embedding. Paying for it
|
||||
//! once per account rather than once per device is the whole point of this
|
||||
//! module, and it is the same bargain the thumbnail store already makes.
|
||||
//!
|
||||
//! # Why shards, and not the catalog snapshot
|
||||
//!
|
||||
//! The catalog is uploaded whole on every sync. Embeddings are 1 KB each, so a
|
||||
//! fully indexed 23,500-image library carries roughly 30 MB of them — which
|
||||
//! would be re-uploaded in its entirety every time anything in the catalog
|
||||
//! changed, and re-downloaded by every device. That is precisely the cost
|
||||
//! `dr_thumbs`'s 25 MB shard cap exists to bound.
|
||||
//!
|
||||
//! So face data splits the way thumbnails and collections already split:
|
||||
//!
|
||||
//! - **Here, in sealed shards:** the faces themselves — box, landmarks,
|
||||
//! embedding, and the run marker. Bulk, immutable once written, and identical
|
||||
//! everywhere. A sealed shard is never rewritten, so a client downloads each
|
||||
//! one exactly once and never asks about it again.
|
||||
//! - **In the catalog snapshot, which already syncs:** people, their names, and
|
||||
//! which face belongs to whom. Small, changes constantly, and merges by uuid.
|
||||
//!
|
||||
//! # Keyed on `file_id`, never on `image_id`
|
||||
//!
|
||||
//! A row id is local and means nothing on another device. `oc:fileid` is the
|
||||
//! identity every client agrees on (FR-NC-5), and it survives a server-side
|
||||
//! move, so a shard written before a reorganisation still applies after it.
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
use rusqlite::{Connection, OptionalExtension};
|
||||
|
||||
use crate::error::CatalogError;
|
||||
|
||||
/// Maximum bytes in a shard before a new one is started.
|
||||
///
|
||||
/// The same 25 MB `dr_thumbs` uses, for the same reason and deliberately not a
|
||||
/// separate number: it bounds what a client re-downloads when the open shard
|
||||
/// changes, and keeps a stalled transfer cheap to retry. At ~1.1 KB per face
|
||||
/// that is roughly 22,000 faces per shard, so a large personal library lands in
|
||||
/// one or two.
|
||||
pub const SHARD_MAX_BYTES: u64 = dr_thumbs::SHARD_MAX_BYTES;
|
||||
|
||||
/// Bytes one stored face occupies, near enough to bound a shard by.
|
||||
///
|
||||
/// Counted rather than measured: the embedding is fixed at 512 × f16, the
|
||||
/// landmarks at 5 × 2 × f32, and the rest is a handful of numbers. Measuring
|
||||
/// the file after each insert would mean a `VACUUM` to get an honest answer.
|
||||
const BYTES_PER_FACE: u64 = 1024 + 40 + 64;
|
||||
|
||||
/// One face as it travels between devices.
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct SharedFace {
|
||||
/// `oc:fileid` of the photograph — the cross-device image identity.
|
||||
pub file_id: u64,
|
||||
pub model_id: String,
|
||||
pub x: f32,
|
||||
pub y: f32,
|
||||
pub w: f32,
|
||||
pub h: f32,
|
||||
pub landmarks: Vec<u8>,
|
||||
pub confidence: f32,
|
||||
pub embedding: Vec<u8>,
|
||||
pub crop_px: f32,
|
||||
}
|
||||
|
||||
/// A store of face shards, beside the thumbnail store.
|
||||
pub struct FaceShardStore {
|
||||
dir: PathBuf,
|
||||
index: Connection,
|
||||
client: String,
|
||||
}
|
||||
|
||||
impl FaceShardStore {
|
||||
pub fn open(dir: &Path) -> Result<Self, CatalogError> {
|
||||
std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?;
|
||||
let index = Connection::open(dir.join("index.sqlite"))?;
|
||||
index.execute_batch(INDEX_SCHEMA)?;
|
||||
let client = mint_client_id(&index)?;
|
||||
Ok(Self {
|
||||
dir: dir.to_path_buf(),
|
||||
index,
|
||||
client,
|
||||
})
|
||||
}
|
||||
|
||||
/// This device's id, which qualifies its shard numbering.
|
||||
///
|
||||
/// Two devices both writing `shard-0000` would collide on the server; the
|
||||
/// name carries the writer so they cannot.
|
||||
pub fn client_id(&self) -> &str {
|
||||
&self.client
|
||||
}
|
||||
|
||||
pub fn contains(&self, file_id: u64, model_id: &str) -> bool {
|
||||
self.index
|
||||
.query_row(
|
||||
"SELECT 1 FROM entries WHERE file_id = ?1 AND model_id = ?2",
|
||||
rusqlite::params![file_id as i64, model_id],
|
||||
|_| Ok(()),
|
||||
)
|
||||
.optional()
|
||||
.ok()
|
||||
.flatten()
|
||||
.is_some()
|
||||
}
|
||||
|
||||
/// Number of faces held.
|
||||
pub fn len(&self) -> u64 {
|
||||
self.index
|
||||
.query_row("SELECT COUNT(*) FROM faces_meta", [], |r| {
|
||||
r.get::<_, i64>(0)
|
||||
})
|
||||
.map(|n| n as u64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.len() == 0
|
||||
}
|
||||
|
||||
/// Write a batch of faces for one image into the open shard.
|
||||
///
|
||||
/// Per image rather than per face, because the run marker is per image and
|
||||
/// the two must land together: a shard holding faces but no marker would
|
||||
/// make the receiving device re-detect the image it just adopted.
|
||||
pub fn put_image(
|
||||
&mut self,
|
||||
file_id: u64,
|
||||
model_id: &str,
|
||||
source_edge: u32,
|
||||
faces: &[SharedFace],
|
||||
) -> Result<u32, CatalogError> {
|
||||
let incoming = faces.len() as u64 * BYTES_PER_FACE + 64;
|
||||
let shard = self.active_shard(incoming)?;
|
||||
let conn = self.open_shard(shard, true)?;
|
||||
|
||||
let tx = conn.unchecked_transaction()?;
|
||||
// Replace rather than append: re-indexing an image must not double its
|
||||
// faces in the shard, exactly as it must not in the catalog.
|
||||
tx.execute(
|
||||
"DELETE FROM faces WHERE file_id = ?1 AND model_id = ?2",
|
||||
rusqlite::params![file_id as i64, model_id],
|
||||
)?;
|
||||
for f in faces {
|
||||
tx.execute(
|
||||
"INSERT INTO faces
|
||||
(file_id, model_id, x, y, w, h, landmarks, confidence,
|
||||
embedding, crop_px)
|
||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
||||
rusqlite::params![
|
||||
f.file_id as i64,
|
||||
f.model_id,
|
||||
f.x as f64,
|
||||
f.y as f64,
|
||||
f.w as f64,
|
||||
f.h as f64,
|
||||
f.landmarks,
|
||||
f.confidence as f64,
|
||||
f.embedding,
|
||||
f.crop_px as f64,
|
||||
],
|
||||
)?;
|
||||
}
|
||||
tx.execute(
|
||||
"INSERT INTO indexed (file_id, model_id, faces_found, source_edge)
|
||||
VALUES (?1, ?2, ?3, ?4)
|
||||
ON CONFLICT(file_id, model_id) DO UPDATE SET
|
||||
faces_found = excluded.faces_found,
|
||||
source_edge = excluded.source_edge",
|
||||
rusqlite::params![
|
||||
file_id as i64,
|
||||
model_id,
|
||||
faces.len() as i64,
|
||||
source_edge as i64
|
||||
],
|
||||
)?;
|
||||
tx.commit()?;
|
||||
|
||||
self.index.execute(
|
||||
"INSERT INTO entries (file_id, model_id, shard, bytes)
|
||||
VALUES (?1, ?2, ?3, ?4)
|
||||
ON CONFLICT(file_id, model_id) DO UPDATE SET
|
||||
shard = excluded.shard, bytes = excluded.bytes",
|
||||
rusqlite::params![file_id as i64, model_id, shard as i64, incoming as i64],
|
||||
)?;
|
||||
self.index.execute(
|
||||
"INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2)
|
||||
ON CONFLICT(file_id, model_id) DO NOTHING",
|
||||
rusqlite::params![file_id as i64, model_id],
|
||||
)?;
|
||||
self.index.execute(
|
||||
"UPDATE shards SET bytes = bytes + ?2 WHERE id = ?1",
|
||||
rusqlite::params![shard as i64, incoming as i64],
|
||||
)?;
|
||||
Ok(shard)
|
||||
}
|
||||
|
||||
/// The shard currently open for writing, sealing and starting a new one
|
||||
/// where this batch would overflow the cap.
|
||||
fn active_shard(&self, incoming: u64) -> Result<u32, CatalogError> {
|
||||
let open: Option<(i64, i64)> = self
|
||||
.index
|
||||
.query_row(
|
||||
"SELECT id, bytes FROM shards WHERE sealed = 0 ORDER BY id DESC LIMIT 1",
|
||||
[],
|
||||
|r| Ok((r.get(0)?, r.get(1)?)),
|
||||
)
|
||||
.optional()?;
|
||||
|
||||
match open {
|
||||
Some((id, bytes)) if (bytes as u64) + incoming <= SHARD_MAX_BYTES => Ok(id as u32),
|
||||
Some((id, _)) => {
|
||||
// Sealed, and never written to again. That immutability is what
|
||||
// lets a peer download it once and stop asking.
|
||||
self.index
|
||||
.execute("UPDATE shards SET sealed = 1 WHERE id = ?1", [id])?;
|
||||
let next = (id + 1) as u32;
|
||||
self.index.execute(
|
||||
"INSERT INTO shards (id, bytes, sealed) VALUES (?1, 0, 0)",
|
||||
[next as i64],
|
||||
)?;
|
||||
Ok(next)
|
||||
}
|
||||
None => {
|
||||
self.index.execute(
|
||||
"INSERT INTO shards (id, bytes, sealed) VALUES (0, 0, 0)",
|
||||
[],
|
||||
)?;
|
||||
Ok(0)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn open_shard(&self, shard: u32, create: bool) -> Result<Connection, CatalogError> {
|
||||
let path = self.shard_path(shard);
|
||||
if !create && !path.is_file() {
|
||||
return Err(CatalogError::Io(format!("face shard {shard} is missing")));
|
||||
}
|
||||
let conn = Connection::open(&path)?;
|
||||
conn.execute_batch(SHARD_SCHEMA)?;
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
pub fn shard_path(&self, shard: u32) -> PathBuf {
|
||||
self.dir.join(format!("shard-{}-{shard:04}.sqlite", self.client))
|
||||
}
|
||||
|
||||
/// Every shard, with whether it is sealed.
|
||||
pub fn shards(&self) -> Result<Vec<ShardInfo>, CatalogError> {
|
||||
let mut q = self
|
||||
.index
|
||||
.prepare("SELECT id, bytes, sealed FROM shards ORDER BY id")?;
|
||||
let rows = q.query_map([], |r| {
|
||||
Ok(ShardInfo {
|
||||
id: r.get::<_, i64>(0)? as u32,
|
||||
bytes: r.get::<_, i64>(1)? as u64,
|
||||
sealed: r.get(2)?,
|
||||
})
|
||||
})?;
|
||||
rows.collect::<Result<_, _>>().map_err(Into::into)
|
||||
}
|
||||
|
||||
/// Whether a peer's shard has already been taken in.
|
||||
///
|
||||
/// Nothing in an adopted face records where it came from, so without this
|
||||
/// ledger a client re-downloads every peer's shard on every sync.
|
||||
pub fn has_adopted(&self, name: &str, size: u64) -> bool {
|
||||
self.index
|
||||
.query_row(
|
||||
"SELECT 1 FROM adopted WHERE name = ?1 AND size = ?2",
|
||||
rusqlite::params![name, size as i64],
|
||||
|_| Ok(()),
|
||||
)
|
||||
.optional()
|
||||
.ok()
|
||||
.flatten()
|
||||
.is_some()
|
||||
}
|
||||
|
||||
pub fn mark_adopted(&self, name: &str, size: u64) -> Result<(), CatalogError> {
|
||||
self.index.execute(
|
||||
"INSERT INTO adopted (name, size) VALUES (?1, ?2)
|
||||
ON CONFLICT(name) DO UPDATE SET size = excluded.size",
|
||||
rusqlite::params![name, size as i64],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Take in a peer's shard, skipping images this store already holds.
|
||||
///
|
||||
/// Returns how many images were adopted.
|
||||
pub fn merge_shard(&mut self, downloaded: &Path) -> Result<usize, CatalogError> {
|
||||
let src = Connection::open_with_flags(
|
||||
downloaded,
|
||||
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
|
||||
)?;
|
||||
|
||||
let mut q = src.prepare("SELECT file_id, model_id, faces_found, source_edge FROM indexed")?;
|
||||
let images: Vec<(i64, String, i64, i64)> = q
|
||||
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))?
|
||||
.collect::<Result<_, _>>()?;
|
||||
|
||||
let mut adopted = 0;
|
||||
for (file_id, model_id, _found, edge) in images {
|
||||
if self.contains(file_id as u64, &model_id) {
|
||||
continue;
|
||||
}
|
||||
let mut fq = src.prepare(
|
||||
"SELECT file_id, model_id, x, y, w, h, landmarks, confidence,
|
||||
embedding, crop_px
|
||||
FROM faces WHERE file_id = ?1 AND model_id = ?2",
|
||||
)?;
|
||||
let faces: Vec<SharedFace> = fq
|
||||
.query_map(rusqlite::params![file_id, &model_id], read_shared_face)?
|
||||
.collect::<Result<_, _>>()?;
|
||||
self.put_image(file_id as u64, &model_id, edge as u32, &faces)?;
|
||||
adopted += 1;
|
||||
}
|
||||
Ok(adopted)
|
||||
}
|
||||
|
||||
/// Read back everything held for one image.
|
||||
pub fn get_image(
|
||||
&self,
|
||||
file_id: u64,
|
||||
model_id: &str,
|
||||
) -> Result<Option<(Vec<SharedFace>, u32)>, CatalogError> {
|
||||
let Some(shard): Option<i64> = self
|
||||
.index
|
||||
.query_row(
|
||||
"SELECT shard FROM entries WHERE file_id = ?1 AND model_id = ?2",
|
||||
rusqlite::params![file_id as i64, model_id],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.optional()?
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let conn = self.open_shard(shard as u32, false)?;
|
||||
let edge: Option<i64> = conn
|
||||
.query_row(
|
||||
"SELECT source_edge FROM indexed WHERE file_id = ?1 AND model_id = ?2",
|
||||
rusqlite::params![file_id as i64, model_id],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.optional()?;
|
||||
let Some(edge) = edge else { return Ok(None) };
|
||||
|
||||
let mut q = conn.prepare(
|
||||
"SELECT file_id, model_id, x, y, w, h, landmarks, confidence, embedding, crop_px
|
||||
FROM faces WHERE file_id = ?1 AND model_id = ?2",
|
||||
)?;
|
||||
let faces: Vec<SharedFace> = q
|
||||
.query_map(rusqlite::params![file_id as i64, model_id], read_shared_face)?
|
||||
.collect::<Result<_, _>>()?;
|
||||
Ok(Some((faces, edge as u32)))
|
||||
}
|
||||
}
|
||||
|
||||
/// Copy this device's indexed faces into the shard store, ready to upload.
|
||||
///
|
||||
/// Only what the store does not already hold, so a call after a partial sweep
|
||||
/// writes only the new images. Returns how many images were exported.
|
||||
///
|
||||
/// Images with no `oc:fileid` are skipped rather than failing the pass: a
|
||||
/// purely local file has no identity another device could match it by, so there
|
||||
/// is nothing useful to send.
|
||||
pub fn export_to_shards(
|
||||
conn: &Connection,
|
||||
store: &mut FaceShardStore,
|
||||
model_id: &str,
|
||||
) -> Result<usize, CatalogError> {
|
||||
let mut q = conn.prepare(
|
||||
"SELECT r.file_id, fi.image_id, fi.source_edge
|
||||
FROM face_index fi
|
||||
JOIN remote r ON r.image_id = fi.image_id
|
||||
WHERE fi.model_id = ?1
|
||||
ORDER BY fi.image_id",
|
||||
)?;
|
||||
let rows: Vec<(i64, i64, i64)> = q
|
||||
.query_map([model_id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
|
||||
.collect::<Result<_, _>>()?;
|
||||
|
||||
let mut exported = 0;
|
||||
for (file_id, image_id, edge) in rows {
|
||||
if store.contains(file_id as u64, model_id) {
|
||||
continue;
|
||||
}
|
||||
let mut fq = conn.prepare(
|
||||
"SELECT x, y, w, h, landmarks, detector_confidence, embedding, crop_px
|
||||
FROM faces WHERE image_id = ?1 AND model_id = ?2",
|
||||
)?;
|
||||
let faces: Vec<SharedFace> = fq
|
||||
.query_map(rusqlite::params![image_id, model_id], |r| {
|
||||
Ok(SharedFace {
|
||||
file_id: file_id as u64,
|
||||
model_id: model_id.to_string(),
|
||||
x: r.get::<_, f64>(0)? as f32,
|
||||
y: r.get::<_, f64>(1)? as f32,
|
||||
w: r.get::<_, f64>(2)? as f32,
|
||||
h: r.get::<_, f64>(3)? as f32,
|
||||
landmarks: r.get(4)?,
|
||||
confidence: r.get::<_, f64>(5)? as f32,
|
||||
embedding: r.get(6)?,
|
||||
crop_px: r.get::<_, f64>(7)? as f32,
|
||||
})
|
||||
})?
|
||||
.collect::<Result<_, _>>()?;
|
||||
|
||||
store.put_image(file_id as u64, model_id, edge as u32, &faces)?;
|
||||
exported += 1;
|
||||
}
|
||||
Ok(exported)
|
||||
}
|
||||
|
||||
/// Take faces out of the shard store and into this device's catalog.
|
||||
///
|
||||
/// **The point of the whole module.** An image a peer already examined is
|
||||
/// adopted rather than re-detected, which is the difference between a new
|
||||
/// device being useful in a minute and in two hours.
|
||||
///
|
||||
/// Skips any image this device has already indexed itself. Local work is not
|
||||
/// second-guessed by a peer's — the two should agree, since the same model over
|
||||
/// the same proxy is deterministic, but where they do not, the copy this device
|
||||
/// computed is the one it can vouch for.
|
||||
///
|
||||
/// Returns how many images were adopted.
|
||||
pub fn import_from_shards(
|
||||
conn: &Connection,
|
||||
store: &FaceShardStore,
|
||||
model_id: &str,
|
||||
) -> Result<usize, CatalogError> {
|
||||
// Only images this device actually has. A shard covers the whole account,
|
||||
// and a device holding a subset of the library should take only its own
|
||||
// part rather than accumulating faces for photographs it cannot show.
|
||||
let mut q = conn.prepare(
|
||||
"SELECT r.file_id, r.image_id
|
||||
FROM remote r
|
||||
JOIN images i ON i.id = r.image_id
|
||||
WHERE i.trashed_at IS NULL
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM face_index fi
|
||||
WHERE fi.image_id = r.image_id AND fi.model_id = ?1
|
||||
)",
|
||||
)?;
|
||||
let candidates: Vec<(i64, i64)> = q
|
||||
.query_map([model_id], |r| Ok((r.get(0)?, r.get(1)?)))?
|
||||
.collect::<Result<_, _>>()?;
|
||||
|
||||
let mut adopted = 0;
|
||||
for (file_id, image_id) in candidates {
|
||||
let Some((faces, edge)) = store.get_image(file_id as u64, model_id)? else {
|
||||
continue;
|
||||
};
|
||||
let local: Vec<crate::faces::DetectedFace> = faces
|
||||
.into_iter()
|
||||
.map(|f| crate::faces::DetectedFace {
|
||||
x: f.x,
|
||||
y: f.y,
|
||||
w: f.w,
|
||||
h: f.h,
|
||||
landmarks: blob_to_landmarks(&f.landmarks),
|
||||
confidence: f.confidence,
|
||||
embedding: f.embedding,
|
||||
crop_px: f.crop_px,
|
||||
model_id: f.model_id,
|
||||
})
|
||||
.collect();
|
||||
|
||||
crate::faces::record_detections(
|
||||
conn,
|
||||
dr_types::ImageId(image_id as u64),
|
||||
model_id,
|
||||
edge,
|
||||
&local,
|
||||
)?;
|
||||
adopted += 1;
|
||||
}
|
||||
Ok(adopted)
|
||||
}
|
||||
|
||||
fn blob_to_landmarks(b: &[u8]) -> [(f32, f32); 5] {
|
||||
let mut out = [(0.0_f32, 0.0_f32); 5];
|
||||
for (i, o) in out.iter_mut().enumerate() {
|
||||
let at = i * 8;
|
||||
if at + 8 <= b.len() {
|
||||
o.0 = f32::from_le_bytes([b[at], b[at + 1], b[at + 2], b[at + 3]]);
|
||||
o.1 = f32::from_le_bytes([b[at + 4], b[at + 5], b[at + 6], b[at + 7]]);
|
||||
}
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
fn read_shared_face(r: &rusqlite::Row<'_>) -> rusqlite::Result<SharedFace> {
|
||||
Ok(SharedFace {
|
||||
file_id: r.get::<_, i64>(0)? as u64,
|
||||
model_id: r.get(1)?,
|
||||
x: r.get::<_, f64>(2)? as f32,
|
||||
y: r.get::<_, f64>(3)? as f32,
|
||||
w: r.get::<_, f64>(4)? as f32,
|
||||
h: r.get::<_, f64>(5)? as f32,
|
||||
landmarks: r.get(6)?,
|
||||
confidence: r.get::<_, f64>(7)? as f32,
|
||||
embedding: r.get(8)?,
|
||||
crop_px: r.get::<_, f64>(9)? as f32,
|
||||
})
|
||||
}
|
||||
|
||||
/// One shard's sync-relevant state.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub struct ShardInfo {
|
||||
pub id: u32,
|
||||
pub bytes: u64,
|
||||
/// Sealed shards are immutable; once downloaded they never need re-checking.
|
||||
pub sealed: bool,
|
||||
}
|
||||
|
||||
fn mint_client_id(conn: &Connection) -> Result<String, CatalogError> {
|
||||
conn.execute(
|
||||
"INSERT INTO meta(key, value) VALUES('client_id', lower(hex(randomblob(6))))
|
||||
ON CONFLICT(key) DO NOTHING",
|
||||
[],
|
||||
)?;
|
||||
Ok(
|
||||
conn.query_row("SELECT value FROM meta WHERE key = 'client_id'", [], |r| {
|
||||
r.get(0)
|
||||
})?,
|
||||
)
|
||||
}
|
||||
|
||||
const INDEX_SCHEMA: &str = r#"
|
||||
CREATE TABLE IF NOT EXISTS entries (
|
||||
file_id INTEGER NOT NULL,
|
||||
model_id TEXT NOT NULL,
|
||||
shard INTEGER NOT NULL,
|
||||
bytes INTEGER NOT NULL,
|
||||
PRIMARY KEY(file_id, model_id)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS entries_shard ON entries(shard);
|
||||
|
||||
-- One row per (image, model) held, so a count is a count rather than a scan
|
||||
-- across every shard file.
|
||||
CREATE TABLE IF NOT EXISTS faces_meta (
|
||||
file_id INTEGER NOT NULL,
|
||||
model_id TEXT NOT NULL,
|
||||
PRIMARY KEY(file_id, model_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS shards (
|
||||
id INTEGER PRIMARY KEY,
|
||||
bytes INTEGER NOT NULL DEFAULT 0,
|
||||
-- Recorded rather than inferred from file size: a vacuum could shrink a
|
||||
-- sealed shard below the cap and it must still stay closed, or clients
|
||||
-- holding it would see it change.
|
||||
sealed INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS meta (
|
||||
key TEXT PRIMARY KEY,
|
||||
value TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS adopted (
|
||||
name TEXT PRIMARY KEY,
|
||||
size INTEGER NOT NULL
|
||||
);
|
||||
"#;
|
||||
|
||||
const SHARD_SCHEMA: &str = r#"
|
||||
CREATE TABLE IF NOT EXISTS faces (
|
||||
file_id INTEGER NOT NULL,
|
||||
model_id TEXT NOT NULL,
|
||||
x REAL NOT NULL, y REAL NOT NULL, w REAL NOT NULL, h REAL NOT NULL,
|
||||
landmarks BLOB NOT NULL,
|
||||
confidence REAL NOT NULL,
|
||||
embedding BLOB NOT NULL,
|
||||
crop_px REAL NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS faces_file ON faces(file_id, model_id);
|
||||
|
||||
-- The run marker, travelling with the faces it describes. Without it a
|
||||
-- receiving device cannot tell an image with no faces from one never examined,
|
||||
-- and would re-detect every landscape it just adopted.
|
||||
CREATE TABLE IF NOT EXISTS indexed (
|
||||
file_id INTEGER NOT NULL,
|
||||
model_id TEXT NOT NULL,
|
||||
faces_found INTEGER NOT NULL,
|
||||
source_edge INTEGER NOT NULL,
|
||||
PRIMARY KEY(file_id, model_id)
|
||||
);
|
||||
"#;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// A fresh directory per call.
|
||||
///
|
||||
/// Per *call*, not per test: several of these tests stand up two stores to
|
||||
/// play the two devices against each other, and a helper keyed on the
|
||||
/// thread would have the second call delete the first store's files.
|
||||
fn tempdir() -> PathBuf {
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
static N: AtomicU64 = AtomicU64::new(0);
|
||||
let d = std::env::temp_dir().join(format!(
|
||||
"dr-face-shard-{}-{}",
|
||||
std::process::id(),
|
||||
N.fetch_add(1, Ordering::Relaxed)
|
||||
));
|
||||
let _ = std::fs::remove_dir_all(&d);
|
||||
std::fs::create_dir_all(&d).unwrap();
|
||||
d
|
||||
}
|
||||
|
||||
fn face(file_id: u64, seed: u8) -> SharedFace {
|
||||
SharedFace {
|
||||
file_id,
|
||||
model_id: "w600k_mbf".into(),
|
||||
x: 0.1,
|
||||
y: 0.2,
|
||||
w: 0.15,
|
||||
h: 0.2,
|
||||
landmarks: vec![seed; 40],
|
||||
confidence: 0.87,
|
||||
embedding: vec![seed; 1024],
|
||||
crop_px: 180.0,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_stored_image_round_trips() {
|
||||
let dir = tempdir();
|
||||
let mut s = FaceShardStore::open(&dir).unwrap();
|
||||
s.put_image(42, "w600k_mbf", 1024, &[face(42, 1), face(42, 2)])
|
||||
.unwrap();
|
||||
|
||||
assert!(s.contains(42, "w600k_mbf"));
|
||||
let (faces, edge) = s.get_image(42, "w600k_mbf").unwrap().unwrap();
|
||||
assert_eq!(faces.len(), 2);
|
||||
assert_eq!(edge, 1024);
|
||||
assert_eq!(faces[0].embedding.len(), 1024);
|
||||
assert!((faces[0].crop_px - 180.0).abs() < 1e-3);
|
||||
}
|
||||
|
||||
/// The case the run marker exists for, carried across the wire: an image
|
||||
/// with no faces must still read as *indexed* on the receiving device.
|
||||
#[test]
|
||||
fn an_image_with_no_faces_is_still_recorded() {
|
||||
let dir = tempdir();
|
||||
let mut s = FaceShardStore::open(&dir).unwrap();
|
||||
s.put_image(7, "w600k_mbf", 1024, &[]).unwrap();
|
||||
|
||||
assert!(s.contains(7, "w600k_mbf"));
|
||||
let (faces, _) = s.get_image(7, "w600k_mbf").unwrap().unwrap();
|
||||
assert!(faces.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_different_model_is_a_different_entry() {
|
||||
let dir = tempdir();
|
||||
let mut s = FaceShardStore::open(&dir).unwrap();
|
||||
s.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap();
|
||||
assert!(s.contains(1, "w600k_mbf"));
|
||||
assert!(!s.contains(1, "lvface"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn re_storing_an_image_replaces_rather_than_doubling_it() {
|
||||
let dir = tempdir();
|
||||
let mut s = FaceShardStore::open(&dir).unwrap();
|
||||
s.put_image(1, "w600k_mbf", 1024, &[face(1, 1), face(1, 2)])
|
||||
.unwrap();
|
||||
s.put_image(1, "w600k_mbf", 2048, &[face(1, 3)]).unwrap();
|
||||
|
||||
let (faces, edge) = s.get_image(1, "w600k_mbf").unwrap().unwrap();
|
||||
assert_eq!(faces.len(), 1);
|
||||
assert_eq!(edge, 2048);
|
||||
}
|
||||
|
||||
/// The shard cap is the whole reason this module exists rather than the
|
||||
/// data riding in the catalog snapshot.
|
||||
#[test]
|
||||
fn shards_seal_at_the_cap_and_a_new_one_opens() {
|
||||
let dir = tempdir();
|
||||
let mut s = FaceShardStore::open(&dir).unwrap();
|
||||
|
||||
// Enough faces to overflow one shard several times over.
|
||||
let per_image = 64;
|
||||
let images = (SHARD_MAX_BYTES / (per_image as u64 * BYTES_PER_FACE)) + 2;
|
||||
for i in 0..images {
|
||||
let faces: Vec<SharedFace> = (0..per_image).map(|k| face(i, k as u8)).collect();
|
||||
s.put_image(i, "w600k_mbf", 1024, &faces).unwrap();
|
||||
}
|
||||
|
||||
let shards = s.shards().unwrap();
|
||||
assert!(shards.len() >= 2, "expected a seal, got {} shard(s)", shards.len());
|
||||
assert!(shards[0].sealed, "the first shard should be sealed");
|
||||
assert!(
|
||||
shards[0].bytes <= SHARD_MAX_BYTES,
|
||||
"sealed shard is {} bytes, over the {SHARD_MAX_BYTES} cap",
|
||||
shards[0].bytes
|
||||
);
|
||||
assert!(!shards.last().unwrap().sealed, "the last shard stays open");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn merging_a_peers_shard_adopts_what_this_device_lacks() {
|
||||
let peer_dir = tempdir();
|
||||
let mut peer = FaceShardStore::open(&peer_dir).unwrap();
|
||||
peer.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap();
|
||||
peer.put_image(2, "w600k_mbf", 1024, &[]).unwrap();
|
||||
peer.put_image(3, "w600k_mbf", 1024, &[face(3, 3)]).unwrap();
|
||||
let peer_shard = peer.shard_path(0);
|
||||
|
||||
let mine_dir = tempdir();
|
||||
let mut mine = FaceShardStore::open(&mine_dir).unwrap();
|
||||
// Already have image 1, by our own indexing.
|
||||
mine.put_image(1, "w600k_mbf", 1024, &[face(1, 9)]).unwrap();
|
||||
|
||||
assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 2);
|
||||
assert!(mine.contains(2, "w600k_mbf"));
|
||||
assert!(mine.contains(3, "w600k_mbf"));
|
||||
|
||||
// Ours was not overwritten by theirs.
|
||||
let (faces, _) = mine.get_image(1, "w600k_mbf").unwrap().unwrap();
|
||||
assert_eq!(faces[0].landmarks[0], 9);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn merging_the_same_shard_twice_adopts_nothing_the_second_time() {
|
||||
let peer_dir = tempdir();
|
||||
let mut peer = FaceShardStore::open(&peer_dir).unwrap();
|
||||
peer.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap();
|
||||
let peer_shard = peer.shard_path(0);
|
||||
|
||||
let mine_dir = tempdir();
|
||||
let mut mine = FaceShardStore::open(&mine_dir).unwrap();
|
||||
assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 1);
|
||||
assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_adopted_ledger_stops_a_shard_being_fetched_twice() {
|
||||
let dir = tempdir();
|
||||
let s = FaceShardStore::open(&dir).unwrap();
|
||||
assert!(!s.has_adopted("shard-abc-0000.sqlite", 100));
|
||||
s.mark_adopted("shard-abc-0000.sqlite", 100).unwrap();
|
||||
assert!(s.has_adopted("shard-abc-0000.sqlite", 100));
|
||||
// A shard that grew is the *open* one and must be re-read.
|
||||
assert!(!s.has_adopted("shard-abc-0000.sqlite", 200));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn two_stores_get_different_client_ids() {
|
||||
let a = FaceShardStore::open(&tempdir()).unwrap();
|
||||
let b = FaceShardStore::open(&tempdir()).unwrap();
|
||||
assert_ne!(a.client_id(), b.client_id());
|
||||
assert!(a.shard_path(0).to_string_lossy().contains(a.client_id()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_client_id_survives_reopening() {
|
||||
let dir = tempdir();
|
||||
let first = FaceShardStore::open(&dir).unwrap().client_id().to_string();
|
||||
let second = FaceShardStore::open(&dir).unwrap().client_id().to_string();
|
||||
assert_eq!(first, second);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod catalog_round_trip {
|
||||
use super::*;
|
||||
use crate::{faces, schema};
|
||||
|
||||
fn tempdir(tag: &str) -> PathBuf {
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
static N: AtomicU64 = AtomicU64::new(0);
|
||||
let d = std::env::temp_dir().join(format!(
|
||||
"dr-face-rt-{tag}-{}-{}",
|
||||
std::process::id(),
|
||||
N.fetch_add(1, Ordering::Relaxed)
|
||||
));
|
||||
let _ = std::fs::remove_dir_all(&d);
|
||||
std::fs::create_dir_all(&d).unwrap();
|
||||
d
|
||||
}
|
||||
|
||||
/// A device holding photographs keyed by `oc:fileid`, as they are on the
|
||||
/// server. Row ids differ between devices; the file id is what agrees.
|
||||
fn device(images: &[(i64, i64)]) -> Connection {
|
||||
let c = Connection::open_in_memory().unwrap();
|
||||
schema::migrate(&c).unwrap();
|
||||
c.execute(
|
||||
"INSERT INTO roots(id, kind, label) VALUES (1, 'local', 'lib')",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
for (image_id, file_id) in images {
|
||||
c.execute(
|
||||
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (?1, 1, ?2, 0)",
|
||||
rusqlite::params![image_id, format!("IMG_{file_id}.CR3")],
|
||||
)
|
||||
.unwrap();
|
||||
c.execute(
|
||||
"INSERT INTO remote(image_id, file_id) VALUES (?1, ?2)",
|
||||
rusqlite::params![image_id, file_id],
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
c
|
||||
}
|
||||
|
||||
fn detected(seed: u8) -> faces::DetectedFace {
|
||||
faces::DetectedFace {
|
||||
x: 0.1,
|
||||
y: 0.2,
|
||||
w: 0.15,
|
||||
h: 0.2,
|
||||
landmarks: [
|
||||
(0.11, 0.21),
|
||||
(0.2, 0.21),
|
||||
(0.15, 0.26),
|
||||
(0.12, 0.3),
|
||||
(0.19, 0.3),
|
||||
],
|
||||
confidence: 0.9,
|
||||
embedding: vec![seed; 1024],
|
||||
crop_px: 180.0,
|
||||
model_id: "w600k_mbf".into(),
|
||||
}
|
||||
}
|
||||
|
||||
/// The whole point of the module: device B adopts A's work rather than
|
||||
/// spending two hours re-detecting the same library.
|
||||
#[test]
|
||||
fn a_second_device_adopts_the_first_devices_indexing() {
|
||||
let a = device(&[(1, 5001), (2, 5002), (3, 5003)]);
|
||||
let b = device(&[(90, 5001), (91, 5002), (92, 5003)]);
|
||||
|
||||
// A indexes everything: two images with faces, one without.
|
||||
faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)])
|
||||
.unwrap();
|
||||
faces::record_detections(
|
||||
&a,
|
||||
dr_types::ImageId(2),
|
||||
"w600k_mbf",
|
||||
1024,
|
||||
&[detected(2), detected(3)],
|
||||
)
|
||||
.unwrap();
|
||||
faces::record_detections(&a, dr_types::ImageId(3), "w600k_mbf", 1024, &[]).unwrap();
|
||||
|
||||
let mut store_a = FaceShardStore::open(&tempdir("a")).unwrap();
|
||||
assert_eq!(export_to_shards(&a, &mut store_a, "w600k_mbf").unwrap(), 3);
|
||||
|
||||
assert_eq!(faces::coverage(&b, "w600k_mbf").unwrap().indexed, 0);
|
||||
|
||||
let mut store_b = FaceShardStore::open(&tempdir("b")).unwrap();
|
||||
store_b.merge_shard(&store_a.shard_path(0)).unwrap();
|
||||
assert_eq!(import_from_shards(&b, &store_b, "w600k_mbf").unwrap(), 3);
|
||||
|
||||
let cov = faces::coverage(&b, "w600k_mbf").unwrap();
|
||||
assert_eq!(cov.indexed, 3, "B did not adopt A's coverage");
|
||||
assert_eq!(cov.faces, 3);
|
||||
// The faceless image must arrive as *indexed with none*, or B
|
||||
// re-detects it for ever.
|
||||
assert_eq!(cov.without_faces, 1);
|
||||
assert_eq!(cov.outstanding(), 0);
|
||||
|
||||
// The embeddings survived intact, or the two devices would cluster the
|
||||
// same library differently.
|
||||
let got = faces::for_image(&b, dr_types::ImageId(90)).unwrap();
|
||||
assert_eq!(got.len(), 1);
|
||||
assert!((got[0].crop_px - 180.0).abs() < 1e-3);
|
||||
assert!((got[0].landmarks[2].0 - 0.15).abs() < 1e-5);
|
||||
let emb = faces::embeddings(&b, "w600k_mbf").unwrap();
|
||||
assert!(emb.iter().any(|(_, _, blob, _)| blob[0] == 1));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn importing_does_not_overwrite_work_this_device_already_did() {
|
||||
let a = device(&[(1, 5001)]);
|
||||
let b = device(&[(50, 5001)]);
|
||||
|
||||
faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(7)])
|
||||
.unwrap();
|
||||
// B indexed the same photograph itself and got a different embedding.
|
||||
faces::record_detections(&b, dr_types::ImageId(50), "w600k_mbf", 1024, &[detected(9)])
|
||||
.unwrap();
|
||||
|
||||
let mut store = FaceShardStore::open(&tempdir("keep")).unwrap();
|
||||
export_to_shards(&a, &mut store, "w600k_mbf").unwrap();
|
||||
|
||||
assert_eq!(
|
||||
import_from_shards(&b, &store, "w600k_mbf").unwrap(),
|
||||
0,
|
||||
"a peer's copy replaced work this device had already done"
|
||||
);
|
||||
let emb = faces::embeddings(&b, "w600k_mbf").unwrap();
|
||||
assert_eq!(emb[0].2[0], 9, "B's own embedding was overwritten");
|
||||
}
|
||||
|
||||
/// A device holding a subset of the library takes only its own part.
|
||||
#[test]
|
||||
fn a_device_ignores_faces_for_photographs_it_does_not_have() {
|
||||
let a = device(&[(1, 5001), (2, 5002)]);
|
||||
let b = device(&[(80, 5002)]);
|
||||
|
||||
faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)])
|
||||
.unwrap();
|
||||
faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)])
|
||||
.unwrap();
|
||||
|
||||
let mut store = FaceShardStore::open(&tempdir("subset")).unwrap();
|
||||
export_to_shards(&a, &mut store, "w600k_mbf").unwrap();
|
||||
|
||||
assert_eq!(import_from_shards(&b, &store, "w600k_mbf").unwrap(), 1);
|
||||
assert_eq!(faces::coverage(&b, "w600k_mbf").unwrap().indexed, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn exporting_twice_sends_only_what_is_new() {
|
||||
let a = device(&[(1, 5001), (2, 5002)]);
|
||||
faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)])
|
||||
.unwrap();
|
||||
|
||||
let mut store = FaceShardStore::open(&tempdir("incr")).unwrap();
|
||||
assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1);
|
||||
assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 0);
|
||||
|
||||
faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)])
|
||||
.unwrap();
|
||||
assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_local_only_image_with_no_file_id_is_skipped_rather_than_failing() {
|
||||
let a = device(&[(1, 5001)]);
|
||||
a.execute(
|
||||
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (2, 1, 'local.CR3', 0)",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)])
|
||||
.unwrap();
|
||||
faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)])
|
||||
.unwrap();
|
||||
|
||||
let mut store = FaceShardStore::open(&tempdir("localonly")).unwrap();
|
||||
assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1);
|
||||
}
|
||||
}
|
||||
@@ -37,6 +37,7 @@ pub mod cache;
|
||||
pub mod collections;
|
||||
pub mod dedup;
|
||||
pub mod error;
|
||||
pub mod face_shard;
|
||||
pub mod faces;
|
||||
pub mod jobs;
|
||||
pub mod keywords;
|
||||
@@ -53,6 +54,7 @@ pub use cache::{Budget, Cache, DEFAULT_BUDGET_BYTES};
|
||||
pub use collections::{Collection, CollectionKind, TreeRow};
|
||||
pub use dedup::{seen_by_content, seen_by_metadata, set_content_hash};
|
||||
pub use error::CatalogError;
|
||||
pub use face_shard::{FaceShardStore, SharedFace};
|
||||
pub use faces::{Calibration, DetectedFace, Face, FaceId, Person, PersonId};
|
||||
pub use jobs::{Job, JobKind, Priority};
|
||||
pub use keywords::{Coverage, Keyword, KeywordId, SelectionKeyword};
|
||||
|
||||
Reference in New Issue
Block a user