Skip to content

Commit a3755d0

Browse files
committed
feat: add streamReadChunk to StreamIO for bulk-read throughput
Every consumer of StreamIO read one byte at a time, which bounded the perf test app's download direction at under 5 Mbps (issue #276) while uploads ran at 0.49-0.80 Gbps. Add a chunk-level read to StreamIO: streamReadChunk :: Int -> IO ByteString returning between 1 and n bytes (whatever is buffered or arrives next), with EOF surfacing as an IOException exactly like streamReadByte. The max-length argument (not in the issue's sketch) is what lets readExactBounded use it safely: with no push-back mechanism, an unbounded chunk read would consume bytes past a message boundary. Implemented in the yamux adapter and Noise session wrapper (hand back the buffered chunk / decrypted frame), the TCP socket (recv n), the in-memory test pairs, and a mkByteStreamIO helper that derives a one-byte-per-call chunk read for byte-queue test mocks. readExactBounded now reads chunks, which moves the whole receive path off byte-at-a-time reads: Noise frame reads from the raw socket, the yamux read callback, and every length-delimited protocol reader. The relay's forwardWithLimit also forwards at chunk granularity, still never consuming a byte beyond the circuit's limit. Two DCUtR upgrade tests asserted that no direct connection exists 500ms after the circuit dial; the faster relayed path now lets the automatic DCUtR upgrade pool a direct connection inside that window (on loopback the handler-side dial is an ordinary client dial and succeeds). Those tests now run with the automatic upgrade disabled via zero-length timeout windows, making their pool preconditions deterministic.
1 parent 2f222e8 commit a3755d0

23 files changed

Lines changed: 436 additions & 253 deletions

File tree

‎src/LibP2P/MultistreamSelect/Negotiation.hs‎

Lines changed: 53 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -8,13 +8,13 @@ module LibP2P.MultistreamSelect.Negotiation
88
, StreamIO (..)
99
, negotiateInitiator
1010
, negotiateResponder
11+
, mkByteStreamIO
1112
, mkMemoryStreamPair
1213
, readExactBounded
1314
) where
1415

1516
import Control.Concurrent.STM
1617
import Control.Exception (IOException, catch)
17-
import Control.Monad (replicateM)
1818
import Data.ByteString (ByteString)
1919
import qualified Data.ByteString as BS
2020
import Data.Text (Text)
@@ -41,9 +41,28 @@ data NegotiationResult
4141

