Dart 用 Completer 等待长连接响应

Aug 13·6 min
AI 生成的摘要
商业 IM 软件把长连接事件包装成 Future,供认证、绑定和心跳等待响应。记录监听顺序、ID 匹配、超时与订阅清理。

写 HTTP 时通常是发请求,再 await 返回值。商业 IM 软件的长连接没有这么直接:发送走队列,收到的内容都进同一条事件流。认证、绑定和心跳需要各自从里面找响应。

项目封装了 sendAndWaitForResponseOptimized,把这件事包装成一个 Future。名字有点长,顺着实现看下来,主要是监听的时机和结束等待的条件需要注意。

先监听,是为了不漏掉很快的响应

假设按下面的顺序处理:先发送,再等待发送方法结束,最后才注册响应监听。

如果服务端的响应在中间已经到达,而这一层没有为请求保存响应,那么后注册的监听器就可能错过它。

当前项目先准备好接收入口,再发出请求。这里可以把顺序简化成:

snippet
text
创建等待结果的对象
  → 注册响应监听
  → 发送请求
  → 匹配响应并结束等待

先监听并不能解决全部并发问题,但它让请求发出去之前,就已经有了接住响应的位置。

Completer 把事件变成一次结果

接收流会不断产生事件,而一次请求通常只需要一个最终结果。

项目创建 Completer<String>,把它的 future 返回给等待方;监听器收到符合条件的数据时,再调用 complete。当前源码中的完成分支是:

snippet
dart
if (!completer.isCompleted) {
  completer.complete(data);
  cleanup();
}

这段保护用于避免重复完成。与此同时,调用方不需要自己管理整条接收流,只需要等待这一次请求。

这种封装真正需要定义的是完成规则:哪些事件算成功,哪些事件算失败,连接关闭算什么,以及没有结果时要等多久。只创建一个 Completer,这些规则不会自动出现。

收到了响应,还要确认属于谁

同一个连接上可能同时有心跳和其他请求。只判断收到的是不是 <iq>,范围太宽。

当前方法支持 elementNames、额外的 predicateexpectedId。其中 ID 的判断如下:

snippet
dart
if (expectedId != null && !data.contains(expectedId)) {
  return;
}

这能排除不包含目标 ID 的数据,但它仍然是字符串包含判断,不是解析 XML 后严格比较响应的 id 属性。

例如,request-1request-10 存在字符串包含关系;同一个数据块里也可能同时出现两条响应。字符串里出现了某个 ID,并不能证明与之匹配的元素就是这次请求需要的元素。

还有一个容易忽略的前提:这个方法监听的是 incomingStream。当前接收入口会先把每块解码后的字符串推入这个流,再追加到消息缓冲区。因此这里收到的事件,不保证已经是完整的协议元素。

如果开始标签或 ID 被拆开,后面的完整消息解析能力也不会自动修复这一条等待路径。

进一步整理时,我会考虑让等待层消费已经解析好的协议事件,再按请求 ID、类型和命名空间关联。目前还没有走这条处理路径。

发送失败了,还要继续等响应吗?

当前方法调用 sendXml(xmlToSend) 后,没有等待它返回的 Future<bool>,而是直接等待响应。

不等待发送结束,可以让接收逻辑继续工作,但也意味着发送失败没有直接结束这一次响应等待。假如发送层已经返回 false,这一层仍可能一直等到超时才报告失败。

这里要协调两个并发结果:发送操作的结果,以及接收流中的匹配响应。

一套更明确的契约应该说明:发送明确失败时怎样结束等待;如果响应已经先到达,迟到的发送错误又怎样处理。不能让两个分支各自随意修改最终结果。

isCompleted 可以避免重复完成的调用,但“哪个结果应该优先”仍然要由业务规则决定。

timeout 结束等待,不代表请求被撤销

当前实现给等待结果设置了超时,超时分支会取消监听并抛出异常:

snippet
dart
return await completer.future.timeout(
  timeout,
  onTimeout: () {
    cleanup();
    throw TimeoutException('等待 ${elementNames.join(", ")} 超时');
  },
);

Dart 的 Future.timeout 会产生一个受时间限制的新 Future;源 Future 仍可能在之后完成。这一点在 官方文档 中有明确说明。

因此,超时不能推导出服务端没有处理请求,也不会自动让已进入发送队列的内容消失。

超时后如果允许重试,还要继续考虑请求身份和迟到响应。否则上一轮的结果回来,可能被下一轮用过宽的匹配条件接收。

这和前面写过的消息发送确认是相通的:客户端停止等待,只描述了客户端这一侧的决定。

清理监听器,也有自己的完成时间

当前实现的清理函数比较短:

snippet
dart
void cleanup() {
  subscription?.cancel();
}

成功、超时和捕获异常的路径都会调用它。

cancel() 调用后,订阅不再接收事件;它返回的 Future 表示底层清理何时结束。如果后续操作必须等待资源释放,就需要等待这个 Future,而不只是发起取消。参见 StreamSubscription.cancel

如果下一步要复用底层资源,就不能漏掉这个等待。

当前监听也没有显式为 onErroronDone 定义如何完成请求等待。如果希望连接结束时立即让所有待处理请求失败,就需要在封装中安排对应的结束路径。

还需要补的情况

这类函数只测“发送后拿到响应”,很容易漏掉关键分支。我会重点检查:响应很快到达、两条响应乱序、ID 相似、响应分片、发送明确失败,以及超时之后才返回结果。

连接关闭也是一个结束条件。后面整理这个方法时,我会先补 onError 和 onDone,避免连接已经没了,调用方还要等到 timeout。

最后修改时间: Sep 15
cd ..