feat(sync): add typed acknowledgement effect

req: sync/006
This commit is contained in:
slhx agent
2026-07-13 23:15:57 +02:00
parent 88598f181f
commit 5a82956c08
6 changed files with 163 additions and 7 deletions
+35 -1
View File
@@ -15,7 +15,7 @@ use hemx_axum::{
use hemx_kanban_example::ui::board::{self as board};
use hemx_kanban_example::ui::board_card as card_board;
use hemx_kanban_example::ui::{self, board as board_ui};
use hemx_sync::{Channel, PresenceScope, PresenceTracker, PresenceUpdate};
use hemx_sync::{Channel, PresenceScope, PresenceTracker, PresenceUpdate, SyncEffect};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::convert::Infallible;
@@ -502,6 +502,7 @@ async fn main() {
.route("/", get(home).post(interact))
.route("/events", get(events))
.route("/sync/broadcast", get(sync_broadcast))
.route("/sync/ack", get(sync_ack))
.route("/sync-demo", get(sync_demo))
.route("/sync.js", get(sync_js))
.route("/sync/context", get(sync_context))
@@ -601,6 +602,39 @@ fn presence_changed(signal: PresenceSignal) -> impl hemx::IntoEffect {
})
}
async fn sync_ack(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
Query(params): Query<BTreeMap<String, String>>,
) -> Result<Response, SyncRejection> {
let principal = current_sync_principal(&headers, &state.sync_sessions)?;
let command_id = CommandId::parse(params.get("command_id"))?;
let acknowledgement = {
let sync = state.sync.lock().unwrap();
let acknowledgement =
sync.acknowledgements
.get(&command_id)
.ok_or(SyncRejection::Conflict(
"command has no canonical acknowledgement",
))?;
if !visible_acknowledgement(principal, acknowledgement) {
return Err(SyncRejection::Forbidden(
"current tenant cannot access command acknowledgement",
));
}
acknowledgement.clone()
};
let batch = (
SyncEffect::ack(board::atoms::sync_ack),
board::sync_status.text(format!(
"Canonical acknowledgement {} at server sequence {}",
acknowledgement.command_id, acknowledgement.server_sequence
)),
)
.into_batch(ui::BUILD_FINGERPRINT);
Ok(sse(stream::iter([Ok::<_, Infallible>(batch)]).boxed()).into_response())
}
async fn sync_broadcast(
State(state): State<Arc<AppState>>,
Query(params): Query<BTreeMap<String, String>>,
+2
View File
@@ -11,5 +11,7 @@
<div data-hemx-slot="board">{+= self.board =+}</div>
<aside data-hemx-slot="presence">Waiting for presence…</aside>
<output id="sync-ack" data-hemx-atom="sync_ack" aria-live="polite">pending</output>
<output data-hemx-slot="sync_status" aria-live="polite">Waiting for acknowledgement…</output>
</section>
+91
View File
@@ -2200,6 +2200,97 @@ async fn adversarial_wire_inputs_are_rejected_before_partial_application() -> We
result.and(quit)
}
#[tokio::test]
async fn canonical_acknowledgement_updates_generated_atom_over_ordinary_batch(
) -> WebDriverResult<()> {
// test req: sync/006
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command.env("HEMX_KANBAN_ADDR", &app_addr);
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
wait_until(
&driver,
"return document.querySelector('#sync-ack')?.textContent.trim() === 'pending'",
)
.await?;
let proof = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
(async () => {
const commandId = 'ack-proof:1';
const accepted = await fetch(`/sync/commands?command_id=${encodeURIComponent(commandId)}&card_id=1&column=done`, { method: 'POST' });
const canonical = await accepted.json();
const root = document.querySelector('[data-hemx-root]');
let acknowledgementEvent;
root.addEventListener('hemx:sync-ack', (event) => { acknowledgementEvent = event.detail; }, { once: true });
const source = new EventSource(`/sync/ack?command_id=${encodeURIComponent(commandId)}`);
await new Promise((resolve, reject) => {
const timeout = setTimeout(() => {
source.close();
reject(new Error('ack batch timed out'));
}, 5000);
source.addEventListener('hemx', (event) => {
clearTimeout(timeout);
const normalized = event.data.replace(/-/g, '+').replace(/_/g, '/');
const padded = normalized + '='.repeat((4 - normalized.length % 4) % 4);
const raw = atob(padded);
const bytes = Uint8Array.from(raw, (character) => character.charCodeAt(0));
window.hemx.applyBatch(bytes.buffer, root);
source.close();
resolve();
});
source.onerror = () => {
clearTimeout(timeout);
source.close();
reject(new Error('ack batch failed'));
};
});
done({
acceptedStatus: accepted.status,
canonical,
atom: document.querySelector('#sync-ack').textContent.trim(),
status: document.body.textContent,
acknowledgementEvent,
});
})().catch((error) => done({ error: String(error), stack: error?.stack }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert!(proof["error"].is_null(), "typed acknowledgement failed: {proof}");
assert_eq!(proof["acceptedStatus"], 200);
assert_eq!(proof["canonical"]["commandId"], "ack-proof:1");
assert_eq!(proof["canonical"]["serverSequence"], 1);
assert_eq!(proof["atom"], "acknowledged");
assert!(proof["status"]
.as_str()
.is_some_and(|status| status.contains("ack-proof:1 at server sequence 1")));
assert!(proof["acknowledgementEvent"].as_str().is_some_and(|payload| payload.contains("atomId")));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn typed_presence_join_leave_updates_generated_atom_over_sse() -> WebDriverResult<()> {
// test req: sync/001 req: sync/004 req: sync/005