mirror of https://github.com/vacp2p/nim-libp2p.git
add support for setting protocol handlers with `{.raises.}` annotation (#1064)
This commit is contained in:
parent
bb97a9de79
commit
03f67d3db5
|
@ -20,7 +20,7 @@ proc new(T: typedesc[TestProto]): T =
|
||||||
# We must close the connections ourselves when we're done with it
|
# We must close the connections ourselves when we're done with it
|
||||||
await conn.close()
|
await conn.close()
|
||||||
|
|
||||||
return T(codecs: @[TestCodec], handler: handle)
|
return T.new(codecs = @[TestCodec], handler = handle)
|
||||||
|
|
||||||
##
|
##
|
||||||
# Helper to create a switch/node
|
# Helper to create a switch/node
|
||||||
|
|
|
@ -246,9 +246,13 @@ proc addHandler*(m: MultistreamSelect,
|
||||||
matcher: Matcher = nil) =
|
matcher: Matcher = nil) =
|
||||||
addHandler(m, @[codec], protocol, matcher)
|
addHandler(m, @[codec], protocol, matcher)
|
||||||
|
|
||||||
proc addHandler*(m: MultistreamSelect,
|
proc addHandler*[E](
|
||||||
|
m: MultistreamSelect,
|
||||||
codec: string,
|
codec: string,
|
||||||
handler: LPProtoHandler,
|
handler: LPProtoHandler |
|
||||||
|
proc (
|
||||||
|
conn: Connection,
|
||||||
|
proto: string): InternalRaisesFuture[void, E],
|
||||||
matcher: Matcher = nil) =
|
matcher: Matcher = nil) =
|
||||||
## helper to allow registering pure handlers
|
## helper to allow registering pure handlers
|
||||||
trace "registering proto handler", proto = codec
|
trace "registering proto handler", proto = codec
|
||||||
|
|
|
@ -20,13 +20,11 @@ const
|
||||||
type
|
type
|
||||||
LPProtoHandler* = proc (
|
LPProtoHandler* = proc (
|
||||||
conn: Connection,
|
conn: Connection,
|
||||||
proto: string):
|
proto: string): Future[void] {.async.}
|
||||||
Future[void]
|
|
||||||
{.gcsafe, raises: [].}
|
|
||||||
|
|
||||||
LPProtocol* = ref object of RootObj
|
LPProtocol* = ref object of RootObj
|
||||||
codecs*: seq[string]
|
codecs*: seq[string]
|
||||||
handler*: LPProtoHandler ## this handler gets invoked by the protocol negotiator
|
handlerImpl: LPProtoHandler ## invoked by the protocol negotiator
|
||||||
started*: bool
|
started*: bool
|
||||||
maxIncomingStreams: Opt[int]
|
maxIncomingStreams: Opt[int]
|
||||||
|
|
||||||
|
@ -52,7 +50,7 @@ proc `maxIncomingStreams=`*(p: LPProtocol, val: int) =
|
||||||
p.maxIncomingStreams = Opt.some(val)
|
p.maxIncomingStreams = Opt.some(val)
|
||||||
|
|
||||||
func codec*(p: LPProtocol): string =
|
func codec*(p: LPProtocol): string =
|
||||||
assert(p.codecs.len > 0, "Codecs sequence was empty!")
|
doAssert(p.codecs.len > 0, "Codecs sequence was empty!")
|
||||||
p.codecs[0]
|
p.codecs[0]
|
||||||
|
|
||||||
func `codec=`*(p: LPProtocol, codec: string) =
|
func `codec=`*(p: LPProtocol, codec: string) =
|
||||||
|
@ -60,6 +58,31 @@ func `codec=`*(p: LPProtocol, codec: string) =
|
||||||
# if we use this abstraction
|
# if we use this abstraction
|
||||||
p.codecs.insert(codec, 0)
|
p.codecs.insert(codec, 0)
|
||||||
|
|
||||||
|
template `handler`*(p: LPProtocol): LPProtoHandler =
|
||||||
|
p.handlerImpl
|
||||||
|
|
||||||
|
template `handler`*(
|
||||||
|
p: LPProtocol, conn: Connection, proto: string): Future[void] =
|
||||||
|
p.handlerImpl(conn, proto)
|
||||||
|
|
||||||
|
func `handler=`*(p: LPProtocol, handler: LPProtoHandler) =
|
||||||
|
p.handlerImpl = handler
|
||||||
|
|
||||||
|
# Callbacks that are annotated with `{.async: (raises).}` explicitly
|
||||||
|
# document the types of errors that they may raise, but are not compatible
|
||||||
|
# with `LPProtoHandler` and need to use a custom `proc` type.
|
||||||
|
# They are internally wrapped into a `LPProtoHandler`, but still allow the
|
||||||
|
# compiler to check that their `{.async: (raises).}` annotation is correct.
|
||||||
|
# https://github.com/nim-lang/Nim/issues/23432
|
||||||
|
func `handler=`*[E](
|
||||||
|
p: LPProtocol,
|
||||||
|
handler: proc (
|
||||||
|
conn: Connection,
|
||||||
|
proto: string): InternalRaisesFuture[void, E]) =
|
||||||
|
proc wrap(conn: Connection, proto: string): Future[void] {.async.} =
|
||||||
|
await handler(conn, proto)
|
||||||
|
p.handlerImpl = wrap
|
||||||
|
|
||||||
proc new*(
|
proc new*(
|
||||||
T: type LPProtocol,
|
T: type LPProtocol,
|
||||||
codecs: seq[string],
|
codecs: seq[string],
|
||||||
|
@ -67,8 +90,19 @@ proc new*(
|
||||||
maxIncomingStreams: Opt[int] | int = Opt.none(int)): T =
|
maxIncomingStreams: Opt[int] | int = Opt.none(int)): T =
|
||||||
T(
|
T(
|
||||||
codecs: codecs,
|
codecs: codecs,
|
||||||
handler: handler,
|
handlerImpl: handler,
|
||||||
maxIncomingStreams:
|
maxIncomingStreams:
|
||||||
when maxIncomingStreams is int: Opt.some(maxIncomingStreams)
|
when maxIncomingStreams is int: Opt.some(maxIncomingStreams)
|
||||||
else: maxIncomingStreams
|
else: maxIncomingStreams
|
||||||
)
|
)
|
||||||
|
|
||||||
|
proc new*[E](
|
||||||
|
T: type LPProtocol,
|
||||||
|
codecs: seq[string],
|
||||||
|
handler: proc (
|
||||||
|
conn: Connection,
|
||||||
|
proto: string): InternalRaisesFuture[void, E],
|
||||||
|
maxIncomingStreams: Opt[int] | int = Opt.none(int)): T =
|
||||||
|
proc wrap(conn: Connection, proto: string): Future[void] {.async.} =
|
||||||
|
await handler(conn, proto)
|
||||||
|
T.new(codec, wrap, maxIncomingStreams)
|
||||||
|
|
|
@ -24,6 +24,6 @@ proc allFuturesThrowing*(args: varargs[FutureBase]): Future[void] =
|
||||||
proc allFuturesThrowing*[T](futs: varargs[Future[T]]): Future[void] =
|
proc allFuturesThrowing*[T](futs: varargs[Future[T]]): Future[void] =
|
||||||
allFuturesThrowing(futs.mapIt(FutureBase(it)))
|
allFuturesThrowing(futs.mapIt(FutureBase(it)))
|
||||||
|
|
||||||
proc allFuturesThrowing*[T, E](
|
proc allFuturesThrowing*[T, E]( # https://github.com/nim-lang/Nim/issues/23432
|
||||||
futs: varargs[InternalRaisesFuture[T, E]]): Future[void] =
|
futs: varargs[InternalRaisesFuture[T, E]]): Future[void] =
|
||||||
allFuturesThrowing(futs.mapIt(FutureBase(it)))
|
allFuturesThrowing(futs.mapIt(FutureBase(it)))
|
||||||
|
|
|
@ -1,7 +1,7 @@
|
||||||
{.used.}
|
{.used.}
|
||||||
|
|
||||||
# Nim-Libp2p
|
# Nim-Libp2p
|
||||||
# Copyright (c) 2023 Status Research & Development GmbH
|
# Copyright (c) 2023-2024 Status Research & Development GmbH
|
||||||
# Licensed under either of
|
# Licensed under either of
|
||||||
# * Apache License, version 2.0, ([LICENSE-APACHE](LICENSE-APACHE))
|
# * Apache License, version 2.0, ([LICENSE-APACHE](LICENSE-APACHE))
|
||||||
# * MIT license ([LICENSE-MIT](LICENSE-MIT))
|
# * MIT license ([LICENSE-MIT](LICENSE-MIT))
|
||||||
|
@ -88,7 +88,6 @@ suite "Tor transport":
|
||||||
|
|
||||||
# every incoming connections will be in handled in this closure
|
# every incoming connections will be in handled in this closure
|
||||||
proc handle(conn: Connection, proto: string) {.async.} =
|
proc handle(conn: Connection, proto: string) {.async.} =
|
||||||
|
|
||||||
var resp: array[6, byte]
|
var resp: array[6, byte]
|
||||||
await conn.readExactly(addr resp, 6)
|
await conn.readExactly(addr resp, 6)
|
||||||
check string.fromBytes(resp) == "client"
|
check string.fromBytes(resp) == "client"
|
||||||
|
@ -97,7 +96,7 @@ suite "Tor transport":
|
||||||
# We must close the connections ourselves when we're done with it
|
# We must close the connections ourselves when we're done with it
|
||||||
await conn.close()
|
await conn.close()
|
||||||
|
|
||||||
return T(codecs: @[TestCodec], handler: handle)
|
return T.new(codecs = @[TestCodec], handler = handle)
|
||||||
|
|
||||||
let rng = newRng()
|
let rng = newRng()
|
||||||
|
|
||||||
|
|
Loading…
Reference in New Issue