Store documents as Automerge state
This commit is contained in:
parent
606641dd4a
commit
b01da7be4b
9 changed files with 490 additions and 73 deletions
|
|
@ -6,6 +6,9 @@ rust-version.workspace = true
|
|||
license.workspace = true
|
||||
|
||||
[dependencies]
|
||||
automerge.workspace = true
|
||||
base64.workspace = true
|
||||
hexane.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
thiserror.workspace = true
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
use geth_types::{DocumentId, ResourceId, UnixMillis};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
pub const AUTOMERGE_STATE_FORMAT: &str = "geth.automerge-document.v1";
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct DocumentResource {
|
||||
pub id: DocumentId,
|
||||
|
|
@ -20,9 +22,18 @@ pub struct DocumentState {
|
|||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum DocumentError {
|
||||
#[error("invalid document name: {0}")]
|
||||
InvalidName(String),
|
||||
Name(String),
|
||||
#[error("invalid document JSON state: {0}")]
|
||||
InvalidState(#[from] serde_json::Error),
|
||||
JsonState(#[from] serde_json::Error),
|
||||
#[error("invalid Automerge document state: {0}")]
|
||||
AutomergeState(String),
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
struct AutomergeStateEnvelope {
|
||||
format: String,
|
||||
automerge_base64: String,
|
||||
view_json: String,
|
||||
}
|
||||
|
||||
pub fn validate_document_name(name: &str) -> Result<(), DocumentError> {
|
||||
|
|
@ -31,7 +42,7 @@ pub fn validate_document_name(name: &str) -> Result<(), DocumentError> {
|
|||
.bytes()
|
||||
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
|
||||
{
|
||||
return Err(DocumentError::InvalidName(name.to_owned()));
|
||||
return Err(DocumentError::Name(name.to_owned()));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -41,9 +52,201 @@ pub fn normalize_document_state(state_json: &str) -> Result<String, DocumentErro
|
|||
Ok(serde_json::to_string(&value)?)
|
||||
}
|
||||
|
||||
pub fn create_automerge_state(state_json: &str) -> Result<String, DocumentError> {
|
||||
let value: serde_json::Value = serde_json::from_str(state_json)?;
|
||||
let normalized = serde_json::to_string(&value)?;
|
||||
let mut doc = automerge::Automerge::new();
|
||||
{
|
||||
let mut tx = doc.transaction();
|
||||
write_json_object(&mut tx, &automerge::ROOT, &value)?;
|
||||
tx.commit();
|
||||
}
|
||||
envelope_from_doc(&doc, normalized)
|
||||
}
|
||||
|
||||
pub fn document_view_json(stored_state: &str) -> Result<String, DocumentError> {
|
||||
if let Some(envelope) = parse_envelope(stored_state)? {
|
||||
return Ok(envelope.view_json);
|
||||
}
|
||||
normalize_document_state(stored_state)
|
||||
}
|
||||
|
||||
pub fn document_state_bytes(stored_state: &str) -> u64 {
|
||||
match parse_envelope(stored_state) {
|
||||
Ok(Some(envelope)) => {
|
||||
use base64::Engine;
|
||||
base64::engine::general_purpose::STANDARD
|
||||
.decode(envelope.automerge_base64)
|
||||
.map(|bytes| bytes.len() as u64)
|
||||
.unwrap_or(stored_state.len() as u64)
|
||||
}
|
||||
Ok(None) | Err(_) => stored_state.len() as u64,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn merge_automerge_states(local: &str, remote: &str) -> Result<String, DocumentError> {
|
||||
let Some(local_envelope) = parse_envelope(local)? else {
|
||||
return create_automerge_state(remote);
|
||||
};
|
||||
let Some(remote_envelope) = parse_envelope(remote)? else {
|
||||
return create_automerge_state(&document_view_json(remote)?);
|
||||
};
|
||||
let mut local_doc = load_envelope_doc(&local_envelope)?;
|
||||
let mut remote_doc = load_envelope_doc(&remote_envelope)?;
|
||||
local_doc
|
||||
.merge(&mut remote_doc)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
let view_json = serde_json::to_string(&automerge::AutoSerde::from(&local_doc))?;
|
||||
envelope_from_doc(&local_doc, view_json)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn automerge_roadmap() -> &'static str {
|
||||
"future documents use Automerge sync over Iroh with resource-local authorization"
|
||||
"documents are stored as Automerge save bytes with a JSON view; remote sync is resource-authorized and exchanges the durable Automerge state"
|
||||
}
|
||||
|
||||
fn parse_envelope(state: &str) -> Result<Option<AutomergeStateEnvelope>, DocumentError> {
|
||||
let value: serde_json::Value = serde_json::from_str(state)?;
|
||||
if value
|
||||
.get("format")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.is_some_and(|format| format == AUTOMERGE_STATE_FORMAT)
|
||||
{
|
||||
return serde_json::from_value(value)
|
||||
.map(Some)
|
||||
.map_err(DocumentError::from);
|
||||
}
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
fn load_envelope_doc(
|
||||
envelope: &AutomergeStateEnvelope,
|
||||
) -> Result<automerge::Automerge, DocumentError> {
|
||||
use base64::Engine;
|
||||
let bytes = base64::engine::general_purpose::STANDARD
|
||||
.decode(&envelope.automerge_base64)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
automerge::Automerge::load(&bytes)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))
|
||||
}
|
||||
|
||||
fn envelope_from_doc(
|
||||
doc: &automerge::Automerge,
|
||||
view_json: String,
|
||||
) -> Result<String, DocumentError> {
|
||||
use base64::Engine;
|
||||
let envelope = AutomergeStateEnvelope {
|
||||
format: AUTOMERGE_STATE_FORMAT.to_owned(),
|
||||
automerge_base64: base64::engine::general_purpose::STANDARD.encode(doc.save()),
|
||||
view_json,
|
||||
};
|
||||
serde_json::to_string(&envelope).map_err(DocumentError::from)
|
||||
}
|
||||
|
||||
fn write_json_object(
|
||||
tx: &mut automerge::transaction::Transaction<'_>,
|
||||
obj: &automerge::ObjId,
|
||||
value: &serde_json::Value,
|
||||
) -> Result<(), DocumentError> {
|
||||
match value {
|
||||
serde_json::Value::Object(map) => {
|
||||
for (key, value) in map {
|
||||
put_json_value(tx, obj, key, value)?;
|
||||
}
|
||||
}
|
||||
other => {
|
||||
put_json_value(tx, obj, "value", other)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn put_json_value(
|
||||
tx: &mut automerge::transaction::Transaction<'_>,
|
||||
obj: &automerge::ObjId,
|
||||
key: &str,
|
||||
value: &serde_json::Value,
|
||||
) -> Result<(), DocumentError> {
|
||||
use automerge::transaction::Transactable;
|
||||
match value {
|
||||
serde_json::Value::Null => tx.put(obj, key, ()),
|
||||
serde_json::Value::Bool(value) => tx.put(obj, key, *value),
|
||||
serde_json::Value::Number(value) => {
|
||||
if let Some(value) = value.as_i64() {
|
||||
tx.put(obj, key, value)
|
||||
} else if let Some(value) = value.as_u64() {
|
||||
tx.put(obj, key, value)
|
||||
} else if let Some(value) = value.as_f64() {
|
||||
tx.put(obj, key, value)
|
||||
} else {
|
||||
Err(automerge::AutomergeError::Fail)
|
||||
}
|
||||
}
|
||||
serde_json::Value::String(value) => tx.put(obj, key, value.as_str()),
|
||||
serde_json::Value::Array(values) => {
|
||||
let list = tx
|
||||
.put_object(obj, key, automerge::ObjType::List)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
for (index, value) in values.iter().enumerate() {
|
||||
insert_json_value(tx, &list, index, value)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
serde_json::Value::Object(map) => {
|
||||
let map_obj = tx
|
||||
.put_object(obj, key, automerge::ObjType::Map)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
for (child_key, value) in map {
|
||||
put_json_value(tx, &map_obj, child_key, value)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))
|
||||
}
|
||||
|
||||
fn insert_json_value(
|
||||
tx: &mut automerge::transaction::Transaction<'_>,
|
||||
obj: &automerge::ObjId,
|
||||
index: usize,
|
||||
value: &serde_json::Value,
|
||||
) -> Result<(), DocumentError> {
|
||||
use automerge::transaction::Transactable;
|
||||
match value {
|
||||
serde_json::Value::Null => tx.insert(obj, index, ()),
|
||||
serde_json::Value::Bool(value) => tx.insert(obj, index, *value),
|
||||
serde_json::Value::Number(value) => {
|
||||
if let Some(value) = value.as_i64() {
|
||||
tx.insert(obj, index, value)
|
||||
} else if let Some(value) = value.as_u64() {
|
||||
tx.insert(obj, index, value)
|
||||
} else if let Some(value) = value.as_f64() {
|
||||
tx.insert(obj, index, value)
|
||||
} else {
|
||||
Err(automerge::AutomergeError::Fail)
|
||||
}
|
||||
}
|
||||
serde_json::Value::String(value) => tx.insert(obj, index, value.as_str()),
|
||||
serde_json::Value::Array(values) => {
|
||||
let list = tx
|
||||
.insert_object(obj, index, automerge::ObjType::List)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
for (index, value) in values.iter().enumerate() {
|
||||
insert_json_value(tx, &list, index, value)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
serde_json::Value::Object(map) => {
|
||||
let map_obj = tx
|
||||
.insert_object(obj, index, automerge::ObjType::Map)
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))?;
|
||||
for (child_key, value) in map {
|
||||
put_json_value(tx, &map_obj, child_key, value)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
.map_err(|error| DocumentError::AutomergeState(error.to_string()))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
|
@ -68,4 +271,29 @@ mod tests {
|
|||
);
|
||||
assert!(normalize_document_state("{").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn automerge_state_preserves_json_view_and_durable_bytes() {
|
||||
let state = create_automerge_state(r#"{ "title": "notes", "items": [1, true, null] }"#)
|
||||
.expect("automerge state");
|
||||
|
||||
assert_eq!(
|
||||
document_view_json(&state).expect("view"),
|
||||
r#"{"items":[1,true,null],"title":"notes"}"#
|
||||
);
|
||||
assert!(document_state_bytes(&state) > 2);
|
||||
assert!(state.contains(AUTOMERGE_STATE_FORMAT));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn automerge_states_can_merge_and_render_json_view() {
|
||||
let left = create_automerge_state(r#"{ "left": true }"#).expect("left state");
|
||||
let right = create_automerge_state(r#"{ "right": true }"#).expect("right state");
|
||||
|
||||
let merged = merge_automerge_states(&left, &right).expect("merge");
|
||||
let view = document_view_json(&merged).expect("merged view");
|
||||
|
||||
assert!(view.contains(r#""left":true"#));
|
||||
assert!(view.contains(r#""right":true"#));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2537,12 +2537,15 @@ async fn document_sync_from_peer(
|
|||
if let Some(state) = state {
|
||||
let local = ensure_local_document(&store, name)?;
|
||||
if state.updated_at.0 >= local.updated_at_ms {
|
||||
let normalized = geth_document::normalize_document_state(&state.state_json)?;
|
||||
let merged = geth_document::merge_automerge_states(
|
||||
&local.state_json,
|
||||
&state.state_json,
|
||||
)?;
|
||||
store.insert_document_resource(&StoredDocumentResource {
|
||||
document_id: local.document_id,
|
||||
resource_id: local.resource_id,
|
||||
name: local.name,
|
||||
state_json: normalized,
|
||||
state_json: merged,
|
||||
updated_at_ms: state.updated_at.0,
|
||||
})?;
|
||||
updated = true;
|
||||
|
|
@ -5242,7 +5245,7 @@ async fn handle_iroh_control_connection(
|
|||
bearer_proof.as_ref(),
|
||||
)?;
|
||||
let state = if explanation.allowed && document.updated_at_ms >= since_ms {
|
||||
Some(document_state_from_stored(&document))
|
||||
Some(document_sync_state_from_stored(&document)?)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
|
@ -5258,7 +5261,7 @@ async fn handle_iroh_control_connection(
|
|||
reason: explanation.reason,
|
||||
evaluated_ops: explanation.evaluated_ops,
|
||||
nonce,
|
||||
note: "document sync authenticated endpoint/card binding and required document.read on the remote document resource; JSON LWW sync is a bootstrap before Automerge".to_owned(),
|
||||
note: "document sync authenticated endpoint/card binding and required document.read on the remote document resource; durable Automerge state is exchanged with a JSON view for CLI output".to_owned(),
|
||||
}
|
||||
} else {
|
||||
PeerControlResponse::Error {
|
||||
|
|
@ -7430,7 +7433,7 @@ pub fn handle_request(
|
|||
document_id,
|
||||
resource_id,
|
||||
name,
|
||||
state_json: "{}".to_owned(),
|
||||
state_json: geth_document::create_automerge_state("{}")?,
|
||||
updated_at_ms: geth_store::now_ms(),
|
||||
};
|
||||
store.insert_document_resource(&stored)?;
|
||||
|
|
@ -7454,11 +7457,11 @@ pub fn handle_request(
|
|||
let mut stored = store
|
||||
.get_document_resource_by_name(&name)?
|
||||
.ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?;
|
||||
stored.state_json = geth_document::normalize_document_state(&state_json)?;
|
||||
stored.state_json = geth_document::create_automerge_state(&state_json)?;
|
||||
stored.updated_at_ms = geth_store::now_ms();
|
||||
store.insert_document_resource(&stored)?;
|
||||
Ok(ControlResponse::DocumentSet {
|
||||
state: document_state_from_stored(&stored),
|
||||
state: document_state_from_stored(&stored)?,
|
||||
})
|
||||
}
|
||||
ControlRequest::DocumentGet { name } => {
|
||||
|
|
@ -7468,7 +7471,7 @@ pub fn handle_request(
|
|||
.get_document_resource_by_name(&name)?
|
||||
.ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?;
|
||||
Ok(ControlResponse::DocumentGet {
|
||||
state: document_state_from_stored(&stored),
|
||||
state: document_state_from_stored(&stored)?,
|
||||
})
|
||||
}
|
||||
ControlRequest::PubsubPub {
|
||||
|
|
@ -7820,7 +7823,7 @@ fn ensure_local_document(store: &Store, name: &str) -> Result<StoredDocumentReso
|
|||
document_id: format!("document:{name}"),
|
||||
resource_id: format!("resource:document:{name}"),
|
||||
name: name.to_owned(),
|
||||
state_json: "{}".to_owned(),
|
||||
state_json: geth_document::create_automerge_state("{}")?,
|
||||
updated_at_ms: 0,
|
||||
};
|
||||
store.insert_document_resource(&stored)?;
|
||||
|
|
@ -7832,17 +7835,27 @@ fn document_resource_from_stored(stored: &StoredDocumentResource) -> DocumentRes
|
|||
id: stored.document_id.clone().into(),
|
||||
resource: stored.resource_id.clone().into(),
|
||||
name: stored.name.clone(),
|
||||
sync_status: "local-only".to_owned(),
|
||||
state_bytes: stored.state_json.len() as u64,
|
||||
sync_status: "automerge-local".to_owned(),
|
||||
state_bytes: geth_document::document_state_bytes(&stored.state_json),
|
||||
}
|
||||
}
|
||||
|
||||
fn document_state_from_stored(stored: &StoredDocumentResource) -> DocumentState {
|
||||
DocumentState {
|
||||
fn document_state_from_stored(stored: &StoredDocumentResource) -> Result<DocumentState, NodeError> {
|
||||
Ok(DocumentState {
|
||||
document: document_resource_from_stored(stored),
|
||||
state_json: geth_document::document_view_json(&stored.state_json)?,
|
||||
updated_at: UnixMillis(stored.updated_at_ms),
|
||||
})
|
||||
}
|
||||
|
||||
fn document_sync_state_from_stored(
|
||||
stored: &StoredDocumentResource,
|
||||
) -> Result<DocumentState, NodeError> {
|
||||
Ok(DocumentState {
|
||||
document: document_resource_from_stored(stored),
|
||||
state_json: stored.state_json.clone(),
|
||||
updated_at: UnixMillis(stored.updated_at_ms),
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
fn file_root_from_stored(stored: &StoredFileRoot) -> FileRoot {
|
||||
|
|
@ -11528,7 +11541,7 @@ mod tests {
|
|||
assert!(allowed);
|
||||
assert!(updated);
|
||||
assert!(reason.contains("direct grant"));
|
||||
assert!(note.contains("JSON LWW"));
|
||||
assert!(note.contains("durable Automerge"));
|
||||
}
|
||||
other => panic!("unexpected allowed document sync response: {other:?}"),
|
||||
}
|
||||
|
|
@ -11804,7 +11817,10 @@ mod tests {
|
|||
left_store
|
||||
.get_document_resource_by_name("notes")
|
||||
.expect("get live-synced document")
|
||||
.map(|document| document.state_json),
|
||||
.map(|document| {
|
||||
geth_document::document_view_json(&document.state_json)
|
||||
.expect("document view json")
|
||||
}),
|
||||
Some(r#"{"title":"live"}"#.to_owned())
|
||||
);
|
||||
let live_remote_root = left_store
|
||||
|
|
|
|||
|
|
@ -2889,8 +2889,8 @@ fn document_create_and_status_use_local_store() {
|
|||
match response {
|
||||
geth_control::ControlResponse::DocumentCreated { document } => {
|
||||
assert_eq!(document.name, "notes");
|
||||
assert_eq!(document.sync_status, "local-only");
|
||||
assert_eq!(document.state_bytes, 2);
|
||||
assert_eq!(document.sync_status, "automerge-local");
|
||||
assert!(document.state_bytes > 2);
|
||||
}
|
||||
other => panic!("unexpected response: {other:?}"),
|
||||
}
|
||||
|
|
@ -2906,7 +2906,7 @@ fn document_create_and_status_use_local_store() {
|
|||
geth_control::ControlResponse::DocumentStatus { document } => {
|
||||
assert_eq!(document.id.to_string(), "document:notes");
|
||||
assert_eq!(document.resource.to_string(), "resource:document:notes");
|
||||
assert_eq!(document.state_bytes, 2);
|
||||
assert!(document.state_bytes > 2);
|
||||
}
|
||||
other => panic!("unexpected response: {other:?}"),
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue