feat(drive): start implementation of drive

- add storage.drives
    - prepare migration phase
    - add created_by and updated_by on storage.folders
This commit is contained in:
Edouard Vanbelle
2026-06-18 13:29:41 +02:00
parent 77545aee05
commit eab7a609b9
43 changed files with 2434 additions and 154 deletions
+44 -23
View File
@@ -2776,12 +2776,22 @@ mod rechunk_integration_tests {
Arc::new(pool)
}
async fn seed_user(pool: &PgPool) -> Uuid {
sqlx::query("SELECT id FROM auth.users LIMIT 1")
.fetch_one(pool)
.await
.map(|r| r.get::<Uuid, _>("id"))
.expect("auth.users must be seeded (init-test-schema.sh)")
/// Returns `(user_id, drive_id)`. Post-D0 every internal user has a
/// default Personal drive (provisioned by `PersonalDriveLifecycleHook`
/// during init-test-schema.sh's user seeding); the JOIN below picks
/// the user-drive pair atomically so test fixtures can insert into
/// `storage.files` with both `user_id` and `drive_id` populated.
async fn seed_user(pool: &PgPool) -> (Uuid, Uuid) {
sqlx::query(
"SELECT u.id AS user_id, d.id AS drive_id
FROM auth.users u
JOIN storage.drives d ON d.default_for_user = u.id
LIMIT 1",
)
.fetch_one(pool)
.await
.map(|r| (r.get::<Uuid, _>("user_id"), r.get::<Uuid, _>("drive_id")))
.expect("auth.users + storage.drives must be seeded (init-test-schema.sh)")
}
/// Plain local backend in a fresh temp dir.
@@ -2848,7 +2858,7 @@ mod rechunk_integration_tests {
.await
.expect("insert legacy blob row");
let user_id = seed_user(pool).await;
let (user_id, drive_id) = seed_user(pool).await;
let mut file_ids = Vec::new();
for i in 0..n_files {
let name = format!(
@@ -2856,11 +2866,12 @@ mod rechunk_integration_tests {
&Uuid::new_v4().to_string()[..8]
);
let id: Uuid = sqlx::query_scalar(
"INSERT INTO storage.files (name, user_id, blob_hash, size)
VALUES ($1, $2, $3, $4) RETURNING id",
"INSERT INTO storage.files (name, user_id, drive_id, blob_hash, size)
VALUES ($1, $2, $3, $4, $5) RETURNING id",
)
.bind(&name)
.bind(user_id)
.bind(drive_id)
.bind(&hash)
.bind(data.len() as i64)
.fetch_one(pool)
@@ -3121,12 +3132,20 @@ mod delta_upload_integration_tests {
Arc::new(pool)
}
async fn seed_user(pool: &PgPool) -> Uuid {
sqlx::query("SELECT id FROM auth.users LIMIT 1")
.fetch_one(pool)
.await
.map(|r| r.get::<Uuid, _>("id"))
.expect("auth.users must be seeded (init-test-schema.sh)")
/// Returns `(user_id, drive_id)` — same shape as the rechunk tests'
/// `seed_user`. Post-D0 every internal user has a default Personal
/// drive provisioned by `PersonalDriveLifecycleHook`.
async fn seed_user(pool: &PgPool) -> (Uuid, Uuid) {
sqlx::query(
"SELECT u.id AS user_id, d.id AS drive_id
FROM auth.users u
JOIN storage.drives d ON d.default_for_user = u.id
LIMIT 1",
)
.fetch_one(pool)
.await
.map(|r| (r.get::<Uuid, _>("user_id"), r.get::<Uuid, _>("drive_id")))
.expect("auth.users + storage.drives must be seeded (init-test-schema.sh)")
}
async fn local_svc(pool: &Arc<PgPool>, dir: &TempDir) -> DedupService {
@@ -3141,6 +3160,7 @@ mod delta_upload_integration_tests {
svc: &DedupService,
pool: &PgPool,
user_id: Uuid,
drive_id: Uuid,
data: &[u8],
label: &str,
) -> (String, Vec<String>, Uuid) {
@@ -3159,14 +3179,15 @@ mod delta_upload_integration_tests {
.expect("chunks");
let file_id: Uuid = sqlx::query_scalar(
"INSERT INTO storage.files (name, user_id, blob_hash, size)
VALUES ($1, $2, $3, $4) RETURNING id",
"INSERT INTO storage.files (name, user_id, drive_id, blob_hash, size)
VALUES ($1, $2, $3, $4, $5) RETURNING id",
)
.bind(format!(
"rust-test-delta-{label}-{}",
&Uuid::new_v4().to_string()[..8]
))
.bind(user_id)
.bind(drive_id)
.bind(&file_hash)
.bind(data.len() as i64)
.fetch_one(pool)
@@ -3226,13 +3247,13 @@ mod delta_upload_integration_tests {
let pool = test_pool().await;
let dir = TempDir::new().unwrap();
let svc = local_svc(&pool, &dir).await;
let user = seed_user(&pool).await;
let (user, drive_id) = seed_user(&pool).await;
// Owned content (multi-chunk), one foreign chunk (ref 1, no file
// row for this user), one orphan (ref 0), one unknown hash.
let data = content(3 * 1024 * 1024, 21);
let (file_hash, owned_chunks, file_id) =
seed_owned_content(&svc, &pool, user, &data, "claim").await;
seed_owned_content(&svc, &pool, user, drive_id, &data, "claim").await;
assert!(owned_chunks.len() >= 3, "3 MiB must split into ≥3 chunks");
let foreign = blake3::hash(format!("foreign-{}", Uuid::new_v4()).as_bytes())
@@ -3316,12 +3337,12 @@ mod delta_upload_integration_tests {
let pool = test_pool().await;
let dir = TempDir::new().unwrap();
let svc = local_svc(&pool, &dir).await;
let user = seed_user(&pool).await;
let (user, drive_id) = seed_user(&pool).await;
// An owned chunk that the client redundantly re-uploads.
let data = content(100 * 1024, 22);
let (file_hash, owned_chunks, file_id) =
seed_owned_content(&svc, &pool, user, &data, "loose").await;
seed_owned_content(&svc, &pool, user, drive_id, &data, "loose").await;
let owned_chunk_bytes = {
let mut stream = svc.read_blob_stream(&file_hash).await.expect("stream");
let mut out = Vec::new();
@@ -3540,11 +3561,11 @@ mod delta_upload_integration_tests {
let pool = test_pool().await;
let dir = TempDir::new().unwrap();
let svc = local_svc(&pool, &dir).await;
let user = seed_user(&pool).await;
let (user, drive_id) = seed_user(&pool).await;
let data = content(2 * 1024 * 1024 + 137, 24);
let (file_hash, _chunks, file_id) =
seed_owned_content(&svc, &pool, user, &data, "verify").await;
seed_owned_content(&svc, &pool, user, drive_id, &data, "verify").await;
let manifest: (Vec<String>, Vec<i64>) = sqlx::query_as(
"SELECT chunk_hashes, chunk_sizes FROM storage.chunk_manifests WHERE file_hash = $1",
+68 -1
View File
@@ -230,11 +230,32 @@ impl PgAclEngine {
}
}
/// Public wrapper around `subject_match_set` for callers that need
/// the expanded `(subject_types, subject_ids)` pair without invoking
/// the engine's full `check`/`require` pipeline. Used by
/// `GET /api/drives` (and future drive-aware listing surfaces) to
/// ask the `DriveRepository` for every drive the caller can read,
/// reusing the engine's cached group-expansion logic.
pub async fn expand_subject_for_listing(
&self,
subject: Subject,
) -> Result<(Vec<&'static str>, Vec<Uuid>), DomainError> {
let counters = QueryCounters::default();
self.subject_match_set(subject, &counters).await
}
/// Returns the owner UUID for any resource type.
async fn owner_of(&self, resource: Resource) -> Result<Uuid, DomainError> {
match resource {
Resource::Folder(id) => self.folder_repo.get_folder_user_id(&id.to_string()).await,
Resource::File(id) => self.file_repo.get_file_user_id(&id.to_string()).await,
// Drive owner resolution wires up in D0-6 once `DriveRepository`
// lands (D0-5). Drive entity carries `default_for_user` for
// `kind='personal'`; shared drives resolve through role_grants
// (Owner role). Returning NotFound here means a permission
// check that reached owner_of on a Drive falls through to the
// grant-lookup path — safe default during D0-1.
Resource::Drive(_) => Err(DomainError::not_found("Drive", resource.id().to_string())),
}
}
@@ -360,6 +381,43 @@ impl PgAclEngine {
Ok(exists.is_some())
}
/// Direct grant lookup for a drive — no ltree cascade (drives have
/// no ancestors). Mirrors the cascade helpers above but with a
/// straight `resource_type='drive' AND resource_id=$4` filter.
async fn drive_grant_exists(
&self,
subject_types: &[&str],
subject_ids: &[Uuid],
permission: Permission,
drive_id: Uuid,
counters: &QueryCounters,
) -> Result<bool, DomainError> {
counters.sql_queries.fetch_add(1, Ordering::Relaxed);
let roles = Self::roles_implying_strings(permission);
let exists: Option<i32> = sqlx::query_scalar(
r#"
SELECT 1
FROM storage.role_grants g
WHERE g.subject_type = ANY($1)
AND g.subject_id = ANY($2)
AND g.role = ANY($3::storage.grant_role[])
AND g.resource_type = 'drive'
AND g.resource_id = $4
AND (g.expires_at IS NULL OR g.expires_at > NOW())
LIMIT 1
"#,
)
.bind(subject_types)
.bind(subject_ids)
.bind(&roles)
.bind(drive_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("PgAcl", format!("drive grant: {e}")))?;
Ok(exists.is_some())
}
/// Look up a single role grant by id, returning the actors a revoke /
/// notify handler needs to make a decision without a second round-trip.
/// Returns `(subject, resource, granted_by)` or `None` if no such row.
@@ -459,7 +517,12 @@ impl PgAclEngine {
) -> Result<bool, DomainError> {
// Owner short-circuit (only for User subjects — groups/tokens/external
// are never owners of resources).
if let Subject::User(uid) = subject {
// Owner short-circuit applies to Folder/File only — they carry a
// single-owner `user_id` column in their respective tables. Drives
// model ownership through the `Owner` role in `role_grants`, so
// there's no analogous fast path: the grant lookup below resolves
// a drive owner via the same query that resolves any drive role.
if let (Subject::User(uid), Resource::Folder(_) | Resource::File(_)) = (subject, resource) {
counters.sql_queries.fetch_add(1, Ordering::Relaxed);
match self.owner_of(resource).await {
Ok(owner) if owner == uid => return Ok(true),
@@ -500,6 +563,10 @@ impl PgAclEngine {
)
.await
}
Resource::Drive(id) => {
self.drive_grant_exists(&subject_types, &subject_ids, permission, id, counters)
.await
}
}
}
}
@@ -242,14 +242,15 @@ impl ContentIndexWorker {
// Authoritative state re-read: a queued 'upsert' whose row vanished
// or got trashed in the meantime becomes a delete.
let files: Vec<(Uuid, String, String, String, String, i64)> =
let files: Vec<(Uuid, String, String, String, String, String, i64)> =
if upsert_candidates.is_empty() {
Vec::new()
} else {
sqlx::query_as(
"SELECT fi.id, fi.user_id::text, fi.name, fi.blob_hash, fi.mime_type, fi.size
FROM storage.files fi
WHERE fi.id = ANY($1) AND NOT fi.is_trashed",
"SELECT fi.id, fi.user_id::text, fi.drive_id::text, fi.name,
fi.blob_hash, fi.mime_type, fi.size
FROM storage.files fi
WHERE fi.id = ANY($1) AND NOT fi.is_trashed",
)
.bind(&upsert_candidates)
.fetch_all(self.maintenance_pool.as_ref())
@@ -261,10 +262,10 @@ impl ContentIndexWorker {
// Per-blob text: batch-read the extraction cache, extract misses.
let wanted_hashes: Vec<String> = files
.iter()
.filter(|(_, _, name, _, mime, size)| {
.filter(|(_, _, _, name, _, mime, size)| {
text_extractor::supports(name, mime) && *size as u64 <= self.max_extract_file_bytes
})
.map(|f| f.3.clone())
.map(|f| f.4.clone())
.collect();
let mut text_by_hash: HashMap<String, Option<String>> = HashMap::new();
if !wanted_hashes.is_empty() {
@@ -281,7 +282,7 @@ impl ContentIndexWorker {
}
let mut records = Vec::with_capacity(files.len());
for (file_id, user_id, name, blob_hash, mime, size) in files {
for (file_id, user_id, drive_id, name, blob_hash, mime, size) in files {
let supported = text_extractor::supports(&name, &mime);
let content = if !supported {
None
@@ -301,6 +302,7 @@ impl ContentIndexWorker {
records.push(IndexDocRecord {
file_id: file_id.to_string(),
user_id,
drive_id,
name,
content,
preview,
@@ -37,7 +37,15 @@ use crate::common::errors::DomainError;
/// Bump whenever the Tantivy schema OR the text extractor output changes in a
/// way that requires re-indexing. A mismatch with the on-disk marker wipes the
/// index directory and reseeds the dirty queue with every live file.
pub const INDEX_SCHEMA_VERSION: &str = "1";
///
/// Version history:
/// 1 — initial schema (file_id, user_id, name, content, preview)
/// 2 — D0 added `drive_id` field; query filter pivots from user_id
/// to a `drive_id ∈ accessible_drives` set membership clause. On
/// deploy, every operator's index is wiped and reseeded against
/// the post-D0 schema (the worker drains the dirty queue with
/// drive_id-aware records).
pub const INDEX_SCHEMA_VERSION: &str = "2";
/// Recorded in `storage.blob_extracted_text.extractor`; rows from another
/// version are dropped at worker startup (the reseed re-extracts them).
@@ -73,6 +81,11 @@ const PREFIX_MIN_CHARS: usize = 3;
pub struct IndexDocRecord {
pub file_id: String,
pub user_id: String,
/// Owning drive — written verbatim into the `drive_id` STRING field
/// for set-membership filtering at query time. The user_id field is
/// kept during the D0 dual-write window for rollback safety; the
/// query filter no longer reads it.
pub drive_id: String,
pub name: String,
pub content: Option<String>,
pub preview: Option<String>,
@@ -82,6 +95,7 @@ pub struct IndexDocRecord {
struct IndexFields {
file_id: Field,
user_id: Field,
drive_id: Field,
name: Field,
content: Field,
preview: Field,
@@ -105,6 +119,7 @@ impl TantivyContentIndex {
let fields = IndexFields {
file_id: builder.add_text_field("file_id", STRING | STORED),
user_id: builder.add_text_field("user_id", STRING),
drive_id: builder.add_text_field("drive_id", STRING),
name: builder.add_text_field("name", TEXT),
content: builder.add_text_field("content", TEXT),
preview: builder.add_text_field("preview", STORED),
@@ -197,6 +212,7 @@ impl TantivyContentIndex {
let mut document = doc!(
self.fields.file_id => record.file_id,
self.fields.user_id => record.user_id,
self.fields.drive_id => record.drive_id,
self.fields.name => record.name,
);
if let Some(content) = record.content {
@@ -234,15 +250,26 @@ impl TantivyContentIndex {
/// Build the scored query: every token must match (in name OR content,
/// exact OR fuzzy OR — for the last token — prefix), and the whole thing
/// is `Must`-scoped to the user.
fn build_query(fields: IndexFields, user_id: &str, tokens: &[String]) -> Box<dyn Query> {
let mut clauses: Vec<(Occur, Box<dyn Query>)> = vec![(
Occur::Must,
Box::new(TermQuery::new(
Term::from_field_text(fields.user_id, user_id),
IndexRecordOption::Basic,
)),
)];
/// is `Must`-scoped to the caller's accessible drives.
///
/// The drive filter is expressed as a BoolQuery with `Should` arms —
/// at least one drive_id must match — wrapped under an outer `Must`.
/// Equivalent to a TermSetQuery; this form avoids the API churn of
/// rebuilding the same shape across Tantivy versions.
fn build_query(fields: IndexFields, drive_ids: &[String], tokens: &[String]) -> Box<dyn Query> {
// Drive-membership Must clause: union of Term(drive_id = $each).
let drive_alternatives: Vec<(Occur, Box<dyn Query>)> = drive_ids
.iter()
.map(|d| {
let q: Box<dyn Query> = Box::new(TermQuery::new(
Term::from_field_text(fields.drive_id, d),
IndexRecordOption::Basic,
));
(Occur::Should, q)
})
.collect();
let mut clauses: Vec<(Occur, Box<dyn Query>)> =
vec![(Occur::Must, Box::new(BooleanQuery::new(drive_alternatives)))];
let last = tokens.len().saturating_sub(1);
for (i, token) in tokens.iter().enumerate() {
@@ -306,7 +333,7 @@ impl TantivyContentIndex {
searcher: tantivy::Searcher,
analyzer: TextAnalyzer,
fields: IndexFields,
user_id: &str,
drive_ids: &[String],
raw_query: &str,
limit: usize,
) -> Result<Vec<ContentHitDto>, DomainError> {
@@ -315,7 +342,7 @@ impl TantivyContentIndex {
return Ok(Vec::new());
}
let query = Self::build_query(fields, user_id, &tokens);
let query = Self::build_query(fields, drive_ids, &tokens);
let top_docs = searcher
.search(&query, &TopDocs::with_limit(limit.max(1)).order_by_score())
.map_err(|e| DomainError::internal_error("ContentIndex", format!("search: {e}")))?;
@@ -365,18 +392,25 @@ impl TantivyContentIndex {
impl ContentIndexPort for TantivyContentIndex {
async fn search_content(
&self,
user_id: Uuid,
accessible_drive_ids: &[Uuid],
query: &str,
limit: usize,
) -> Result<Vec<ContentHitDto>, DomainError> {
// No accessible drives → no hits, no Tantivy work. Matches the
// anti-enumeration semantics (empty filter set returns empty
// results without any side channel).
if accessible_drive_ids.is_empty() {
return Ok(Vec::new());
}
let searcher = self.reader.searcher();
let analyzer = self.analyzer.clone();
let fields = self.fields;
let user_id = user_id.to_string();
let drive_ids: Vec<String> = accessible_drive_ids.iter().map(|d| d.to_string()).collect();
let query = query.to_owned();
tokio::task::spawn_blocking(move || {
Self::search_blocking(searcher, analyzer, fields, &user_id, &query, limit)
Self::search_blocking(searcher, analyzer, fields, &drive_ids, &query, limit)
})
.await
.map_err(|e| DomainError::internal_error("ContentIndex", format!("join: {e}")))?
@@ -391,6 +425,10 @@ mod tests {
IndexDocRecord {
file_id: file_id.to_owned(),
user_id: user_id.to_owned(),
// Tests stamp a placeholder drive_id derived from user_id so the
// record satisfies the post-D0 schema. Query-side filtering by
// drive_id is exercised in D0-12's integration tests, not here.
drive_id: format!("{user_id}-drive"),
name: name.to_owned(),
content: content.map(str::to_owned),
preview: content.map(str::to_owned),
@@ -401,11 +439,16 @@ mod tests {
// Force a reader reload — OnCommitWithDelay is asynchronous and tests
// must observe the commit immediately.
index.reader.reload().unwrap();
// Test records derive `drive_id = format!("{user_id}-drive")` —
// the same convention used by `record()`. Filtering by that
// single drive id exercises the same path the production
// search uses.
let drive_ids = vec![format!("{user_id}-drive")];
TantivyContentIndex::search_blocking(
index.reader.searcher(),
index.analyzer.clone(),
index.fields,
user_id,
&drive_ids,
query,
32,
)
@@ -113,13 +113,15 @@ impl TreeEtagFlushService {
FROM storage.tree_etag_dirty
ORDER BY id
LIMIT $1)
RETURNING lpath, folder_id
RETURNING lpath, folder_id, drive_id
),
targets AS (
-- Captured chain: covers target folders deleted or
-- moved away since enqueue (the old location's
-- surviving ancestors still get their bump).
SELECT lpath FROM drained
-- surviving ancestors still get their bump). drive_id
-- comes along so the victims walk can enforce
-- cross-drive isolation (D0-13).
SELECT lpath, drive_id FROM drained
UNION
-- Flush-time resolution: a folder MOVED since
-- enqueue had its subtree's lpaths rewritten, so
@@ -128,19 +130,29 @@ impl TreeEtagFlushService {
-- this, a bump queued just before a move would be
-- silently lost and sync clients would never
-- discover the change.
SELECT fo.lpath
SELECT fo.lpath, fo.drive_id
FROM storage.folders fo
JOIN drained d ON fo.id = d.folder_id
),
victims AS (
-- `lpath @> target` = the target folder itself plus
-- every ancestor (GiST-indexed). Folder rows deleted
-- since enqueue simply don't match. Lock in id order
-- so overlapping closures cannot deadlock.
-- every ancestor (GiST-indexed). The `drive_id`
-- predicate prevents a numerically-overlapping
-- lpath in a SIBLING drive from spuriously matching
-- (D0-13). Rows from old queue entries (pre-M4) have
-- NULL `drive_id` — `IS NOT DISTINCT FROM` falls
-- back to pure lpath matching for those, preserving
-- the rollover semantics for any rows enqueued
-- between this migration committing and the
-- service restart.
SELECT f.id
FROM storage.folders f
WHERE EXISTS (SELECT 1 FROM targets t
WHERE f.lpath @> t.lpath)
WHERE EXISTS (
SELECT 1 FROM targets t
WHERE f.lpath @> t.lpath
AND (t.drive_id IS NULL
OR f.drive_id = t.drive_id)
)
ORDER BY f.id
FOR NO KEY UPDATE
),