mirror of
https://github.com/codex-storage/nim-libp2p.git
synced 2025-01-23 01:09:32 +00:00
don't timeout in pubsub
This commit is contained in:
parent
1bd933cd5a
commit
2232ca598e
@ -32,19 +32,19 @@ type
|
||||
protocol*: LPProtocol
|
||||
match*: Matcher
|
||||
|
||||
MultisteamSelect* = ref object of RootObj
|
||||
MultistreamSelect* = ref object of RootObj
|
||||
handlers*: seq[HandlerHolder]
|
||||
codec*: string
|
||||
na: string
|
||||
ls: string
|
||||
|
||||
proc newMultistream*(): MultisteamSelect =
|
||||
proc newMultistream*(): MultistreamSelect =
|
||||
new result
|
||||
result.codec = MSCodec
|
||||
result.ls = Ls
|
||||
result.na = Na
|
||||
|
||||
proc select*(m: MultisteamSelect,
|
||||
proc select*(m: MultistreamSelect,
|
||||
conn: Connection,
|
||||
proto: seq[string]):
|
||||
Future[string] {.async.} =
|
||||
@ -81,7 +81,7 @@ proc select*(m: MultisteamSelect,
|
||||
trace "selected protocol", protocol = result
|
||||
break
|
||||
|
||||
proc select*(m: MultisteamSelect,
|
||||
proc select*(m: MultistreamSelect,
|
||||
conn: Connection,
|
||||
proto: string): Future[bool] {.async.} =
|
||||
if proto.len > 0:
|
||||
@ -89,10 +89,10 @@ proc select*(m: MultisteamSelect,
|
||||
else:
|
||||
result = (await m.select(conn, @[])) == Codec
|
||||
|
||||
proc select*(m: MultisteamSelect, conn: Connection): Future[bool] =
|
||||
proc select*(m: MultistreamSelect, conn: Connection): Future[bool] =
|
||||
m.select(conn, "")
|
||||
|
||||
proc list*(m: MultisteamSelect,
|
||||
proc list*(m: MultistreamSelect,
|
||||
conn: Connection): Future[seq[string]] {.async.} =
|
||||
## list remote protos requests on connection
|
||||
if not await m.select(conn):
|
||||
@ -108,7 +108,7 @@ proc list*(m: MultisteamSelect,
|
||||
|
||||
result = list
|
||||
|
||||
proc handle*(m: MultisteamSelect, conn: Connection) {.async, gcsafe.} =
|
||||
proc handle*(m: MultistreamSelect, conn: Connection) {.async, gcsafe.} =
|
||||
trace "handle: starting multistream handling"
|
||||
try:
|
||||
while not conn.closed:
|
||||
@ -152,7 +152,7 @@ proc handle*(m: MultisteamSelect, conn: Connection) {.async, gcsafe.} =
|
||||
finally:
|
||||
trace "leaving multistream loop"
|
||||
|
||||
proc addHandler*[T: LPProtocol](m: MultisteamSelect,
|
||||
proc addHandler*[T: LPProtocol](m: MultistreamSelect,
|
||||
codec: string,
|
||||
protocol: T,
|
||||
matcher: Matcher = nil) =
|
||||
@ -167,7 +167,7 @@ proc addHandler*[T: LPProtocol](m: MultisteamSelect,
|
||||
protocol: protocol,
|
||||
match: matcher))
|
||||
|
||||
proc addHandler*[T: LPProtoHandler](m: MultisteamSelect,
|
||||
proc addHandler*[T: LPProtoHandler](m: MultistreamSelect,
|
||||
codec: string,
|
||||
handler: T,
|
||||
matcher: Matcher = nil) =
|
||||
|
@ -43,6 +43,8 @@ proc isConnected*(p: PubSubPeer): bool =
|
||||
proc `conn=`*(p: PubSubPeer, conn: Connection) =
|
||||
trace "attaching send connection for peer", peer = p.id
|
||||
p.sendConn = conn
|
||||
p.sendConn.timeout = InfiniteDuration
|
||||
|
||||
p.onConnect.fire()
|
||||
|
||||
proc handle*(p: PubSubPeer, conn: Connection) {.async.} =
|
||||
|
@ -40,7 +40,7 @@ type
|
||||
transports*: seq[Transport]
|
||||
protocols*: seq[LPProtocol]
|
||||
muxers*: Table[string, MuxerProvider]
|
||||
ms*: MultisteamSelect
|
||||
ms*: MultistreamSelect
|
||||
identity*: Identify
|
||||
streamHandler*: StreamHandler
|
||||
secureManagers*: Table[string, Secure]
|
||||
|
Loading…
x
Reference in New Issue
Block a user