feat: continue tree rust family cutover
- add rust renderer/state-family scaffolds and inline compat host thinning for page tree, file tree, and picker - route tree/filetree preflight, file projection, resource artifact, and stream delta contracts through rust plans - preserve canonical move-order validation, file-tree search projection, and related frontend/runtime regression coverage
This commit is contained in:
@@ -1,11 +1,13 @@
|
||||
use crate::app::AppConfig;
|
||||
use crate::context::RequestContext;
|
||||
use crate::error::WebError;
|
||||
use crate::transport::convex::execute_convex_command_plan;
|
||||
use crate::transport::convex::{
|
||||
ConvexCommandExecution, execute_convex_command_plan, execute_convex_command_plan_with_artifacts,
|
||||
};
|
||||
use bridge_runtime::{
|
||||
execute_runtime_input, RuntimeActorWire, RuntimeBridgeContextWire, RuntimeCommandEnvelopeWire,
|
||||
RuntimeActorWire, RuntimeBridgeContextWire, RuntimeCommandEnvelopeWire,
|
||||
RuntimeCommandExecutionPlan, RuntimeExecutionPlan, RuntimeInput, RuntimeSourceWire,
|
||||
RuntimeTargetWire,
|
||||
RuntimeTargetWire, execute_runtime_input,
|
||||
};
|
||||
use serde_json::Value;
|
||||
|
||||
@@ -66,6 +68,27 @@ pub async fn execute_runtime_command_via_convex(
|
||||
execute_convex_command_plan(config, context, &plan).await
|
||||
}
|
||||
|
||||
pub async fn execute_runtime_command_via_convex_with_artifacts(
|
||||
config: &AppConfig,
|
||||
context: &RequestContext,
|
||||
effective_workspace_id: Option<&str>,
|
||||
command: RuntimeCommandEnvelopeWire,
|
||||
) -> Result<ConvexCommandExecution, WebError> {
|
||||
let runtime_context = runtime_context(context, effective_workspace_id);
|
||||
let runtime_input = RuntimeInput::Command {
|
||||
context: runtime_context.clone(),
|
||||
command: command.clone(),
|
||||
};
|
||||
let RuntimeExecutionPlan::Command(plan) = execute_runtime_input(runtime_input)
|
||||
.map_err(|error| WebError::bad_request(error.message).with_context(context))?
|
||||
else {
|
||||
return Err(WebError::internal("runtime command 未返回 command plan").with_context(context));
|
||||
};
|
||||
|
||||
execute_convex_command_plan_with_artifacts(config, context, &runtime_context, &command, &plan)
|
||||
.await
|
||||
}
|
||||
|
||||
pub fn build_tree_target(
|
||||
workspace_id: &str,
|
||||
page_id: Option<&str>,
|
||||
|
||||
@@ -64,6 +64,8 @@ pub async fn next_sidebar(
|
||||
workspace_id: &effective_workspace_id,
|
||||
root_node_id: None,
|
||||
depth: None,
|
||||
query: None,
|
||||
max_results: None,
|
||||
projection: KernelProjectionKind::SidebarTree,
|
||||
},
|
||||
)
|
||||
|
||||
@@ -20,6 +20,8 @@ pub struct KernelProjectionQuery {
|
||||
pub workspace_id: Option<String>,
|
||||
pub root_node_id: Option<String>,
|
||||
pub depth: Option<u32>,
|
||||
pub query: Option<String>,
|
||||
pub max_results: Option<usize>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
@@ -73,6 +75,8 @@ async fn project_projection(
|
||||
workspace_id: &effective_workspace_id,
|
||||
root_node_id: query.root_node_id.as_deref(),
|
||||
depth: query.depth,
|
||||
query: query.query.as_deref(),
|
||||
max_results: query.max_results,
|
||||
projection,
|
||||
},
|
||||
)
|
||||
@@ -388,4 +392,64 @@ mod tests {
|
||||
.iter()
|
||||
.any(|value| value == "expand"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn file_tree_projection_query_returns_matches_and_ancestors() {
|
||||
let response = app()
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri("/api/tree/projections/file?workspaceId=ws_demo&rootNodeId=page_root&query=%E9%A2%84%E7%AE%97")
|
||||
.body(Body::empty())
|
||||
.expect("request"),
|
||||
)
|
||||
.await
|
||||
.expect("response");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body");
|
||||
let payload: Value = serde_json::from_slice(&body).expect("json");
|
||||
let items = payload["result"]["items"].as_array().expect("items");
|
||||
let row_ids = items
|
||||
.iter()
|
||||
.filter_map(|item| item["rowId"].as_str())
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
assert_eq!(row_ids, vec!["doc:page_root", "asset:table_1"]);
|
||||
assert_eq!(items[0]["expandedByDefault"], true);
|
||||
assert_eq!(items[1]["resourceMeta"]["resourceKind"], "table");
|
||||
let edges = payload["result"]["edges"].as_array().expect("edges");
|
||||
assert_eq!(edges.len(), 1);
|
||||
assert_eq!(edges[0]["fromNodeId"], "page_root");
|
||||
assert_eq!(edges[0]["toNodeId"], "asset:table_1");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn file_tree_projection_query_honors_max_results_and_keeps_ancestors() {
|
||||
let response = app()
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri("/api/tree/projections/file?workspaceId=ws_demo&rootNodeId=page_root&query=png&maxResults=1")
|
||||
.body(Body::empty())
|
||||
.expect("request"),
|
||||
)
|
||||
.await
|
||||
.expect("response");
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body");
|
||||
let payload: Value = serde_json::from_slice(&body).expect("json");
|
||||
let items = payload["result"]["items"].as_array().expect("items");
|
||||
let row_ids = items
|
||||
.iter()
|
||||
.filter_map(|item| item["rowId"].as_str())
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
assert_eq!(row_ids, vec!["doc:page_root", "asset:asset_file_1"]);
|
||||
assert_eq!(items[0]["expandedByDefault"], true);
|
||||
assert_eq!(items[1]["resourceMeta"]["assetKind"], "image");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,8 @@ pub struct ProjectionSnapshotSpec<'a> {
|
||||
pub workspace_id: &'a str,
|
||||
pub root_node_id: Option<&'a str>,
|
||||
pub depth: Option<u32>,
|
||||
pub query: Option<&'a str>,
|
||||
pub max_results: Option<usize>,
|
||||
pub projection: KernelProjectionKind,
|
||||
}
|
||||
|
||||
@@ -53,6 +55,8 @@ pub fn projection_query(spec: &ProjectionSnapshotSpec<'_>) -> RuntimeQueryEnvelo
|
||||
"workspaceId": spec.workspace_id,
|
||||
"rootNodeId": spec.root_node_id,
|
||||
"depth": spec.depth,
|
||||
"query": spec.query,
|
||||
"maxResults": spec.max_results,
|
||||
"includeEdges": true,
|
||||
"includeContent": false,
|
||||
"nodeTypes": [KernelNodeType::Page],
|
||||
|
||||
@@ -78,7 +78,9 @@ pub async fn events(
|
||||
&workspace_id,
|
||||
&overview,
|
||||
change.cursor,
|
||||
change.delta.unwrap_or_else(|| serde_json::json!({ "op": "noop" })),
|
||||
change
|
||||
.delta
|
||||
.unwrap_or_else(|| serde_json::json!({ "op": "noop" })),
|
||||
);
|
||||
return Some((Ok(stream_event("delta", &payload)), Some(state)));
|
||||
}
|
||||
@@ -95,8 +97,7 @@ pub async fn events(
|
||||
return None;
|
||||
};
|
||||
state.query = next_query;
|
||||
state.current_cursor =
|
||||
read_stream_cursor_from_payload(&snapshot_payload);
|
||||
state.current_cursor = read_stream_cursor_from_payload(&snapshot_payload);
|
||||
return Some((
|
||||
Ok(stream_event(
|
||||
"resync",
|
||||
|
||||
@@ -125,7 +125,11 @@ fn is_record(value: &Value) -> bool {
|
||||
fn read_string_field(value: &Value, keys: &[&str]) -> Option<String> {
|
||||
let map = value.as_object()?;
|
||||
for key in keys {
|
||||
let candidate = map.get(*key).and_then(Value::as_str).map(str::trim).unwrap_or("");
|
||||
let candidate = map
|
||||
.get(*key)
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.unwrap_or("");
|
||||
if !candidate.is_empty() {
|
||||
return Some(candidate.to_string());
|
||||
}
|
||||
@@ -149,22 +153,30 @@ fn encode_stream_cursor(id: &str, created_at: &str) -> Option<String> {
|
||||
if id.is_empty() || created_at.is_empty() {
|
||||
return None;
|
||||
}
|
||||
Some(json!({
|
||||
"createdAt": created_at,
|
||||
"id": id,
|
||||
})
|
||||
.to_string())
|
||||
Some(
|
||||
json!({
|
||||
"createdAt": created_at,
|
||||
"id": id,
|
||||
})
|
||||
.to_string(),
|
||||
)
|
||||
}
|
||||
|
||||
fn encode_command_cursor(row: &Value) -> Option<String> {
|
||||
let id = read_string_field(row, &["id", "command_id", "commandId"])?;
|
||||
let created_at = read_string_field(row, &["created_at", "createdAt", "finished_at", "finishedAt"])?;
|
||||
let created_at = read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)?;
|
||||
encode_stream_cursor(&id, &created_at)
|
||||
}
|
||||
|
||||
fn encode_domain_event_cursor(row: &Value) -> Option<String> {
|
||||
let id = read_string_field(row, &["event_id", "eventId", "id"])?;
|
||||
let created_at = read_string_field(row, &["created_at", "createdAt", "finished_at", "finishedAt"])?;
|
||||
let created_at = read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)?;
|
||||
encode_stream_cursor(&format!("domain_event:{id}"), &created_at)
|
||||
}
|
||||
|
||||
@@ -176,10 +188,7 @@ fn decode_stream_cursor(raw: &str) -> Option<DecodedStreamCursor> {
|
||||
})
|
||||
}
|
||||
|
||||
pub fn resolve_stream_cursor(
|
||||
overview: Option<&Value>,
|
||||
fallback: Option<&str>,
|
||||
) -> Option<String> {
|
||||
pub fn resolve_stream_cursor(overview: Option<&Value>, fallback: Option<&str>) -> Option<String> {
|
||||
let fallback = fallback
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
@@ -218,29 +227,98 @@ pub fn resolve_stream_cursor(
|
||||
}
|
||||
}
|
||||
|
||||
fn collect_new_command_logs(
|
||||
rows: &[Value],
|
||||
previous_cursor: Option<&str>,
|
||||
) -> (Vec<Value>, bool) {
|
||||
fn collect_new_command_logs(rows: &[Value], previous_cursor: Option<&str>) -> (Vec<Value>, bool) {
|
||||
let Some(previous_cursor) = previous_cursor.and_then(decode_stream_cursor) else {
|
||||
return (rows.to_vec(), false);
|
||||
};
|
||||
|
||||
let previous_index = rows.iter().position(|row| {
|
||||
let id = read_string_field(row, &["id", "command_id", "commandId"]).unwrap_or_default();
|
||||
let created_at =
|
||||
read_string_field(row, &["created_at", "createdAt", "finished_at", "finishedAt"])
|
||||
.unwrap_or_default();
|
||||
let created_at = read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)
|
||||
.unwrap_or_default();
|
||||
id == previous_cursor.id && created_at == previous_cursor.created_at
|
||||
});
|
||||
|
||||
if let Some(index) = previous_index {
|
||||
(rows.iter().take(index).cloned().collect(), false)
|
||||
} else {
|
||||
let newer_rows = rows
|
||||
.iter()
|
||||
.filter(|row| {
|
||||
read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)
|
||||
.map(|created_at| created_at > previous_cursor.created_at)
|
||||
.unwrap_or(false)
|
||||
})
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
if newer_rows.len() < rows.len() {
|
||||
return (newer_rows, false);
|
||||
}
|
||||
(rows.to_vec(), !rows.is_empty())
|
||||
}
|
||||
}
|
||||
|
||||
fn collect_new_domain_events(rows: &[Value], previous_cursor: Option<&str>) -> (Vec<Value>, bool) {
|
||||
let Some(previous_cursor) = previous_cursor.and_then(decode_stream_cursor) else {
|
||||
return (rows.to_vec(), false);
|
||||
};
|
||||
let previous_id = previous_cursor
|
||||
.id
|
||||
.strip_prefix("domain_event:")
|
||||
.unwrap_or(previous_cursor.id.as_str());
|
||||
|
||||
let previous_index = rows.iter().position(|row| {
|
||||
let id = read_string_field(row, &["event_id", "eventId", "id"]).unwrap_or_default();
|
||||
let created_at = read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)
|
||||
.unwrap_or_default();
|
||||
id == previous_id && created_at == previous_cursor.created_at
|
||||
});
|
||||
|
||||
if let Some(index) = previous_index {
|
||||
(rows.iter().take(index).cloned().collect(), false)
|
||||
} else {
|
||||
let newer_rows = rows
|
||||
.iter()
|
||||
.filter(|row| {
|
||||
read_string_field(
|
||||
row,
|
||||
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
||||
)
|
||||
.map(|created_at| created_at > previous_cursor.created_at)
|
||||
.unwrap_or(false)
|
||||
})
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
if newer_rows.len() < rows.len() {
|
||||
return (newer_rows, false);
|
||||
}
|
||||
(rows.to_vec(), !rows.is_empty())
|
||||
}
|
||||
}
|
||||
|
||||
fn read_stream_delta_candidate(candidate: &Value) -> Option<Value> {
|
||||
if candidate
|
||||
.as_object()
|
||||
.and_then(|map| map.get("op"))
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.is_some()
|
||||
{
|
||||
return Some(candidate.clone());
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn read_command_payload_delta(row: &Value) -> Option<Value> {
|
||||
let command_name = read_string_field(row, &["command_name", "commandName"]).unwrap_or_default();
|
||||
if TREE_STREAM_NOOP_COMMANDS.contains(&command_name.as_str()) {
|
||||
@@ -255,17 +333,52 @@ fn read_command_payload_delta(row: &Value) -> Option<Value> {
|
||||
let candidate = payload
|
||||
.as_object()
|
||||
.and_then(|map| map.get("streamDelta").or_else(|| map.get("stream_delta")))?;
|
||||
if candidate
|
||||
read_stream_delta_candidate(candidate)
|
||||
}
|
||||
|
||||
fn read_domain_event_payload_delta(row: &Value) -> Option<Value> {
|
||||
let payload = row.as_object()?.get("payload")?;
|
||||
if !is_record(payload) {
|
||||
return None;
|
||||
}
|
||||
|
||||
let candidate = payload
|
||||
.as_object()
|
||||
.and_then(|map| map.get("op"))
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.is_some()
|
||||
{
|
||||
return Some(candidate.clone());
|
||||
.and_then(|map| map.get("streamDelta").or_else(|| map.get("stream_delta")))?;
|
||||
read_stream_delta_candidate(candidate)
|
||||
}
|
||||
|
||||
fn read_command_row_command_id(row: &Value) -> Option<String> {
|
||||
read_string_field(row, &["command_id", "commandId", "id"])
|
||||
}
|
||||
|
||||
fn read_domain_event_command_id(row: &Value) -> Option<String> {
|
||||
read_string_field(row, &["command_id", "commandId"])
|
||||
}
|
||||
|
||||
fn resolve_matching_command_domain_event_delta(
|
||||
command_rows: &[Value],
|
||||
command_drifted: bool,
|
||||
domain_event_rows: &[Value],
|
||||
_domain_event_drifted: bool,
|
||||
) -> Option<Value> {
|
||||
if command_drifted || command_rows.len() != 1 || domain_event_rows.len() != 1 {
|
||||
return None;
|
||||
}
|
||||
|
||||
let command_id = read_command_row_command_id(&command_rows[0])?;
|
||||
let event_command_id = read_domain_event_command_id(&domain_event_rows[0])?;
|
||||
if command_id != event_command_id {
|
||||
return None;
|
||||
}
|
||||
|
||||
let command_delta = read_command_payload_delta(&command_rows[0])?;
|
||||
let event_delta = read_domain_event_payload_delta(&domain_event_rows[0])?;
|
||||
if command_delta == event_delta {
|
||||
Some(event_delta)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
pub fn resolve_stream_change(
|
||||
@@ -281,8 +394,36 @@ pub fn resolve_stream_change(
|
||||
return None;
|
||||
}
|
||||
|
||||
let rows = read_array_field(overview, &["command_logs", "commandLogs"]).cloned().unwrap_or_default();
|
||||
let rows = read_array_field(overview, &["command_logs", "commandLogs"])
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let (new_rows, drifted) = collect_new_command_logs(&rows, previous_cursor.as_deref());
|
||||
let event_rows = read_array_field(overview, &["domain_events", "domainEvents"])
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let (new_event_rows, event_drifted) =
|
||||
collect_new_domain_events(&event_rows, previous_cursor.as_deref());
|
||||
if !new_rows.is_empty() && !new_event_rows.is_empty() {
|
||||
if let Some(delta) = resolve_matching_command_domain_event_delta(
|
||||
&new_rows,
|
||||
drifted,
|
||||
&new_event_rows,
|
||||
event_drifted,
|
||||
) {
|
||||
return Some(StreamChange {
|
||||
kind: StreamChangeKind::Delta,
|
||||
cursor: next_cursor,
|
||||
delta: Some(delta),
|
||||
});
|
||||
}
|
||||
|
||||
return Some(StreamChange {
|
||||
kind: StreamChangeKind::Resync,
|
||||
cursor: next_cursor,
|
||||
delta: None,
|
||||
});
|
||||
}
|
||||
|
||||
if !drifted && new_rows.len() == 1 {
|
||||
if let Some(delta) = read_command_payload_delta(&new_rows[0]) {
|
||||
return Some(StreamChange {
|
||||
@@ -293,6 +434,16 @@ pub fn resolve_stream_change(
|
||||
}
|
||||
}
|
||||
|
||||
if !event_drifted && new_rows.is_empty() && new_event_rows.len() == 1 {
|
||||
if let Some(delta) = read_domain_event_payload_delta(&new_event_rows[0]) {
|
||||
return Some(StreamChange {
|
||||
kind: StreamChangeKind::Delta,
|
||||
cursor: next_cursor,
|
||||
delta: Some(delta),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
Some(StreamChange {
|
||||
kind: StreamChangeKind::Resync,
|
||||
cursor: next_cursor,
|
||||
@@ -377,6 +528,8 @@ pub async fn load_stream_snapshot(
|
||||
workspace_id: &effective_workspace_id,
|
||||
root_node_id: None,
|
||||
depth: query.depth,
|
||||
query: None,
|
||||
max_results: None,
|
||||
projection: KernelProjectionKind::SidebarTree,
|
||||
},
|
||||
)
|
||||
@@ -388,8 +541,8 @@ pub async fn load_stream_snapshot(
|
||||
})
|
||||
}
|
||||
StreamSnapshotScope::Subtree => {
|
||||
let root_node_id = normalize_root_node_id(query)
|
||||
.expect("subtree scope 已确保 rootNodeId 存在");
|
||||
let root_node_id =
|
||||
normalize_root_node_id(query).expect("subtree scope 已确保 rootNodeId 存在");
|
||||
let dataset = load_sidebar_dataset(config, context, &effective_workspace_id).await?;
|
||||
let tree = execute_kernel_query(
|
||||
context,
|
||||
@@ -435,8 +588,8 @@ pub async fn load_stream_snapshot(
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
resolve_stream_change, resolve_stream_cursor, resolve_stream_scope,
|
||||
StreamChangeKind, StreamSnapshotQuery, StreamSnapshotScope,
|
||||
resolve_stream_change, resolve_stream_cursor, resolve_stream_scope, StreamChangeKind,
|
||||
StreamSnapshotQuery, StreamSnapshotScope,
|
||||
};
|
||||
use serde_json::json;
|
||||
|
||||
@@ -487,6 +640,7 @@ mod tests {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.node.archive",
|
||||
@@ -498,6 +652,7 @@ mod tests {
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
@@ -507,14 +662,14 @@ mod tests {
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"cmd_1"}"#),
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.cursor,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:02Z","id":"cmd_2"}"#.into())
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:02Z","id":"clog_2"}"#.into())
|
||||
);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
@@ -525,11 +680,195 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_preserves_move_document_delta_fields() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.subtree.move",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "move_document",
|
||||
"documentId": "page_2",
|
||||
"parentId": "page_1",
|
||||
"sortOrder": 3,
|
||||
"updatedAt": "2026-04-25T10:00:02Z"
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": []
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "move_document",
|
||||
"documentId": "page_2",
|
||||
"parentId": "page_1",
|
||||
"sortOrder": 3,
|
||||
"updatedAt": "2026-04-25T10:00:02Z"
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_preserves_upsert_documents_delta_fields() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.subtree.copy",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "upsert_documents",
|
||||
"upsertDocuments": [
|
||||
{
|
||||
"id": "copy_1",
|
||||
"workspace_id": "ws_1",
|
||||
"title": "Copy",
|
||||
"parent_id": null,
|
||||
"sort_order": 2,
|
||||
"is_starred": false,
|
||||
"access_scope": "private",
|
||||
"is_template": false,
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"updated_at": "2026-04-25T10:00:02Z"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": []
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "upsert_documents",
|
||||
"upsertDocuments": [
|
||||
{
|
||||
"id": "copy_1",
|
||||
"workspace_id": "ws_1",
|
||||
"title": "Copy",
|
||||
"parent_id": null,
|
||||
"sort_order": 2,
|
||||
"is_starred": false,
|
||||
"access_scope": "private",
|
||||
"is_template": false,
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"updated_at": "2026-04-25T10:00:02Z"
|
||||
}
|
||||
]
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_preserves_upsert_assets_delta_fields() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.resource.move",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "upsert_assets",
|
||||
"upsertAssets": [
|
||||
{
|
||||
"id": "asset_1",
|
||||
"workspace_id": "ws_1",
|
||||
"document_id": "doc_target",
|
||||
"asset_type": "file",
|
||||
"file_url": "/file.pdf",
|
||||
"thumbnail_url": "/file.pdf",
|
||||
"file_name": "file.pdf",
|
||||
"file_size": 1024,
|
||||
"mime_type": "application/pdf",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"updated_at": "2026-04-25T10:00:02Z"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": []
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "upsert_assets",
|
||||
"upsertAssets": [
|
||||
{
|
||||
"id": "asset_1",
|
||||
"workspace_id": "ws_1",
|
||||
"document_id": "doc_target",
|
||||
"asset_type": "file",
|
||||
"file_url": "/file.pdf",
|
||||
"thumbnail_url": "/file.pdf",
|
||||
"file_name": "file.pdf",
|
||||
"file_size": 1024,
|
||||
"mime_type": "application/pdf",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"updated_at": "2026-04-25T10:00:02Z"
|
||||
}
|
||||
]
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_detects_noop_delta_for_non_tree_mutating_command() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "page.body.save",
|
||||
@@ -538,6 +877,7 @@ mod tests {
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
@@ -547,7 +887,7 @@ mod tests {
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"cmd_1"}"#),
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
@@ -555,11 +895,276 @@ mod tests {
|
||||
assert_eq!(change.delta, Some(json!({ "op": "noop" })));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_detects_delta_from_single_new_domain_event() {
|
||||
let overview = json!({
|
||||
"command_logs": [],
|
||||
"domain_events": [
|
||||
{
|
||||
"event_id": "evt_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"payload": {
|
||||
"command_name": "tree.node.rename",
|
||||
"streamDelta": {
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "新标题"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"event_id": "evt_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
]
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"domain_event:evt_1"}"#),
|
||||
)
|
||||
.expect("应识别到 domain event delta");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.cursor,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:02Z","id":"domain_event:evt_2"}"#.into())
|
||||
);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "新标题"
|
||||
}
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_falls_back_to_resync_for_unknown_domain_event_payload() {
|
||||
let overview = json!({
|
||||
"command_logs": [],
|
||||
"domain_events": [
|
||||
{
|
||||
"event_id": "evt_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"payload": {
|
||||
"schema": "mnote.tree.domain_event",
|
||||
"schemaVersion": 1,
|
||||
"eventType": "tree.node.unknown"
|
||||
}
|
||||
},
|
||||
{
|
||||
"event_id": "evt_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
]
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"domain_event:evt_1"}"#),
|
||||
)
|
||||
.expect("应识别到未知 domain event 推进");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Resync);
|
||||
assert_eq!(
|
||||
change.cursor,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:02Z","id":"domain_event:evt_2"}"#.into())
|
||||
);
|
||||
assert_eq!(change.delta, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_falls_back_to_resync_when_command_and_domain_event_both_advance() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:03Z",
|
||||
"command_name": "tree.node.rename",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "命令标题"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": [
|
||||
{
|
||||
"event_id": "evt_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "remove_document",
|
||||
"documentId": "page_3"
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到混合变化");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Resync);
|
||||
assert_eq!(change.delta, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_dedupes_matching_command_and_domain_event_delta() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.node.rename",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "同一标题"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": [
|
||||
{
|
||||
"event_id": "evt_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "同一标题"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到同命令去重 delta");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "upsert_document",
|
||||
"document": {
|
||||
"id": "page_2",
|
||||
"title": "同一标题"
|
||||
}
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_dedupes_matching_delta_when_previous_cursor_is_command_log() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.subtree.move",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "move_document",
|
||||
"documentId": "page_2",
|
||||
"parentId": "page_1",
|
||||
"sortOrder": 2
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
],
|
||||
"domain_events": [
|
||||
{
|
||||
"event_id": "evt_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "move_document",
|
||||
"documentId": "page_2",
|
||||
"parentId": "page_1",
|
||||
"sortOrder": 2
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"event_id": "evt_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z",
|
||||
"payload": {
|
||||
"streamDelta": {
|
||||
"op": "noop"
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
});
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到同命令去重 delta");
|
||||
|
||||
assert_eq!(change.kind, StreamChangeKind::Delta);
|
||||
assert_eq!(
|
||||
change.delta,
|
||||
Some(json!({
|
||||
"op": "move_document",
|
||||
"documentId": "page_2",
|
||||
"parentId": "page_1",
|
||||
"sortOrder": 2
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_change_falls_back_to_resync_when_delta_is_unstable() {
|
||||
let overview = json!({
|
||||
"command_logs": [
|
||||
{
|
||||
"id": "clog_2",
|
||||
"command_id": "cmd_2",
|
||||
"created_at": "2026-04-25T10:00:02Z",
|
||||
"command_name": "tree.subtree.move",
|
||||
@@ -568,6 +1173,7 @@ mod tests {
|
||||
}
|
||||
},
|
||||
{
|
||||
"id": "clog_1",
|
||||
"command_id": "cmd_1",
|
||||
"created_at": "2026-04-25T10:00:01Z"
|
||||
}
|
||||
@@ -577,7 +1183,7 @@ mod tests {
|
||||
|
||||
let change = resolve_stream_change(
|
||||
&overview,
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"cmd_1"}"#),
|
||||
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
||||
)
|
||||
.expect("应识别到变化");
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user