4242
-- | Abstraction for stream I/O to enable testing with in-memory buffers.
4343
data StreamIO = StreamIO
44-
{ streamWrite :: ByteString -> IO ()
45-
, streamReadByte :: IO Word8 -- ^ Read exactly one byte (blocks until available)
46-
, streamClose :: IO () -- ^ Close/half-close the stream (signals EOF to remote)
44+
{ streamWrite :: ByteString -> IO ()
45+
, streamReadByte :: IO Word8 -- ^ Read exactly one byte (blocks until available)
46+
, streamReadChunk :: Int -> IO ByteString
47+
-- ^ Read between 1 and @n@ bytes (@n >= 1@): whatever is already
48+
-- buffered or arrives next, without waiting for the full @n@.
49+
-- Blocks until at least one byte is available and never returns an
50+
-- empty ByteString; EOF and failures surface as 'IOException',
51+
-- exactly like 'streamReadByte'. Bulk readers use this to move
52+
-- data at chunk granularity instead of byte-at-a-time (#276).
53+
, streamClose :: IO () -- ^ Close/half-close the stream (signals EOF to remote)
54+
}
55+
56+
-- | Build a 'StreamIO' from byte-level primitives: 'streamReadChunk'
57+
-- falls back to one byte per call. Correct for any consumer (chunk
58+
-- reads promise at least one byte, not @n@), just not fast — intended
59+
-- for tests and mocks built on byte queues.
60+
mkByteStreamIO :: (ByteString -> IO ()) -> IO Word8 -> IO () -> StreamIO
61+
mkByteStreamIO write readByte close = StreamIO
62+
{ streamWrite = write
63+
, streamReadByte = readByte
64+
, streamReadChunk = \_ -> BS.singleton <$> readByte
65+
, streamClose = close
4766
}
4867

4968
-- | Create an in-memory stream pair for testing using STM TQueue.
@@ -54,13 +73,28 @@ mkMemoryStreamPair = do
5473
queueBtoA <- newTQueueIO :: IO (TQueue Word8)
5574
let writeToQueue q bs = mapM_ (atomically . writeTQueue q) (BS.unpack bs)
5675
readFromQueue q = atomically (readTQueue q)
76+
-- Chunk read: block for the first byte, then drain whatever else
77+
-- is already queued (up to the requested length) in the same
78+
-- transaction.
79+
drainUpTo q k
80+
| k <= (0 :: Int) = pure []
81+
| otherwise = do
82+
mb <- tryReadTQueue q
83+
case mb of
84+
Nothing -> pure []
85+
Just b -> (b :) <$> drainUpTo q (k - 1)
86+
readChunkFromQueue q n = atomically $ do
87+
b <- readTQueue q
88+
rest <- drainUpTo q (n - 1)
89+
pure (BS.pack (b : rest))
5790
pure
58-
( StreamIO (writeToQueue queueAtoB) (readFromQueue queueBtoA) (pure ())
59-
, StreamIO (writeToQueue queueBtoA) (readFromQueue queueAtoB) (pure ())
91+
( StreamIO (writeToQueue queueAtoB) (readFromQueue queueBtoA) (readChunkFromQueue queueBtoA) (pure ())
92+
, StreamIO (writeToQueue queueBtoA) (readFromQueue queueAtoB) (readChunkFromQueue queueAtoB) (pure ())
6093
)
6194

62-
-- | Chunk size for 'readExactBounded'. Bounds the transient boxed-list
63-
-- allocation per read step regardless of the requested length.
95+
-- | Maximum bytes requested per 'streamReadChunk' call in
96+
-- 'readExactBounded'. Bounds transient allocation per read step
97+
-- regardless of the requested length.
6498
readChunkSize :: Int
6599
readChunkSize = 32768
66100

@@ -70,8 +104,10 @@ readChunkSize = 32768
70104
-- #169): the declared length is validated against the caller's
71105
-- protocol-defined cap before a single byte is read or allocated, so a
72106
-- hostile length prefix cannot trigger an unbounded allocation. Bytes
73-
-- are accumulated in chunks of at most 'readChunkSize', keeping
74-
-- transient memory use proportional to the chunk size, not to @n@.
107+
-- are read via 'streamReadChunk' in requests of at most
108+
-- 'readChunkSize', keeping transient memory use proportional to the
109+
-- chunk size, not to @n@. A chunk request never exceeds the bytes
110+
-- still owed, so no byte beyond @n@ is consumed from the stream.
75111
--
76112
-- I/O failures during the read (stream reset, EOF) are returned as
77113
-- 'Left' instead of propagating as 'IOException's.
@@ -92,11 +128,13 @@ readExactBounded stream maxLen n
92128
pure (Left ("readExactBounded: read failed: " <> show e))
93129
where
94130
go :: Int -> IO [ByteString]
95-
go 0 = pure []
96-
go remaining = do
97-
let m = min readChunkSize remaining
98-
chunk <- BS.pack <$> replicateM m (streamReadByte stream)
99-
(chunk :) <$> go (remaining - m)
131+
go remaining
132+
| remaining <= 0 = pure []
133+
| otherwise = do
134+
chunk <- streamReadChunk stream (min readChunkSize remaining)
135+
if BS.null chunk
136+
then fail "readExactBounded: streamReadChunk returned no bytes"
137+
else (chunk :) <$> go (remaining - BS.length chunk)
100138

101139
-- | Read a complete multistream-select message from a stream.
102140
-- Reads varint length byte-by-byte, then reads the full payload.

‎src/LibP2P/NAT/Relay.hs‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -312,20 +312,22 @@ bridgeStreams mLimit streamA streamB = do
312312
streamClose streamB
313313

314314
-- | Forward bytes from source to destination with a byte limit.
315-
-- The limit is checked before each read, so the circuit terminates as soon
316-
-- as exactly @limit@ bytes have been forwarded — no byte beyond the limit
317-
-- is consumed from the source.
315+
-- Data moves at chunk granularity ('streamReadChunk'), but a chunk
316+
-- request never exceeds the bytes still allowed, so the circuit
317+
-- terminates as soon as exactly @limit@ bytes have been forwarded —
318+
-- no byte beyond the limit is consumed from the source.
318319
forwardWithLimit :: StreamIO -> StreamIO -> IORef Int -> Int -> IO ()
319320
forwardWithLimit src dst countRef limit = go
320321
where
322+
forwardChunkSize = 32768
321323
go = do
322324
count <- readIORef countRef
323325
if count >= limit
324326
then pure () -- limit reached, stop forwarding
325327
else do
326-
b <- streamReadByte src
327-
modifyIORef' countRef (+ 1)
328-
streamWrite dst (BS.singleton b)
328+
chunk <- streamReadChunk src (min forwardChunkSize (limit - count))
329+
modifyIORef' countRef (+ BS.length chunk)
330+
streamWrite dst chunk
329331
go
330332

331333
-- | Build a relay multiaddr in binary format.

‎src/LibP2P/Switch/Upgrade.hs‎

Lines changed: 60 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -214,6 +214,7 @@ noiseSessionToStreamIO
214214
noiseSessionToStreamIO sendRef recvRef bufRef rawIO = StreamIO
215215
{ streamWrite = encryptAndWrite sendRef rawIO
216216
, streamReadByte = decryptAndReadByte recvRef bufRef rawIO
217+
, streamReadChunk = decryptAndReadChunk recvRef bufRef rawIO
217218
, streamClose = pure () -- Encryption layer does not own the connection
218219
}
219220

@@ -233,37 +234,47 @@ encryptAndWrite sendRef rawIO plaintext =
233234
writeIORef sendRef sess'
234235
writeFramedMessage rawIO ct
235236

236-
-- | Read and decrypt a byte from the Noise channel.
237-
-- If the buffer has bytes, return the first. Otherwise, read Noise
238-
-- frames from the raw stream until one decrypts to a non-empty
239-
-- plaintext, and buffer the result. A transport message with an empty
240-
-- plaintext (a frame carrying only the AEAD tag) is legal — some
241-
-- implementations send it as a keepalive — and a zero-length frame
242-
-- carries no Noise message at all; both yield zero application bytes,
243-
-- so reading continues at the next frame.
237+
-- | Read Noise frames from the raw stream until one decrypts to a
238+
-- non-empty plaintext, and return that plaintext. A transport message
239+
-- with an empty plaintext (a frame carrying only the AEAD tag) is
240+
-- legal — some implementations send it as a keepalive — and a
241+
-- zero-length frame carries no Noise message at all; both yield zero
242+
-- application bytes, so reading continues at the next frame.
243+
nextPlaintext :: IORef NoiseSession -> StreamIO -> IO ByteString
244+
nextPlaintext recvRef rawIO = do
245+
ct <- readFramedMessage rawIO
246+
if BS.null ct
247+
then nextPlaintext recvRef rawIO -- zero-length frame: no message to decrypt
248+
else do
249+
sess <- readIORef recvRef
250+
case decryptMessage sess ct of
251+
Left err -> fail $ "nextPlaintext: decrypt failed: " <> err
252+
Right (pt, sess') -> do
253+
writeIORef recvRef sess'
254+
if BS.null pt
255+
then nextPlaintext recvRef rawIO -- empty transport message (keepalive)
256+
else pure pt
257+
258+
-- | Read and decrypt a byte from the Noise channel: pop the buffer if
259+
-- it has bytes, otherwise decrypt the next frame and buffer the rest.
244260
decryptAndReadByte :: IORef NoiseSession -> IORef ByteString -> StreamIO -> IO Word8
245261
decryptAndReadByte recvRef bufRef rawIO = do
246262
buf <- readIORef bufRef
247-
if BS.null buf
248-
then fillFromNextFrame
249-
else popByte buf
250-
where
251-
popByte bs = do
252-
writeIORef bufRef (BS.tail bs)
253-
pure (BS.head bs)
254-
fillFromNextFrame = do
255-
ct <- readFramedMessage rawIO
256-
if BS.null ct
257-
then fillFromNextFrame -- zero-length frame: no message to decrypt
258-
else do
259-
sess <- readIORef recvRef
260-
case decryptMessage sess ct of
261-
Left err -> fail $ "decryptAndReadByte: " <> err
262-
Right (pt, sess') -> do
263-
writeIORef recvRef sess'
264-
if BS.null pt
265-
then fillFromNextFrame -- empty transport message (keepalive)
266-
else popByte pt
263+
bs <- if BS.null buf then nextPlaintext recvRef rawIO else pure buf
264+
writeIORef bufRef (BS.tail bs)
265+
pure (BS.head bs)
266+
267+
-- | Chunk-level read from the Noise channel: hand back up to @n@ bytes
268+
-- of the buffered plaintext (a decrypted frame is already a chunk),
269+
-- decrypting the next frame only when the buffer is empty. Bytes
270+
-- beyond @n@ stay buffered for the next read.
271+
decryptAndReadChunk :: IORef NoiseSession -> IORef ByteString -> StreamIO -> Int -> IO ByteString
272+
decryptAndReadChunk recvRef bufRef rawIO n = do
273+
buf <- readIORef bufRef
274+
bs <- if BS.null buf then nextPlaintext recvRef rawIO else pure buf
275+
let (front, rest) = BS.splitAt n bs
276+
writeIORef bufRef rest
277+
pure front
267278

268279
-- | Bounded window given to the send loop to flush the GoAway frame
269280
-- before the transport is closed underneath it.
@@ -312,18 +323,13 @@ yamuxToMuxerSession yamuxSess closeTransport = do
312323
}
313324

314325
-- | Convert a YamuxStream to StreamIO with a read buffer.
315-
-- Yamux delivers data in chunks via streamRead, but StreamIO requires
316-
-- byte-by-byte reads. An IORef buffer bridges this gap.
326+
-- Yamux delivers data in chunks via streamRead; an IORef buffer holds
327+
-- the bytes a byte- or chunk-level read did not consume.
317328
yamuxStreamToStreamIO :: YamuxStream -> IO StreamIO
318329
yamuxStreamToStreamIO yamuxStream = do
319330
readBuf <- newIORef BS.empty
320-
pure StreamIO
321-
{ streamWrite = \bs -> do
322-
result <- YS.streamWrite yamuxStream bs
323-
case result of
324-
Right () -> pure ()
325-
Left err -> fail $ "yamuxStreamWrite: " <> show err
326-
, streamReadByte = do
331+
let -- Buffered bytes if any, otherwise the next yamux chunk.
332+
nextChunk = do
327333
buf <- readIORef readBuf
328334
if BS.null buf
329335
then do
@@ -332,13 +338,23 @@ yamuxStreamToStreamIO yamuxStream = do
332338
Left err -> fail $ "yamuxStreamRead: " <> show err
333339
Right chunk
334340
| BS.null chunk -> fail "yamuxStreamRead: empty chunk"
335-
| BS.length chunk == 1 -> pure (BS.head chunk)
336-
| otherwise -> do
337-
writeIORef readBuf (BS.tail chunk)
338-
pure (BS.head chunk)
339-
else do
340-
writeIORef readBuf (BS.tail buf)
341-
pure (BS.head buf)
341+
| otherwise -> pure chunk
342+
else pure buf
343+
pure StreamIO
344+
{ streamWrite = \bs -> do
345+
result <- YS.streamWrite yamuxStream bs
346+
case result of
347+
Right () -> pure ()
348+
Left err -> fail $ "yamuxStreamWrite: " <> show err
349+
, streamReadByte = do
350+
chunk <- nextChunk
351+
writeIORef readBuf (BS.tail chunk)
352+
pure (BS.head chunk)
353+
, streamReadChunk = \n -> do
354+
chunk <- nextChunk
355+
let (front, rest) = BS.splitAt n chunk
356+
writeIORef readBuf rest
357+
pure front
342358
, streamClose = do
343359
_ <- YS.streamClose yamuxStream -- Sends FIN flag
344360
pure ()

‎src/LibP2P/Transport/TCP.hs‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -134,13 +134,16 @@ tcpListen _ = fail "tcpListen: unsupported multiaddr"
134134
socketToStreamIO :: NS.Socket -> StreamIO
135135
socketToStreamIO sock = StreamIO
136136
{ streamWrite = NSB.sendAll sock
137-
, streamReadByte = do
138-
bs <- NSB.recv sock 1
139-
if BS.null bs
140-
then fail "socketToStreamIO: connection closed"
141-
else pure (BS.head bs)
137+
, streamReadByte = BS.head <$> recvChunk 1
138+
, streamReadChunk = recvChunk
142139
, streamClose = NS.close sock
143140
}
141+
where
142+
recvChunk n = do
143+
bs <- NSB.recv sock n
144+
if BS.null bs
145+
then fail "socketToStreamIO: connection closed"
146+
else pure bs
144147

145148
-- | Convert a SockAddr to a Multiaddr.
146149
sockAddrToMultiaddr :: NS.SockAddr -> IO Multiaddr

‎test/LibP2P/ConformanceSpec.hs‎

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import LibP2P.Crypto.PeerId (PeerId, fromPublicKey)
2424
import LibP2P.MultistreamSelect.Negotiation
2525
( NegotiationResult (..)
2626
, StreamIO (..)
27+
, mkByteStreamIO
2728
, mkMemoryStreamPair
2829
, negotiateInitiator
2930
, negotiateResponder
@@ -55,15 +56,15 @@ mkScriptedStream :: ByteString -> IO (StreamIO, IO ByteString)
5556
mkScriptedStream canned = do
5657
writtenRef <- newIORef BS.empty
5758
readRef <- newIORef canned
58-
let stream = StreamIO
59-
{ streamWrite = \bs -> modifyIORef' writtenRef (`BS.append` bs)
60-
, streamReadByte = do
61-
buf <- readIORef readRef
62-
case BS.uncons buf of
63-
Nothing -> ioError (userError "scripted stream: EOF")
64-
Just (b, rest) -> writeIORef readRef rest >> pure b
65-
, streamClose = pure ()
66-
}
59+
let readB = do
60+
buf <- readIORef readRef
61+
case BS.uncons buf of
62+
Nothing -> ioError (userError "scripted stream: EOF")
63+
Just (b, rest) -> writeIORef readRef rest >> pure b
64+
stream = mkByteStreamIO
65+
(\bs -> modifyIORef' writtenRef (`BS.append` bs))
66+
readB
67+
(pure ())
6768
pure (stream, readIORef writtenRef)
6869

6970
-- multistream-select messages, hand-derived from the spec:

‎test/LibP2P/DHT/DHTSpec.hs‎

Lines changed: 9 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ import LibP2P.DHT.Message
2727
import LibP2P.DHT.RoutingTable (allPeers, bucketForPeer, insertPeer, newRoutingTable)
2828
import LibP2P.DHT.Types
2929
import LibP2P.Multiaddr (Multiaddr, fromText, toBytes)
30-
import LibP2P.MultistreamSelect.Negotiation (StreamIO (..), negotiateResponder)
30+
import LibP2P.MultistreamSelect.Negotiation (StreamIO (..), mkByteStreamIO, negotiateResponder)
3131
import LibP2P.Switch.ConnPool (addConn)
3232
import LibP2P.Switch.Types
3333
( ConnState (..)
@@ -151,16 +151,14 @@ mkStreamPair = do
151151
Nothing -> do
152152
closed <- readTVar closedVar
153153
if closed then throwSTM (userError "stream closed") else retry
154-
streamA = StreamIO
155-
{ streamWrite = writeAll q1
156-
, streamReadByte = readOrEOF q2 closedBtoA
157-
, streamClose = atomically (writeTVar closedAtoB True)
158-
}
159-
streamB = StreamIO
160-
{ streamWrite = writeAll q2
161-
, streamReadByte = readOrEOF q1 closedAtoB
162-
, streamClose = atomically (writeTVar closedBtoA True)
163-
}
154+
streamA = mkByteStreamIO
155+
(writeAll q1)
156+
(readOrEOF q2 closedBtoA)
157+
(atomically (writeTVar closedAtoB True))
158+
streamB = mkByteStreamIO
159+
(writeAll q2)
160+
(readOrEOF q1 closedAtoB)
161+
(atomically (writeTVar closedBtoA True))
164162
pure (streamA, streamB)
165163

166164
-- | A mock Connection that hands out the given stream on the first

0 commit comments

Comments
 (0)