nim-libp2p-experimental/tests/testmultistream.nim

479 lines
14 KiB
Nim
Raw Normal View History

import strutils, strformat, stew/byteutils
2019-08-30 05:16:55 +00:00
import chronos
import ../libp2p/errors,
2019-09-06 07:13:03 +00:00
../libp2p/multistream,
2020-02-12 14:43:42 +00:00
../libp2p/stream/bufferstream,
../libp2p/stream/connection,
../libp2p/multiaddress,
../libp2p/transports/transport,
../libp2p/transports/tcptransport,
../libp2p/protocols/protocol,
../libp2p/upgrademngrs/upgrade
2019-08-30 05:16:55 +00:00
when (NimMajor, NimMinor) < (1, 4):
{.push raises: [Defect].}
else:
{.push raises: [].}
import ./helpers
when defined(nimHasUsed): {.used.}
2019-08-30 05:16:55 +00:00
## Mock stream for select test
type
TestSelectStream = ref object of Connection
2019-08-30 05:16:55 +00:00
step*: int
method readOnce*(s: TestSelectStream,
pbytes: pointer,
nbytes: int): Future[int] {.async, gcsafe.} =
2019-08-30 05:16:55 +00:00
case s.step:
of 1:
var buf = newSeq[byte](1)
buf[0] = 19
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 2
return buf.len
2019-08-30 05:16:55 +00:00
of 2:
var buf = "/multistream/1.0.0\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 3
return buf.len
2019-08-30 05:16:55 +00:00
of 3:
var buf = newSeq[byte](1)
buf[0] = 18
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 4
return buf.len
2019-08-30 05:16:55 +00:00
of 4:
var buf = "/test/proto/1.0.0\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
return buf.len
2019-08-30 05:16:55 +00:00
else:
2019-09-01 00:23:11 +00:00
copyMem(pbytes,
2019-08-30 05:16:55 +00:00
cstring("\0x3na\n"),
"\0x3na\n".len())
return "\0x3na\n".len()
method write*(s: TestSelectStream, msg: seq[byte]) {.async, gcsafe.} = discard
2019-09-08 07:59:14 +00:00
method close(s: TestSelectStream) {.async, gcsafe.} =
s.isClosed = true
s.isEof = true
2019-09-06 07:13:03 +00:00
2019-08-30 05:16:55 +00:00
proc newTestSelectStream(): TestSelectStream =
new result
result.step = 1
## Mock stream for handles `ls` test
type
LsHandler = proc(procs: seq[byte]): Future[void] {.gcsafe, raises: [Defect].}
2019-08-30 05:16:55 +00:00
TestLsStream = ref object of Connection
2019-08-30 05:16:55 +00:00
step*: int
ls*: LsHandler
method readOnce*(s: TestLsStream,
pbytes: pointer,
nbytes: int):
Future[int] {.async.} =
2019-08-30 05:16:55 +00:00
case s.step:
of 1:
var buf = newSeq[byte](1)
buf[0] = 19
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 2
return buf.len()
2019-08-30 05:16:55 +00:00
of 2:
var buf = "/multistream/1.0.0\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 3
return buf.len()
2019-08-30 05:16:55 +00:00
of 3:
var buf = newSeq[byte](1)
buf[0] = 3
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 4
return buf.len()
2019-08-30 05:16:55 +00:00
of 4:
var buf = "ls\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
return buf.len()
2019-08-30 05:16:55 +00:00
else:
var buf = "na\n"
copyMem(pbytes, addr buf[0], buf.len())
return buf.len()
2019-08-30 05:16:55 +00:00
method write*(s: TestLsStream, msg: seq[byte]) {.async, gcsafe.} =
2019-08-30 05:16:55 +00:00
if s.step == 4:
await s.ls(msg)
2020-02-12 14:43:42 +00:00
method close(s: TestLsStream) {.async, gcsafe.} =
s.isClosed = true
s.isEof = true
2019-09-01 17:31:24 +00:00
2019-10-04 20:10:01 +00:00
proc newTestLsStream(ls: LsHandler): TestLsStream {.gcsafe.} =
2019-08-30 05:16:55 +00:00
new result
result.ls = ls
result.step = 1
## Mock stream for handles `na` test
type
NaHandler = proc(procs: string): Future[void] {.gcsafe, raises: [Defect].}
2019-08-30 05:16:55 +00:00
TestNaStream = ref object of Connection
2019-08-30 05:16:55 +00:00
step*: int
na*: NaHandler
method readOnce*(s: TestNaStream,
pbytes: pointer,
nbytes: int):
Future[int] {.async, gcsafe.} =
2019-08-30 05:16:55 +00:00
case s.step:
of 1:
var buf = newSeq[byte](1)
buf[0] = 19
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 2
return buf.len()
2019-08-30 05:16:55 +00:00
of 2:
var buf = "/multistream/1.0.0\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 3
return buf.len()
2019-08-30 05:16:55 +00:00
of 3:
var buf = newSeq[byte](1)
buf[0] = 18
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
2019-08-30 05:16:55 +00:00
s.step = 4
return buf.len()
2019-08-30 05:16:55 +00:00
of 4:
var buf = "/test/proto/1.0.0\n"
2019-09-01 00:23:11 +00:00
copyMem(pbytes, addr buf[0], buf.len())
return buf.len()
2019-08-30 05:16:55 +00:00
else:
2019-09-01 00:23:11 +00:00
copyMem(pbytes,
2019-08-30 05:16:55 +00:00
cstring("\0x3na\n"),
"\0x3na\n".len())
return "\0x3na\n".len()
method write*(s: TestNaStream, msg: seq[byte]) {.async, gcsafe.} =
2019-08-30 05:16:55 +00:00
if s.step == 4:
await s.na(string.fromBytes(msg))
2019-08-30 05:16:55 +00:00
2020-02-12 14:43:42 +00:00
method close(s: TestNaStream) {.async, gcsafe.} =
s.isClosed = true
s.isEof = true
2019-09-01 17:31:24 +00:00
2019-08-30 05:16:55 +00:00
proc newTestNaStream(na: NaHandler): TestNaStream =
new result
result.na = na
result.step = 1
suite "Multistream select":
teardown:
checkTrackers()
asyncTest "test select custom proto":
let ms = MultistreamSelect.new()
let conn = newTestSelectStream()
check (await ms.select(conn, @["/test/proto/1.0.0"])) == "/test/proto/1.0.0"
await conn.close()
asyncTest "test handle custom proto":
let ms = MultistreamSelect.new()
let conn = newTestSelectStream()
var protocol: LPProtocol = new LPProtocol
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} =
check proto == "/test/proto/1.0.0"
await conn.close()
2019-08-30 05:16:55 +00:00
protocol.handler = testHandler
ms.addHandler("/test/proto/1.0.0", protocol)
await ms.handle(conn)
asyncTest "test handle `ls`":
let ms = MultistreamSelect.new()
var conn: Connection = nil
let done = newFuture[void]()
proc testLsHandler(proto: seq[byte]) {.async, gcsafe.} =
var strProto: string = string.fromBytes(proto)
check strProto == "\x26/test/proto1/1.0.0\n/test/proto2/1.0.0\n"
await conn.close()
done.complete()
conn = Connection(newTestLsStream(testLsHandler))
proc testHandler(conn: Connection, proto: string): Future[void]
{.async, gcsafe.} = discard
var protocol: LPProtocol = new LPProtocol
protocol.handler = testHandler
ms.addHandler("/test/proto1/1.0.0", protocol)
ms.addHandler("/test/proto2/1.0.0", protocol)
await ms.handle(conn)
await done.wait(5.seconds)
asyncTest "test handle `na`":
let ms = MultistreamSelect.new()
var conn: Connection = nil
proc testNaHandler(msg: string): Future[void] {.async, gcsafe.} =
2023-03-08 11:30:19 +00:00
check msg == "\x03na\n"
await conn.close()
conn = newTestNaStream(testNaHandler)
2019-08-30 05:16:55 +00:00
var protocol: LPProtocol = new LPProtocol
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} = discard
protocol.handler = testHandler
ms.addHandler("/unabvailable/proto/1.0.0", protocol)
2019-08-30 05:16:55 +00:00
await ms.handle(conn)
2019-08-30 05:16:55 +00:00
asyncTest "e2e - handle":
2021-12-16 10:05:20 +00:00
let ma = @[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()]
var protocol: LPProtocol = new LPProtocol
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} =
check proto == "/test/proto/1.0.0"
await conn.writeLp("Hello!")
await conn.close()
protocol.handler = testHandler
let msListen = MultistreamSelect.new()
msListen.addHandler("/test/proto/1.0.0", protocol)
2019-08-30 05:16:55 +00:00
let transport1 = TcpTransport.new(upgrade = Upgrade())
Remove asynccheck (#590) * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * cleanup Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im>
2021-06-14 23:21:44 +00:00
asyncSpawn transport1.start(ma)
proc acceptHandler(): Future[void] {.async, gcsafe.} =
let conn = await transport1.accept()
await msListen.handle(conn)
await conn.close()
let handlerWait = acceptHandler()
2019-08-30 05:16:55 +00:00
let msDial = MultistreamSelect.new()
let transport2 = TcpTransport.new(upgrade = Upgrade())
let conn = await transport2.dial(transport1.addrs[0])
2019-08-30 05:16:55 +00:00
check (await msDial.select(conn, "/test/proto/1.0.0")) == true
let hello = string.fromBytes(await conn.readLp(1024))
check hello == "Hello!"
await conn.close()
2019-08-30 05:16:55 +00:00
await transport2.stop()
await transport1.stop()
await handlerWait.wait(30.seconds)
asyncTest "e2e - streams limit":
let ma = @[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()]
let blocker = newFuture[void]()
# Start 5 streams which are blocked by `blocker`
# Try to start a new one, which should fail
# Unblock the 5 streams, check that we can open a new one
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} =
await blocker
await conn.writeLp("Hello!")
await conn.close()
var protocol: LPProtocol = LPProtocol.new(
@["/test/proto/1.0.0"],
testHandler,
maxIncomingStreams = 5
)
protocol.handler = testHandler
let msListen = MultistreamSelect.new()
msListen.addHandler("/test/proto/1.0.0", protocol)
let transport1 = TcpTransport.new(upgrade = Upgrade())
await transport1.start(ma)
proc acceptedOne(c: Connection) {.async.} =
await msListen.handle(c)
await c.close()
proc acceptHandler() {.async, gcsafe.} =
while true:
let conn = await transport1.accept()
asyncSpawn acceptedOne(conn)
var handlerWait = acceptHandler()
let msDial = MultistreamSelect.new()
let transport2 = TcpTransport.new(upgrade = Upgrade())
proc connector {.async.} =
let conn = await transport2.dial(transport1.addrs[0])
check: (await msDial.select(conn, "/test/proto/1.0.0")) == true
check: string.fromBytes(await conn.readLp(1024)) == "Hello!"
await conn.close()
# Fill up the 5 allowed streams
var dialers: seq[Future[void]]
for _ in 0..<5:
dialers.add(connector())
# This one will fail during negotiation
expect(CatchableError):
try: waitFor(connector().wait(1.seconds))
except AsyncTimeoutError as exc:
check false
raise exc
# check that the dialers aren't finished
check: (await dialers[0].withTimeout(10.milliseconds)) == false
# unblock the dialers
blocker.complete()
await allFutures(dialers)
# now must work
waitFor(connector())
await transport2.stop()
await transport1.stop()
await handlerWait.cancelAndWait()
asyncTest "e2e - ls":
2021-12-16 10:05:20 +00:00
let ma = @[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()]
let msListen = MultistreamSelect.new()
var protocol: LPProtocol = new LPProtocol
protocol.handler = proc(conn: Connection, proto: string) {.async, gcsafe.} =
# never reached
discard
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async.} =
# never reached
discard
protocol.handler = testHandler
msListen.addHandler("/test/proto1/1.0.0", protocol)
msListen.addHandler("/test/proto2/1.0.0", protocol)
let transport1: TcpTransport = TcpTransport.new(upgrade = Upgrade())
let listenFut = transport1.start(ma)
proc acceptHandler(): Future[void] {.async, gcsafe.} =
let conn = await transport1.accept()
try:
await msListen.handle(conn)
except LPStreamEOFError:
discard
except LPStreamClosedError:
discard
finally:
await conn.close()
let acceptFut = acceptHandler()
let msDial = MultistreamSelect.new()
let transport2: TcpTransport = TcpTransport.new(upgrade = Upgrade())
let conn = await transport2.dial(transport1.addrs[0])
let ls = await msDial.list(conn)
let protos: seq[string] = @["/test/proto1/1.0.0", "/test/proto2/1.0.0"]
check ls == protos
await conn.close()
await acceptFut
await transport2.stop()
await transport1.stop()
await listenFut.wait(5.seconds)
asyncTest "e2e - select one from a list with unsupported protos":
2021-12-16 10:05:20 +00:00
let ma = @[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()]
var protocol: LPProtocol = new LPProtocol
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} =
check proto == "/test/proto/1.0.0"
await conn.writeLp("Hello!")
await conn.close()
protocol.handler = testHandler
let msListen = MultistreamSelect.new()
msListen.addHandler("/test/proto/1.0.0", protocol)
let transport1: TcpTransport = TcpTransport.new(upgrade = Upgrade())
Remove asynccheck (#590) * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * cleanup Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im>
2021-06-14 23:21:44 +00:00
asyncSpawn transport1.start(ma)
proc acceptHandler(): Future[void] {.async, gcsafe.} =
let conn = await transport1.accept()
await msListen.handle(conn)
let acceptFut = acceptHandler()
let msDial = MultistreamSelect.new()
let transport2: TcpTransport = TcpTransport.new(upgrade = Upgrade())
let conn = await transport2.dial(transport1.addrs[0])
check (await msDial.select(conn,
@["/test/proto/1.0.0", "/test/no/proto/1.0.0"])) == "/test/proto/1.0.0"
let hello = string.fromBytes(await conn.readLp(1024))
check hello == "Hello!"
await conn.close()
await acceptFut
await transport2.stop()
await transport1.stop()
asyncTest "e2e - select one with both valid":
2021-12-16 10:05:20 +00:00
let ma = @[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()]
var protocol: LPProtocol = new LPProtocol
proc testHandler(conn: Connection,
proto: string):
Future[void] {.async, gcsafe.} =
await conn.writeLp(&"Hello from {proto}!")
await conn.close()
protocol.handler = testHandler
let msListen = MultistreamSelect.new()
msListen.addHandler("/test/proto1/1.0.0", protocol)
msListen.addHandler("/test/proto2/1.0.0", protocol)
let transport1: TcpTransport = TcpTransport.new(upgrade = Upgrade())
Remove asynccheck (#590) * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * Merge master (#555) * Revisit Floodsub (#543) Fixes #525 add coverage to unsubscribeAll and testing * add mounted protos to identify message (#546) * add stable/unstable auto bumps * fix auto-bump CI * merge nbc auto bump with CI in order to bump only on CI success * put conditional locks on nbc bump (#549) * Fix minor exception issues (#550) Makes code compatible with https://github.com/status-im/nim-chronos/pull/166 without requiring it. * fix nimbus ref for auto-bump stable's PR * Split dialer (#542) * extracting dialing logic to dialer * exposing upgrade methods on transport * cleanup * fixing tests to use new interfaces * add comments * add base exception class and fix hierarchy * fix imports * `doAssert` is `ValueError` not `AssertionError`? * revert back to `AssertionError` Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im> * cleanup Co-authored-by: Giovanni Petrantoni <7008900+sinkingsugar@users.noreply.github.com> Co-authored-by: Jacek Sieka <jacek@status.im>
2021-06-14 23:21:44 +00:00
asyncSpawn transport1.start(ma)
proc acceptHandler(): Future[void] {.async, gcsafe.} =
let conn = await transport1.accept()
await msListen.handle(conn)
let acceptFut = acceptHandler()
let msDial = MultistreamSelect.new()
let transport2: TcpTransport = TcpTransport.new(upgrade = Upgrade())
let conn = await transport2.dial(transport1.addrs[0])
check (await msDial.select(conn,
@[
"/test/proto2/1.0.0",
"/test/proto1/1.0.0"
])) == "/test/proto2/1.0.0"
check string.fromBytes(await conn.readLp(1024)) == "Hello from /test/proto2/1.0.0!"
await conn.close()
await acceptFut
await transport2.stop()
await transport1.stop()