Exchange DB changes over Iroh
This commit is contained in:
parent
284207eb01
commit
20fd773a42
10 changed files with 563 additions and 15 deletions
|
|
@ -263,6 +263,11 @@ pub async fn handle_request_async(
|
|||
node: peer_node,
|
||||
name,
|
||||
} => kv_sync_from_peer(node, &peer_node, &name).await,
|
||||
ControlRequest::DbSync {
|
||||
node: peer_node,
|
||||
name,
|
||||
limit,
|
||||
} => db_sync_from_peer(node, &peer_node, &name, limit).await,
|
||||
ControlRequest::PubsubPub {
|
||||
topic,
|
||||
message,
|
||||
|
|
@ -555,7 +560,8 @@ async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result<ControlResponse,
|
|||
| PeerControlResponse::KvSynced { .. }
|
||||
| PeerControlResponse::PubsubPublished { .. }
|
||||
| PeerControlResponse::PipeConnected { .. }
|
||||
| PeerControlResponse::DocumentSynced { .. } => Err(NodeError::IrohPeer(
|
||||
| PeerControlResponse::DocumentSynced { .. }
|
||||
| PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer(
|
||||
"peer returned wrong response type to ping request".to_owned(),
|
||||
)),
|
||||
}
|
||||
|
|
@ -666,7 +672,8 @@ async fn peer_auth_check(
|
|||
| PeerControlResponse::KvSynced { .. }
|
||||
| PeerControlResponse::PubsubPublished { .. }
|
||||
| PeerControlResponse::PipeConnected { .. }
|
||||
| PeerControlResponse::DocumentSynced { .. } => Err(NodeError::IrohPeer(
|
||||
| PeerControlResponse::DocumentSynced { .. }
|
||||
| PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer(
|
||||
"peer returned wrong response type to auth-check request".to_owned(),
|
||||
)),
|
||||
}
|
||||
|
|
@ -806,7 +813,8 @@ async fn cas_fetch_from_peer(
|
|||
| PeerControlResponse::KvSynced { .. }
|
||||
| PeerControlResponse::PubsubPublished { .. }
|
||||
| PeerControlResponse::PipeConnected { .. }
|
||||
| PeerControlResponse::DocumentSynced { .. } => Err(NodeError::IrohPeer(
|
||||
| PeerControlResponse::DocumentSynced { .. }
|
||||
| PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer(
|
||||
"peer returned wrong response type to CAS fetch".to_owned(),
|
||||
)),
|
||||
}
|
||||
|
|
@ -1212,6 +1220,94 @@ async fn document_sync_from_peer(
|
|||
}
|
||||
}
|
||||
|
||||
async fn db_sync_from_peer(
|
||||
node: &LocalNode,
|
||||
peer_node: &str,
|
||||
name: &str,
|
||||
limit: u32,
|
||||
) -> Result<ControlResponse, NodeError> {
|
||||
geth_db::validate_db_name(name).map_err(|_| NodeError::InvalidDbName(name.to_owned()))?;
|
||||
let store = Store::open(&node.paths.metadata_db())?;
|
||||
let local = store
|
||||
.get_db_resource_by_name(name)?
|
||||
.ok_or_else(|| NodeError::DbNotFound(name.to_owned()))?;
|
||||
let stream = format!("db:{name}");
|
||||
let cursor = load_live_sync_cursor(&store, peer_node, &stream)?;
|
||||
let after_db_version = (cursor > 0).then_some(cursor);
|
||||
let response = request_peer_control(node, peer_node, "db-sync", |peer_card, nonce| {
|
||||
PeerControlRequest::DbSync {
|
||||
peer_card,
|
||||
name: name.to_owned(),
|
||||
after_db_version,
|
||||
limit,
|
||||
nonce,
|
||||
}
|
||||
})
|
||||
.await?;
|
||||
match response {
|
||||
PeerControlResponse::DbSynced {
|
||||
node_id,
|
||||
agent_id,
|
||||
endpoint_id,
|
||||
name: response_name,
|
||||
batch,
|
||||
high_water_db_version,
|
||||
allowed,
|
||||
reason,
|
||||
note,
|
||||
..
|
||||
} if response_name == name => {
|
||||
if !allowed {
|
||||
return Ok(ControlResponse::DbSynced {
|
||||
peer_node_id: node_id,
|
||||
peer_agent_id: agent_id,
|
||||
endpoint_id,
|
||||
name: response_name,
|
||||
changes_received: 0,
|
||||
max_db_version: None,
|
||||
schema_match: false,
|
||||
allowed,
|
||||
reason,
|
||||
note,
|
||||
});
|
||||
}
|
||||
let local_schema = geth_db::schema_metadata(Path::new(&local.path))?;
|
||||
let (changes_received, max_db_version, schema_match) = if let Some(batch) = batch {
|
||||
let schema_match = batch.schema_metadata == local_schema;
|
||||
let changes_received = batch.changes.len();
|
||||
let max_db_version = batch.max_db_version;
|
||||
if schema_match {
|
||||
if let Some(next_cursor) = high_water_db_version.or(max_db_version) {
|
||||
store_live_sync_cursor(&store, peer_node, &stream, next_cursor)?;
|
||||
}
|
||||
}
|
||||
(changes_received, max_db_version, schema_match)
|
||||
} else {
|
||||
(0, high_water_db_version, false)
|
||||
};
|
||||
Ok(ControlResponse::DbSynced {
|
||||
peer_node_id: node_id,
|
||||
peer_agent_id: agent_id,
|
||||
endpoint_id,
|
||||
name: response_name,
|
||||
changes_received,
|
||||
max_db_version,
|
||||
schema_match,
|
||||
allowed,
|
||||
reason,
|
||||
note,
|
||||
})
|
||||
}
|
||||
PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer(
|
||||
"peer DB sync response did not match request".to_owned(),
|
||||
)),
|
||||
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
|
||||
_ => Err(NodeError::IrohPeer(
|
||||
"peer returned wrong response type to DB sync".to_owned(),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
fn live_sync_cursor_key(peer_node: &str, stream: &str) -> String {
|
||||
format!("live-sync:{peer_node}:{stream}")
|
||||
}
|
||||
|
|
@ -1320,6 +1416,10 @@ async fn request_peer_control(
|
|||
| PeerControlResponse::DocumentSynced {
|
||||
nonce: response_nonce,
|
||||
..
|
||||
}
|
||||
| PeerControlResponse::DbSynced {
|
||||
nonce: response_nonce,
|
||||
..
|
||||
} if response_nonce == &nonce => Ok(response),
|
||||
PeerControlResponse::Error { .. } => Ok(response),
|
||||
_ => Err(NodeError::IrohPeer(format!(
|
||||
|
|
@ -1369,6 +1469,7 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> {
|
|||
let peers = Store::open(&node.paths.metadata_db())?.list_peer_cards()?;
|
||||
let kv_stores = Store::open(&node.paths.metadata_db())?.list_kv_stores()?;
|
||||
let documents = Store::open(&node.paths.metadata_db())?.list_document_resources()?;
|
||||
let dbs = Store::open(&node.paths.metadata_db())?.list_db_resources()?;
|
||||
for peer in peers {
|
||||
if let Err(error) = ssh_cert_sync_from_peer(node, &peer.peer_id).await {
|
||||
tracing::debug!(peer = %peer.peer_id, %error, "SSH cert live sync failed");
|
||||
|
|
@ -1386,6 +1487,11 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> {
|
|||
tracing::debug!(peer = %peer.peer_id, document = %document.name, %error, "document live sync failed");
|
||||
}
|
||||
}
|
||||
for db in &dbs {
|
||||
if let Err(error) = db_sync_from_peer(node, &peer.peer_id, &db.name, 100).await {
|
||||
tracing::debug!(peer = %peer.peer_id, db = %db.name, %error, "DB live sync failed");
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -1868,6 +1974,78 @@ async fn handle_iroh_control_connection(
|
|||
}
|
||||
}
|
||||
}
|
||||
PeerControlRequest::DbSync {
|
||||
peer_card,
|
||||
name,
|
||||
after_db_version,
|
||||
limit,
|
||||
nonce,
|
||||
} => {
|
||||
geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?;
|
||||
peer_card.validate_candidate()?;
|
||||
ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?;
|
||||
let discovered = DiscoveredPeer::candidate(
|
||||
peer_card.clone(),
|
||||
UnixMillis(geth_store::now_ms()),
|
||||
DiscoverySource::PeerExchange,
|
||||
)?;
|
||||
let store = Store::open(&node.paths.metadata_db())?;
|
||||
store.upsert_peer_card(&StoredPeerCard {
|
||||
peer_id: peer_card.node_id.to_string(),
|
||||
card_json: serde_json::to_string(&peer_card)?,
|
||||
updated_at_ms: discovered.discovered_at.0,
|
||||
})?;
|
||||
if let Some(db) = store.get_db_resource_by_name(&name)? {
|
||||
let capability = "db.sync".to_owned();
|
||||
let explanation = geth_auth::explain_auth_ops(
|
||||
&load_auth_ops_for_resource(&store, &db.resource_id)?,
|
||||
PrincipalId::new(peer_card.node_id.to_string()),
|
||||
ResourceId::new(db.resource_id.clone()),
|
||||
Capability::new(capability),
|
||||
);
|
||||
let sync_result: Result<_, String> = if explanation.allowed {
|
||||
let path = Path::new(&db.path);
|
||||
match geth_db::extract_crsqlite_changes(path, after_db_version, limit) {
|
||||
Ok(batch) => match geth_db::crsqlite_change_metadata(path) {
|
||||
Ok(metadata) => {
|
||||
let high_water_db_version =
|
||||
metadata.max_db_version.or(batch.max_db_version);
|
||||
Ok((Some(batch), high_water_db_version))
|
||||
}
|
||||
Err(error) => Err(format!(
|
||||
"db sync cannot read crsql_changes metadata for {name}: {error}"
|
||||
)),
|
||||
},
|
||||
Err(error) => Err(format!(
|
||||
"db sync cannot read crsql_changes for {name}: {error}"
|
||||
)),
|
||||
}
|
||||
} else {
|
||||
Ok((None, None))
|
||||
};
|
||||
match sync_result {
|
||||
Ok((batch, high_water_db_version)) => PeerControlResponse::DbSynced {
|
||||
node_id: node.node_id.clone(),
|
||||
agent_id: node.agent_id.clone(),
|
||||
endpoint_id: node.iroh_status.endpoint_id.clone().unwrap_or_default(),
|
||||
remote_endpoint_id,
|
||||
name,
|
||||
batch,
|
||||
high_water_db_version,
|
||||
allowed: explanation.allowed,
|
||||
reason: explanation.reason,
|
||||
evaluated_ops: explanation.evaluated_ops,
|
||||
nonce,
|
||||
note: "DB sync authenticated endpoint/card binding and required db.sync on the remote DB resource; bootstrap exchanges typed crsql_changes and cursors them, but applying remote changes is not implemented yet".to_owned(),
|
||||
},
|
||||
Err(message) => PeerControlResponse::Error { message },
|
||||
}
|
||||
} else {
|
||||
PeerControlResponse::Error {
|
||||
message: format!("db resource not found: {name}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
send.write_all(geth_control::encode_peer_response(&response)?.as_bytes())
|
||||
.await
|
||||
|
|
@ -2007,6 +2185,7 @@ pub fn handle_request(
|
|||
ControlRequest::SshCertSync { .. } => Err(NodeError::IrohEndpointUnavailable),
|
||||
ControlRequest::SshRevocationSync { .. } => Err(NodeError::IrohEndpointUnavailable),
|
||||
ControlRequest::KvSync { .. } => Err(NodeError::IrohEndpointUnavailable),
|
||||
ControlRequest::DbSync { .. } => Err(NodeError::IrohEndpointUnavailable),
|
||||
ControlRequest::DocumentSync { .. } => Err(NodeError::IrohEndpointUnavailable),
|
||||
ControlRequest::ResourceList => Ok(ControlResponse::ResourceList {
|
||||
resources: store
|
||||
|
|
@ -3471,6 +3650,63 @@ mod tests {
|
|||
.expect("write config");
|
||||
}
|
||||
|
||||
fn create_mock_crsqlite_db(path: &Path, change_db_version: Option<i64>) {
|
||||
let conn = rusqlite::Connection::open(path).expect("open sqlite");
|
||||
conn.execute("CREATE TABLE notes(id INTEGER PRIMARY KEY, body TEXT)", [])
|
||||
.expect("create notes");
|
||||
conn.execute(
|
||||
r#"CREATE TABLE crsql_changes(
|
||||
table_name TEXT NOT NULL,
|
||||
pk BLOB NOT NULL,
|
||||
cid TEXT NOT NULL,
|
||||
val BLOB,
|
||||
col_version INTEGER NOT NULL,
|
||||
db_version INTEGER NOT NULL,
|
||||
site_id BLOB,
|
||||
cl INTEGER,
|
||||
seq INTEGER
|
||||
)"#,
|
||||
[],
|
||||
)
|
||||
.expect("create crsql_changes");
|
||||
if let Some(db_version) = change_db_version {
|
||||
conn.execute(
|
||||
"INSERT INTO crsql_changes(table_name, pk, cid, val, col_version, db_version, site_id, cl, seq) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
|
||||
(
|
||||
"notes",
|
||||
vec![db_version as u8],
|
||||
"body",
|
||||
Vec::from("hello".as_bytes()),
|
||||
1_i64,
|
||||
db_version,
|
||||
vec![1_u8],
|
||||
1_i64,
|
||||
db_version,
|
||||
),
|
||||
)
|
||||
.expect("insert initial mock crsqlite change");
|
||||
}
|
||||
}
|
||||
|
||||
fn insert_mock_crsqlite_change(path: &Path, db_version: i64, body: &str) {
|
||||
let conn = rusqlite::Connection::open(path).expect("open sqlite");
|
||||
conn.execute(
|
||||
"INSERT INTO crsql_changes(table_name, pk, cid, val, col_version, db_version, site_id, cl, seq) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
|
||||
(
|
||||
"notes",
|
||||
vec![db_version as u8],
|
||||
"body",
|
||||
body.as_bytes().to_vec(),
|
||||
1_i64,
|
||||
db_version,
|
||||
vec![1_u8],
|
||||
1_i64,
|
||||
db_version,
|
||||
),
|
||||
)
|
||||
.expect("insert mock crsqlite change");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lan_discovery_address_selection_uses_iroh_direct_addresses() {
|
||||
let key = AgentKey::generate();
|
||||
|
|
@ -3618,6 +3854,26 @@ mod tests {
|
|||
},
|
||||
)
|
||||
.expect("right document set");
|
||||
let left_db_path = left_home.path().join("notes.sqlite");
|
||||
let right_db_path = right_home.path().join("notes.sqlite");
|
||||
create_mock_crsqlite_db(&left_db_path, None);
|
||||
create_mock_crsqlite_db(&right_db_path, Some(7));
|
||||
handle_request(
|
||||
&left,
|
||||
ControlRequest::DbAdd {
|
||||
name: "notes".to_owned(),
|
||||
path: left_db_path.clone(),
|
||||
},
|
||||
)
|
||||
.expect("left db add");
|
||||
handle_request(
|
||||
&right,
|
||||
ControlRequest::DbAdd {
|
||||
name: "notes".to_owned(),
|
||||
path: right_db_path.clone(),
|
||||
},
|
||||
)
|
||||
.expect("right db add");
|
||||
|
||||
let ping = handle_request_async(
|
||||
&left,
|
||||
|
|
@ -3789,6 +4045,32 @@ mod tests {
|
|||
other => panic!("unexpected denied document sync response: {other:?}"),
|
||||
}
|
||||
|
||||
let denied_db_sync = handle_request_async(
|
||||
&left,
|
||||
ControlRequest::DbSync {
|
||||
node: right_card.node_id.to_string(),
|
||||
name: "notes".to_owned(),
|
||||
limit: 10,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("denied db sync");
|
||||
match denied_db_sync {
|
||||
ControlResponse::DbSynced {
|
||||
allowed,
|
||||
changes_received,
|
||||
schema_match,
|
||||
reason,
|
||||
..
|
||||
} => {
|
||||
assert!(!allowed);
|
||||
assert_eq!(changes_received, 0);
|
||||
assert!(!schema_match);
|
||||
assert!(reason.contains("no active direct or group grant"));
|
||||
}
|
||||
other => panic!("unexpected denied DB sync response: {other:?}"),
|
||||
}
|
||||
|
||||
handle_request(
|
||||
&right,
|
||||
ControlRequest::AuthGrant {
|
||||
|
|
@ -3839,6 +4121,16 @@ mod tests {
|
|||
},
|
||||
)
|
||||
.expect("grant left document read");
|
||||
handle_request(
|
||||
&right,
|
||||
ControlRequest::AuthGrant {
|
||||
subject: left.node_id.clone(),
|
||||
resource: "resource:db:notes".to_owned(),
|
||||
capability: "db.sync".to_owned(),
|
||||
grant_id: Some("grant:left-db-sync".to_owned()),
|
||||
},
|
||||
)
|
||||
.expect("grant left db sync");
|
||||
|
||||
let allowed = handle_request_async(
|
||||
&left,
|
||||
|
|
@ -4040,6 +4332,37 @@ mod tests {
|
|||
other => panic!("unexpected synced document get response: {other:?}"),
|
||||
}
|
||||
|
||||
let db_sync = handle_request_async(
|
||||
&left,
|
||||
ControlRequest::DbSync {
|
||||
node: right_card.node_id.to_string(),
|
||||
name: "notes".to_owned(),
|
||||
limit: 10,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("allowed db sync");
|
||||
match db_sync {
|
||||
ControlResponse::DbSynced {
|
||||
allowed,
|
||||
changes_received,
|
||||
max_db_version,
|
||||
schema_match,
|
||||
reason,
|
||||
note,
|
||||
..
|
||||
} => {
|
||||
assert!(allowed);
|
||||
assert_eq!(changes_received, 1);
|
||||
assert_eq!(max_db_version, Some(7));
|
||||
assert!(schema_match);
|
||||
assert!(reason.contains("direct grant"));
|
||||
assert!(note.contains("db.sync"));
|
||||
assert!(note.contains("applying remote changes is not implemented yet"));
|
||||
}
|
||||
other => panic!("unexpected allowed DB sync response: {other:?}"),
|
||||
}
|
||||
|
||||
let denied_cert_sync = handle_request_async(
|
||||
&left,
|
||||
ControlRequest::SshCertSync {
|
||||
|
|
@ -4207,6 +4530,7 @@ mod tests {
|
|||
},
|
||||
)
|
||||
.expect("right live document set");
|
||||
insert_mock_crsqlite_change(&right_db_path, 8, "live");
|
||||
|
||||
run_live_sync_once(&left)
|
||||
.await
|
||||
|
|
@ -4240,6 +4564,19 @@ mod tests {
|
|||
.map(|document| document.state_json),
|
||||
Some(r#"{"title":"live"}"#.to_owned())
|
||||
);
|
||||
let db_cursor = left_store
|
||||
.get_module_state(&live_sync_cursor_key(
|
||||
right_card.node_id.as_str(),
|
||||
"db:notes",
|
||||
))
|
||||
.expect("get live-synced db cursor")
|
||||
.expect("db cursor exists");
|
||||
assert_eq!(
|
||||
serde_json::from_str::<LiveSyncCursor>(&db_cursor.state_json)
|
||||
.expect("parse db cursor")
|
||||
.cursor_ms,
|
||||
8
|
||||
);
|
||||
|
||||
left_endpoint.shutdown().await;
|
||||
right_endpoint.shutdown().await;
|
||||
|
|
|
|||
Loading…
Reference in a new issue