fix: reconnect listener instead of waiting for dispatcher restart - #259
Conversation
… events Add non-blocking event_receiver_receive_nonblock() and use it in a draining loop in event_listener to avoid GMainContext re-dispatch overhead per event. Enlarge SO_SNDBUF/SO_RCVBUF to 1 MiB on both dispatcher and receiver so bursts of 4 KB dispatch_event_t messages no longer hit EAGAIN and get the client kicked. 新增非阻塞接收接口 event_receiver_receive_nonblock(),event_listener 改为循环排空内核缓冲区,减少每事件 GMainContext 重派发开销;分发端与 接收端均将 socket 缓冲区扩大至 1 MiB,避免 4 KB 事件突发时 EAGAIN 导致服务端踢掉客户端。 Log: 修复高吞吐场景下事件接收端被踢与逐事件派发开销问题 PMS: BUG-370779, BUG-370781 Influence: 提升高吞吐场景下事件分发链路的稳定性,避免客户端因缓冲区不足被断开,并降低每事件派发开销。
Reviewer's GuideThe PR changes listener recovery from waiting for a dispatcher restart to actively retrying a delayed reconnect, while draining socket bursts non-blockingly and enlarging Unix socket buffers to reduce disconnects. It also makes subscription path resolution resilient by falling back to expanded paths when mount metadata is unavailable. Sequence diagram for delayed listener reconnect and event drainingsequenceDiagram
participant Dispatcher
participant Listener
participant Timer
participant EventReceiver
participant MainLoop
Dispatcher-->>Listener: Connection closes or receive error
Listener->>EventReceiver: event_receiver_free()
Listener->>Timer: g_timeout_source_new(RECONNECT_DELAY_MS)
Timer-->>Listener: on_reconnect_attempt()
Listener->>EventReceiver: event_receiver_new(DISPATCHER_SOCKET_PATH)
alt reconnect succeeds
Listener->>MainLoop: g_unix_fd_source_new(socket_fd)
Listener-->>Dispatcher: Reconnected
Dispatcher-->>MainLoop: Socket readable
MainLoop->>Listener: on_fd_readable()
loop pending events
Listener->>EventReceiver: event_receiver_receive_nonblock()
EventReceiver-->>Listener: EVENT_RECEIVE_OK
Listener-->>Listener: callback(event)
end
EventReceiver-->>Listener: EVENT_RECEIVE_WOULD_BLOCK
else reconnect fails
Listener->>Timer: schedule next attempt after RECONNECT_DELAY_MS
end
Flow diagram for resilient subscription path resolutionflowchart TD
P[User subscription path] --> E[Expand path]
E --> R{get_full_path succeeds?}
R -->|yes| F[Resolved full path]
R -->|no| B[Use expanded path fallback]
F --> T[Ensure trailing slash]
B --> T
T --> O[Add subscription prefix]
File-Level Changes
Tips and commandsInteracting with Sourcery
Customizing Your ExperienceAccess your dashboard to:
Getting Help
|
There was a problem hiding this comment.
Hey - I've found 1 issue
Prompt for AI Agents
Please address the comments from this code review:
## Individual Comments
### Comment 1
<location path="src/daemon/src/core/event_listener.cpp" line_range="195-197" />
<code_context>
- default:
- return G_SOURCE_CONTINUE;
+ /* Drain all pending events in one callback to avoid GMainContext
+ * re-dispatch overhead per event under high throughput. The fd source
+ * is level-triggered, so it will fire again if the kernel buffer still
+ * has data after we return. */
+ for (;;) {
+ dispatch_event_t dispatch_evt;
</code_context>
<issue_to_address>
**issue (broader_impact):** When the socket reports both readable data and `G_IO_HUP` or `G_IO_ERR`, the callback enters the HUP/ERR branch before draining the socket, frees the receiver, and reconnects, so events already queued in the kernel are discarded.
**Triggers:** When the dispatcher closes or errors the connection while one or more complete event packets remain unread.
**Suggested fix:** Drain readable packets before treating HUP/ERR as a connection loss, or explicitly document and preserve the queued packets before reconnecting.
</issue_to_address>
deepin pr auto review🤖 AI 代码审查报告📊 总体评价
🔍 详细分析1. 语法逻辑 ❌评价: 存在逻辑缺陷 ❌ 不通过 (17/25分) 潜在问题:
建议: 将 reconnect_attempts = 0 的重置操作移到 fd source 成功 attach 之后(第115行 return G_SOURCE_REMOVE 之前),确保只有完全成功的重连才重置计数器。 2. 代码质量 ✅评价: 代码结构清晰,注释完整 ✅ 通过 (23/25分) 潜在问题:
建议: 函数指针转换是 GLib 惯用模式,可保持现状。goto 模式可考虑使用辅助函数进一步简化错误处理路径。 3. 代码性能 ✅评价: 性能良好,资源使用合理 ✅ 通过 (20/20分) 潜在问题: 无 建议: 性能改进显著:非阻塞排空循环避免高频事件下GMainContext重调度开销;1MiB缓冲区扩大减少EAGAIN断连;10s重连延迟避免事件突发期间频繁重连。 4. 代码安全 🔒评价: 存在0个安全漏洞 ✅ 通过 (30/30分)
安全漏洞详情: 无安全漏洞 建议: 无安全风险。SO_PEERCRED认证保持不变,recv使用sizeof(*event)防止缓冲区溢出,g_strlcpy确保路径拷贝安全,get_full_path失败时使用expanded路径作为fallback是合理的降级行为。 漏洞对比统计:新增漏洞 0 个,减少漏洞 0 个,持平 0 个 💡 改进建议代码示例// 修复方案:仅在整条重连路径成功后重置 reconnect_attempts
static gboolean on_reconnect_attempt(gpointer data)
{
EventListener *listener = (EventListener *)data;
listener->reconnect_timer_id = 0;
EventReceiver *new_receiver = event_receiver_new(DISPATCHER_SOCKET_PATH);
if (new_receiver != NULL) {
listener->receiver = new_receiver;
// 删除此处的 reconnect_attempts = 0;
int sock_fd = event_receiver_get_socket(listener->receiver);
if (sock_fd < 0) {
spdlog::error("Invalid receiver socket fd after reconnect");
event_receiver_free(listener->receiver);
listener->receiver = NULL;
listener->reconnect_attempts++;
goto schedule_retry;
}
GSource *fd_source = g_unix_fd_source_new(sock_fd,
(GIOCondition)(G_IO_IN | G_IO_HUP | G_IO_ERR));
if (fd_source == NULL) {
spdlog::error("Failed to create fd watch source");
event_receiver_free(listener->receiver);
listener->receiver = NULL;
listener->reconnect_attempts++;
goto schedule_retry;
}
g_source_set_priority(fd_source, G_PRIORITY_DEFAULT);
g_source_set_callback(fd_source,
(GSourceFunc)(void (*)(void))on_fd_readable, listener, NULL);
listener->fd_source_id = g_source_attach(fd_source, listener->context);
g_source_unref(fd_source);
if (listener->fd_source_id == 0) {
spdlog::error("Failed to attach fd watch to listener context");
event_receiver_free(listener->receiver);
listener->receiver = NULL;
listener->reconnect_attempts++;
goto schedule_retry;
}
// 仅在全部成功后重置计数器
listener->reconnect_attempts = 0;
spdlog::info("Reconnected to dispatcher: {}", DISPATCHER_SOCKET_PATH);
return G_SOURCE_REMOVE;
}
// ...
}📋 审查信息
本报告由 AI 代码审查工具自动生成 |
Replace inode-polling restart check with a reconnect state machine (10s interval, 30 attempts, then quit). Also fall back to the expanded path when get_full_path fails in the dispatcher so user subscriptions are not silently dropped. 将基于inode轮询的重启检测替换为重连状态机(每10秒一次,最多 30次,超限后退出)。同时修复调度器中get_full_path失败时订阅 路径被静默丢弃的问题,回退使用展开路径。 Log: 重连逻辑重写及订阅路径回退修复 PMS: TASK-395673 Influence: 守护进程与调度器断连后可自动重连恢复事件接收,最长约5分钟后放弃并退出;用户订阅路径解析失败时不再被丢弃,事件分发更可靠。
ceae4d8 to
c54db10
Compare
|
[APPROVALNOTIFIER] This PR is NOT APPROVED This pull-request has been approved by: lzwind, wangrong1069 The full list of commands accepted by this bot can be found here. DetailsNeeds approval from an approver in each of these files:Approvers can indicate their approval by writing |
|
/merge |
Summary by Sourcery
Make event delivery resilient to dispatcher restarts and high-volume bursts by reconnecting the listener directly, buffering socket traffic, and preserving unresolved subscriptions.
New Features:
Bug Fixes:
Enhancements: