mirror of
https://github.com/toeverything/AFFiNE.git
synced 2026-08-05 03:19:52 +08:00
feat(server): improve doc gc (#15363)
#### PR Dependency Tree * **PR #15363** 👈 This tree was auto-generated by [Charcoal](https://github.com/danerwilliams/charcoal) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Bug Fixes** * Improved validation of workspace roots and document projections, with clearer failures for malformed or incomplete data. * Improved document reference rebuilding and cleanup reliability. * Updated document update merging to better handle invalid binary data. * **Performance** * Avoided unnecessary document reconstruction when no updates are pending. * **Tests** * Updated coverage for malformed workspace roots and document snapshot parsing. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
Generated
+13
-23
@@ -76,11 +76,11 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "affine_doc_loader"
|
name = "affine_doc_loader"
|
||||||
version = "0.1.3"
|
version = "0.1.4"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "1bd208d52725dd0c63583b171ffc3121a60d7686bacedfee5c135f9c0c58d00b"
|
checksum = "1344b7af4cfa7e4c17c676281db8f4244914779a59250fbe71c06b5219b15952"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"nanoid 0.5.0",
|
"nanoid",
|
||||||
"pulldown-cmark 0.13.1",
|
"pulldown-cmark 0.13.1",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
@@ -96,7 +96,7 @@ checksum = "5826d670f6d43faa7809ef6a06d4a8c5ff493b970399c5c3c5659d25e572746b"
|
|||||||
dependencies = [
|
dependencies = [
|
||||||
"affine_doc_loader",
|
"affine_doc_loader",
|
||||||
"chrono",
|
"chrono",
|
||||||
"nanoid 0.5.0",
|
"nanoid",
|
||||||
"pulldown-cmark 0.13.1",
|
"pulldown-cmark 0.13.1",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
@@ -2930,7 +2930,7 @@ dependencies = [
|
|||||||
"libc",
|
"libc",
|
||||||
"log",
|
"log",
|
||||||
"rustversion",
|
"rustversion",
|
||||||
"windows-link 0.1.3",
|
"windows-link 0.2.1",
|
||||||
"windows-result 0.4.1",
|
"windows-result 0.4.1",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -5209,15 +5209,6 @@ dependencies = [
|
|||||||
"typenum",
|
"typenum",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
|
||||||
name = "nanoid"
|
|
||||||
version = "0.4.0"
|
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
|
||||||
checksum = "3ffa00dec017b5b1a8b7cf5e2c008bfda1aa7e0697ac1508b491fdf2622fb4d8"
|
|
||||||
dependencies = [
|
|
||||||
"rand 0.8.6",
|
|
||||||
]
|
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "nanoid"
|
name = "nanoid"
|
||||||
version = "0.5.0"
|
version = "0.5.0"
|
||||||
@@ -6658,12 +6649,12 @@ checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69"
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "rand_distr"
|
name = "rand_distr"
|
||||||
version = "0.5.1"
|
version = "0.6.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "6a8615d50dcf34fa31f7ab52692afec947c4dd0ab803cc87cb3b0b4570ff7463"
|
checksum = "4d431c2703ccf129de4d45253c03f49ebb22b97d6ad79ee3ecfc7e3f4862c1d8"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"num-traits",
|
"num-traits",
|
||||||
"rand 0.9.4",
|
"rand 0.10.1",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
@@ -10167,7 +10158,7 @@ version = "0.1.11"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
|
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"windows-sys 0.48.0",
|
"windows-sys 0.61.2",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
@@ -10740,9 +10731,9 @@ checksum = "ec7a2a501ed189703dba8b08142f057e887dfc4b2cc4db2d343ac6376ba3e0b9"
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "y-octo"
|
name = "y-octo"
|
||||||
version = "0.0.3"
|
version = "0.1.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "bb412cb21b56fe2b48f460b009b230343f3e0f30b18aaaf97e6a67739a6ec0e7"
|
checksum = "b72d63c74097d8d9f79237787d714a18ef0203a5da7a8bdb75e99db2db26cf78"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"ahash",
|
"ahash",
|
||||||
"arbitrary",
|
"arbitrary",
|
||||||
@@ -10750,11 +10741,10 @@ dependencies = [
|
|||||||
"byteorder",
|
"byteorder",
|
||||||
"log",
|
"log",
|
||||||
"loom",
|
"loom",
|
||||||
"nanoid 0.4.0",
|
"nanoid",
|
||||||
"nom 8.0.0",
|
"nom 8.0.0",
|
||||||
"ordered-float",
|
"ordered-float",
|
||||||
"rand 0.9.4",
|
"rand 0.10.1",
|
||||||
"rand_chacha 0.9.0",
|
|
||||||
"rand_distr",
|
"rand_distr",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
|
|||||||
+2
-2
@@ -16,7 +16,7 @@ resolver = "3"
|
|||||||
[workspace.dependencies]
|
[workspace.dependencies]
|
||||||
aes-gcm = "0.10"
|
aes-gcm = "0.10"
|
||||||
affine_common = { path = "./packages/common/native" }
|
affine_common = { path = "./packages/common/native" }
|
||||||
affine_doc_loader = "0.1.3"
|
affine_doc_loader = "0.1.4"
|
||||||
affine_importer = "0.1.2"
|
affine_importer = "0.1.2"
|
||||||
affine_nbstore = { path = "./packages/frontend/native/nbstore" }
|
affine_nbstore = { path = "./packages/frontend/native/nbstore" }
|
||||||
affine_preview = { version = "0.1.0", default-features = false }
|
affine_preview = { version = "0.1.0", default-features = false }
|
||||||
@@ -114,7 +114,7 @@ resolver = "3"
|
|||||||
"Win32_UI_Shell_PropertiesSystem",
|
"Win32_UI_Shell_PropertiesSystem",
|
||||||
] }
|
] }
|
||||||
windows-core = { version = "0.61" }
|
windows-core = { version = "0.61" }
|
||||||
y-octo = "0.0.3"
|
y-octo = "0.1.0"
|
||||||
zip = "8.6"
|
zip = "8.6"
|
||||||
|
|
||||||
[profile.dev.package.sqlx-macros]
|
[profile.dev.package.sqlx-macros]
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
use affine_doc_loader as doc_loader;
|
use affine_doc_loader as doc_loader;
|
||||||
|
use chrono::{DateTime, Utc};
|
||||||
use sqlx::PgPool;
|
use sqlx::PgPool;
|
||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
@@ -8,11 +9,7 @@ use super::{
|
|||||||
|
|
||||||
const PARSER_VERSION: i32 = 1;
|
const PARSER_VERSION: i32 = 1;
|
||||||
|
|
||||||
struct ExtractedRef {
|
type ExtractedRef = doc_loader::BlobRef;
|
||||||
blob_key: String,
|
|
||||||
block_id: String,
|
|
||||||
flavour: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
struct ProjectionState {
|
struct ProjectionState {
|
||||||
@@ -153,23 +150,9 @@ async fn purge_removed_doc_refs(pool: &PgPool, workspace_id: &str, current_doc_i
|
|||||||
Ok(result.rows_affected() as i64)
|
Ok(result.rows_affected() as i64)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn extract_refs(snapshot: &CurrentDoc) -> RuntimeResult<Vec<ExtractedRef>> {
|
fn extract_refs(blob: Vec<u8>) -> RuntimeResult<Vec<ExtractedRef>> {
|
||||||
let parsed = doc_loader::parse_doc_from_binary(snapshot.blob.clone(), snapshot.doc_id.clone())
|
doc_loader::get_blob_refs_from_binary(blob)
|
||||||
.map_err(|err| RuntimeError::invalid_state(format!("Doc blob refs parse failed: {err}")))?;
|
.map_err(|err| RuntimeError::invalid_state(format!("Doc blob refs parse failed: {err}")))
|
||||||
let mut refs = Vec::new();
|
|
||||||
for block in parsed.blocks {
|
|
||||||
let Some(blob_keys) = block.blob else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
for blob_key in blob_keys {
|
|
||||||
refs.push(ExtractedRef {
|
|
||||||
blob_key,
|
|
||||||
block_id: block.block_id.clone(),
|
|
||||||
flavour: block.flavour.clone(),
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(refs)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -191,7 +174,7 @@ mod tests {
|
|||||||
updated_at: Utc::now(),
|
updated_at: Utc::now(),
|
||||||
};
|
};
|
||||||
|
|
||||||
let refs = extract_refs(&snapshot).expect("refs should parse");
|
let refs = extract_refs(snapshot.blob).expect("refs should parse");
|
||||||
|
|
||||||
assert!(
|
assert!(
|
||||||
refs
|
refs
|
||||||
@@ -233,19 +216,25 @@ mod tests {
|
|||||||
updated_at: Utc::now(),
|
updated_at: Utc::now(),
|
||||||
};
|
};
|
||||||
|
|
||||||
assert!(extract_refs(&snapshot).is_err());
|
assert!(extract_refs(snapshot.blob).is_err());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn replace_doc_refs(pool: &PgPool, snapshot: &CurrentDoc, refs: Vec<ExtractedRef>) -> RuntimeResult<(i64, i64)> {
|
async fn replace_doc_refs(
|
||||||
|
pool: &PgPool,
|
||||||
|
workspace_id: &str,
|
||||||
|
doc_id: &str,
|
||||||
|
updated_at: DateTime<Utc>,
|
||||||
|
refs: Vec<ExtractedRef>,
|
||||||
|
) -> RuntimeResult<(i64, i64)> {
|
||||||
let mut tx = pool
|
let mut tx = pool
|
||||||
.begin()
|
.begin()
|
||||||
.await
|
.await
|
||||||
.map_err(|err| RuntimeError::database("Doc blob refs transaction failed", err))?;
|
.map_err(|err| RuntimeError::database("Doc blob refs transaction failed", err))?;
|
||||||
|
|
||||||
let deleted = sqlx::query("DELETE FROM doc_blob_refs WHERE workspace_id = $1 AND doc_id = $2")
|
let deleted = sqlx::query("DELETE FROM doc_blob_refs WHERE workspace_id = $1 AND doc_id = $2")
|
||||||
.bind(&snapshot.workspace_id)
|
.bind(workspace_id)
|
||||||
.bind(&snapshot.doc_id)
|
.bind(doc_id)
|
||||||
.execute(&mut *tx)
|
.execute(&mut *tx)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| RuntimeError::database("Doc blob refs delete failed", err))?
|
.map_err(|err| RuntimeError::database("Doc blob refs delete failed", err))?
|
||||||
@@ -267,12 +256,12 @@ async fn replace_doc_refs(pool: &PgPool, snapshot: &CurrentDoc, refs: Vec<Extrac
|
|||||||
error = NULL
|
error = NULL
|
||||||
"#,
|
"#,
|
||||||
)
|
)
|
||||||
.bind(&snapshot.workspace_id)
|
.bind(workspace_id)
|
||||||
.bind(&snapshot.doc_id)
|
.bind(doc_id)
|
||||||
.bind(reference.blob_key)
|
.bind(reference.blob_key)
|
||||||
.bind(reference.block_id)
|
.bind(reference.block_id)
|
||||||
.bind(reference.flavour)
|
.bind(reference.flavour)
|
||||||
.bind(snapshot.updated_at)
|
.bind(updated_at)
|
||||||
.bind(PARSER_VERSION)
|
.bind(PARSER_VERSION)
|
||||||
.execute(&mut *tx)
|
.execute(&mut *tx)
|
||||||
.await
|
.await
|
||||||
@@ -330,9 +319,15 @@ async fn rebuild_doc_blob_refs_inner(
|
|||||||
return Ok(result);
|
return Ok(result);
|
||||||
};
|
};
|
||||||
|
|
||||||
match extract_refs(&snapshot) {
|
let CurrentDoc {
|
||||||
|
workspace_id,
|
||||||
|
doc_id,
|
||||||
|
blob,
|
||||||
|
updated_at,
|
||||||
|
} = snapshot;
|
||||||
|
match extract_refs(blob) {
|
||||||
Ok(refs) => {
|
Ok(refs) => {
|
||||||
let (written, deleted) = replace_doc_refs(&pool, &snapshot, refs).await?;
|
let (written, deleted) = replace_doc_refs(&pool, &workspace_id, &doc_id, updated_at, refs).await?;
|
||||||
result.parsed_docs = 1;
|
result.parsed_docs = 1;
|
||||||
result.refs_written = written;
|
result.refs_written = written;
|
||||||
result.refs_deleted = deleted;
|
result.refs_deleted = deleted;
|
||||||
|
|||||||
@@ -311,9 +311,12 @@ async fn current_activity(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn root_contains(root: CurrentDoc, doc_id: &str) -> RuntimeResult<bool> {
|
fn root_contains(root: CurrentDoc, doc_id: &str) -> RuntimeResult<bool> {
|
||||||
let ids = affine_doc_loader::get_doc_ids_from_binary(root.blob, true)
|
let projection = affine_doc_loader::project_workspace_root(root.blob, true)
|
||||||
.map_err(|err| RuntimeError::invalid_state(format!("Document cleanup root parse failed: {err}")))?;
|
.map_err(|err| RuntimeError::invalid_state(format!("Document cleanup root parse failed: {err}")))?;
|
||||||
Ok(ids.iter().any(|id| id == doc_id))
|
if !projection.complete {
|
||||||
|
return Err(RuntimeError::invalid_state("Document cleanup root doc is incomplete"));
|
||||||
|
}
|
||||||
|
Ok(projection.doc_ids.iter().any(|id| id == doc_id))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_doc_rows(tx: &mut Transaction<'_, Postgres>, candidate: &Candidate) -> RuntimeResult<i64> {
|
async fn delete_doc_rows(tx: &mut Transaction<'_, Postgres>, candidate: &Candidate) -> RuntimeResult<i64> {
|
||||||
|
|||||||
@@ -101,6 +101,9 @@ fn merge_current_doc(
|
|||||||
if snapshot.is_none() && updates.is_empty() {
|
if snapshot.is_none() && updates.is_empty() {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
}
|
}
|
||||||
|
if updates.is_empty() {
|
||||||
|
return Ok(snapshot);
|
||||||
|
}
|
||||||
let mut doc = Doc::default();
|
let mut doc = Doc::default();
|
||||||
let mut updated_at = snapshot
|
let mut updated_at = snapshot
|
||||||
.as_ref()
|
.as_ref()
|
||||||
@@ -136,8 +139,12 @@ async fn load_workspace_live_doc_ids(pool: &PgPool, workspace_id: &str) -> Runti
|
|||||||
|
|
||||||
fn workspace_live_doc_ids(root: Option<CurrentDoc>) -> RuntimeResult<Vec<String>> {
|
fn workspace_live_doc_ids(root: Option<CurrentDoc>) -> RuntimeResult<Vec<String>> {
|
||||||
let root = root.ok_or_else(|| RuntimeError::invalid_state("Workspace root doc is missing"))?;
|
let root = root.ok_or_else(|| RuntimeError::invalid_state("Workspace root doc is missing"))?;
|
||||||
let mut ids = affine_doc_loader::get_doc_ids_from_binary(root.blob, true)
|
let projection = affine_doc_loader::project_workspace_root(root.blob, true)
|
||||||
.map_err(|err| RuntimeError::invalid_state(format!("Workspace root doc parse failed: {err}")))?;
|
.map_err(|err| RuntimeError::invalid_state(format!("Workspace root doc parse failed: {err}")))?;
|
||||||
|
if !projection.complete {
|
||||||
|
return Err(RuntimeError::invalid_state("Workspace root doc is incomplete"));
|
||||||
|
}
|
||||||
|
let mut ids = projection.doc_ids;
|
||||||
ids.sort();
|
ids.sort();
|
||||||
ids.dedup();
|
ids.dedup();
|
||||||
Ok(ids)
|
Ok(ids)
|
||||||
@@ -1556,6 +1563,18 @@ mod tests {
|
|||||||
}))
|
}))
|
||||||
.is_err()
|
.is_err()
|
||||||
);
|
);
|
||||||
|
assert!(
|
||||||
|
workspace_live_doc_ids(Some(CurrentDoc {
|
||||||
|
workspace_id: "workspace".to_string(),
|
||||||
|
doc_id: "workspace".to_string(),
|
||||||
|
blob: vec![
|
||||||
|
1, 1, 1, 1, 40, 0, 1, 0, 11, 115, 117, 98, 95, 109, 97, 112, 95, 107, 101, 121, 1, 119, 13, 115, 117, 98, 95,
|
||||||
|
109, 97, 112, 95, 118, 97, 108, 117, 101, 0,
|
||||||
|
],
|
||||||
|
updated_at: Utc::now(),
|
||||||
|
}))
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ use memory_indexer::{SearchHit, SnapshotData};
|
|||||||
use napi_derive::napi;
|
use napi_derive::napi;
|
||||||
use serde::Serialize;
|
use serde::Serialize;
|
||||||
use sqlx::Row;
|
use sqlx::Row;
|
||||||
use y_octo::DocOptions;
|
use y_octo::merge_updates_v1;
|
||||||
|
|
||||||
// Increment this whenever there is a breaking change in the index format or how
|
// Increment this whenever there is a breaking change in the index format or how
|
||||||
// updates are applied
|
// updates are applied
|
||||||
@@ -120,7 +120,7 @@ impl SqliteDocStorage {
|
|||||||
}
|
}
|
||||||
segments.extend(updates.into_iter().map(|update| update.bin.to_vec()));
|
segments.extend(updates.into_iter().map(|update| update.bin.to_vec()));
|
||||||
|
|
||||||
merge_updates(segments, doc_id).map(Some)
|
merge_updates(segments).map(Some)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn init_index(&self) -> Result<()> {
|
pub async fn init_index(&self) -> Result<()> {
|
||||||
@@ -249,7 +249,7 @@ impl SqliteDocStorage {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn merge_updates(mut segments: Vec<Vec<u8>>, guid: &str) -> Result<Vec<u8>> {
|
fn merge_updates(mut segments: Vec<Vec<u8>>) -> Result<Vec<u8>> {
|
||||||
if segments.is_empty() {
|
if segments.is_empty() {
|
||||||
return Err(ParseError::DocNotFound.into());
|
return Err(ParseError::DocNotFound.into());
|
||||||
}
|
}
|
||||||
@@ -258,15 +258,9 @@ fn merge_updates(mut segments: Vec<Vec<u8>>, guid: &str) -> Result<Vec<u8>> {
|
|||||||
return segments.pop().ok_or(ParseError::DocNotFound.into());
|
return segments.pop().ok_or(ParseError::DocNotFound.into());
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut doc = DocOptions::new().with_guid(guid.to_string()).build();
|
let update = merge_updates_v1(segments).map_err(|_| ParseError::InvalidBinary)?;
|
||||||
for update in segments.iter() {
|
let buffer = update
|
||||||
doc
|
.encode_v1()
|
||||||
.apply_update_from_binary_v1(update)
|
|
||||||
.map_err(|_| ParseError::InvalidBinary)?;
|
|
||||||
}
|
|
||||||
|
|
||||||
let buffer = doc
|
|
||||||
.encode_update_v1()
|
|
||||||
.map_err(|err| ParseError::ParserError(err.to_string()))?;
|
.map_err(|err| ParseError::ParserError(err.to_string()))?;
|
||||||
|
|
||||||
Ok(buffer)
|
Ok(buffer)
|
||||||
@@ -317,6 +311,13 @@ mod tests {
|
|||||||
.execute(&storage.pool)
|
.execute(&storage.pool)
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
sqlx::query(r#"INSERT INTO updates (doc_id, data, created_at) VALUES (?, ?, ?)"#)
|
||||||
|
.bind("demo-doc")
|
||||||
|
.bind(&[0, 0][..])
|
||||||
|
.bind(Utc::now().naive_utc())
|
||||||
|
.execute(&storage.pool)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
let result = storage.crawl_doc_data("demo-doc").await.unwrap();
|
let result = storage.crawl_doc_data("demo-doc").await.unwrap();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user