Skip to content

fix: reconnect listener instead of waiting for dispatcher restart - #259

Merged
deepin-bot[bot] merged 2 commits into
linuxdeepin:develop/snipefrom
wangrong1069:pr0921-2
Sep 21, 2026
Merged

deepin-bot[bot] merged 2 commits into
linuxdeepin:develop/snipefrom
wangrong1069:pr0921-2

Conversation

@wangrong1069

@wangrong1069 wangrong1069 commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

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:

  • Reconnect the event listener automatically after dispatcher connection loss, with delayed retries and a bounded failure limit.
  • Drain all available dispatcher events per readable callback using non-blocking receives.
  • Add socket buffering and subscription path fallbacks to better handle event bursts and unresolved mount paths.

Bug Fixes:

  • Prevent the daemon from restarting immediately when the listener exhausts reconnection attempts; exit instead.
  • Reduce disconnects caused by socket buffer exhaustion during high-volume event bursts.
  • Preserve subscriptions when full path resolution fails by falling back to the expanded path.

Enhancements:

  • Replace dispatcher socket inode polling with direct connection recreation and fd monitoring.

… 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: 提升高吞吐场景下事件分发链路的稳定性,避免客户端因缓冲区不足被断开,并降低每事件派发开销。
@sourcery-ai

sourcery-ai Bot commented Sep 21, 2026

Copy link
Copy Markdown

Reviewer's Guide

The 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 draining

sequenceDiagram
    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
Loading

Flow diagram for resilient subscription path resolution

flowchart 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]
Loading

File-Level Changes

Change Details Files
Replace dispatcher restart detection with a bounded delayed reconnect state machine.
  • Remove socket inode polling and restart-check limits.
  • Schedule one-shot reconnect attempts every 10 seconds, recreating the receiver and GLib fd source on success.
  • Quit after 30 failed attempts or timer/source allocation failures.
  • Update shutdown cleanup and quit handling for reconnect timers and permanent failure.
src/daemon/src/core/event_listener.cpp
src/daemon/src/main.cpp
Drain dispatcher events non-blockingly during each fd callback.
  • Add a non-blocking receive API and WOULD_BLOCK result.
  • Loop over pending events while continuing on interruptions and return when the socket is drained.
  • Trigger reconnect on disconnect or receive errors.
src/daemon/src/core/event_listener.cpp
src/dispatcher/event_dispatcher.h
src/dispatcher/event_receiver.c
Increase Unix socket buffering to tolerate event bursts.
  • Configure a 1 MiB send buffer for accepted dispatcher clients.
  • Configure a 1 MiB receive buffer for event listeners.
  • Log buffer configuration failures without aborting the connection.
src/dispatcher/event_dispatcher.c
src/dispatcher/event_receiver.c
Preserve subscriptions when full path resolution is unavailable.
  • Fall back to the expanded subscription path when mount/path resolution fails.
  • Continue enforcing trailing-slash prefix matching for fallback paths.
src/server/event-dispatcher.c

Tips and commands

Interacting with Sourcery

  • Trigger a new review: Comment @sourcery-ai review on the pull request.
  • Continue discussions: Reply directly to Sourcery's review comments.
  • Generate a GitHub issue from a review comment: Ask Sourcery to create an
    issue from a review comment by replying to it. You can also reply to a
    review comment with @sourcery-ai issue to create an issue from it.
  • Generate a pull request title: Write @sourcery-ai anywhere in the pull
    request title to generate a title at any time. You can also comment
    @sourcery-ai title on the pull request to (re-)generate the title at any time.
  • Generate a pull request summary: Write @sourcery-ai summary anywhere in
    the pull request body to generate a PR summary at any time exactly where you
    want it. You can also comment @sourcery-ai summary on the pull request to
    (re-)generate the summary at any time.
  • Generate reviewer's guide: Comment @sourcery-ai guide on the pull
    request to (re-)generate the reviewer's guide at any time.
  • Resolve all Sourcery comments: Comment @sourcery-ai resolve on the
    pull request to resolve all Sourcery comments. Useful if you've already
    addressed all the comments and don't want to see them anymore.
  • Dismiss all Sourcery reviews: Comment @sourcery-ai dismiss on the pull
    request to dismiss all existing Sourcery reviews. Especially useful if you
    want to start fresh with a new review - don't forget to comment
    @sourcery-ai review to trigger a new review!

