2026-04-18 05:43:49 +08:00
|
|
|
use crate::app::AppConfig;
|
|
|
|
|
use crate::context::RequestContext;
|
|
|
|
|
use crate::error::WebError;
|
|
|
|
|
use crate::routes::query_support::{
|
|
|
|
|
execute_runtime_query_via_convex, resolve_effective_workspace_id,
|
|
|
|
|
};
|
|
|
|
|
use crate::routes::snapshot_support::{
|
2026-04-29 12:24:44 +08:00
|
|
|
execute_kernel_query, load_projection_snapshot, load_sidebar_dataset, subtree_query,
|
|
|
|
|
ProjectionSnapshotSpec,
|
2026-04-18 05:43:49 +08:00
|
|
|
};
|
|
|
|
|
use bridge_runtime::RuntimeQueryEnvelopeWire;
|
|
|
|
|
use core_protocol::KernelProjectionKind;
|
|
|
|
|
use serde::Deserialize;
|
2026-04-29 12:24:44 +08:00
|
|
|
use serde_json::{json, Value};
|
2026-04-18 05:43:49 +08:00
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
const TREE_STREAM_NOOP_COMMANDS: [&str; 6] = [
|
|
|
|
|
"page.body.save",
|
|
|
|
|
"page.layout.updateOptions",
|
|
|
|
|
"documents.stats.update",
|
|
|
|
|
"blocks.patch",
|
|
|
|
|
"blocks.move",
|
|
|
|
|
"blocks.embed",
|
|
|
|
|
];
|
|
|
|
|
|
2026-04-18 05:43:49 +08:00
|
|
|
#[derive(Debug, Clone, Deserialize, Default)]
|
|
|
|
|
#[serde(rename_all = "camelCase")]
|
|
|
|
|
pub struct StreamSnapshotQuery {
|
|
|
|
|
pub workspace_id: Option<String>,
|
|
|
|
|
pub root_node_id: Option<String>,
|
|
|
|
|
pub depth: Option<u32>,
|
|
|
|
|
pub cursor: Option<String>,
|
|
|
|
|
pub limit: Option<u32>,
|
2026-04-26 04:29:23 +08:00
|
|
|
pub poll_ms: Option<u64>,
|
|
|
|
|
pub max_polls: Option<u32>,
|
2026-04-18 05:43:49 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
|
|
|
pub enum StreamSnapshotScope {
|
|
|
|
|
Workspace,
|
|
|
|
|
Subtree,
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
|
|
|
|
pub enum StreamChangeKind {
|
|
|
|
|
Delta,
|
|
|
|
|
Resync,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Debug, Clone, PartialEq)]
|
|
|
|
|
pub struct StreamChange {
|
|
|
|
|
pub kind: StreamChangeKind,
|
|
|
|
|
pub cursor: Option<String>,
|
|
|
|
|
pub delta: Option<Value>,
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-18 05:43:49 +08:00
|
|
|
impl StreamSnapshotScope {
|
|
|
|
|
pub fn as_str(self) -> &'static str {
|
|
|
|
|
match self {
|
|
|
|
|
Self::Workspace => "workspace",
|
|
|
|
|
Self::Subtree => "subtree",
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-04-26 04:29:23 +08:00
|
|
|
|
|
|
|
|
pub fn projection(self) -> &'static str {
|
|
|
|
|
match self {
|
|
|
|
|
Self::Workspace => "sidebar_tree",
|
|
|
|
|
Self::Subtree => "page_tree",
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
|
|
|
|
struct DecodedStreamCursor {
|
|
|
|
|
created_at: String,
|
|
|
|
|
id: String,
|
2026-04-18 05:43:49 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn resolve_stream_scope(query: &StreamSnapshotQuery) -> StreamSnapshotScope {
|
|
|
|
|
if query
|
|
|
|
|
.root_node_id
|
|
|
|
|
.as_deref()
|
|
|
|
|
.map(str::trim)
|
|
|
|
|
.filter(|value| !value.is_empty())
|
|
|
|
|
.is_some()
|
|
|
|
|
{
|
|
|
|
|
StreamSnapshotScope::Subtree
|
|
|
|
|
} else {
|
|
|
|
|
StreamSnapshotScope::Workspace
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
fn normalize_root_node_id(query: &StreamSnapshotQuery) -> Option<String> {
|
|
|
|
|
query
|
|
|
|
|
.root_node_id
|
|
|
|
|
.as_deref()
|
|
|
|
|
.map(str::trim)
|
|
|
|
|
.filter(|value| !value.is_empty())
|
|
|
|
|
.map(ToOwned::to_owned)
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-18 05:43:49 +08:00
|
|
|
fn workspace_overview_query(
|
|
|
|
|
workspace_id: &str,
|
|
|
|
|
query: &StreamSnapshotQuery,
|
|
|
|
|
) -> RuntimeQueryEnvelopeWire {
|
|
|
|
|
RuntimeQueryEnvelopeWire {
|
|
|
|
|
name: "bridge.workspace.overview".into(),
|
|
|
|
|
payload: json!({
|
|
|
|
|
"workspaceId": workspace_id,
|
|
|
|
|
"cursor": query.cursor,
|
|
|
|
|
"limit": query.limit.unwrap_or(20),
|
|
|
|
|
"commandStatus": Value::Null,
|
|
|
|
|
"eventStatus": Value::Null,
|
|
|
|
|
"targetPageId": Value::Null,
|
|
|
|
|
"targetBlockId": Value::Null,
|
|
|
|
|
"aggregateType": Value::Null,
|
|
|
|
|
"aggregateId": query.root_node_id,
|
|
|
|
|
}),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
fn is_record(value: &Value) -> bool {
|
|
|
|
|
value.is_object()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn read_string_field(value: &Value, keys: &[&str]) -> Option<String> {
|
|
|
|
|
let map = value.as_object()?;
|
|
|
|
|
for key in keys {
|
2026-04-26 19:35:52 +08:00
|
|
|
let candidate = map
|
|
|
|
|
.get(*key)
|
|
|
|
|
.and_then(Value::as_str)
|
|
|
|
|
.map(str::trim)
|
|
|
|
|
.unwrap_or("");
|
2026-04-26 04:29:23 +08:00
|
|
|
if !candidate.is_empty() {
|
|
|
|
|
return Some(candidate.to_string());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
None
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn read_array_field<'a>(value: &'a Value, keys: &[&str]) -> Option<&'a Vec<Value>> {
|
|
|
|
|
let map = value.as_object()?;
|
|
|
|
|
for key in keys {
|
|
|
|
|
if let Some(items) = map.get(*key).and_then(Value::as_array) {
|
|
|
|
|
return Some(items);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
None
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn encode_stream_cursor(id: &str, created_at: &str) -> Option<String> {
|
|
|
|
|
let id = id.trim();
|
|
|
|
|
let created_at = created_at.trim();
|
|
|
|
|
if id.is_empty() || created_at.is_empty() {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
2026-04-26 19:35:52 +08:00
|
|
|
Some(
|
|
|
|
|
json!({
|
|
|
|
|
"createdAt": created_at,
|
|
|
|
|
"id": id,
|
|
|
|
|
})
|
|
|
|
|
.to_string(),
|
|
|
|
|
)
|
2026-04-26 04:29:23 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn encode_command_cursor(row: &Value) -> Option<String> {
|
|
|
|
|
let id = read_string_field(row, &["id", "command_id", "commandId"])?;
|
2026-04-26 19:35:52 +08:00
|
|
|
let created_at = read_string_field(
|
|
|
|
|
row,
|
|
|
|
|
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
|
|
|
|
)?;
|
2026-04-26 04:29:23 +08:00
|
|
|
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"])?;
|
2026-04-26 19:35:52 +08:00
|
|
|
let created_at = read_string_field(
|
|
|
|
|
row,
|
|
|
|
|
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
|
|
|
|
)?;
|
2026-04-26 04:29:23 +08:00
|
|
|
encode_stream_cursor(&format!("domain_event:{id}"), &created_at)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn decode_stream_cursor(raw: &str) -> Option<DecodedStreamCursor> {
|
|
|
|
|
let parsed = serde_json::from_str::<Value>(raw).ok()?;
|
|
|
|
|
Some(DecodedStreamCursor {
|
|
|
|
|
created_at: read_string_field(&parsed, &["createdAt"])?,
|
|
|
|
|
id: read_string_field(&parsed, &["id"])?,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
pub fn resolve_stream_cursor(overview: Option<&Value>, fallback: Option<&str>) -> Option<String> {
|
2026-04-26 04:29:23 +08:00
|
|
|
let fallback = fallback
|
|
|
|
|
.map(str::trim)
|
|
|
|
|
.filter(|value| !value.is_empty())
|
|
|
|
|
.map(ToOwned::to_owned);
|
|
|
|
|
let Some(overview) = overview else {
|
|
|
|
|
return fallback;
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
let command_cursor = read_array_field(overview, &["command_logs", "commandLogs"])
|
|
|
|
|
.and_then(|rows| rows.first())
|
|
|
|
|
.and_then(encode_command_cursor);
|
|
|
|
|
let domain_event_cursor = read_array_field(overview, &["domain_events", "domainEvents"])
|
|
|
|
|
.and_then(|rows| rows.first())
|
|
|
|
|
.and_then(encode_domain_event_cursor);
|
|
|
|
|
|
|
|
|
|
match (command_cursor, domain_event_cursor) {
|
|
|
|
|
(None, None) => fallback,
|
|
|
|
|
(Some(cursor), None) => Some(cursor),
|
|
|
|
|
(None, Some(cursor)) => Some(cursor),
|
|
|
|
|
(Some(command_cursor), Some(domain_event_cursor)) => {
|
|
|
|
|
let decoded_command = decode_stream_cursor(&command_cursor);
|
|
|
|
|
let decoded_domain_event = decode_stream_cursor(&domain_event_cursor);
|
|
|
|
|
match (decoded_command, decoded_domain_event) {
|
|
|
|
|
(Some(command), Some(event)) => {
|
|
|
|
|
if event.created_at > command.created_at {
|
|
|
|
|
Some(domain_event_cursor)
|
|
|
|
|
} else {
|
|
|
|
|
Some(command_cursor)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
(Some(_), None) => Some(command_cursor),
|
|
|
|
|
(None, Some(_)) => Some(domain_event_cursor),
|
|
|
|
|
(None, None) => fallback,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
fn collect_new_command_logs(rows: &[Value], previous_cursor: Option<&str>) -> (Vec<Value>, bool) {
|
2026-04-26 04:29:23 +08:00
|
|
|
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();
|
2026-04-26 19:35:52 +08:00
|
|
|
let created_at = read_string_field(
|
|
|
|
|
row,
|
|
|
|
|
&["created_at", "createdAt", "finished_at", "finishedAt"],
|
|
|
|
|
)
|
|
|
|
|
.unwrap_or_default();
|
2026-04-26 04:29:23 +08:00
|
|
|
id == previous_cursor.id && created_at == previous_cursor.created_at
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
if let Some(index) = previous_index {
|
|
|
|
|
(rows.iter().take(index).cloned().collect(), false)
|
|
|
|
|
} else {
|
2026-04-26 19:35:52 +08:00
|
|
|
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);
|
|
|
|
|
}
|
2026-04-26 04:29:23 +08:00
|
|
|
(rows.to_vec(), !rows.is_empty())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
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()) {
|
|
|
|
|
return Some(json!({ "op": "noop" }));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let payload = row.as_object()?.get("payload")?;
|
|
|
|
|
if !is_record(payload) {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let candidate = payload
|
|
|
|
|
.as_object()
|
|
|
|
|
.and_then(|map| map.get("streamDelta").or_else(|| map.get("stream_delta")))?;
|
2026-04-26 19:35:52 +08:00
|
|
|
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
|
2026-04-26 04:29:23 +08:00
|
|
|
.as_object()
|
2026-04-26 19:35:52 +08:00
|
|
|
.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
|
2026-04-26 04:29:23 +08:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn resolve_stream_change(
|
|
|
|
|
overview: &Value,
|
|
|
|
|
previous_cursor: Option<&str>,
|
|
|
|
|
) -> Option<StreamChange> {
|
|
|
|
|
let next_cursor = resolve_stream_cursor(Some(overview), previous_cursor);
|
|
|
|
|
let previous_cursor = previous_cursor
|
|
|
|
|
.map(str::trim)
|
|
|
|
|
.filter(|value| !value.is_empty())
|
|
|
|
|
.map(ToOwned::to_owned);
|
|
|
|
|
if next_cursor == previous_cursor {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
let rows = read_array_field(overview, &["command_logs", "commandLogs"])
|
|
|
|
|
.cloned()
|
|
|
|
|
.unwrap_or_default();
|
2026-04-26 04:29:23 +08:00
|
|
|
let (new_rows, drifted) = collect_new_command_logs(&rows, previous_cursor.as_deref());
|
2026-04-26 19:35:52 +08:00
|
|
|
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,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
if !drifted && new_rows.len() == 1 {
|
|
|
|
|
if let Some(delta) = read_command_payload_delta(&new_rows[0]) {
|
|
|
|
|
return Some(StreamChange {
|
|
|
|
|
kind: StreamChangeKind::Delta,
|
|
|
|
|
cursor: next_cursor,
|
|
|
|
|
delta: Some(delta),
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
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),
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
Some(StreamChange {
|
|
|
|
|
kind: StreamChangeKind::Resync,
|
|
|
|
|
cursor: next_cursor,
|
|
|
|
|
delta: None,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn read_stream_cursor_from_payload(payload: &Value) -> Option<String> {
|
|
|
|
|
read_string_field(payload, &["cursor"])
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn with_stream_kind(payload: &Value, kind: &str) -> Value {
|
|
|
|
|
if let Some(mut map) = payload.as_object().cloned() {
|
|
|
|
|
map.insert("kind".into(), Value::String(kind.into()));
|
|
|
|
|
return Value::Object(map);
|
|
|
|
|
}
|
|
|
|
|
json!({
|
|
|
|
|
"kind": kind,
|
|
|
|
|
"data": payload,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn build_stream_delta_payload(
|
|
|
|
|
context: &RequestContext,
|
|
|
|
|
query: &StreamSnapshotQuery,
|
|
|
|
|
workspace_id: &str,
|
|
|
|
|
overview: &Value,
|
|
|
|
|
cursor: Option<String>,
|
|
|
|
|
delta: Value,
|
|
|
|
|
) -> Value {
|
|
|
|
|
let scope = resolve_stream_scope(query);
|
|
|
|
|
json!({
|
|
|
|
|
"kind": "delta",
|
|
|
|
|
"stream": scope.as_str(),
|
|
|
|
|
"projection": scope.projection(),
|
|
|
|
|
"requestId": context.trace.request_id,
|
|
|
|
|
"traceId": context.trace.trace_id,
|
|
|
|
|
"workspaceId": workspace_id,
|
|
|
|
|
"rootNodeId": normalize_root_node_id(query),
|
|
|
|
|
"depth": query.depth,
|
|
|
|
|
"cursor": cursor,
|
|
|
|
|
"data": delta,
|
|
|
|
|
"snapshot": Value::Null,
|
|
|
|
|
"overview": overview,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub async fn load_stream_overview(
|
|
|
|
|
config: &AppConfig,
|
|
|
|
|
context: &RequestContext,
|
|
|
|
|
query: &StreamSnapshotQuery,
|
|
|
|
|
) -> Result<(String, Value), WebError> {
|
|
|
|
|
let effective_workspace_id =
|
|
|
|
|
resolve_effective_workspace_id(context, query.workspace_id.as_deref(), true)?
|
|
|
|
|
.expect("workspace_required 已确保存在");
|
|
|
|
|
let overview = execute_runtime_query_via_convex(
|
|
|
|
|
config,
|
|
|
|
|
context,
|
|
|
|
|
Some(&effective_workspace_id),
|
|
|
|
|
workspace_overview_query(&effective_workspace_id, query),
|
|
|
|
|
)
|
|
|
|
|
.await?;
|
|
|
|
|
Ok((effective_workspace_id, overview))
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-18 05:43:49 +08:00
|
|
|
pub async fn load_stream_snapshot(
|
|
|
|
|
config: &AppConfig,
|
|
|
|
|
context: &RequestContext,
|
|
|
|
|
query: &StreamSnapshotQuery,
|
|
|
|
|
) -> Result<Value, WebError> {
|
|
|
|
|
let effective_workspace_id =
|
|
|
|
|
resolve_effective_workspace_id(context, query.workspace_id.as_deref(), true)?
|
|
|
|
|
.expect("workspace_required 已确保存在");
|
|
|
|
|
let scope = resolve_stream_scope(query);
|
|
|
|
|
|
|
|
|
|
let snapshot = match scope {
|
|
|
|
|
StreamSnapshotScope::Workspace => {
|
|
|
|
|
let loaded = load_projection_snapshot(
|
|
|
|
|
config,
|
|
|
|
|
context,
|
|
|
|
|
&ProjectionSnapshotSpec {
|
|
|
|
|
workspace_id: &effective_workspace_id,
|
|
|
|
|
root_node_id: None,
|
|
|
|
|
depth: query.depth,
|
2026-04-26 19:35:52 +08:00
|
|
|
query: None,
|
|
|
|
|
max_results: None,
|
2026-04-18 05:43:49 +08:00
|
|
|
projection: KernelProjectionKind::SidebarTree,
|
|
|
|
|
},
|
|
|
|
|
)
|
|
|
|
|
.await?;
|
|
|
|
|
|
|
|
|
|
json!({
|
|
|
|
|
"dataset": loaded.dataset,
|
|
|
|
|
"tree": loaded.projection,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
StreamSnapshotScope::Subtree => {
|
2026-04-26 19:35:52 +08:00
|
|
|
let root_node_id =
|
|
|
|
|
normalize_root_node_id(query).expect("subtree scope 已确保 rootNodeId 存在");
|
2026-04-18 05:43:49 +08:00
|
|
|
let dataset = load_sidebar_dataset(config, context, &effective_workspace_id).await?;
|
|
|
|
|
let tree = execute_kernel_query(
|
|
|
|
|
context,
|
|
|
|
|
&effective_workspace_id,
|
2026-04-26 04:29:23 +08:00
|
|
|
subtree_query(&effective_workspace_id, &root_node_id, query.depth),
|
2026-04-18 05:43:49 +08:00
|
|
|
dataset.clone(),
|
|
|
|
|
)?;
|
|
|
|
|
|
|
|
|
|
json!({
|
|
|
|
|
"dataset": dataset,
|
|
|
|
|
"tree": tree,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
let overview = execute_runtime_query_via_convex(
|
|
|
|
|
config,
|
|
|
|
|
context,
|
|
|
|
|
Some(&effective_workspace_id),
|
|
|
|
|
workspace_overview_query(&effective_workspace_id, query),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.ok();
|
2026-04-26 04:29:23 +08:00
|
|
|
let cursor = resolve_stream_cursor(overview.as_ref(), query.cursor.as_deref());
|
2026-04-18 05:43:49 +08:00
|
|
|
|
|
|
|
|
Ok(json!({
|
|
|
|
|
"kind": "snapshot",
|
|
|
|
|
"scope": scope.as_str(),
|
2026-04-26 04:29:23 +08:00
|
|
|
"stream": scope.as_str(),
|
|
|
|
|
"projection": scope.projection(),
|
|
|
|
|
"cursor": cursor,
|
2026-04-18 05:43:49 +08:00
|
|
|
"requestId": context.trace.request_id,
|
|
|
|
|
"traceId": context.trace.trace_id,
|
|
|
|
|
"workspaceId": effective_workspace_id,
|
2026-04-26 04:29:23 +08:00
|
|
|
"rootNodeId": normalize_root_node_id(query),
|
2026-04-18 05:43:49 +08:00
|
|
|
"depth": query.depth,
|
2026-04-26 04:29:23 +08:00
|
|
|
"data": snapshot,
|
2026-04-18 05:43:49 +08:00
|
|
|
"snapshot": snapshot,
|
|
|
|
|
"overview": overview,
|
|
|
|
|
}))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[cfg(test)]
|
|
|
|
|
mod tests {
|
2026-04-26 04:29:23 +08:00
|
|
|
use super::{
|
2026-04-29 12:24:44 +08:00
|
|
|
resolve_stream_change, resolve_stream_cursor, resolve_stream_scope, StreamChangeKind,
|
|
|
|
|
StreamSnapshotQuery, StreamSnapshotScope,
|
2026-04-26 04:29:23 +08:00
|
|
|
};
|
|
|
|
|
use serde_json::json;
|
2026-04-18 05:43:49 +08:00
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn stream_scope_defaults_to_workspace() {
|
|
|
|
|
assert_eq!(
|
|
|
|
|
resolve_stream_scope(&StreamSnapshotQuery::default()),
|
|
|
|
|
StreamSnapshotScope::Workspace
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn stream_scope_switches_to_subtree_when_root_exists() {
|
|
|
|
|
assert_eq!(
|
|
|
|
|
resolve_stream_scope(&StreamSnapshotQuery {
|
|
|
|
|
root_node_id: Some("page_root".into()),
|
|
|
|
|
..StreamSnapshotQuery::default()
|
|
|
|
|
}),
|
|
|
|
|
StreamSnapshotScope::Subtree
|
|
|
|
|
);
|
|
|
|
|
}
|
2026-04-26 04:29:23 +08:00
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn stream_cursor_prefers_newer_domain_event() {
|
|
|
|
|
let overview = json!({
|
|
|
|
|
"command_logs": [
|
|
|
|
|
{
|
|
|
|
|
"command_id": "cmd_1",
|
|
|
|
|
"created_at": "2026-04-25T10:00:00Z"
|
|
|
|
|
}
|
|
|
|
|
],
|
|
|
|
|
"domain_events": [
|
|
|
|
|
{
|
|
|
|
|
"event_id": "evt_2",
|
|
|
|
|
"created_at": "2026-04-25T10:00:01Z"
|
|
|
|
|
}
|
|
|
|
|
]
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
assert_eq!(
|
|
|
|
|
resolve_stream_cursor(Some(&overview), None),
|
|
|
|
|
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"domain_event:evt_2"}"#.into())
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn stream_change_detects_delta_from_single_new_command() {
|
|
|
|
|
let overview = json!({
|
|
|
|
|
"command_logs": [
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_2",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_2",
|
|
|
|
|
"created_at": "2026-04-25T10:00:02Z",
|
|
|
|
|
"command_name": "tree.node.archive",
|
|
|
|
|
"payload": {
|
|
|
|
|
"streamDelta": {
|
|
|
|
|
"op": "remove_document",
|
|
|
|
|
"documentId": "page_2"
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_1",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_1",
|
|
|
|
|
"created_at": "2026-04-25T10:00:01Z"
|
|
|
|
|
}
|
|
|
|
|
],
|
|
|
|
|
"domain_events": []
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let change = resolve_stream_change(
|
|
|
|
|
&overview,
|
2026-04-26 19:35:52 +08:00
|
|
|
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
2026-04-26 04:29:23 +08:00
|
|
|
)
|
|
|
|
|
.expect("应识别到变化");
|
|
|
|
|
|
|
|
|
|
assert_eq!(change.kind, StreamChangeKind::Delta);
|
|
|
|
|
assert_eq!(
|
|
|
|
|
change.cursor,
|
2026-04-26 19:35:52 +08:00
|
|
|
Some(r#"{"createdAt":"2026-04-25T10:00:02Z","id":"clog_2"}"#.into())
|
2026-04-26 04:29:23 +08:00
|
|
|
);
|
|
|
|
|
assert_eq!(
|
|
|
|
|
change.delta,
|
|
|
|
|
Some(json!({
|
|
|
|
|
"op": "remove_document",
|
|
|
|
|
"documentId": "page_2"
|
|
|
|
|
}))
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
#[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"
|
|
|
|
|
}
|
|
|
|
|
]
|
|
|
|
|
}))
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
#[test]
|
|
|
|
|
fn stream_change_detects_noop_delta_for_non_tree_mutating_command() {
|
|
|
|
|
let overview = json!({
|
|
|
|
|
"command_logs": [
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_2",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_2",
|
|
|
|
|
"created_at": "2026-04-25T10:00:02Z",
|
|
|
|
|
"command_name": "page.body.save",
|
|
|
|
|
"payload": {
|
|
|
|
|
"documentId": "page_1"
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_1",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_1",
|
|
|
|
|
"created_at": "2026-04-25T10:00:01Z"
|
|
|
|
|
}
|
|
|
|
|
],
|
|
|
|
|
"domain_events": []
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let change = resolve_stream_change(
|
|
|
|
|
&overview,
|
2026-04-26 19:35:52 +08:00
|
|
|
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
2026-04-26 04:29:23 +08:00
|
|
|
)
|
|
|
|
|
.expect("应识别到变化");
|
|
|
|
|
|
|
|
|
|
assert_eq!(change.kind, StreamChangeKind::Delta);
|
|
|
|
|
assert_eq!(change.delta, Some(json!({ "op": "noop" })));
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 19:35:52 +08:00
|
|
|
#[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
|
|
|
|
|
}))
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-26 04:29:23 +08:00
|
|
|
#[test]
|
|
|
|
|
fn stream_change_falls_back_to_resync_when_delta_is_unstable() {
|
|
|
|
|
let overview = json!({
|
|
|
|
|
"command_logs": [
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_2",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_2",
|
|
|
|
|
"created_at": "2026-04-25T10:00:02Z",
|
|
|
|
|
"command_name": "tree.subtree.move",
|
|
|
|
|
"payload": {
|
|
|
|
|
"documentId": "page_2"
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
{
|
2026-04-26 19:35:52 +08:00
|
|
|
"id": "clog_1",
|
2026-04-26 04:29:23 +08:00
|
|
|
"command_id": "cmd_1",
|
|
|
|
|
"created_at": "2026-04-25T10:00:01Z"
|
|
|
|
|
}
|
|
|
|
|
],
|
|
|
|
|
"domain_events": []
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let change = resolve_stream_change(
|
|
|
|
|
&overview,
|
2026-04-26 19:35:52 +08:00
|
|
|
Some(r#"{"createdAt":"2026-04-25T10:00:01Z","id":"clog_1"}"#),
|
2026-04-26 04:29:23 +08:00
|
|
|
)
|
|
|
|
|
.expect("应识别到变化");
|
|
|
|
|
|
|
|
|
|
assert_eq!(change.kind, StreamChangeKind::Resync);
|
|
|
|
|
assert_eq!(change.delta, None);
|
|
|
|
|
}
|
2026-04-18 05:43:49 +08:00
|
|
|
}
|