fix quicClient
This commit is contained in:
parent
270e2ac81b
commit
1d19f727d3
@ -268,36 +268,18 @@ actor SDLContextActor {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 这里必须等待quic的协商完成
|
// 这里必须等待quic的协商完成
|
||||||
try await Task.sleep(for: .seconds(0.3))
|
try await Task.sleep(for: .seconds(0.5))
|
||||||
SDLLogger.log("[SDLContext] start quic client: \(self.config.serverHost)")
|
SDLLogger.log("[SDLContext] start quic client: \(self.config.serverHost)")
|
||||||
|
|
||||||
try await withThrowingTaskGroup { group in
|
try await withThrowingTaskGroup { group in
|
||||||
defer {
|
// 创建一个简单的异步状态等待机制(可以用一个 Actor 或者 AsyncStream 模拟)
|
||||||
group.cancelAll()
|
let (readyStream, readyContinuation) = AsyncStream<Void>.makeStream()
|
||||||
}
|
|
||||||
|
|
||||||
group.addTask {
|
|
||||||
for await message in quicClient.messageStream {
|
|
||||||
await self.handleQUICMessage(message: message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
group.addTask {
|
|
||||||
let timerStream = SDLAsyncTimerStream()
|
|
||||||
timerStream.start(interval: .seconds(5))
|
|
||||||
|
|
||||||
for await _ in timerStream.stream {
|
|
||||||
if Task.isCancelled {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
quicClient.send(type: .ping, data: Data())
|
|
||||||
}
|
|
||||||
SDLLogger.log("[SDLQUICClient] udp pingTask cancel", for: .debug)
|
|
||||||
}
|
|
||||||
|
|
||||||
group.addTask {
|
group.addTask {
|
||||||
for await event in quicClient.eventStream {
|
for await event in quicClient.eventStream {
|
||||||
switch event {
|
switch event {
|
||||||
|
case .ready:
|
||||||
|
readyContinuation.yield()
|
||||||
case .failed(let error):
|
case .failed(let error):
|
||||||
throw error
|
throw error
|
||||||
case .cancelled:
|
case .cancelled:
|
||||||
@ -308,9 +290,36 @@ actor SDLContextActor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
group.addTask {
|
||||||
|
// 等待信号
|
||||||
|
var it = readyStream.makeAsyncIterator()
|
||||||
|
await it.next()
|
||||||
|
|
||||||
|
try Task.checkCancellation()
|
||||||
|
|
||||||
|
await withThrowingTaskGroup { workerGroup in
|
||||||
|
workerGroup.addTask {
|
||||||
|
for await message in quicClient.messageStream {
|
||||||
|
try Task.checkCancellation()
|
||||||
|
await self.handleQUICMessage(message: message)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
workerGroup.addTask {
|
||||||
|
let timerStream = SDLAsyncTimerStream()
|
||||||
|
timerStream.start(interval: .seconds(5))
|
||||||
|
|
||||||
|
for await _ in timerStream.stream {
|
||||||
|
try Task.checkCancellation()
|
||||||
|
quicClient.send(type: .ping, data: Data())
|
||||||
|
}
|
||||||
|
SDLLogger.log("[SDLQUICClient] udp pingTask cancel", for: .debug)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
try await group.next()
|
try await group.next()
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private func handleQUICMessage(message: SDLQUICInboundMessage) async {
|
private func handleQUICMessage(message: SDLQUICInboundMessage) async {
|
||||||
|
|||||||
@ -22,6 +22,7 @@ enum SDLQUICError: Error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
enum SDLQUICEvent: Error {
|
enum SDLQUICEvent: Error {
|
||||||
|
case ready
|
||||||
case failed(Error)
|
case failed(Error)
|
||||||
case cancelled
|
case cancelled
|
||||||
case writeFailed(Error)
|
case writeFailed(Error)
|
||||||
@ -71,6 +72,7 @@ final class SDLQUICClient {
|
|||||||
switch state {
|
switch state {
|
||||||
case .ready:
|
case .ready:
|
||||||
self?.startReadTask()
|
self?.startReadTask()
|
||||||
|
self?.eventCont.yield(.ready)
|
||||||
case .failed(let error):
|
case .failed(let error):
|
||||||
self?.eventCont.yield(.failed(error))
|
self?.eventCont.yield(.failed(error))
|
||||||
case .cancelled:
|
case .cancelled:
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user