From 65d1c71ba7fa075aab248ae5b4c429c3cdfc3dd7 Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 30 Sep 2026 12:26:20 +0300 Subject: [PATCH] knot-xrpc,interop: drain birth commit frames before #syncy-syncy Lewis: May this revision serve well! --- knot2/crates/knot-xrpc/tests/atproto.rs | 59 ++++++++++++++++--------- knot2/interop/atproto_firehose_test.go | 12 ++++- 2 files changed, 49 insertions(+), 22 deletions(-) diff --git a/knot2/crates/knot-xrpc/tests/atproto.rs b/knot2/crates/knot-xrpc/tests/atproto.rs index 06eafbd50..216cdd572 100644 --- a/knot2/crates/knot-xrpc/tests/atproto.rs +++ b/knot2/crates/knot-xrpc/tests/atproto.rs @@ -6,7 +6,9 @@ use std::time::Duration; use axum::body::Body; use cid::Cid; -use common::{FIXTURE_COLLECTION as COLLECTION, KNOT_HOST, OWNER, World, percent_encode, post_authed}; +use common::{ + FIXTURE_COLLECTION as COLLECTION, KNOT_HOST, OWNER, World, percent_encode, post_authed, +}; use futures::StreamExt; use http::{HeaderValue, StatusCode, header}; use jacquard_repo::Mst; @@ -462,25 +464,42 @@ async fn account_creation_announces_identity_account_and_an_empty_commit() { identity_seq < account_seq && account_seq < payload.seq, "creation goes identity, account, commit" ); - let genesis_rev = payload.rev.clone(); - let commit_seq = payload.seq; - let sync = next_binary(&mut ws).await; - let (header, rest) = frame_header(&sync); - assert_eq!(header["t"], "#sync"); - let payload: SyncProbe = frame_payload(rest); - assert!( - commit_seq < payload.seq, - "the sync assertion follows the commit it names" - ); - assert_eq!(payload.did, repo_did); - assert_eq!( - payload.rev, genesis_rev, - "the sync frame asserts the genesis commit just announced" - ); - assert!( - !payload.blocks.is_empty(), - "the sync frame carries the genesis commit in its CAR" - ); + let mut last_seq = payload.seq; + let mut last_rev = payload.rev; + loop { + let frame = next_binary(&mut ws).await; + let (header, rest) = frame_header(&frame); + match header["t"].as_str() { + Some("#commit") => { + let payload: CommitProbe = frame_payload(rest); + assert_eq!(payload.repo, repo_did); + assert!( + last_seq < payload.seq, + "creation's commits stay in seq order" + ); + last_seq = payload.seq; + last_rev = payload.rev; + } + Some("#sync") => { + let payload: SyncProbe = frame_payload(rest); + assert!( + last_seq < payload.seq, + "the sync assertion follows the commit it names" + ); + assert_eq!(payload.did, repo_did); + assert_eq!( + payload.rev, last_rev, + "the sync frame asserts the last commit creation announced" + ); + assert!( + !payload.blocks.is_empty(), + "the sync frame carries that commit in its CAR" + ); + break; + } + other => panic!("creation announced an unexpected frame: {other:?}"), + } + } let (status, body) = post_authed( &world, diff --git a/knot2/interop/atproto_firehose_test.go b/knot2/interop/atproto_firehose_test.go index a09c1fe6e..ed45df92e 100644 --- a/knot2/interop/atproto_firehose_test.go +++ b/knot2/interop/atproto_firehose_test.go @@ -182,10 +182,18 @@ func TestKnotServesSubscribeRepos(t *testing.T) { assert.NotEmpty(t, genesis.Blocks, "the commit frame includes its CAR") assert.NotEmpty(t, genesis.Commit.String()) + lastRev := genesis.Rev + lastSeq := genesis.Seq sync := nextFrame(t, conn) - require.NotNil(t, sync.Sync) + for sync.Sync == nil { + require.NotNil(t, sync.Commit, "creation only announces commits before the sync assertion") + assert.Equal(t, repoDid, sync.Commit.Repo) + assert.Less(t, lastSeq, sync.Commit.Seq, "creation's commits stay in seq order") + lastRev, lastSeq = sync.Commit.Rev, sync.Commit.Seq + sync = nextFrame(t, conn) + } assert.Equal(t, repoDid, sync.Sync.Did) - assert.Equal(t, genesis.Rev, sync.Sync.Rev, "the sync frame asserts the genesis commit") + assert.Equal(t, lastRev, sync.Sync.Rev, "the sync frame asserts the last commit creation announced") assert.NotEmpty(t, sync.Sync.Blocks, "the sync frame includes the commit block") assert.Less(t, identity.Identity.Seq, account.Account.Seq, "identity precedes account") -- 2.51.2