From 128cabe2bc369edbf4b5f8cf0f063358e8ba984a Mon Sep 17 00:00:00 2001 From: Blake Griffith Date: Thu, 4 Jun 2026 16:00:33 -0400 Subject: [PATCH 1/4] update CHANGELOG.md --- CHANGELOG.md | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d017ea5..a723de2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,12 +2,17 @@ All notable changes to this Rust implementation of hypercore-protocol will be documented here. -### unreleased +### 7.0.1 + +* Rewrite serveral async methods to return owned futures and fix a busy loop ([PR 148](https://github.com/datrs/hypercore-protocol-rs/pull/24)). + +### 7.0.0 BIG CHANGES: * Encryption and framing of streams has been moved out of this crate into `hypercore_handshake` and `uint24le_framing` respectively. This had big impacts on the public API. Now `Protocol::new` just takes a `impl CipherTrait` argument. -* Remove dependence on `hypercore` instead we use `hypercore_schema`. +* Remove dependence on `hypercore` instead we use `hypercore_schema` (so hypercore related features have been removed). * Bumped to edition 2024. +* Dropped support for async-std (and its feature flag) ### 0.6.1 From 37f3d682d40408fa8af67fadb8e018f53017a6d9 Mon Sep 17 00:00:00 2001 From: Blake Griffith Date: Thu, 4 Jun 2026 16:09:11 -0400 Subject: [PATCH 2/4] chore: Release hypercore-protocol version 0.7.1 --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index dbf4c86..5e7a3cc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "hypercore-protocol" -version = "0.7.0" +version = "0.7.1" license = "MIT OR Apache-2.0" description = "Replication protocol for Hypercore feeds" authors = [ From 51a9bf3a3e0c65f20e41ee9a45e1c4bef1db153c Mon Sep 17 00:00:00 2001 From: Blake Griffith Date: Tue, 7 Jul 2026 21:35:19 +0200 Subject: [PATCH 3/4] Fix Open/Close messages corrupting multi-message batches Vec's multi-message encoding routes every message through the generic Message::encode, which deliberately omits the type-tag byte for Open/Close (they're framed via their own dedicated 2-byte prefix on the single-message path instead). But the multi-message decode loop always goes through the generic Message::decode, which has no case for Open/Close at all, so it misread the untagged bytes as a bogus type and errored ("Invalid message type to decode"). This only ever showed up once something started opening 2+ channels close enough together to get batched into one write - every existing test only ever sent one Open at a time. Fixed by never batching an Open/Close with anything else in MessageIo::poll_outbound, rather than changing the wire encoding itself. --- src/mqueue.rs | 22 ++++++++++++++++++---- tests/basic.rs | 35 +++++++++++++++++++++++++++++++++++ 2 files changed, 53 insertions(+), 4 deletions(-) diff --git a/src/mqueue.rs b/src/mqueue.rs index 1bd912c..348f503 100644 --- a/src/mqueue.rs +++ b/src/mqueue.rs @@ -16,7 +16,7 @@ use futures::{Sink, Stream}; use hypercore_handshake::{CipherTrait, state_machine::PUBLIC_KEYLEN}; use tracing::{error, instrument, trace}; -use crate::message::ChannelMessage; +use crate::message::{ChannelMessage, Message}; /// Message IO layer that encodes/decodes `ChannelMessage` over a byte stream. /// @@ -73,10 +73,24 @@ impl MessageIo { break; } - // Batch all queued messages + // Batch queued messages, but never batch an `Open`/`Close` message together + // with anything else. `Vec`'s multi-message wire encoding + // groups messages by channel number to avoid repeating it per message, but + // `Open`/`Close` don't have an outer channel number at all (it's embedded in + // their own payload, framed via a dedicated 2-byte prefix) — only the + // single-message path encodes them correctly. So an `Open`/`Close` at the + // front of the queue is sent alone; a batch otherwise stops right before one. let mut messages = vec![]; - while let Some(msg) = self.write_queue.pop_front() { - messages.push(msg); + while let Some(front) = self.write_queue.front() { + let front_is_open_or_close = + matches!(front.message, Message::Open(_) | Message::Close(_)); + if front_is_open_or_close && !messages.is_empty() { + break; + } + messages.push(self.write_queue.pop_front().expect("front just checked")); + if front_is_open_or_close { + break; + } } let buf = match messages.to_encoded_bytes() { diff --git a/tests/basic.rs b/tests/basic.rs index 0618d74..22c4c35 100644 --- a/tests/basic.rs +++ b/tests/basic.rs @@ -176,3 +176,38 @@ async fn open_close_channels() -> anyhow::Result<()> { fn want(start: u64, length: u64) -> Message { Message::Want(Want { start, length }) } + +/// Regression test: two `Open` messages queued before the first flush must not get batched +/// into one multi-message write. `Vec`'s multi-message encoding groups +/// messages by channel number to save repeating it, but `Open`/`Close` don't have an outer +/// channel number at all (it's embedded in their own payload, framed via a dedicated 2-byte +/// prefix) — only the single-message path encodes them correctly. Unlike `open_close_channels` +/// above (which fully establishes key1 before ever opening key2, so each `Open` is always +/// flushed alone), this opens both keys back-to-back with no driving in between, so they're +/// still queued together when the first flush happens. +#[tokio::test] +async fn two_opens_queued_before_first_flush() -> anyhow::Result<()> { + let (proto_a, proto_b) = create_pair(); + + let key1 = [4u8; 32]; + let key2 = [5u8; 32]; + + proto_a.open(key1).await?; + proto_a.open(key2).await?; + proto_b.open(key1).await?; + proto_b.open(key2).await?; + + let next_a = drive_until_channel(proto_a); + let next_b = drive_until_channel(proto_b); + let (proto_a, _channel_a1) = next_a.await??; + let (proto_b, _channel_b1) = next_b.await??; + + let next_a = drive_until_channel(proto_a); + let next_b = drive_until_channel(proto_b); + let (proto_a, _channel_a2) = next_a.await??; + let (proto_b, _channel_b2) = next_b.await??; + + assert_eq!(proto_a.channels().count(), 2); + assert_eq!(proto_b.channels().count(), 2); + Ok(()) +} From 9b1caf4c9fd7e49a42330ba7e669bfd6b57ab05b Mon Sep 17 00:00:00 2001 From: Blake Griffith Date: Mon, 3 Aug 2026 13:35:32 -0400 Subject: [PATCH 4/4] No benchmarks in CI The numberes generated by benchmarks in CI are pointless bc the runtime's performance is so variable. Also, it's really slow. --- .github/workflows/ci.yml | 1 - 1 file changed, 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f8589d1..8e7cfb4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -34,7 +34,6 @@ jobs: cargo check --all-targets --no-default-features cargo test --features js_tests cargo test --no-default-features --features js_tests - cargo test --benches build-extra: runs-on: ubuntu-latest