Customizing Your Experience

Access your dashboard to:

  • Enable or disable review features such as the Sourcery-generated pull request
    summary, the reviewer's guide, and others.
  • Change the review language.
  • Add, remove or edit custom review instructions.
  • Adjust other review settings.

Getting Help

@sourcery-ai sourcery-ai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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>

Sourcery is free for open source - if you like our reviews please consider sharing them ✨

Comment thread src/daemon/src/core/event_listener.cpp
@deepin-ci-robot

Copy link
Copy Markdown

deepin pr auto review

🤖 AI 代码审查报告

总体评分: 78 分 (通过阈值: 70分)

Fail


📊 总体评价

项目 结果
审查结论 代码审查通过(有条件)
评分详情 代码实现了重连机制替代重启检测、非阻塞事件排空和缓冲区扩大,整体设计合理且注释充分。但存在1个逻辑缺陷:reconnect_attempts重置时机不当可能导致无限重连循环。无安全漏洞。
漏洞对比统计 新增漏洞 0 个,减少漏洞 0 个,持平 0 个

🔍 详细分析

1. 语法逻辑 ❌

评价: 存在逻辑缺陷 ❌ 不通过 (17/25分)

潜在问题:

  1. src/daemon/src/core/event_listener.cpp 第79-82行,函数 on_reconnect_attempt:reconnect_attempts 在 event_receiver_new 成功时立即重置为0(第82行),但如果后续 fd source 设置失败(无效fd、source创建失败、attach失败),计数器仅递增到1。下次重连时若 event_receiver_new 再次成功,计数器又重置为0,导致无限重连循环——监听器永远无法达到 MAX_RECONNECT_ATTEMPTS 并退出。正确做法是仅在整条重连路径(包括fd source设置)全部成功后才重置 reconnect_attempts。

建议: 将 reconnect_attempts = 0 的重置操作移到 fd source 成功 attach 之后(第115行 return G_SOURCE_REMOVE 之前),确保只有完全成功的重连才重置计数器。


2. 代码质量 ✅

评价: 代码结构清晰,注释完整 ✅ 通过 (23/25分)

潜在问题:

  1. src/daemon/src/core/event_listener.cpp 第103-104行、第281-282行,函数 on_reconnect_attempt/event_listener_thread_func:函数指针转换 (GSourceFunc)(void (*)(void))on_fd_readable 虽然是 GLib 常用模式,但技术上属于未定义行为。goto schedule_retry 模式虽然可读性尚可,但增加了控制流复杂度。

建议: 函数指针转换是 GLib 惯用模式,可保持现状。goto 模式可考虑使用辅助函数进一步简化错误处理路径。


3. 代码性能 ✅

评价: 性能良好,资源使用合理 ✅ 通过 (20/20分)

潜在问题: 无

建议: 性能改进显著:非阻塞排空循环避免高频事件下GMainContext重调度开销;1MiB缓冲区扩大减少EAGAIN断连;10s重连延迟避免事件突发期间频繁重连。


4. 代码安全 🔒

评价: 存在0个安全漏洞 ✅ 通过 (30/30分)

🔐 发现 0 个安全漏洞

安全漏洞详情: 无安全漏洞

建议: 无安全风险。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;
    }
    // ...
}

📋 审查信息

项目 内容
PR #259
标题 fix: reconnect listener instead of waiting for dispatcher restart
作者 wangrong1069
分支 pr0921-2 → develop/snipe
修改文件 6 个文件 (+230/-88)
分析模式 全量分析
OCR 审查 已完成(发现1个bug)

本报告由 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分钟后放弃并退出;用户订阅路径解析失败时不再被丢弃,事件分发更可靠。
@deepin-ci-robot

Copy link
Copy Markdown

[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.

Details Needs approval from an approver in each of these files:

Approvers can indicate their approval by writing /approve in a comment
Approvers can cancel approval by writing /approve cancel in a comment

@wangrong1069

Copy link
Copy Markdown
Contributor Author

/merge

@deepin-bot
deepin-bot Bot merged commit 97fbe73 into linuxdeepin:develop/snipe Sep 21, 2026
16 checks passed
@wangrong1069
wangrong1069 deleted the pr0921-2 branch September 21, 2026 07:34
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants