From 05c2b27d598af63c282b23495a822ce0a63fda68 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 00:49:39 -0500 Subject: [PATCH 01/35] ci: capture the retained Share recovery fault on hosted runners --- .github/workflows/rc2-share-fault.yml | 67 +++++++++++++++++++++++++++ 1 file changed, 67 insertions(+) create mode 100644 .github/workflows/rc2-share-fault.yml diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml new file mode 100644 index 0000000..65163d8 --- /dev/null +++ b/.github/workflows/rc2-share-fault.yml @@ -0,0 +1,67 @@ +name: RC2 hosted Share fault diagnostics + +on: + pull_request: + +permissions: + contents: read + +jobs: + share-fault: + runs-on: ubuntu-latest + timeout-minutes: 30 + steps: + - uses: actions/checkout@v4 + with: + path: testlab + persist-credentials: false + - uses: actions/checkout@v4 + with: + repository: kafkars/kafkars + ref: 6eb8e644da164eaddf546e1090fb2011fa775e20 + path: kafkars-candidate + persist-credentials: false + - name: Instrument only the disposable candidate and select unchanged diagnostics + run: | + python3 - <<'PY' + from pathlib import Path + source = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer/share/registry_delivery.rs") + code = source.read_text() + changes = [ + ('return Err(ShareConsumerDeliveryError::MembershipFault);', 'eprintln!("RC2_SHARE_MEMBERSHIP fault={:?} fatal={:?}", entry.fault, entry.membership.as_ref().and_then(|membership| membership.machine().fatal()));\n return Err(ShareConsumerDeliveryError::MembershipFault);'), + ('return Err(ShareConsumerDeliveryError::FetchFault);', 'eprintln!("RC2_SHARE_FETCH routing={:?} session={:?}", entry.fetch().fault(), entry.fetch().session_fault());\n return Err(ShareConsumerDeliveryError::FetchFault);'), + ] + for before, after in changes: + assert code.count(before) == 1 + code = code.replace(before, after) + source.write_text(code) + Path("testlab/packs/rc2-share-fault.toml").write_text('''schema_version = 1 + id = "rc2-share-fault" + title = "RC2 Share fault diagnostics, not release evidence" + scenarios = ["scenarios/kafka/share-group-session-recovery.toml"] + ''') + Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 + id = "rc2-share-fault" + title = "RC2 Share fault diagnostics, not release evidence" + [[cells]] + id = "apache-kafka-4-3-1-three-plaintext" + environment = "clusters/apache-kafka/4.3.1/three-plaintext.toml" + pack = "packs/rc2-share-fault.toml" + attempts = 6 + gating = true + ''') + PY + git -C kafkars-candidate diff > instrumented-candidate.patch + - uses: ./testlab + with: + kafkars-path: kafkars-candidate + allow-dirty: "true" + evidence-directory: rc2-share-fault-evidence + - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a + if: ${{ always() }} + with: + name: rc2-share-fault-evidence + path: | + rc2-share-fault-evidence + instrumented-candidate.patch + if-no-files-found: error From eb0a39464127379e1a36aaf7301d1ac33d72a46d Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 00:51:33 -0500 Subject: [PATCH 02/35] ci: expose bounded adapter stderr in temporary Share diagnostics --- .github/workflows/rc2-share-fault.yml | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index 65163d8..ee64dde 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -35,6 +35,11 @@ jobs: assert code.count(before) == 1 code = code.replace(before, after) source.write_text(code) + reader = Path("testlab/crates/testctl/src/process_io.rs") + code = reader.read_text() + before = " let _ = sender.send(value);" + assert code.count(before) == 1 + reader.write_text(code.replace(before, ' eprintln!("{value}");\n' + before)) Path("testlab/packs/rc2-share-fault.toml").write_text('''schema_version = 1 id = "rc2-share-fault" title = "RC2 Share fault diagnostics, not release evidence" @@ -52,6 +57,7 @@ jobs: ''') PY git -C kafkars-candidate diff > instrumented-candidate.patch + git -C testlab diff > instrumented-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate @@ -64,4 +70,5 @@ jobs: path: | rc2-share-fault-evidence instrumented-candidate.patch + instrumented-harness.patch if-no-files-found: error From cf14134586b028de1cdeb12d6223d3bb4bb6c00f Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 01:05:35 -0500 Subject: [PATCH 03/35] ci: capture exact Share partition rejection and heartbeat retry facts --- .github/workflows/rc2-share-fault.yml | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index ee64dde..6d546ca 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -35,6 +35,23 @@ jobs: assert code.count(before) == 1 code = code.replace(before, after) source.write_text(code) + base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer/share") + patches = [ + ("fetch_session_set/execution.rs", '.map_err(|error| ShareFetchExecutionError::Settlement(error.kind()))?', '.map_err(|error| { eprintln!("RC2_SHARE_SETTLEMENT {error:?}"); ShareFetchExecutionError::Settlement(error.kind()) })?'), + ("fetch_session_settlement/settlement.rs", 'ShareFetchResolution::Succeeded(success) => {', 'ShareFetchResolution::Succeeded(success) => {\n for topic in &success.topics { for partition in &topic.partitions { if let Some(rejection) = &partition.rejection { eprintln!("RC2_SHARE_PARTITION {rejection:?}"); } } }'), + ("settlement/failure.rs", 'let prepared = self.prepared.ok_or(ShareMembershipError::EffectShape)?;', 'let prepared = self.prepared.ok_or(ShareMembershipError::EffectShape)?;\n eprintln!("RC2_SHARE_HEARTBEAT now={now:?} prepared={prepared:?} failure={failure:?}");'), + ("transition/retry.rs", 'let transition = self\n .machine\n .apply(ShareGroupHeartbeatInput::RetryDue { schedule, now })', 'eprintln!("RC2_SHARE_RETRY now={now:?} schedule={schedule:?}");\n let transition = self\n .machine\n .apply(ShareGroupHeartbeatInput::RetryDue { schedule, now })'), + ] + for relative, before, after in patches: + source = base / relative + code = source.read_text() + if relative == "settlement/failure.rs": + assert code.count(before) == 2 + code = code.replace(before, after, 1) + else: + assert code.count(before) == 1, relative + code = code.replace(before, after) + source.write_text(code) reader = Path("testlab/crates/testctl/src/process_io.rs") code = reader.read_text() before = " let _ = sender.send(value);" From 8d53ec1835ace9268f26c23ca3e900740764b603 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 01:18:45 -0500 Subject: [PATCH 04/35] ci: verify route-less Share expiry recovery against broker faults --- .github/workflows/rc2-share-fault.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index 6d546ca..d112d8c 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -18,7 +18,7 @@ jobs: - uses: actions/checkout@v4 with: repository: kafkars/kafkars - ref: 6eb8e644da164eaddf546e1090fb2011fa775e20 + ref: e15a33cc9c6057d19e8715f483fd0b8d4b38233e path: kafkars-candidate persist-credentials: false - name: Instrument only the disposable candidate and select unchanged diagnostics From 5bf4bfee62fbc02030edee59d462c804c20e8df3 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 01:32:37 -0500 Subject: [PATCH 05/35] ci: qualify Share leader-epoch recovery against broker faults --- .github/workflows/rc2-share-fault.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index d112d8c..22f3c74 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -18,7 +18,7 @@ jobs: - uses: actions/checkout@v4 with: repository: kafkars/kafkars - ref: e15a33cc9c6057d19e8715f483fd0b8d4b38233e + ref: 7e54c80c654f4228f2e84f3b583a6671d671c70a path: kafkars-candidate persist-credentials: false - name: Instrument only the disposable candidate and select unchanged diagnostics From 07f3975947480d08b4e01b163175fd26ce33e066 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 01:46:17 -0500 Subject: [PATCH 06/35] ci: identify exact Share invalidation failure during recovery --- .github/workflows/rc2-share-fault.yml | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index 22f3c74..00d74e9 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -52,6 +52,16 @@ jobs: assert code.count(before) == 1, relative code = code.replace(before, after) source.write_text(code) + source = Path("kafkars-candidate/crates/kafka-client-engine/src/driver/rpc/share_group_heartbeat/invalidation_drive.rs") + code = source.read_text() + changes = [ + (' let result = match result {', ' eprintln!("RC2_SHARE_INVALIDATION_TERMINAL group={group_id:?} result={result:?}");\n let result = match result {'), + ('let (source, token) = rejected.into_parts();', 'let (source, token) = rejected.into_parts();\n eprintln!("RC2_SHARE_INVALIDATION_SUBMIT {source:?}");'), + ] + for before, after in changes: + assert code.count(before) == 1 + code = code.replace(before, after) + source.write_text(code) reader = Path("testlab/crates/testctl/src/process_io.rs") code = reader.read_text() before = " let _ = sender.send(value);" From beaa39bcf8a397252b3c6397c79fb45c17a0b583 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 01:51:49 -0500 Subject: [PATCH 07/35] ci: compare bounded Share invalidation recovery with control --- .github/workflows/rc2-share-fault.yml | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml index 00d74e9..54f295f 100644 --- a/.github/workflows/rc2-share-fault.yml +++ b/.github/workflows/rc2-share-fault.yml @@ -10,6 +10,13 @@ jobs: share-fault: runs-on: ubuntu-latest timeout-minutes: 30 + strategy: + fail-fast: false + max-parallel: 2 + matrix: + candidate: + - 7e54c80c654f4228f2e84f3b583a6671d671c70a + - 192ed9c11ba4382a744440fb3b2df51b54e4d8bf steps: - uses: actions/checkout@v4 with: @@ -18,7 +25,7 @@ jobs: - uses: actions/checkout@v4 with: repository: kafkars/kafkars - ref: 7e54c80c654f4228f2e84f3b583a6671d671c70a + ref: ${{ matrix.candidate }} path: kafkars-candidate persist-credentials: false - name: Instrument only the disposable candidate and select unchanged diagnostics @@ -55,7 +62,7 @@ jobs: source = Path("kafkars-candidate/crates/kafka-client-engine/src/driver/rpc/share_group_heartbeat/invalidation_drive.rs") code = source.read_text() changes = [ - (' let result = match result {', ' eprintln!("RC2_SHARE_INVALIDATION_TERMINAL group={group_id:?} result={result:?}");\n let result = match result {'), + ('return Ok(terminal(group_id, result));', 'eprintln!("RC2_SHARE_INVALIDATION_TERMINAL group={group_id:?} result={result:?}");\n return Ok(terminal(group_id, result));'), ('let (source, token) = rejected.into_parts();', 'let (source, token) = rejected.into_parts();\n eprintln!("RC2_SHARE_INVALIDATION_SUBMIT {source:?}");'), ] for before, after in changes: @@ -93,7 +100,7 @@ jobs: - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-share-fault-evidence + name: rc2-share-fault-evidence-${{ matrix.candidate }} path: | rc2-share-fault-evidence instrumented-candidate.patch From b9bab1419aefad57fc29d26083e50861cf75984b Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 02:05:11 -0500 Subject: [PATCH 08/35] ci: retire archived RC2 Share diagnostics --- .github/workflows/rc2-share-fault.yml | 108 -------------------------- 1 file changed, 108 deletions(-) delete mode 100644 .github/workflows/rc2-share-fault.yml diff --git a/.github/workflows/rc2-share-fault.yml b/.github/workflows/rc2-share-fault.yml deleted file mode 100644 index 54f295f..0000000 --- a/.github/workflows/rc2-share-fault.yml +++ /dev/null @@ -1,108 +0,0 @@ -name: RC2 hosted Share fault diagnostics - -on: - pull_request: - -permissions: - contents: read - -jobs: - share-fault: - runs-on: ubuntu-latest - timeout-minutes: 30 - strategy: - fail-fast: false - max-parallel: 2 - matrix: - candidate: - - 7e54c80c654f4228f2e84f3b583a6671d671c70a - - 192ed9c11ba4382a744440fb3b2df51b54e4d8bf - steps: - - uses: actions/checkout@v4 - with: - path: testlab - persist-credentials: false - - uses: actions/checkout@v4 - with: - repository: kafkars/kafkars - ref: ${{ matrix.candidate }} - path: kafkars-candidate - persist-credentials: false - - name: Instrument only the disposable candidate and select unchanged diagnostics - run: | - python3 - <<'PY' - from pathlib import Path - source = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer/share/registry_delivery.rs") - code = source.read_text() - changes = [ - ('return Err(ShareConsumerDeliveryError::MembershipFault);', 'eprintln!("RC2_SHARE_MEMBERSHIP fault={:?} fatal={:?}", entry.fault, entry.membership.as_ref().and_then(|membership| membership.machine().fatal()));\n return Err(ShareConsumerDeliveryError::MembershipFault);'), - ('return Err(ShareConsumerDeliveryError::FetchFault);', 'eprintln!("RC2_SHARE_FETCH routing={:?} session={:?}", entry.fetch().fault(), entry.fetch().session_fault());\n return Err(ShareConsumerDeliveryError::FetchFault);'), - ] - for before, after in changes: - assert code.count(before) == 1 - code = code.replace(before, after) - source.write_text(code) - base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer/share") - patches = [ - ("fetch_session_set/execution.rs", '.map_err(|error| ShareFetchExecutionError::Settlement(error.kind()))?', '.map_err(|error| { eprintln!("RC2_SHARE_SETTLEMENT {error:?}"); ShareFetchExecutionError::Settlement(error.kind()) })?'), - ("fetch_session_settlement/settlement.rs", 'ShareFetchResolution::Succeeded(success) => {', 'ShareFetchResolution::Succeeded(success) => {\n for topic in &success.topics { for partition in &topic.partitions { if let Some(rejection) = &partition.rejection { eprintln!("RC2_SHARE_PARTITION {rejection:?}"); } } }'), - ("settlement/failure.rs", 'let prepared = self.prepared.ok_or(ShareMembershipError::EffectShape)?;', 'let prepared = self.prepared.ok_or(ShareMembershipError::EffectShape)?;\n eprintln!("RC2_SHARE_HEARTBEAT now={now:?} prepared={prepared:?} failure={failure:?}");'), - ("transition/retry.rs", 'let transition = self\n .machine\n .apply(ShareGroupHeartbeatInput::RetryDue { schedule, now })', 'eprintln!("RC2_SHARE_RETRY now={now:?} schedule={schedule:?}");\n let transition = self\n .machine\n .apply(ShareGroupHeartbeatInput::RetryDue { schedule, now })'), - ] - for relative, before, after in patches: - source = base / relative - code = source.read_text() - if relative == "settlement/failure.rs": - assert code.count(before) == 2 - code = code.replace(before, after, 1) - else: - assert code.count(before) == 1, relative - code = code.replace(before, after) - source.write_text(code) - source = Path("kafkars-candidate/crates/kafka-client-engine/src/driver/rpc/share_group_heartbeat/invalidation_drive.rs") - code = source.read_text() - changes = [ - ('return Ok(terminal(group_id, result));', 'eprintln!("RC2_SHARE_INVALIDATION_TERMINAL group={group_id:?} result={result:?}");\n return Ok(terminal(group_id, result));'), - ('let (source, token) = rejected.into_parts();', 'let (source, token) = rejected.into_parts();\n eprintln!("RC2_SHARE_INVALIDATION_SUBMIT {source:?}");'), - ] - for before, after in changes: - assert code.count(before) == 1 - code = code.replace(before, after) - source.write_text(code) - reader = Path("testlab/crates/testctl/src/process_io.rs") - code = reader.read_text() - before = " let _ = sender.send(value);" - assert code.count(before) == 1 - reader.write_text(code.replace(before, ' eprintln!("{value}");\n' + before)) - Path("testlab/packs/rc2-share-fault.toml").write_text('''schema_version = 1 - id = "rc2-share-fault" - title = "RC2 Share fault diagnostics, not release evidence" - scenarios = ["scenarios/kafka/share-group-session-recovery.toml"] - ''') - Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 - id = "rc2-share-fault" - title = "RC2 Share fault diagnostics, not release evidence" - [[cells]] - id = "apache-kafka-4-3-1-three-plaintext" - environment = "clusters/apache-kafka/4.3.1/three-plaintext.toml" - pack = "packs/rc2-share-fault.toml" - attempts = 6 - gating = true - ''') - PY - git -C kafkars-candidate diff > instrumented-candidate.patch - git -C testlab diff > instrumented-harness.patch - - uses: ./testlab - with: - kafkars-path: kafkars-candidate - allow-dirty: "true" - evidence-directory: rc2-share-fault-evidence - - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a - if: ${{ always() }} - with: - name: rc2-share-fault-evidence-${{ matrix.candidate }} - path: | - rc2-share-fault-evidence - instrumented-candidate.patch - instrumented-harness.patch - if-no-files-found: error From c877c96200db733be46c91e067517cb32ae8dd41 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 04:22:39 -0500 Subject: [PATCH 09/35] ci: diagnose RC2 group recovery failures --- .github/workflows/rc2-group-fault.yml | 106 ++++++++++++++++++++++++++ 1 file changed, 106 insertions(+) create mode 100644 .github/workflows/rc2-group-fault.yml diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml new file mode 100644 index 0000000..1093edd --- /dev/null +++ b/.github/workflows/rc2-group-fault.yml @@ -0,0 +1,106 @@ +name: RC2 hosted group recovery diagnostics + +on: + pull_request: + +permissions: + contents: read + +jobs: + group-fault: + runs-on: ubuntu-latest + timeout-minutes: 30 + steps: + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 + with: + path: testlab + persist-credentials: false + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 + with: + repository: kafkars/kafkars + ref: d4f7eaf1a68e62bd4623bd3644f7c8a285924321 + path: kafkars-candidate + persist-credentials: false + - name: Instrument only disposable copies and retain unchanged scenarios + run: | + python3 - <<'PY' + from pathlib import Path + base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer") + def patch(relative, before, after): + source = base / relative + code = source.read_text() + assert code.count(before) == 1, relative + source.write_text(code.replace(before, after)) + patch("group/registry_graceful_revocation.rs", + ".map_err(GroupConsumerRevocationPortError::Acknowledge)", + '''.map_err(|error| { + static CALLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if CALLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_GROUP_ACK group={group_id:?} public={assignment_epoch} error={error:?}"); + } + GroupConsumerRevocationPortError::Acknowledge(error) + })''') + patch("group/registry_graceful_revocation.rs", + "return Err(GroupConsumerRevocationPortError::GroupUnavailable);", + '''static UNAVAILABLE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if UNAVAILABLE.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_GROUP_UNAVAILABLE group={group_id:?} public={assignment_epoch} active={} fault={}", entry.state == GroupConsumerEntryState::Active, entry.fault.is_some()); + } + return Err(GroupConsumerRevocationPortError::GroupUnavailable);''') + patch("group/classic_group_fetch/delivery.rs", + ") -> Result, ClassicGroupFetchDeliveryError> {", + ''') -> Result, ClassicGroupFetchDeliveryError> { + static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_GROUP_FETCH machine={:?} effects={} pending={} sessions={}", self.machine, self.effects.len(), self.pending_fetches.len(), self.fetches.retained_broker_sessions()); + }''') + patch("fetch_execution/broker_session_begin.rs", + "self.entries[index].in_flight = true;\n Ok(BrokerSessionPlan {", + '''static PLANS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if PLANS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_GROUP_SESSION broker={broker_id:?} session={:?} active={active:?} forgotten={forgotten:?}", self.entries[index].metadata); + } + self.entries[index].in_flight = true; + Ok(BrokerSessionPlan {''') + reader = Path("testlab/crates/testctl/src/process_io.rs") + code = reader.read_text() + for before, after in [ + ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), + ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), + (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), + ]: + assert code.count(before) == 1 + code = code.replace(before, after) + reader.write_text(code) + Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 + id = "rc2-group-fault" + title = "RC2 group diagnostics, not release evidence" + scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml"] + ''') + Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 + id = "rc2-group-fault" + title = "RC2 group diagnostics, not release evidence" + [[cells]] + id = "apache-kafka-4-3-1-three-sasl-plain" + environment = "clusters/apache-kafka/4.3.1/three-sasl-plain.toml" + pack = "packs/rc2-group-fault.toml" + attempts = 3 + gating = true + ''') + PY + git -C kafkars-candidate diff > instrumented-candidate.patch + git -C testlab diff > instrumented-harness.patch + - uses: ./testlab + with: + kafkars-path: kafkars-candidate + allow-dirty: "true" + evidence-directory: rc2-group-fault-evidence + - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a + if: ${{ always() }} + with: + name: rc2-group-fault-evidence + path: | + rc2-group-fault-evidence + instrumented-candidate.patch + instrumented-harness.patch + if-no-files-found: error From 24c343074f615412343871ac7f7a127ed8a4ec87 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 04:36:17 -0500 Subject: [PATCH 10/35] ci: trace RC2 TLS membership loss --- .github/workflows/rc2-group-fault.yml | 54 ++++++++++----------------- 1 file changed, 20 insertions(+), 34 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 1093edd..7127e1e 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -25,43 +25,29 @@ jobs: run: | python3 - <<'PY' from pathlib import Path - base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer") + base = Path("kafkars-candidate/crates") def patch(relative, before, after): source = base / relative code = source.read_text() assert code.count(before) == 1, relative source.write_text(code.replace(before, after)) - patch("group/registry_graceful_revocation.rs", - ".map_err(GroupConsumerRevocationPortError::Acknowledge)", - '''.map_err(|error| { - static CALLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if CALLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_GROUP_ACK group={group_id:?} public={assignment_epoch} error={error:?}"); - } - GroupConsumerRevocationPortError::Acknowledge(error) - })''') - patch("group/registry_graceful_revocation.rs", - "return Err(GroupConsumerRevocationPortError::GroupUnavailable);", - '''static UNAVAILABLE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if UNAVAILABLE.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_GROUP_UNAVAILABLE group={group_id:?} public={assignment_epoch} active={} fault={}", entry.state == GroupConsumerEntryState::Active, entry.fault.is_some()); - } - return Err(GroupConsumerRevocationPortError::GroupUnavailable);''') - patch("group/classic_group_fetch/delivery.rs", - ") -> Result, ClassicGroupFetchDeliveryError> {", - ''') -> Result, ClassicGroupFetchDeliveryError> { - static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_GROUP_FETCH machine={:?} effects={} pending={} sessions={}", self.machine, self.effects.len(), self.pending_fetches.len(), self.fetches.retained_broker_sessions()); + patch("kafka-client-core/src/consumer/consumer_group/transition_support.rs", + " let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);", + ''' eprintln!("RC2_GROUP_FATAL failure={failure:?} machine={self:?}"); + let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);''') + patch("kafka-client-engine/src/consumer/group/consumer_group_heartbeat_settlement.rs", + " let (resolution, route) = outcome.into_resolution();", + ''' let (resolution, route) = outcome.into_resolution(); + match &resolution { + ConsumerGroupHeartbeatResolution::Failed(failure) => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} failure={failure:?}"), + ConsumerGroupHeartbeatResolution::BrokerRejected { error_code, .. } => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} broker={error_code}"), + _ => {}, }''') - patch("fetch_execution/broker_session_begin.rs", - "self.entries[index].in_flight = true;\n Ok(BrokerSessionPlan {", - '''static PLANS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if PLANS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_GROUP_SESSION broker={broker_id:?} session={:?} active={active:?} forgotten={forgotten:?}", self.entries[index].metadata); - } - self.entries[index].in_flight = true; - Ok(BrokerSessionPlan {''') + patch("kafka-client-engine/src/driver/rpc/consumer_group_heartbeat_failure.rs", + ") -> ConsumerGroupHeartbeatDriverFailureKind {\n #[allow(", + ''') -> ConsumerGroupHeartbeatDriverFailureKind { + eprintln!("RC2_GROUP_DRIVER error={error:?}"); + #[allow(''') reader = Path("testlab/crates/testctl/src/process_io.rs") code = reader.read_text() for before, after in [ @@ -75,14 +61,14 @@ jobs: Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" title = "RC2 group diagnostics, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml"] + scenarios = ["scenarios/kafka/consumer-protocol-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 id = "rc2-group-fault" title = "RC2 group diagnostics, not release evidence" [[cells]] - id = "apache-kafka-4-3-1-three-sasl-plain" - environment = "clusters/apache-kafka/4.3.1/three-sasl-plain.toml" + id = "apache-kafka-4-3-1-three-tls" + environment = "clusters/apache-kafka/4.3.1/three-tls.toml" pack = "packs/rc2-group-fault.toml" attempts = 3 gating = true From 5c78f6f8c37d24df159b3cc0e340455cfba621de Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 05:01:31 -0500 Subject: [PATCH 11/35] ci: verify clean RC2 group recovery fixes --- .github/workflows/rc2-group-fault.yml | 76 +++++++++------------------ 1 file changed, 25 insertions(+), 51 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 7127e1e..64060d1 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted group recovery diagnostics +name: RC2 hosted group regression checks on: pull_request: @@ -9,7 +9,12 @@ permissions: jobs: group-fault: runs-on: ubuntu-latest - timeout-minutes: 30 + timeout-minutes: 40 + strategy: + fail-fast: false + max-parallel: 2 + matrix: + environment: [three-sasl-plain, three-tls] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: @@ -18,75 +23,44 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: d4f7eaf1a68e62bd4623bd3644f7c8a285924321 + ref: c93e2094de4d70d3af9e8f7cc8bd99d5f7904fee path: kafkars-candidate persist-credentials: false - - name: Instrument only disposable copies and retain unchanged scenarios + - name: Select unchanged scenarios against the clean candidate + env: + TARGET_CELL: ${{ matrix.environment }} run: | python3 - <<'PY' + import os from pathlib import Path - base = Path("kafkars-candidate/crates") - def patch(relative, before, after): - source = base / relative - code = source.read_text() - assert code.count(before) == 1, relative - source.write_text(code.replace(before, after)) - patch("kafka-client-core/src/consumer/consumer_group/transition_support.rs", - " let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);", - ''' eprintln!("RC2_GROUP_FATAL failure={failure:?} machine={self:?}"); - let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);''') - patch("kafka-client-engine/src/consumer/group/consumer_group_heartbeat_settlement.rs", - " let (resolution, route) = outcome.into_resolution();", - ''' let (resolution, route) = outcome.into_resolution(); - match &resolution { - ConsumerGroupHeartbeatResolution::Failed(failure) => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} failure={failure:?}"), - ConsumerGroupHeartbeatResolution::BrokerRejected { error_code, .. } => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} broker={error_code}"), - _ => {}, - }''') - patch("kafka-client-engine/src/driver/rpc/consumer_group_heartbeat_failure.rs", - ") -> ConsumerGroupHeartbeatDriverFailureKind {\n #[allow(", - ''') -> ConsumerGroupHeartbeatDriverFailureKind { - eprintln!("RC2_GROUP_DRIVER error={error:?}"); - #[allow(''') - reader = Path("testlab/crates/testctl/src/process_io.rs") - code = reader.read_text() - for before, after in [ - ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), - ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), - (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), - ]: - assert code.count(before) == 1 - code = code.replace(before, after) - reader.write_text(code) + environment = os.environ["TARGET_CELL"] + assert environment in ("three-sasl-plain", "three-tls") Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 group diagnostics, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-session-recovery.toml"] + title = "RC2 focused group checks, not release evidence" + scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] ''') - Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 + Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 group diagnostics, not release evidence" + title = "RC2 focused group checks, not release evidence" [[cells]] - id = "apache-kafka-4-3-1-three-tls" - environment = "clusters/apache-kafka/4.3.1/three-tls.toml" + id = "apache-kafka-4-3-1-{environment}" + environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" attempts = 3 gating = true ''') PY - git -C kafkars-candidate diff > instrumented-candidate.patch - git -C testlab diff > instrumented-harness.patch + git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - allow-dirty: "true" - evidence-directory: rc2-group-fault-evidence + evidence-directory: rc2-group-check-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-group-fault-evidence + name: rc2-group-check-${{ matrix.environment }} path: | - rc2-group-fault-evidence - instrumented-candidate.patch - instrumented-harness.patch + rc2-group-check-evidence + focused-harness.patch if-no-files-found: error From 21a2458306370fc32a45aa9f897e1accc66531c9 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 05:25:23 -0500 Subject: [PATCH 12/35] ci: diagnose the remaining RC2 checkpoint rejection --- .github/workflows/rc2-group-fault.yml | 84 +++++++++++++++++++-------- 1 file changed, 59 insertions(+), 25 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 64060d1..4fc3567 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted group regression checks +name: RC2 hosted checkpoint recovery diagnostics on: pull_request: @@ -10,11 +10,6 @@ jobs: group-fault: runs-on: ubuntu-latest timeout-minutes: 40 - strategy: - fail-fast: false - max-parallel: 2 - matrix: - environment: [three-sasl-plain, three-tls] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: @@ -23,44 +18,83 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: c93e2094de4d70d3af9e8f7cc8bd99d5f7904fee + ref: c87c5a0eeebba4e676662e2c26ca35bb76970dea path: kafkars-candidate persist-credentials: false - - name: Select unchanged scenarios against the clean candidate - env: - TARGET_CELL: ${{ matrix.environment }} + - name: Diagnose disposable copies with the unchanged recovery scenario run: | python3 - <<'PY' - import os from pathlib import Path - environment = os.environ["TARGET_CELL"] - assert environment in ("three-sasl-plain", "three-tls") + base = Path("kafkars-candidate/crates") + def patch(relative, before, after): + source = base / relative + code = source.read_text() + assert code.count(before) == 1, relative + source.write_text(code.replace(before, after)) + patch("kafka-client-core/src/consumer/consumer_group/transition_support.rs", + " let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);", + ''' eprintln!("RC2_GROUP_FATAL failure={failure:?} machine={self:?}"); + let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);''') + patch("kafka-client-engine/src/consumer/group/consumer_group_heartbeat_settlement.rs", + " let (resolution, route) = outcome.into_resolution();", + ''' let (resolution, route) = outcome.into_resolution(); + match &resolution { + ConsumerGroupHeartbeatResolution::Failed(failure) => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} failure={failure:?}"), + ConsumerGroupHeartbeatResolution::BrokerRejected { error_code, .. } => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} broker={error_code}"), + _ => {}, + }''') + patch("kafka-client-engine/src/driver/rpc/consumer_group_heartbeat_failure.rs", + ") -> ConsumerGroupHeartbeatDriverFailureKind {\n #[allow(", + ''') -> ConsumerGroupHeartbeatDriverFailureKind { + eprintln!("RC2_GROUP_DRIVER error={error:?}"); + #[allow(''') + patch("kafka-client-engine/src/consumer/group/registry_commit.rs", + " let offset_commits = &mut self.offset_commits;", + ''' eprintln!("RC2_COMMIT_ADMIT checkpoint={checkpoint:?} live={:?}", entry.catalog.live_assignment()); + let offset_commits = &mut self.offset_commits;''') + patch("kafka-client-engine/src/consumer/group/registry_commit.rs", + "fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure {", + '''fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure { + eprintln!("RC2_COMMIT_REJECT kind={:?}", failure.kind);''') + reader = Path("testlab/crates/testctl/src/process_io.rs") + code = reader.read_text() + for before, after in [ + ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), + ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), + (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), + ]: + assert code.count(before) == 1 + code = code.replace(before, after) + reader.write_text(code) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 focused group checks, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] + title = "RC2 checkpoint diagnostics, not release evidence" + scenarios = ["scenarios/kafka/consumer-protocol-group-session-recovery.toml"] ''') - Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 + Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 id = "rc2-group-fault" - title = "RC2 focused group checks, not release evidence" + title = "RC2 checkpoint diagnostics, not release evidence" [[cells]] - id = "apache-kafka-4-3-1-{environment}" - environment = "clusters/apache-kafka/4.3.1/{environment}.toml" + id = "apache-kafka-4-3-1-three-sasl-plain" + environment = "clusters/apache-kafka/4.3.1/three-sasl-plain.toml" pack = "packs/rc2-group-fault.toml" - attempts = 3 + attempts = 5 gating = true ''') PY - git -C testlab diff > focused-harness.patch + git -C kafkars-candidate diff > instrumented-candidate.patch + git -C testlab diff > instrumented-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - evidence-directory: rc2-group-check-evidence + allow-dirty: "true" + evidence-directory: rc2-checkpoint-diagnostic-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-group-check-${{ matrix.environment }} + name: rc2-checkpoint-diagnostic-evidence path: | - rc2-group-check-evidence - focused-harness.patch + rc2-checkpoint-diagnostic-evidence + instrumented-candidate.patch + instrumented-harness.patch if-no-files-found: error From 9d30cdbfb1f484780f7098b0ad676dc4769ad03d Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 05:41:24 -0500 Subject: [PATCH 13/35] ci: trace coordinator routing behind RC2 recovery failures --- .github/workflows/rc2-group-fault.yml | 77 +++++++++++++++++++++++++-- 1 file changed, 72 insertions(+), 5 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 4fc3567..2237fd6 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted checkpoint recovery diagnostics +name: RC2 hosted coordinator routing diagnostics on: pull_request: @@ -21,9 +21,16 @@ jobs: ref: c87c5a0eeebba4e676662e2c26ca35bb76970dea path: kafkars-candidate persist-credentials: false + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 + with: + repository: kafkars/kafka-driver + ref: 66c33e2f9de8dc787b8ccd2b4dd3eced39e7c6a4 + path: kafka-driver-diagnostic + persist-credentials: false - name: Diagnose disposable copies with the unchanged recovery scenario run: | python3 - <<'PY' + import json from pathlib import Path base = Path("kafkars-candidate/crates") def patch(relative, before, after): @@ -50,12 +57,69 @@ jobs: #[allow(''') patch("kafka-client-engine/src/consumer/group/registry_commit.rs", " let offset_commits = &mut self.offset_commits;", - ''' eprintln!("RC2_COMMIT_ADMIT checkpoint={checkpoint:?} live={:?}", entry.catalog.live_assignment()); + ''' if entry.catalog.live_assignment().is_none() { + eprintln!("RC2_COMMIT_NO_ASSIGNMENT checkpoint={checkpoint:?}"); + } let offset_commits = &mut self.offset_commits;''') patch("kafka-client-engine/src/consumer/group/registry_commit.rs", "fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure {", '''fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure { eprintln!("RC2_COMMIT_REJECT kind={:?}", failure.kind);''') + base = Path("kafka-driver-diagnostic") + patch("src/reactor/host/coordinator_routing.rs", + " let directory = self.metadata.as_ref()?.current()?.brokers();\n let route = directory.route_to(coordinator.broker_id())?;", + ''' let Some(snapshot) = self.metadata.as_ref().and_then(|metadata| metadata.current()) else { + eprintln!("RC2_COORD_ROUTE no_metadata"); + return None; + }; + let directory = snapshot.brokers(); + let Some(route) = directory.route_to(coordinator.broker_id()) else { + eprintln!("RC2_COORD_ROUTE missing_broker={:?}", coordinator.broker_id()); + return None; + };''') + patch("src/reactor/host/coordinator_routing.rs", + " (entry.endpoint() == coordinator.endpoint()).then_some(route)", + ''' if entry.endpoint() != coordinator.endpoint() { + eprintln!("RC2_COORD_ROUTE endpoint_mismatch broker={:?}", coordinator.broker_id()); + } + (entry.endpoint() == coordinator.endpoint()).then_some(route)''') + patch("src/reactor/host/coordinator_routing.rs", + " let Some(seed) = seed else {\n request.fail(RequestError::RouteUnavailable);", + ''' let Some(seed) = seed else { + eprintln!("RC2_COORD_ROUTE missing_seed"); + request.fail(RequestError::RouteUnavailable);''') + patch("src/reactor/coordinator/routing.rs", + " waiting.request.fail(RequestError::RouteUnavailable);", + ''' eprintln!("RC2_COORD_WAIT no_route key={:?} state={:?}", waiting.key, self.entry(&waiting.key).map(|entry| entry.machine.state())); + waiting.request.fail(RequestError::RouteUnavailable);''') + patch("src/reactor/coordinator/drive.rs", + " let input = match result {", + ''' if let Ok(Err(error)) = &result { + eprintln!("RC2_COORD_DISCOVERY error={error:?}"); + } + let input = match result {''') + manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") + code = manifest.read_text() + support_patches = "".join( + f'{name} = {{ path = {json.dumps(str(path.resolve()))} }}\n' + for name, path in [ + ("kafka-driver", base), + ("kafka-driver-core", base / "crates/kafka-driver-core"), + ("kafka-driver-transport", base / "crates/kafka-driver-transport"), + ] + ) + before = " let support_patches = support_patches(artifacts, &source)?;" + assert code.count(before) == 1 + code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") + before = 'format!("kafkars-{version}-{}", &digest[..16])' + assert code.count(before) == 1 + code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') + manifest.write_text(code) + Path("DIAGNOSTIC-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: the adapter overrides all three driver packages with the " + "logged disposable RC4 checkout. Subject archive hashes describe the original " + "packages, not the instrumented driver runtime. See instrumented-driver.patch.\n" + ) reader = Path("testlab/crates/testctl/src/process_io.rs") code = reader.read_text() for before, after in [ @@ -83,18 +147,21 @@ jobs: ''') PY git -C kafkars-candidate diff > instrumented-candidate.patch + git -C kafka-driver-diagnostic diff > instrumented-driver.patch git -C testlab diff > instrumented-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate allow-dirty: "true" - evidence-directory: rc2-checkpoint-diagnostic-evidence + evidence-directory: rc2-routing-diagnostic-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-checkpoint-diagnostic-evidence + name: rc2-routing-diagnostic-evidence path: | - rc2-checkpoint-diagnostic-evidence + rc2-routing-diagnostic-evidence + DIAGNOSTIC-ONLY.txt instrumented-candidate.patch + instrumented-driver.patch instrumented-harness.patch if-no-files-found: error From 5484865ad43620389f43cadc05be9a34e8a30dff Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 05:47:18 -0500 Subject: [PATCH 14/35] ci: verify explicit driver source in hosted routing diagnostics --- .github/workflows/rc2-group-fault.yml | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 2237fd6..62bd8b2 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -115,6 +115,20 @@ jobs: assert code.count(before) == 1 code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') manifest.write_text(code) + provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") + code = provenance.read_text() + before = " let actual = locked_registry_artifact(&lock, name)?;" + assert code.count(before) == 1 + code = code.replace(before, ''' if matches!(name, "kafka-driver" | "kafka-driver-core" | "kafka-driver-transport") { + let packages: Vec<_> = lock.package.iter().filter(|package| package.name == name).collect(); + if packages.len() != 1 || packages[0].version != expected.version + || packages[0].source.is_some() || packages[0].checksum.is_some() { + return Err(candidate(format!("diagnostic driver source lock mismatch: {name}"))); + } + continue; + } + let actual = locked_registry_artifact(&lock, name)?;''') + provenance.write_text(code) Path("DIAGNOSTIC-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: the adapter overrides all three driver packages with the " "logged disposable RC4 checkout. Subject archive hashes describe the original " From b6f35c689ac89a7683ccf1f06852bc88209528a8 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 06:06:31 -0500 Subject: [PATCH 15/35] ci: verify coordinator directory repair against clean source --- .github/workflows/rc2-group-fault.yml | 147 +++++++------------------- 1 file changed, 40 insertions(+), 107 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 62bd8b2..db64367 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted coordinator routing diagnostics +name: RC2 hosted coordinator recovery source checks on: pull_request: @@ -10,6 +10,11 @@ jobs: group-fault: runs-on: ubuntu-latest timeout-minutes: 40 + strategy: + fail-fast: false + max-parallel: 2 + matrix: + environment: [three-sasl-plain, three-tls] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: @@ -24,88 +29,28 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafka-driver - ref: 66c33e2f9de8dc787b8ccd2b4dd3eced39e7c6a4 - path: kafka-driver-diagnostic + ref: 0a8bc088b93269ff6217737651784dd69e833063 + path: kafka-driver-candidate persist-credentials: false - - name: Diagnose disposable copies with the unchanged recovery scenario + - name: Select unchanged scenarios against clean fixed source + env: + TARGET_CELL: ${{ matrix.environment }} run: | python3 - <<'PY' import json + import os from pathlib import Path - base = Path("kafkars-candidate/crates") - def patch(relative, before, after): - source = base / relative - code = source.read_text() - assert code.count(before) == 1, relative - source.write_text(code.replace(before, after)) - patch("kafka-client-core/src/consumer/consumer_group/transition_support.rs", - " let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);", - ''' eprintln!("RC2_GROUP_FATAL failure={failure:?} machine={self:?}"); - let fatal = ConsumerGroupHeartbeatFatal::new(attempt, failure);''') - patch("kafka-client-engine/src/consumer/group/consumer_group_heartbeat_settlement.rs", - " let (resolution, route) = outcome.into_resolution();", - ''' let (resolution, route) = outcome.into_resolution(); - match &resolution { - ConsumerGroupHeartbeatResolution::Failed(failure) => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} failure={failure:?}"), - ConsumerGroupHeartbeatResolution::BrokerRejected { error_code, .. } => eprintln!("RC2_GROUP_HEARTBEAT kind={kind:?} broker={error_code}"), - _ => {}, - }''') - patch("kafka-client-engine/src/driver/rpc/consumer_group_heartbeat_failure.rs", - ") -> ConsumerGroupHeartbeatDriverFailureKind {\n #[allow(", - ''') -> ConsumerGroupHeartbeatDriverFailureKind { - eprintln!("RC2_GROUP_DRIVER error={error:?}"); - #[allow(''') - patch("kafka-client-engine/src/consumer/group/registry_commit.rs", - " let offset_commits = &mut self.offset_commits;", - ''' if entry.catalog.live_assignment().is_none() { - eprintln!("RC2_COMMIT_NO_ASSIGNMENT checkpoint={checkpoint:?}"); - } - let offset_commits = &mut self.offset_commits;''') - patch("kafka-client-engine/src/consumer/group/registry_commit.rs", - "fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure {", - '''fn delegated_failure(failure: GroupOffsetCommitAdmissionFailure) -> GroupConsumerCommitFailure { - eprintln!("RC2_COMMIT_REJECT kind={:?}", failure.kind);''') - base = Path("kafka-driver-diagnostic") - patch("src/reactor/host/coordinator_routing.rs", - " let directory = self.metadata.as_ref()?.current()?.brokers();\n let route = directory.route_to(coordinator.broker_id())?;", - ''' let Some(snapshot) = self.metadata.as_ref().and_then(|metadata| metadata.current()) else { - eprintln!("RC2_COORD_ROUTE no_metadata"); - return None; - }; - let directory = snapshot.brokers(); - let Some(route) = directory.route_to(coordinator.broker_id()) else { - eprintln!("RC2_COORD_ROUTE missing_broker={:?}", coordinator.broker_id()); - return None; - };''') - patch("src/reactor/host/coordinator_routing.rs", - " (entry.endpoint() == coordinator.endpoint()).then_some(route)", - ''' if entry.endpoint() != coordinator.endpoint() { - eprintln!("RC2_COORD_ROUTE endpoint_mismatch broker={:?}", coordinator.broker_id()); - } - (entry.endpoint() == coordinator.endpoint()).then_some(route)''') - patch("src/reactor/host/coordinator_routing.rs", - " let Some(seed) = seed else {\n request.fail(RequestError::RouteUnavailable);", - ''' let Some(seed) = seed else { - eprintln!("RC2_COORD_ROUTE missing_seed"); - request.fail(RequestError::RouteUnavailable);''') - patch("src/reactor/coordinator/routing.rs", - " waiting.request.fail(RequestError::RouteUnavailable);", - ''' eprintln!("RC2_COORD_WAIT no_route key={:?} state={:?}", waiting.key, self.entry(&waiting.key).map(|entry| entry.machine.state())); - waiting.request.fail(RequestError::RouteUnavailable);''') - patch("src/reactor/coordinator/drive.rs", - " let input = match result {", - ''' if let Ok(Err(error)) = &result { - eprintln!("RC2_COORD_DISCOVERY error={error:?}"); - } - let input = match result {''') + environment = os.environ["TARGET_CELL"] + assert environment in ("three-sasl-plain", "three-tls") + driver = Path("kafka-driver-candidate") manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") code = manifest.read_text() support_patches = "".join( f'{name} = {{ path = {json.dumps(str(path.resolve()))} }}\n' for name, path in [ - ("kafka-driver", base), - ("kafka-driver-core", base / "crates/kafka-driver-core"), - ("kafka-driver-transport", base / "crates/kafka-driver-transport"), + ("kafka-driver", driver), + ("kafka-driver-core", driver / "crates/kafka-driver-core"), + ("kafka-driver-transport", driver / "crates/kafka-driver-transport"), ] ) before = " let support_patches = support_patches(artifacts, &source)?;" @@ -113,7 +58,7 @@ jobs: code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") before = 'format!("kafkars-{version}-{}", &digest[..16])' assert code.count(before) == 1 - code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') + code = code.replace(before, 'format!("source-check-kafkars-{version}-{}", &digest[..16])') manifest.write_text(code) provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") code = provenance.read_text() @@ -123,59 +68,47 @@ jobs: let packages: Vec<_> = lock.package.iter().filter(|package| package.name == name).collect(); if packages.len() != 1 || packages[0].version != expected.version || packages[0].source.is_some() || packages[0].checksum.is_some() { - return Err(candidate(format!("diagnostic driver source lock mismatch: {name}"))); + return Err(candidate(format!("source-check driver lock mismatch: {name}"))); } continue; } let actual = locked_registry_artifact(&lock, name)?;''') provenance.write_text(code) - Path("DIAGNOSTIC-ONLY.txt").write_text( - "NOT RELEASE EVIDENCE: the adapter overrides all three driver packages with the " - "logged disposable RC4 checkout. Subject archive hashes describe the original " - "packages, not the instrumented driver runtime. See instrumented-driver.patch.\n" + Path("SOURCE-CHECK-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: clean native c87c5a0e plus clean unpublished driver " + "0a8bc088b93269ff6217737651784dd69e833063. The three driver crates use " + "explicit source paths. Subject archive hashes for driver RC4 do not qualify " + "the fixed runtime. No native or driver instrumentation is applied.\n" ) - reader = Path("testlab/crates/testctl/src/process_io.rs") - code = reader.read_text() - for before, after in [ - ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), - ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), - (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), - ]: - assert code.count(before) == 1 - code = code.replace(before, after) - reader.write_text(code) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 checkpoint diagnostics, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-session-recovery.toml"] + title = "RC2 focused source checks, not release evidence" + scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] ''') - Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 + Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 checkpoint diagnostics, not release evidence" + title = "RC2 focused source checks, not release evidence" [[cells]] - id = "apache-kafka-4-3-1-three-sasl-plain" - environment = "clusters/apache-kafka/4.3.1/three-sasl-plain.toml" + id = "apache-kafka-4-3-1-{environment}" + environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" - attempts = 5 + attempts = 3 gating = true ''') PY - git -C kafkars-candidate diff > instrumented-candidate.patch - git -C kafka-driver-diagnostic diff > instrumented-driver.patch - git -C testlab diff > instrumented-harness.patch + git -C kafkars-candidate diff --exit-code + git -C kafka-driver-candidate diff --exit-code + git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - allow-dirty: "true" - evidence-directory: rc2-routing-diagnostic-evidence + evidence-directory: rc2-source-check-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-routing-diagnostic-evidence + name: rc2-source-check-${{ matrix.environment }} path: | - rc2-routing-diagnostic-evidence - DIAGNOSTIC-ONLY.txt - instrumented-candidate.patch - instrumented-driver.patch - instrumented-harness.patch + rc2-source-check-evidence + SOURCE-CHECK-ONLY.txt + focused-harness.patch if-no-files-found: error From fd136741e08da58c5442090b340a2b1bb548f25d Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 06:29:46 -0500 Subject: [PATCH 16/35] ci: trace classic revocation and fetch ownership --- .github/workflows/rc2-group-fault.yml | 81 ++++++++++++++++++++------- 1 file changed, 62 insertions(+), 19 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index db64367..80fa1c7 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted coordinator recovery source checks +name: RC2 hosted classic recovery ownership diagnostics on: pull_request: @@ -14,7 +14,7 @@ jobs: fail-fast: false max-parallel: 2 matrix: - environment: [three-sasl-plain, three-tls] + environment: [three-sasl-plain] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: @@ -29,10 +29,10 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafka-driver - ref: 0a8bc088b93269ff6217737651784dd69e833063 + ref: e5b697468eadbd99c27af48cb3653373e4599001 path: kafka-driver-candidate persist-credentials: false - - name: Select unchanged scenarios against clean fixed source + - name: Trace classic recovery without changing scenario behavior env: TARGET_CELL: ${{ matrix.environment }} run: | @@ -41,7 +41,47 @@ jobs: import os from pathlib import Path environment = os.environ["TARGET_CELL"] - assert environment in ("three-sasl-plain", "three-tls") + assert environment == "three-sasl-plain" + base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer") + def patch(relative, before, after): + source = base / relative + code = source.read_text() + assert code.count(before) == 1, relative + source.write_text(code.replace(before, after)) + patch("group/registry_delivery.rs", + " if !entry.revocation.is_dormant() {", + ''' if !entry.revocation.is_dormant() { + static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_CLASSIC_RECV_BLOCKED group={group_id:?} at={:?} deadline={:?}", std::time::Instant::now(), entry.revocation.next_deadline()); + }''') + patch("group/registry_graceful_revocation.rs", + " if entry.revocation.expire_if_due(now)? {", + ''' if entry.revocation.expire_if_due(now)? { + eprintln!("RC2_CLASSIC_REVOCATION_EXPIRED group={:?} now={now:?}", entry.group_id());''') + patch("group/registry_graceful_revocation.rs", + ".map_err(GroupConsumerRevocationPortError::Acknowledge)", + '''.map_err(|error| { + eprintln!("RC2_CLASSIC_ACK group={group_id:?} public={assignment_epoch} error={error:?}"); + GroupConsumerRevocationPortError::Acknowledge(error) + })''') + patch("group/classic_group_fetch/delivery.rs", + ") -> Result, ClassicGroupFetchDeliveryError> {", + ''') -> Result, ClassicGroupFetchDeliveryError> { + static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_CLASSIC_FETCH at={:?} machine={:?} effects={} pending={} sessions={}", std::time::Instant::now(), self.machine, self.effects.len(), self.pending_fetches.len(), self.fetches.retained_broker_sessions()); + }''') + reader = Path("testlab/crates/testctl/src/process_io.rs") + code = reader.read_text() + for before, after in [ + ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), + ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), + (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), + ]: + assert code.count(before) == 1 + code = code.replace(before, after) + reader.write_text(code) driver = Path("kafka-driver-candidate") manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") code = manifest.read_text() @@ -58,7 +98,7 @@ jobs: code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") before = 'format!("kafkars-{version}-{}", &digest[..16])' assert code.count(before) == 1 - code = code.replace(before, 'format!("source-check-kafkars-{version}-{}", &digest[..16])') + code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') manifest.write_text(code) provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") code = provenance.read_text() @@ -74,41 +114,44 @@ jobs: } let actual = locked_registry_artifact(&lock, name)?;''') provenance.write_text(code) - Path("SOURCE-CHECK-ONLY.txt").write_text( - "NOT RELEASE EVIDENCE: clean native c87c5a0e plus clean unpublished driver " - "0a8bc088b93269ff6217737651784dd69e833063. The three driver crates use " + Path("DIAGNOSTIC-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: instrumented native c87c5a0e plus clean unpublished driver " + "e5b697468eadbd99c27af48cb3653373e4599001. The three driver crates use " "explicit source paths. Subject archive hashes for driver RC4 do not qualify " - "the fixed runtime. No native or driver instrumentation is applied.\n" + "the fixed runtime. Native stderr traces diagnose revocation and Fetch ownership; " + "scenario and adapter receive behavior are unchanged.\n" ) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 focused source checks, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] + title = "RC2 classic recovery diagnostics, not release evidence" + scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 focused source checks, not release evidence" + title = "RC2 classic recovery diagnostics, not release evidence" [[cells]] id = "apache-kafka-4-3-1-{environment}" environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" - attempts = 3 + attempts = 5 gating = true ''') PY - git -C kafkars-candidate diff --exit-code + git -C kafkars-candidate diff > instrumented-candidate.patch git -C kafka-driver-candidate diff --exit-code git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - evidence-directory: rc2-source-check-evidence + allow-dirty: "true" + evidence-directory: rc2-classic-diagnostic-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-source-check-${{ matrix.environment }} + name: rc2-classic-diagnostic-${{ matrix.environment }} path: | - rc2-source-check-evidence - SOURCE-CHECK-ONLY.txt + rc2-classic-diagnostic-evidence + DIAGNOSTIC-ONLY.txt + instrumented-candidate.patch focused-harness.patch if-no-files-found: error From 192595cd770753bd8856f1b36c29462b5f2b1d0d Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 07:16:57 -0500 Subject: [PATCH 17/35] test: verify receive recovery with clean fixed sources --- .github/workflows/rc2-group-fault.yml | 79 +++++++-------------------- 1 file changed, 19 insertions(+), 60 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 80fa1c7..e0352cd 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted classic recovery ownership diagnostics +name: RC2 hosted receive recovery source checks on: pull_request: @@ -14,10 +14,11 @@ jobs: fail-fast: false max-parallel: 2 matrix: - environment: [three-sasl-plain] + environment: [three-sasl-plain, three-tls] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: + ref: e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866 path: testlab persist-credentials: false - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 @@ -32,7 +33,7 @@ jobs: ref: e5b697468eadbd99c27af48cb3653373e4599001 path: kafka-driver-candidate persist-credentials: false - - name: Trace classic recovery without changing scenario behavior + - name: Select unchanged scenarios against clean fixed source env: TARGET_CELL: ${{ matrix.environment }} run: | @@ -41,47 +42,7 @@ jobs: import os from pathlib import Path environment = os.environ["TARGET_CELL"] - assert environment == "three-sasl-plain" - base = Path("kafkars-candidate/crates/kafka-client-engine/src/consumer") - def patch(relative, before, after): - source = base / relative - code = source.read_text() - assert code.count(before) == 1, relative - source.write_text(code.replace(before, after)) - patch("group/registry_delivery.rs", - " if !entry.revocation.is_dormant() {", - ''' if !entry.revocation.is_dormant() { - static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_CLASSIC_RECV_BLOCKED group={group_id:?} at={:?} deadline={:?}", std::time::Instant::now(), entry.revocation.next_deadline()); - }''') - patch("group/registry_graceful_revocation.rs", - " if entry.revocation.expire_if_due(now)? {", - ''' if entry.revocation.expire_if_due(now)? { - eprintln!("RC2_CLASSIC_REVOCATION_EXPIRED group={:?} now={now:?}", entry.group_id());''') - patch("group/registry_graceful_revocation.rs", - ".map_err(GroupConsumerRevocationPortError::Acknowledge)", - '''.map_err(|error| { - eprintln!("RC2_CLASSIC_ACK group={group_id:?} public={assignment_epoch} error={error:?}"); - GroupConsumerRevocationPortError::Acknowledge(error) - })''') - patch("group/classic_group_fetch/delivery.rs", - ") -> Result, ClassicGroupFetchDeliveryError> {", - ''') -> Result, ClassicGroupFetchDeliveryError> { - static POLLS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if POLLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_CLASSIC_FETCH at={:?} machine={:?} effects={} pending={} sessions={}", std::time::Instant::now(), self.machine, self.effects.len(), self.pending_fetches.len(), self.fetches.retained_broker_sessions()); - }''') - reader = Path("testlab/crates/testctl/src/process_io.rs") - code = reader.read_text() - for before, after in [ - ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), - ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), - (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), - ]: - assert code.count(before) == 1 - code = code.replace(before, after) - reader.write_text(code) + assert environment in ("three-sasl-plain", "three-tls") driver = Path("kafka-driver-candidate") manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") code = manifest.read_text() @@ -98,7 +59,7 @@ jobs: code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") before = 'format!("kafkars-{version}-{}", &digest[..16])' assert code.count(before) == 1 - code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') + code = code.replace(before, 'format!("source-check-kafkars-{version}-{}", &digest[..16])') manifest.write_text(code) provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") code = provenance.read_text() @@ -114,44 +75,42 @@ jobs: } let actual = locked_registry_artifact(&lock, name)?;''') provenance.write_text(code) - Path("DIAGNOSTIC-ONLY.txt").write_text( - "NOT RELEASE EVIDENCE: instrumented native c87c5a0e plus clean unpublished driver " + Path("SOURCE-CHECK-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: clean native c87c5a0e, Testlab receive fix " + "e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866, and clean driver " "e5b697468eadbd99c27af48cb3653373e4599001. The three driver crates use " "explicit source paths. Subject archive hashes for driver RC4 do not qualify " - "the fixed runtime. Native stderr traces diagnose revocation and Fetch ownership; " - "scenario and adapter receive behavior are unchanged.\n" + "the fixed runtime. No native or driver instrumentation is applied.\n" ) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 classic recovery diagnostics, not release evidence" - scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] + title = "RC2 focused source checks, not release evidence" + scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 classic recovery diagnostics, not release evidence" + title = "RC2 focused source checks, not release evidence" [[cells]] id = "apache-kafka-4-3-1-{environment}" environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" - attempts = 5 + attempts = 3 gating = true ''') PY - git -C kafkars-candidate diff > instrumented-candidate.patch + git -C kafkars-candidate diff --exit-code git -C kafka-driver-candidate diff --exit-code git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - allow-dirty: "true" - evidence-directory: rc2-classic-diagnostic-evidence + evidence-directory: rc2-receive-source-check-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-classic-diagnostic-${{ matrix.environment }} + name: rc2-receive-source-check-${{ matrix.environment }} path: | - rc2-classic-diagnostic-evidence - DIAGNOSTIC-ONLY.txt - instrumented-candidate.patch + rc2-receive-source-check-evidence + SOURCE-CHECK-ONLY.txt focused-harness.patch if-no-files-found: error From dd481aa8bedf8e04401a39a0d2c190296db83d50 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 07:39:21 -0500 Subject: [PATCH 18/35] test: trace classic position recovery ownership --- .github/workflows/rc2-group-fault.yml | 99 ++++++++++++++++++++++----- 1 file changed, 82 insertions(+), 17 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index e0352cd..758b282 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted receive recovery source checks +name: RC2 hosted classic position recovery diagnostics on: pull_request: @@ -14,7 +14,7 @@ jobs: fail-fast: false max-parallel: 2 matrix: - environment: [three-sasl-plain, three-tls] + environment: [three-tls] steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: @@ -33,7 +33,7 @@ jobs: ref: e5b697468eadbd99c27af48cb3653373e4599001 path: kafka-driver-candidate persist-credentials: false - - name: Select unchanged scenarios against clean fixed source + - name: Trace position ownership without changing scenario behavior env: TARGET_CELL: ${{ matrix.environment }} run: | @@ -42,7 +42,69 @@ jobs: import os from pathlib import Path environment = os.environ["TARGET_CELL"] - assert environment in ("three-sasl-plain", "three-tls") + assert environment == "three-tls" + base = Path("kafkars-candidate/crates/kafka-client-engine/src") + def patch(relative, before, after): + source = base / relative + code = source.read_text() + assert code.count(before) == 1, relative + source.write_text(code.replace(before, after)) + patch("consumer/group/classic_group_position/preparation.rs", + " let partitions = copy_core_partitions(assignment.partitions())?;", + ''' eprintln!("RC2_POSITION_PREP now={now:?} fence={fence:?} deadline={:?} partitions={:?}", operation_deadline.core(), assignment.partitions()); + let partitions = copy_core_partitions(assignment.partitions())?;''') + patch("consumer/group/classic_group_position/terminal_application.rs", + " match effect {", + ''' eprintln!("RC2_POSITION_TERMINAL now={observed_at:?} fence={supplied:?} deadline={expected_deadline:?} effect={effect:?}"); + match effect {''') + patch("consumer/group/classic_group_position/registry_turn.rs", + " entry.retain_position_failure_observation(failure.observation_kind());", + ''' eprintln!("RC2_POSITION_FAULT group={:?} kind={:?} owner={:?} phase={:?} current={:?} catalog={:?}", entry.group_id(), failure.observation_kind(), std::mem::discriminant(&failure), entry.classic.machine().phase(), entry.classic.machine().live_assignment(), entry.catalog.live_assignment()); + if let super::ClassicGroupPositionFailure::Bootstrap(completed) = &failure { + eprintln!("RC2_POSITION_FAULT_BOOTSTRAP fence={:?} observed={:?} terminal={:?}", completed.fence(), completed.observed_at(), completed.terminal()); + } + entry.retain_position_failure_observation(failure.observation_kind());''') + patch("consumer/group/classic_group_position/registry_submission.rs", + " match entry.position.expire_prepared_if_due(now) {", + ''' if entry.position.next_deadline().is_some_and(|deadline| deadline.is_elapsed_at(now)) { + eprintln!("RC2_POSITION_PREPARED_EXPIRED now={now:?} deadline={:?} phase={:?} current={:?}", entry.position.next_deadline(), entry.classic.machine().phase(), entry.catalog.live_assignment()); + } + match entry.position.expire_prepared_if_due(now) {''') + patch("consumer/group/classic_group_position_reset/settlement.rs", + " let input = match terminal {", + ''' eprintln!("RC2_POSITION_RESET_TERMINAL now={now:?} fence={:?} deadline={:?} terminal={terminal:?}", reset.fence(), operation_deadline.core()); + let input = match terminal {''') + patch("driver/rpc/group_position_offset_fetch/terminal.rs", + " match &self.result {", + ''' eprintln!("RC2_POSITION_DRIVER fence={:?} deadline={:?} version={:?} error={:?}", self.key.fence(), self.key.operation_deadline().core(), self.selected_version, self.result.as_ref().err()); + match &self.result {''') + patch("consumer/group/classic_group_heartbeat_interpret.rs", + " let key = terminal.key();", + ''' let key = terminal.key(); + eprintln!("RC2_CLASSIC_HEARTBEAT now={now:?} attempt={:?} deadline={:?} route_lost={} result={:?}", key.attempt(), key.deadline().core(), terminal.coordinator_path_lost(), terminal.result());''') + patch("consumer/group/registry_graceful_revocation.rs", + " if entry.revocation.expire_if_due(now)? {", + ''' if entry.revocation.expire_if_due(now)? { + eprintln!("RC2_CLASSIC_REVOCATION_EXPIRED group={:?} now={now:?}", entry.group_id());''') + patch("consumer/group/registry_graceful_revocation.rs", + ".map_err(GroupConsumerRevocationPortError::Acknowledge)", + '''.map_err(|error| { + static ERRORS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + if ERRORS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { + eprintln!("RC2_CLASSIC_ACK at={:?} group={group_id:?} public={assignment_epoch} error={error:?}", std::time::Instant::now()); + } + GroupConsumerRevocationPortError::Acknowledge(error) + })''') + reader = Path("testlab/crates/testctl/src/process_io.rs") + code = reader.read_text() + for before, after in [ + ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), + ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), + (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), + ]: + assert code.count(before) == 1 + code = code.replace(before, after) + reader.write_text(code) driver = Path("kafka-driver-candidate") manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") code = manifest.read_text() @@ -59,7 +121,7 @@ jobs: code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") before = 'format!("kafkars-{version}-{}", &digest[..16])' assert code.count(before) == 1 - code = code.replace(before, 'format!("source-check-kafkars-{version}-{}", &digest[..16])') + code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') manifest.write_text(code) provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") code = provenance.read_text() @@ -75,42 +137,45 @@ jobs: } let actual = locked_registry_artifact(&lock, name)?;''') provenance.write_text(code) - Path("SOURCE-CHECK-ONLY.txt").write_text( - "NOT RELEASE EVIDENCE: clean native c87c5a0e, Testlab receive fix " + Path("DIAGNOSTIC-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: instrumented native c87c5a0e, Testlab receive fix " "e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866, and clean driver " "e5b697468eadbd99c27af48cb3653373e4599001. The three driver crates use " "explicit source paths. Subject archive hashes for driver RC4 do not qualify " - "the fixed runtime. No native or driver instrumentation is applied.\n" + "the fixed runtime. Native stderr traces identify position and membership " + "ownership. Scenario and adapter receive behavior are unchanged.\n" ) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 focused source checks, not release evidence" - scenarios = ["scenarios/kafka/consumer-protocol-group-membership-ownership.toml", "scenarios/kafka/classic-group-session-recovery.toml", "scenarios/kafka/consumer-protocol-group-session-recovery.toml"] + title = "RC2 classic position diagnostics, not release evidence" + scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 focused source checks, not release evidence" + title = "RC2 classic position diagnostics, not release evidence" [[cells]] id = "apache-kafka-4-3-1-{environment}" environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" - attempts = 3 + attempts = 5 gating = true ''') PY - git -C kafkars-candidate diff --exit-code + git -C kafkars-candidate diff > instrumented-candidate.patch git -C kafka-driver-candidate diff --exit-code git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - evidence-directory: rc2-receive-source-check-evidence + allow-dirty: "true" + evidence-directory: rc2-classic-position-diagnostic-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-receive-source-check-${{ matrix.environment }} + name: rc2-classic-position-diagnostic-${{ matrix.environment }} path: | - rc2-receive-source-check-evidence - SOURCE-CHECK-ONLY.txt + rc2-classic-position-diagnostic-evidence + DIAGNOSTIC-ONLY.txt + instrumented-candidate.patch focused-harness.patch if-no-files-found: error From c99992b4d3fd5ea016ca62c08d76b405124f0a99 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 07:59:52 -0500 Subject: [PATCH 19/35] test: distinguish classic heartbeat readiness failures --- .github/workflows/rc2-group-fault.yml | 43 +++++++++++++++++++-------- 1 file changed, 31 insertions(+), 12 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 758b282..0c06131 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted classic position recovery diagnostics +name: RC2 hosted classic readiness origin diagnostics on: pull_request: @@ -33,7 +33,7 @@ jobs: ref: e5b697468eadbd99c27af48cb3653373e4599001 path: kafka-driver-candidate persist-credentials: false - - name: Trace position ownership without changing scenario behavior + - name: Trace readiness origin without changing scenario behavior env: TARGET_CELL: ${{ matrix.environment }} run: | @@ -106,6 +106,23 @@ jobs: code = code.replace(before, after) reader.write_text(code) driver = Path("kafka-driver-candidate") + def driver_patch(relative, before, after): + source = driver / relative + code = source.read_text() + assert code.count(before) == 1, relative + source.write_text(code.replace(before, after)) + driver_patch("src/reactor/direct_plaintext/cluster_runtime/route_admission.rs", + " request.fail_observed(not_ready(), observed_at);", + ''' eprintln!("RC2_NOT_READY_ROUTE at={:?} lane={lane:?} failure_at={observed_at:?} physical_present={}", std::time::Instant::now(), physical.is_some()); + request.fail_observed(not_ready(), observed_at);''') + driver_patch("src/reactor/route_waiting.rs", + " let failure = if deadline <= now {", + ''' eprintln!("RC2_NOT_READY_WAIT now={now:?} deadline={deadline:?} failure_at={observed_at:?} traffic={:?}", waiting.request.traffic_class()); + let failure = if deadline <= now {''') + driver_patch("src/reactor/direct_plaintext/failure_translation.rs", + " not_sent(CallFailure::NotReady)", + ''' eprintln!("RC2_NOT_READY_RESERVE at={:?} reason=AdmissionClosed", std::time::Instant::now()); + not_sent(CallFailure::NotReady)''') manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") code = manifest.read_text() support_patches = "".join( @@ -139,43 +156,45 @@ jobs: provenance.write_text(code) Path("DIAGNOSTIC-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: instrumented native c87c5a0e, Testlab receive fix " - "e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866, and clean driver " + "e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866, and instrumented driver " "e5b697468eadbd99c27af48cb3653373e4599001. The three driver crates use " "explicit source paths. Subject archive hashes for driver RC4 do not qualify " - "the fixed runtime. Native stderr traces identify position and membership " - "ownership. Scenario and adapter receive behavior are unchanged.\n" + "the fixed runtime. Native traces identify position/membership ownership; " + "driver traces distinguish historical-route rejection from AdmissionClosed. " + "Scenario and adapter receive behavior are unchanged.\n" ) Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 id = "rc2-group-fault" - title = "RC2 classic position diagnostics, not release evidence" + title = "RC2 classic readiness diagnostics, not release evidence" scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 id = "rc2-group-fault" - title = "RC2 classic position diagnostics, not release evidence" + title = "RC2 classic readiness diagnostics, not release evidence" [[cells]] id = "apache-kafka-4-3-1-{environment}" environment = "clusters/apache-kafka/4.3.1/{environment}.toml" pack = "packs/rc2-group-fault.toml" - attempts = 5 + attempts = 3 gating = true ''') PY git -C kafkars-candidate diff > instrumented-candidate.patch - git -C kafka-driver-candidate diff --exit-code + git -C kafka-driver-candidate diff > instrumented-driver.patch git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate allow-dirty: "true" - evidence-directory: rc2-classic-position-diagnostic-evidence + evidence-directory: rc2-classic-readiness-diagnostic-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-classic-position-diagnostic-${{ matrix.environment }} + name: rc2-classic-readiness-diagnostic-${{ matrix.environment }} path: | - rc2-classic-position-diagnostic-evidence + rc2-classic-readiness-diagnostic-evidence DIAGNOSTIC-ONLY.txt instrumented-candidate.patch + instrumented-driver.patch focused-harness.patch if-no-files-found: error From 04c8a2aff56cc8a6d7b22098651c07984ad58e68 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 08:21:54 -0500 Subject: [PATCH 20/35] test: focus classic route retention source check --- .github/workflows/rc2-group-fault.yml | 188 ++++---------------------- 1 file changed, 27 insertions(+), 161 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 0c06131..0a99a65 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted classic readiness origin diagnostics +name: RC2 hosted classic route retention check on: pull_request: @@ -7,194 +7,60 @@ permissions: contents: read jobs: - group-fault: + classic-route-retention: runs-on: ubuntu-latest - timeout-minutes: 40 - strategy: - fail-fast: false - max-parallel: 2 - matrix: - environment: [three-tls] + timeout-minutes: 30 steps: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: - ref: e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866 + ref: 508ada91114a8004087b75ae2e724d84a77d7424 path: testlab persist-credentials: false - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: c87c5a0eeebba4e676662e2c26ca35bb76970dea + ref: dc709c17d8d3b1efce6ded1224083ef2f94039d1 path: kafkars-candidate persist-credentials: false - - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 - with: - repository: kafkars/kafka-driver - ref: e5b697468eadbd99c27af48cb3653373e4599001 - path: kafka-driver-candidate - persist-credentials: false - - name: Trace readiness origin without changing scenario behavior - env: - TARGET_CELL: ${{ matrix.environment }} + - name: Select the unchanged classic recovery scenario run: | + test "$(git -C kafkars-candidate rev-parse HEAD)" = "dc709c17d8d3b1efce6ded1224083ef2f94039d1" + test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' - import json - import os from pathlib import Path - environment = os.environ["TARGET_CELL"] - assert environment == "three-tls" - base = Path("kafkars-candidate/crates/kafka-client-engine/src") - def patch(relative, before, after): - source = base / relative - code = source.read_text() - assert code.count(before) == 1, relative - source.write_text(code.replace(before, after)) - patch("consumer/group/classic_group_position/preparation.rs", - " let partitions = copy_core_partitions(assignment.partitions())?;", - ''' eprintln!("RC2_POSITION_PREP now={now:?} fence={fence:?} deadline={:?} partitions={:?}", operation_deadline.core(), assignment.partitions()); - let partitions = copy_core_partitions(assignment.partitions())?;''') - patch("consumer/group/classic_group_position/terminal_application.rs", - " match effect {", - ''' eprintln!("RC2_POSITION_TERMINAL now={observed_at:?} fence={supplied:?} deadline={expected_deadline:?} effect={effect:?}"); - match effect {''') - patch("consumer/group/classic_group_position/registry_turn.rs", - " entry.retain_position_failure_observation(failure.observation_kind());", - ''' eprintln!("RC2_POSITION_FAULT group={:?} kind={:?} owner={:?} phase={:?} current={:?} catalog={:?}", entry.group_id(), failure.observation_kind(), std::mem::discriminant(&failure), entry.classic.machine().phase(), entry.classic.machine().live_assignment(), entry.catalog.live_assignment()); - if let super::ClassicGroupPositionFailure::Bootstrap(completed) = &failure { - eprintln!("RC2_POSITION_FAULT_BOOTSTRAP fence={:?} observed={:?} terminal={:?}", completed.fence(), completed.observed_at(), completed.terminal()); - } - entry.retain_position_failure_observation(failure.observation_kind());''') - patch("consumer/group/classic_group_position/registry_submission.rs", - " match entry.position.expire_prepared_if_due(now) {", - ''' if entry.position.next_deadline().is_some_and(|deadline| deadline.is_elapsed_at(now)) { - eprintln!("RC2_POSITION_PREPARED_EXPIRED now={now:?} deadline={:?} phase={:?} current={:?}", entry.position.next_deadline(), entry.classic.machine().phase(), entry.catalog.live_assignment()); - } - match entry.position.expire_prepared_if_due(now) {''') - patch("consumer/group/classic_group_position_reset/settlement.rs", - " let input = match terminal {", - ''' eprintln!("RC2_POSITION_RESET_TERMINAL now={now:?} fence={:?} deadline={:?} terminal={terminal:?}", reset.fence(), operation_deadline.core()); - let input = match terminal {''') - patch("driver/rpc/group_position_offset_fetch/terminal.rs", - " match &self.result {", - ''' eprintln!("RC2_POSITION_DRIVER fence={:?} deadline={:?} version={:?} error={:?}", self.key.fence(), self.key.operation_deadline().core(), self.selected_version, self.result.as_ref().err()); - match &self.result {''') - patch("consumer/group/classic_group_heartbeat_interpret.rs", - " let key = terminal.key();", - ''' let key = terminal.key(); - eprintln!("RC2_CLASSIC_HEARTBEAT now={now:?} attempt={:?} deadline={:?} route_lost={} result={:?}", key.attempt(), key.deadline().core(), terminal.coordinator_path_lost(), terminal.result());''') - patch("consumer/group/registry_graceful_revocation.rs", - " if entry.revocation.expire_if_due(now)? {", - ''' if entry.revocation.expire_if_due(now)? { - eprintln!("RC2_CLASSIC_REVOCATION_EXPIRED group={:?} now={now:?}", entry.group_id());''') - patch("consumer/group/registry_graceful_revocation.rs", - ".map_err(GroupConsumerRevocationPortError::Acknowledge)", - '''.map_err(|error| { - static ERRORS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); - if ERRORS.fetch_add(1, std::sync::atomic::Ordering::Relaxed) % 100 == 0 { - eprintln!("RC2_CLASSIC_ACK at={:?} group={group_id:?} public={assignment_epoch} error={error:?}", std::time::Instant::now()); - } - GroupConsumerRevocationPortError::Acknowledge(error) - })''') - reader = Path("testlab/crates/testctl/src/process_io.rs") - code = reader.read_text() - for before, after in [ - ("const MAX_STDERR_BYTES: usize = 64 * 1024;", "const MAX_STDERR_BYTES: usize = 2 * 1024 * 1024;"), - ("const MAX_STDERR_READ: u64 = 64 * 1024 + 1;", "const MAX_STDERR_READ: u64 = 2 * 1024 * 1024 + 1;"), - (" let _ = sender.send(value);", ' eprintln!("{value}");\n let _ = sender.send(value);'), - ]: - assert code.count(before) == 1 - code = code.replace(before, after) - reader.write_text(code) - driver = Path("kafka-driver-candidate") - def driver_patch(relative, before, after): - source = driver / relative - code = source.read_text() - assert code.count(before) == 1, relative - source.write_text(code.replace(before, after)) - driver_patch("src/reactor/direct_plaintext/cluster_runtime/route_admission.rs", - " request.fail_observed(not_ready(), observed_at);", - ''' eprintln!("RC2_NOT_READY_ROUTE at={:?} lane={lane:?} failure_at={observed_at:?} physical_present={}", std::time::Instant::now(), physical.is_some()); - request.fail_observed(not_ready(), observed_at);''') - driver_patch("src/reactor/route_waiting.rs", - " let failure = if deadline <= now {", - ''' eprintln!("RC2_NOT_READY_WAIT now={now:?} deadline={deadline:?} failure_at={observed_at:?} traffic={:?}", waiting.request.traffic_class()); - let failure = if deadline <= now {''') - driver_patch("src/reactor/direct_plaintext/failure_translation.rs", - " not_sent(CallFailure::NotReady)", - ''' eprintln!("RC2_NOT_READY_RESERVE at={:?} reason=AdmissionClosed", std::time::Instant::now()); - not_sent(CallFailure::NotReady)''') - manifest = Path("testlab/crates/testctl/src/candidate_manifest.rs") - code = manifest.read_text() - support_patches = "".join( - f'{name} = {{ path = {json.dumps(str(path.resolve()))} }}\n' - for name, path in [ - ("kafka-driver", driver), - ("kafka-driver-core", driver / "crates/kafka-driver-core"), - ("kafka-driver-transport", driver / "crates/kafka-driver-transport"), - ] - ) - before = " let support_patches = support_patches(artifacts, &source)?;" - assert code.count(before) == 1 - code = code.replace(before, f" let support_patches = String::from({json.dumps(support_patches)});") - before = 'format!("kafkars-{version}-{}", &digest[..16])' - assert code.count(before) == 1 - code = code.replace(before, 'format!("diagnostic-kafkars-{version}-{}", &digest[..16])') - manifest.write_text(code) - provenance = Path("testlab/crates/testctl/src/candidate_provenance.rs") - code = provenance.read_text() - before = " let actual = locked_registry_artifact(&lock, name)?;" - assert code.count(before) == 1 - code = code.replace(before, ''' if matches!(name, "kafka-driver" | "kafka-driver-core" | "kafka-driver-transport") { - let packages: Vec<_> = lock.package.iter().filter(|package| package.name == name).collect(); - if packages.len() != 1 || packages[0].version != expected.version - || packages[0].source.is_some() || packages[0].checksum.is_some() { - return Err(candidate(format!("source-check driver lock mismatch: {name}"))); - } - continue; - } - let actual = locked_registry_artifact(&lock, name)?;''') - provenance.write_text(code) - Path("DIAGNOSTIC-ONLY.txt").write_text( - "NOT RELEASE EVIDENCE: instrumented native c87c5a0e, Testlab receive fix " - "e9bc3abe28e43b01c0d5c9d46178c5b1a7e1e866, and instrumented driver " - "e5b697468eadbd99c27af48cb3653373e4599001. The three driver crates use " - "explicit source paths. Subject archive hashes for driver RC4 do not qualify " - "the fixed runtime. Native traces identify position/membership ownership; " - "driver traces distinguish historical-route rejection from AdmissionClosed. " - "Scenario and adapter receive behavior are unchanged.\n" + Path("SOURCE-CHECK-ONLY.txt").write_text( + "NOT RELEASE EVIDENCE: clean native candidate " + "dc709c17d8d3b1efce6ded1224083ef2f94039d1 with Testlab " + "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " + "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) - Path("testlab/packs/rc2-group-fault.toml").write_text('''schema_version = 1 - id = "rc2-group-fault" - title = "RC2 classic readiness diagnostics, not release evidence" + Path("testlab/packs/rc2-classic-route-retention.toml").write_text('''schema_version = 1 + id = "rc2-classic-route-retention" + title = "RC2 classic route retention source check" scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] ''') - Path("testlab/qualifications/kafkars-pr.toml").write_text(f'''schema_version = 2 - id = "rc2-group-fault" - title = "RC2 classic readiness diagnostics, not release evidence" + Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 + id = "rc2-classic-route-retention" + title = "RC2 classic route retention source check" [[cells]] - id = "apache-kafka-4-3-1-{environment}" - environment = "clusters/apache-kafka/4.3.1/{environment}.toml" - pack = "packs/rc2-group-fault.toml" - attempts = 3 + id = "apache-kafka-4-3-1-three-tls" + environment = "clusters/apache-kafka/4.3.1/three-tls.toml" + pack = "packs/rc2-classic-route-retention.toml" + attempts = 5 gating = true ''') PY - git -C kafkars-candidate diff > instrumented-candidate.patch - git -C kafka-driver-candidate diff > instrumented-driver.patch git -C testlab diff > focused-harness.patch - uses: ./testlab with: kafkars-path: kafkars-candidate - allow-dirty: "true" - evidence-directory: rc2-classic-readiness-diagnostic-evidence + evidence-directory: rc2-classic-route-retention-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-classic-readiness-diagnostic-${{ matrix.environment }} + name: rc2-classic-route-retention-source-check path: | - rc2-classic-readiness-diagnostic-evidence - DIAGNOSTIC-ONLY.txt - instrumented-candidate.patch - instrumented-driver.patch + rc2-classic-route-retention-evidence + SOURCE-CHECK-ONLY.txt focused-harness.patch if-no-files-found: error From 90e61e8b05a0be54445a0206cb206175f0cc19a0 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 12:04:13 -0500 Subject: [PATCH 21/35] test: focus classic revocation deadline source check --- .github/workflows/rc2-group-fault.yml | 32 +++++++++++++-------------- 1 file changed, 16 insertions(+), 16 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 0a99a65..f48b646 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -1,4 +1,4 @@ -name: RC2 hosted classic route retention check +name: RC2 hosted classic revocation deadline check on: pull_request: @@ -7,7 +7,7 @@ permissions: contents: read jobs: - classic-route-retention: + classic-revocation-deadline: runs-on: ubuntu-latest timeout-minutes: 30 steps: @@ -19,33 +19,33 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: dc709c17d8d3b1efce6ded1224083ef2f94039d1 + ref: a6861d75face09e15a0e92dadf8239a0668dbeb7 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "dc709c17d8d3b1efce6ded1224083ef2f94039d1" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "a6861d75face09e15a0e92dadf8239a0668dbeb7" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "dc709c17d8d3b1efce6ded1224083ef2f94039d1 with Testlab " + "a6861d75face09e15a0e92dadf8239a0668dbeb7 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) - Path("testlab/packs/rc2-classic-route-retention.toml").write_text('''schema_version = 1 - id = "rc2-classic-route-retention" - title = "RC2 classic route retention source check" + Path("testlab/packs/rc2-classic-revocation-deadline.toml").write_text('''schema_version = 1 + id = "rc2-classic-revocation-deadline" + title = "RC2 classic revocation deadline source check" scenarios = ["scenarios/kafka/classic-group-session-recovery.toml"] ''') Path("testlab/qualifications/kafkars-pr.toml").write_text('''schema_version = 2 - id = "rc2-classic-route-retention" - title = "RC2 classic route retention source check" + id = "rc2-classic-revocation-deadline" + title = "RC2 classic revocation deadline source check" [[cells]] - id = "apache-kafka-4-3-1-three-tls" - environment = "clusters/apache-kafka/4.3.1/three-tls.toml" - pack = "packs/rc2-classic-route-retention.toml" + id = "apache-kafka-4-3-1-three-scram-sha-512" + environment = "clusters/apache-kafka/4.3.1/three-scram-sha-512.toml" + pack = "packs/rc2-classic-revocation-deadline.toml" attempts = 5 gating = true ''') @@ -54,13 +54,13 @@ jobs: - uses: ./testlab with: kafkars-path: kafkars-candidate - evidence-directory: rc2-classic-route-retention-evidence + evidence-directory: rc2-classic-revocation-deadline-evidence - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a if: ${{ always() }} with: - name: rc2-classic-route-retention-source-check + name: rc2-classic-revocation-deadline-source-check path: | - rc2-classic-route-retention-evidence + rc2-classic-revocation-deadline-evidence SOURCE-CHECK-ONLY.txt focused-harness.patch if-no-files-found: error From 8519343bf3e6605505f64a501b8e565881c90c56 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 12:36:17 -0500 Subject: [PATCH 22/35] test: verify recoverable fetch session reuse --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index f48b646..0d4e54b 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: a6861d75face09e15a0e92dadf8239a0668dbeb7 + ref: 818856051c0de27bb6209bd5e3fa4a757800cf5f path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "a6861d75face09e15a0e92dadf8239a0668dbeb7" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "818856051c0de27bb6209bd5e3fa4a757800cf5f" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "a6861d75face09e15a0e92dadf8239a0668dbeb7 with Testlab " + "818856051c0de27bb6209bd5e3fa4a757800cf5f with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 761e2d312a5559b4abde67f1c93a3e23cb4cf34f Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 13:07:07 -0500 Subject: [PATCH 23/35] test: qualify refreshed fetch routes --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 0d4e54b..4906dd0 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 818856051c0de27bb6209bd5e3fa4a757800cf5f + ref: c26835ce9fd1457fafeb8819eb71bbb9e2af8a74 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "818856051c0de27bb6209bd5e3fa4a757800cf5f" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "c26835ce9fd1457fafeb8819eb71bbb9e2af8a74" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "818856051c0de27bb6209bd5e3fa4a757800cf5f with Testlab " + "c26835ce9fd1457fafeb8819eb71bbb9e2af8a74 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From b3d802af81e9e177a23c86491daf0994d781b822 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 13:15:47 -0500 Subject: [PATCH 24/35] test: qualify final fetch route candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 4906dd0..2d90670 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: c26835ce9fd1457fafeb8819eb71bbb9e2af8a74 + ref: aa05bddebd411021fd39fed33e60aa61d2e54455 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "c26835ce9fd1457fafeb8819eb71bbb9e2af8a74" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "aa05bddebd411021fd39fed33e60aa61d2e54455" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "c26835ce9fd1457fafeb8819eb71bbb9e2af8a74 with Testlab " + "aa05bddebd411021fd39fed33e60aa61d2e54455 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From b3832ba8299e989a5f995d9b76fb98d6e0659798 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 13:30:45 -0500 Subject: [PATCH 25/35] test: qualify size-compliant fetch candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 2d90670..fc02b35 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: aa05bddebd411021fd39fed33e60aa61d2e54455 + ref: d8ad58004cb5e4c87bf021dd42d4730d58aed88f path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "aa05bddebd411021fd39fed33e60aa61d2e54455" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "d8ad58004cb5e4c87bf021dd42d4730d58aed88f" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "aa05bddebd411021fd39fed33e60aa61d2e54455 with Testlab " + "d8ad58004cb5e4c87bf021dd42d4730d58aed88f with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 6533ce5e4a9a8ef60a52c87ce0bd2fbc0428f929 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 13:45:29 -0500 Subject: [PATCH 26/35] ci: qualify final rc2 fetch recovery --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index fc02b35..712ac00 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: d8ad58004cb5e4c87bf021dd42d4730d58aed88f + ref: 06d14d15331e39b761daf8881e09aa9a0b87442a path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "d8ad58004cb5e4c87bf021dd42d4730d58aed88f" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "06d14d15331e39b761daf8881e09aa9a0b87442a" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "d8ad58004cb5e4c87bf021dd42d4730d58aed88f with Testlab " + "06d14d15331e39b761daf8881e09aa9a0b87442a with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 1873068b86eb949fd507ee6cb616618b48b6a427 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 13:58:29 -0500 Subject: [PATCH 27/35] ci: qualify corrected rc2 evidence --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 712ac00..3b10146 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 06d14d15331e39b761daf8881e09aa9a0b87442a + ref: 8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "06d14d15331e39b761daf8881e09aa9a0b87442a" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "06d14d15331e39b761daf8881e09aa9a0b87442a with Testlab " + "8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 79c9d38298c4c4a65abd199382582f118239665d Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 14:18:11 -0500 Subject: [PATCH 28/35] test: pin focused recovery to latest candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 3b10146..3e3e075 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd + ref: 742a636f9aac261a1f00d5fa39ad7014d9a17f26 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "742a636f9aac261a1f00d5fa39ad7014d9a17f26" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "8ab1bc3795d9e0b3baca7d9c9f9725c02d0492cd with Testlab " + "742a636f9aac261a1f00d5fa39ad7014d9a17f26 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 5f5773e06e7067843a2cbab1d5e24e492ca3cf57 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 14:24:50 -0500 Subject: [PATCH 29/35] test: repin focused recovery after lint fix --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 3e3e075..c801884 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 742a636f9aac261a1f00d5fa39ad7014d9a17f26 + ref: 5082eda5b77b01ad12fa6929ef611312c7dd8e79 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "742a636f9aac261a1f00d5fa39ad7014d9a17f26" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "5082eda5b77b01ad12fa6929ef611312c7dd8e79" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "742a636f9aac261a1f00d5fa39ad7014d9a17f26 with Testlab " + "5082eda5b77b01ad12fa6929ef611312c7dd8e79 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From e24d5036c338482ddee2474b1132f72fb9a52279 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 15:00:21 -0500 Subject: [PATCH 30/35] test: pin focused rc candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index c801884..9f4b008 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 5082eda5b77b01ad12fa6929ef611312c7dd8e79 + ref: d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "5082eda5b77b01ad12fa6929ef611312c7dd8e79" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "5082eda5b77b01ad12fa6929ef611312c7dd8e79 with Testlab " + "d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 435439fb9c3b2c0bc5bf1272e0160595374fbee7 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 15:04:37 -0500 Subject: [PATCH 31/35] test: repin focused rc candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 9f4b008..ce2e49c 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1 + ref: 7514d81f20b3ef5837c452019cacfdaf022b4ba2 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "7514d81f20b3ef5837c452019cacfdaf022b4ba2" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "d5f9f9b6d5862706b5ec5e6060cc7127fb0c80a1 with Testlab " + "7514d81f20b3ef5837c452019cacfdaf022b4ba2 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 191f169a1fa9b611fe1c79bb402c9f0114aa330c Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 15:12:33 -0500 Subject: [PATCH 32/35] test: update focused rc candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index ce2e49c..ffa8bb1 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 7514d81f20b3ef5837c452019cacfdaf022b4ba2 + ref: 06bbbdae7f366215efff8153c679d91b6c7b3a24 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "7514d81f20b3ef5837c452019cacfdaf022b4ba2" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "06bbbdae7f366215efff8153c679d91b6c7b3a24" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "7514d81f20b3ef5837c452019cacfdaf022b4ba2 with Testlab " + "06bbbdae7f366215efff8153c679d91b6c7b3a24 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From c4fdbc6fc358ee3149f8e5ad0e129d0fccccb4f0 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 15:25:02 -0500 Subject: [PATCH 33/35] ci: qualify exact fetch retry candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index ffa8bb1..1a209f7 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 06bbbdae7f366215efff8153c679d91b6c7b3a24 + ref: 811d3bf430d0aa26a3742787d28f6009005904df path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "06bbbdae7f366215efff8153c679d91b6c7b3a24" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "811d3bf430d0aa26a3742787d28f6009005904df" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "06bbbdae7f366215efff8153c679d91b6c7b3a24 with Testlab " + "811d3bf430d0aa26a3742787d28f6009005904df with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From a455a716cc76d35edd621c30484636424d7866ec Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 15:29:36 -0500 Subject: [PATCH 34/35] ci: repin formatted fetch retry candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index 1a209f7..d05dfea 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 811d3bf430d0aa26a3742787d28f6009005904df + ref: 3fa16a7c2fcc01b27b8c5a3701a3215139fd5332 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "811d3bf430d0aa26a3742787d28f6009005904df" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "3fa16a7c2fcc01b27b8c5a3701a3215139fd5332" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "811d3bf430d0aa26a3742787d28f6009005904df with Testlab " + "3fa16a7c2fcc01b27b8c5a3701a3215139fd5332 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" ) From 40e9ea4d1607ebeeca5fb7cd1f8bd048c6f1f259 Mon Sep 17 00:00:00 2001 From: zsumz Date: Sat, 5 Sep 2026 16:26:11 -0500 Subject: [PATCH 35/35] test: pin causal fetch recovery candidate --- .github/workflows/rc2-group-fault.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/rc2-group-fault.yml b/.github/workflows/rc2-group-fault.yml index d05dfea..68f99f7 100644 --- a/.github/workflows/rc2-group-fault.yml +++ b/.github/workflows/rc2-group-fault.yml @@ -19,18 +19,18 @@ jobs: - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 with: repository: kafkars/kafkars - ref: 3fa16a7c2fcc01b27b8c5a3701a3215139fd5332 + ref: 38f93c31cc9d2c1d4db48b6e485b85ae9270b520 path: kafkars-candidate persist-credentials: false - name: Select the unchanged classic recovery scenario run: | - test "$(git -C kafkars-candidate rev-parse HEAD)" = "3fa16a7c2fcc01b27b8c5a3701a3215139fd5332" + test "$(git -C kafkars-candidate rev-parse HEAD)" = "38f93c31cc9d2c1d4db48b6e485b85ae9270b520" test "$(git -C testlab rev-parse HEAD)" = "508ada91114a8004087b75ae2e724d84a77d7424" python3 - <<'PY' from pathlib import Path Path("SOURCE-CHECK-ONLY.txt").write_text( "NOT RELEASE EVIDENCE: clean native candidate " - "3fa16a7c2fcc01b27b8c5a3701a3215139fd5332 with Testlab " + "38f93c31cc9d2c1d4db48b6e485b85ae9270b520 with Testlab " "508ada91114a8004087b75ae2e724d84a77d7424. This focused check does not " "replace exact-merge code CI or the complete 12-cell release aggregate.\n" )