feat(kanban): authorize canonical sync reads

req: sync/007

req: auth/005

req: security/004
This commit is contained in:
slhx agent
2026-07-13 19:22:13 +02:00
parent 3ef84e7331
commit 1ec8a1368d
3 changed files with 211 additions and 20 deletions
+84 -17
View File
@@ -145,6 +145,8 @@ struct SyncAcknowledgement {
card_id: u64,
canonical_column: CanonicalColumn,
status: &'static str,
#[serde(skip_serializing)]
tenant: String,
}
#[derive(Serialize)]
@@ -213,6 +215,8 @@ struct PersistedAcknowledgement {
server_sequence: u64,
card_id: u64,
canonical_column: CanonicalColumn,
#[serde(default)]
tenant: Option<String>,
}
impl SyncStore {
@@ -227,7 +231,7 @@ impl SyncStore {
fs::read(&self.0).map_err(|error| format!("read {}: {error}", self.0.display()))?;
let persisted: PersistedSync = serde_json::from_slice(&bytes)
.map_err(|error| format!("decode {}: {error}", self.0.display()))?;
if persisted.schema_version != 1 || persisted.next_sequence == 0 {
if !matches!(persisted.schema_version, 1 | 2) || persisted.next_sequence == 0 {
return Err(format!("unsupported sync store {}", self.0.display()));
}
let mut acknowledgements = BTreeMap::new();
@@ -237,12 +241,23 @@ impl SyncStore {
if stored.server_sequence == 0 || stored.card_id == 0 {
return Err(format!("invalid acknowledgement in {}", self.0.display()));
}
let tenant = if persisted.schema_version == 1 {
"demo".to_owned()
} else {
stored
.tenant
.filter(|tenant| matches!(tenant.as_str(), "demo" | "alpha" | "beta"))
.ok_or_else(|| {
format!("invalid acknowledgement tenant in {}", self.0.display())
})?
};
let acknowledgement = SyncAcknowledgement {
command_id: stored.command_id,
server_sequence: stored.server_sequence,
card_id: stored.card_id,
canonical_column: stored.canonical_column,
status: "accepted",
tenant,
};
if acknowledgements
.insert(command_id, acknowledgement)
@@ -271,8 +286,8 @@ impl SyncStore {
acknowledgements.push(PersistedAcknowledgement::from(acknowledgement));
acknowledgements.sort_by_key(|item| item.server_sequence);
let persisted = PersistedSync {
schema_version: 1,
next_sequence: acknowledgement.server_sequence + 1,
schema_version: 2,
next_sequence: sync.next_sequence + 1,
acknowledgements,
};
let bytes = serde_json::to_vec_pretty(&persisted)
@@ -309,6 +324,7 @@ impl From<&SyncAcknowledgement> for PersistedAcknowledgement {
server_sequence: value.server_sequence,
card_id: value.card_id,
canonical_column: value.canonical_column,
tenant: Some(value.tenant.clone()),
}
}
}
@@ -603,6 +619,19 @@ async fn sync_context(
}))
}
fn visible_card(principal: CurrentSyncPrincipal, card_id: u64) -> bool {
principal.tenant == "demo"
|| (principal.tenant == "beta" && card_id == 2)
|| (principal.tenant == "alpha" && card_id != 2)
}
fn visible_acknowledgement(
principal: CurrentSyncPrincipal,
acknowledgement: &SyncAcknowledgement,
) -> bool {
principal.tenant == "demo" || acknowledgement.tenant == principal.tenant
}
fn authorize_sync_replay(
principal: CurrentSyncPrincipal,
card_id: u64,
@@ -642,6 +671,11 @@ async fn sync_command(
let mut sync = state.sync.lock().unwrap();
if let Some(existing) = sync.acknowledgements.get(&command_id) {
if !visible_acknowledgement(principal, existing) {
return Err(SyncRejection::Forbidden(
"current tenant cannot access command acknowledgement",
));
}
if existing.card_id != card_id || existing.canonical_column != canonical_column {
return Err(SyncRejection::Conflict(
"command_id was already used for a different payload",
@@ -667,12 +701,21 @@ async fn sync_command(
.iter()
.position(|card| card.id == card_id)
.ok_or(SyncRejection::BadRequest("unknown card_id"))?;
let server_sequence = sync
.acknowledgements
.values()
.filter(|acknowledgement| visible_acknowledgement(principal, acknowledgement))
.map(|acknowledgement| acknowledgement.server_sequence)
.max()
.unwrap_or_default()
+ 1;
let acknowledgement = SyncAcknowledgement {
command_id: command_id.0.clone(),
server_sequence: sync.next_sequence,
server_sequence,
card_id,
canonical_column,
status: "accepted",
tenant: principal.tenant.to_owned(),
};
if let Some(store) = &state.sync_store {
store.persist(&sync, &acknowledgement).map_err(|error| {
@@ -687,23 +730,35 @@ async fn sync_command(
Ok(Json(acknowledgement))
}
// req: sync/007 req: sync/020
async fn sync_snapshot(State(state): State<Arc<AppState>>) -> Json<SyncSnapshot> {
// req: sync/007 req: sync/020 req: auth/005 req: security/004
async fn sync_snapshot(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
) -> Result<Json<SyncSnapshot>, SyncRejection> {
let principal = current_sync_principal(&headers, &state.sync_sessions)?;
let board = state.board.lock().unwrap();
let sync = state.sync.lock().unwrap();
let cards = board
.cards
.iter()
.filter(|card| visible_card(principal, card.id))
.map(|card| SnapshotCard {
id: card.id,
column: canonical_column(card.column),
})
.collect();
Json(SyncSnapshot {
let server_sequence = sync
.acknowledgements
.values()
.filter(|acknowledgement| visible_acknowledgement(principal, acknowledgement))
.map(|acknowledgement| acknowledgement.server_sequence)
.max()
.unwrap_or_default();
Ok(Json(SyncSnapshot {
schema_version: 1,
server_sequence: sync.next_sequence.saturating_sub(1),
server_sequence,
cards,
})
}))
}
fn canonical_column(column: usize) -> CanonicalColumn {
@@ -714,11 +769,13 @@ fn canonical_column(column: usize) -> CanonicalColumn {
}
}
// req: sync/005 req: sync/006 req: sync/007 req: sync/013
// req: sync/005 req: sync/006 req: sync/007 req: sync/013 req: auth/005 req: security/004
async fn sync_acknowledgements(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
Query(params): Query<BTreeMap<String, String>>,
) -> Sse<impl futures_util::Stream<Item = Result<Event, Infallible>>> {
) -> Result<Sse<impl futures_util::Stream<Item = Result<Event, Infallible>>>, SyncRejection> {
let principal = current_sync_principal(&headers, &state.sync_sessions)?;
let after = params
.get("after")
.and_then(|value| value.parse::<u64>().ok())
@@ -733,22 +790,31 @@ async fn sync_acknowledgements(
let attempts = sync.reconnects.entry(key).or_default();
*attempts += 1;
if *attempts == 1 {
return Sse::new(
return Ok(Sse::new(
stream::iter([Ok(Event::default()
.comment("reconnect")
.retry(Duration::from_millis(25)))])
.boxed(),
)
.keep_alive(KeepAlive::default());
.keep_alive(KeepAlive::default()));
}
}
let first_available = sync
.acknowledgements
.values()
.filter(|acknowledgement| acknowledgement.server_sequence > sync.retained_after)
.filter(|acknowledgement| {
visible_acknowledgement(principal, acknowledgement)
&& acknowledgement.server_sequence > sync.retained_after
})
.map(|acknowledgement| acknowledgement.server_sequence)
.min();
let latest = sync.next_sequence.saturating_sub(1);
let latest = sync
.acknowledgements
.values()
.filter(|acknowledgement| visible_acknowledgement(principal, acknowledgement))
.map(|acknowledgement| acknowledgement.server_sequence)
.max()
.unwrap_or_default();
let history_missing = after < latest
&& first_available.is_none_or(|first_sequence| after.saturating_add(1) < first_sequence);
let events = if history_missing {
@@ -766,7 +832,8 @@ async fn sync_acknowledgements(
sync.acknowledgements
.values()
.filter(|acknowledgement| {
acknowledgement.server_sequence > after
visible_acknowledgement(principal, acknowledgement)
&& acknowledgement.server_sequence > after
&& acknowledgement.server_sequence > sync.retained_after
})
.map(|acknowledgement| {
@@ -778,7 +845,7 @@ async fn sync_acknowledgements(
})
.collect::<Vec<Result<Event, Infallible>>>()
};
Sse::new(stream::iter(events).boxed()).keep_alive(KeepAlive::default())
Ok(Sse::new(stream::iter(events).boxed()).keep_alive(KeepAlive::default()))
}
fn registry(state: Arc<AppState>) -> impl DispatchRegistry {
+125 -1
View File
@@ -598,6 +598,129 @@ async fn account_partition_hides_replay_and_export_until_owner_returns() -> WebD
result.and(quit)
}
#[tokio::test]
async fn canonical_snapshot_and_history_are_tenant_scoped() -> WebDriverResult<()> {
// test req: sync/007 req: security/004 req: auth/005
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)
.env(
"HEMX_KANBAN_SESSION_ALICE_ALPHA_EDITOR",
"test-token-alice-alpha-editor",
)
.env(
"HEMX_KANBAN_SESSION_BOB_ALPHA_VIEWER",
"test-token-bob-alpha-viewer",
)
.env(
"HEMX_KANBAN_SESSION_CAROL_BETA_EDITOR",
"test-token-carol-beta-editor",
);
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready 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?;
let proof = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const session = (token) => { document.cookie = `hemx_kanban_session=${token}; Path=/; SameSite=Strict`; };
const command = (id, card) => fetch(`/sync/commands?command_id=${encodeURIComponent(id)}&card_id=${card}&column=done`, { method: 'POST' });
(async () => {
session('test-token-alice-alpha-editor');
const aliceWrite = await command('alice-history:1', 1);
session('test-token-carol-beta-editor');
const carolWrite = await command('carol-history:1', 2);
const betaSnapshotResponse = await fetch('/sync/snapshot', { cache: 'no-store' });
const betaSnapshot = await betaSnapshotResponse.json();
const betaHistoryResponse = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
const betaHistory = await betaHistoryResponse.text();
session('test-token-bob-alpha-viewer');
const alphaSnapshotResponse = await fetch('/sync/snapshot', { cache: 'no-store' });
const alphaSnapshot = await alphaSnapshotResponse.json();
const alphaHistoryResponse = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
const alphaHistory = await alphaHistoryResponse.text();
document.cookie = 'hemx_kanban_session=; Path=/; Max-Age=0; SameSite=Strict';
const signedOutSnapshot = await fetch('/sync/snapshot', { cache: 'no-store' });
const signedOutHistory = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
done({
aliceWrite: aliceWrite.status,
carolWrite: carolWrite.status,
betaSnapshotStatus: betaSnapshotResponse.status,
betaSnapshot,
betaHistoryStatus: betaHistoryResponse.status,
betaHistory,
alphaSnapshotStatus: alphaSnapshotResponse.status,
alphaSnapshot,
alphaHistoryStatus: alphaHistoryResponse.status,
alphaHistory,
signedOutSnapshotStatus: signedOutSnapshot.status,
signedOutHistoryStatus: signedOutHistory.status,
});
})().catch((error) => done({ error: String(error), stack: error.stack }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert!(proof.get("error").is_none(), "browser proof failed: {proof}");
assert_eq!(proof["aliceWrite"], 200);
assert_eq!(proof["carolWrite"], 200);
assert_eq!(proof["betaSnapshotStatus"], 200);
assert_eq!(
proof["betaSnapshot"],
serde_json::json!({
"schemaVersion": 1,
"serverSequence": 1,
"cards": [{ "id": 2, "column": "done" }]
})
);
assert_eq!(proof["betaHistoryStatus"], 200);
let beta_history = proof["betaHistory"].as_str().unwrap();
assert!(beta_history.contains("carol-history:1"));
assert!(!beta_history.contains("alice-history:1"));
assert!(!beta_history.contains("\"cardId\":1"));
assert_eq!(proof["alphaSnapshotStatus"], 200);
assert_eq!(
proof["alphaSnapshot"],
serde_json::json!({
"schemaVersion": 1,
"serverSequence": 1,
"cards": [
{ "id": 1, "column": "done" },
{ "id": 3, "column": "done" }
]
})
);
assert_eq!(proof["alphaHistoryStatus"], 200);
let alpha_history = proof["alphaHistory"].as_str().unwrap();
assert!(alpha_history.contains("alice-history:1"));
assert!(!alpha_history.contains("carol-history:1"));
assert!(!alpha_history.contains("\"cardId\":2"));
assert_eq!(proof["signedOutSnapshotStatus"], 401);
assert_eq!(proof["signedOutHistoryStatus"], 401);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn schema_upgrade_preserves_queued_order_and_local_intent() -> WebDriverResult<()> {
// test req: sync/004 req: sync/010 req: sync/014
@@ -1592,7 +1715,8 @@ async fn canonical_acknowledgement_survives_server_restart() -> WebDriverResult<
assert!(store.is_file(), "server did not materialize sync store");
let persisted = fs::read_to_string(&store).expect("read sync store");
assert!(persisted.contains("restart-proof:1"));
assert!(persisted.contains("\"schemaVersion\": 1"));
assert!(persisted.contains("\"schemaVersion\": 2"));
assert!(persisted.contains("\"tenant\": \"demo\""));
drop(first_app);
let mut second_app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));