news 2026/9/27 7:13:41

Ajenti Core Push 推送服务解析:基于 Socket.IO 的实时消息广播架构

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Ajenti Core Push 推送服务解析:基于 Socket.IO 的实时消息广播架构
  • 后端
  • 运维

【免费下载链接】ajenti

Ajenti Core and stock plugins

项目地址:https://gitcode.com/gh_mirrors/aj/ajenti
点击查看免费下载

Ajenti 的aj.plugins.core.api.push模块提供了一个向浏览器客户端推送实时消息的服务,是任务进度、系统事件等异步通知得以即时呈现的关键通道。本文以该模块的 API 参考文档为骨架,结合仓库中的核心实现(push.py、broadcast_queue.py 及客户端服务),完整讲解 Push 服务的接口用法、广播队列底层原理、Socket.IO 收发链路,以及如何在插件与后台任务中实际调用它。

Push 服务的定位:服务端到客户端的单向广播

在 Ajenti 的架构中,HTTP 请求/响应模型天然是“请求驱动”的:浏览器发请求,服务端给结果。但很多场景需要服务端主动“说话”——例如后台任务运行到一半报告进度、任务完成或抛出异常、任务列表发生变化。Push服务就是为这类场景设计的单向消息广播通道,定义于 plugins/core/api/push.py:

@service class Push(): """ A service providing push messages to the client. """ def __init__(self, context): self.q = BroadcastQueue() def register(self): return self.q.register() def push(self, plugin, msg): """ Sends a push message to the client. :param plugin: routing ID :param msg: message """ self.q.broadcast((plugin, msg))

它通过 jadi 框架的@service装饰器注册为上下文单例服务,核心职责只有两个:

  • register():为调用方(通常是某个已连接的 Socket 会话)注册一个专属的接收队列,返回该队列供后续读取;
  • push(plugin, msg):向所有已注册的队列广播一条消息。消息被组织为(plugin, msg)二元组,其中plugin是路由 ID,决定这条消息在客户端被分发到哪个监听器。

广播队列 BroadcastQueue:弱引用 + gevent 队列

Push 服务的底层存储是一个BroadcastQueue,实现在 ajenti-core/aj/util/broadcast_queue.py。这是理解 Push 语义的关键,全文只有 20 行:

import weakref from gevent.queue import Queue class BroadcastQueue(): def __init__(self): self._queues = [] def register(self): q = Queue() self._queues.append(weakref.ref(q)) return q def broadcast(self, val): for q in list(self._queues): if q(): q().put(val) else: self._queues.remove(q)

几个值得注意的实现细节:

  • 每个注册者独占一个gevent.queue.Queue:register()每次都会新建独立的队列,因此不同客户端的消息读取互不干扰,先读先得,天然满足实时推送的“只关心最新消息”语义。
  • 弱引用(weakref.ref)管理订阅者生命周期:_queues列表中存放的是队列的弱引用而非强引用。当某个客户端断开、其队列不再被任何强引用持有而被 GC 回收后,broadcast()中q()返回None,该订阅会被自动从列表中移除。这保证了客户端断开后,订阅资源不会泄漏——即使没有显式的 unsubscribe 操作。
  • 广播期间安全迭代:broadcast()对list(self._queues)做快照遍历,一边投递一边清理失效引用,不会因迭代过程中修改列表而报错。

服务端接收端:PushSocket 端点与消息回发

Push 只是“投递队列”,真正把消息送到浏览器的是 core 插件中基于 Socket.IO 的端点 plugins/core/views/push.py。它通过@component(SocketEndpoint)注册,plugin = 'push'是该端点的 Socket 路由 ID:

@component(SocketEndpoint) class PushSocket(SocketEndpoint): plugin = 'push' def on_connect(self, message, *args): self.spawn(self._reader) def _reader(self): q = Push.get(self.context).register() while True: try: plugin, msg = q.get() except gevent.queue.Empty: return except EOFError: return if msg: self.send({ 'plugin': plugin, 'message': msg, })

工作方式可以概括为:

  1. 浏览器与/socket命名空间建立 Socket.IO 连接后,on_connect触发,端点通过self.spawn(self._reader)在 gevent 中派生一个读协程;
  2. _reader调用Push.get(self.context).register()拿到自己的广播队列,进入阻塞读取循环;
  3. 任何插件调用Push.push(plugin, msg)时,(plugin, msg)被广播进该队列,_reader取到后调用self.send(...),将{'plugin': ..., 'message': ...}结构体通过 Socket.IO 回发到客户端。

这里还体现了SocketEndpoint基类(定义于 ajenti-core/aj/api/http.py)的两个重要能力:

  • spawn(target):在端点内派生 greenlet,并记录到self.greenlets,客户端断开时destroy()会统一 kill,避免协程泄漏;
  • send(data, plugin=None):通过self.context.worker.send_to_upstream把消息按 Socket.IO 协议写给上游连接。

完整消息链路:从服务端 push() 到浏览器事件

把服务端与客户端两侧拼接起来,一次推送的完整链路是:

插件调用 Push.push(plugin, msg) │ ▼ BroadcastQueue.broadcast((plugin, msg)) # 投递到每个已注册队列 │ ▼ PushSocket._reader 从自己的队列取出 (plugin, msg) │ ▼ SocketEndpoint.send({'plugin':..., 'message':...}) │ (Socket.IO /socket 命名空间) ▼ 浏览器 socket.service 收到 'message' 事件 │ $rootScope.$broadcast('socket:push', ...) ▼ push.service 广播 'push:<plugin>' 事件 │ ▼ AngularJS 各模块 $on('push:<plugin>') 监听并刷新 UI

客户端侧有两个 AngularJS 服务支撑这条链路(均在 plugins/core/resources/js/core/services/ 下):

socket.service.es负责维护 Socket.IO 连接(支持断线重连,'max reconnection attempts': 999999),收到服务端message事件后解析 JSON 并广播 Angular 事件:

this.socket.on('message', msg => { if (msg[0] === '{') { msg = JSON.parse(msg); } $rootScope.$broadcast(`socket:${msg.plugin}`, msg.data); });

push.service.es再转发一层,把socket:push事件按路由 ID 二次广播为push:<plugin>事件,让各业务模块只需监听自己的路由即可:

$rootScope.$on('socket:push', ($event, msg) => { $rootScope.$broadcast(`push:${msg.plugin}`, msg.message); });

因此,一个完整的推送约定是:服务端Push.push('tasks', {...})→ 客户端$rootScope.$on('push:tasks', ...)。

实战:在插件与后台任务中发送 Push

方式一:直接注入 Push 服务

任何插件代码中,只要通过 jadi 的依赖注入拿到Push服务,即可随时推送:

from aj.plugins.core.api.push import Push Push.get(self.context).push('my-plugin', { 'type': 'info', 'message': 'Something happened', })

其中第一个参数plugin是路由 ID,客户端将收到push:my-plugin事件;第二个参数msg可以是任意可 JSON 序列化的对象(推荐使用字典,便于客户端按type等字段分派)。

方式二:从后台任务进程推送(推荐)

Push 最典型的应用场景是 core 插件的任务系统(Tasks)。任务运行在独立的 gipc 子进程中,无法直接访问主进程的 Push 单例,因此Task基类(plugins/core/api/tasks.py)提供了进程内接口Task.push(plugin, message):

def push(self, plugin, message): """ An interface to :class:`aj.plugins.core.api.push.Push` usable from inside the task's process """ self.pipe.put({ 'type': 'push', 'plugin': plugin, 'message': message, })

子进程把推送请求写入与主进程之间的 gipc 管道,主进程侧的_reader循环收到msg['type'] == 'push'后,才真正调用Push.get(self.context).push(msg['plugin'], msg['message'])完成广播(见 tasks.py)。这种“子进程 → 管道 → 主进程 → Push 广播”的桥接模式,让长时间运行的后台任务也能实时上报状态。

TasksService(同文件 tasks.py)在此基础上封装了两类标准推送负载:

  • notify(message):推送{'type': 'message', 'message': ...}到tasks路由,用于任务完成、异常等事件通知;
  • send_update():推送{'type': 'update', 'tasks': [...]}到tasks路由,用于刷新前端任务列表。

结合 tasks 模块的 API 参考 aj.plugins.core.api.tasks,可以看到任务系统与 Push 服务在设计上是互相配合的整体:任务进程只负责产出事件,Push 负责分发,浏览器端push:tasks监听器负责渲染。

使用约束与注意事项

从实现可以总结出几条实际使用 Push 时需要注意的约束:

  1. 单向通道:Push 只负责服务端 → 客户端方向的广播。客户端 → 服务端的反向消息走socket.service.send(plugin, data)(socket.emit('message', ...)),两者方向不同,不要混用。
  2. 消息不持久化:BroadcastQueue的语义是“广播给当前在线订阅者”,消息发送时未连接或已断开的客户端不会收到补发。需要保证送达的业务应当自行设计持久化或重试机制。
  3. 弱引用自动清理:断开连接的客户端队列会被broadcast()自动清理,无需手动注销;但这也意味着服务端无法通过 Push 获知“谁没收到”,可靠性完全取决于 Socket 连接本身。
  4. 路由 ID 即客户端命名空间:plugin参数在客户端被拼接为push:<plugin>事件名,因此应使用稳定、有语义的 ID(如tasks、push),并避免在消息负载中混入未序列化的对象。
  5. 依赖基础组件:Push 依赖 gevent 队列、Socket.IO(服务端端点与客户端 socket.io.js)以及 jadi 的服务/组件注入机制,这些都属于 Ajenti Core 的运行时底座,插件直接 import 即可,无需额外配置。

小结

aj.plugins.core.api.push是 Ajenti Core 提供的一条简洁而完整的实时推送通道:Push服务负责把消息广播进所有订阅队列,PushSocket端点负责把队列内容经 Socket.IO 送到浏览器,客户端socket.service与push.service负责把事件按路由 ID 分发到各 AngularJS 模块。理解这条链路后,无论是插件界面刷新、后台任务进度上报,还是自定义实时通知,都可以在几十行代码内基于既有基础设施落地。

  • 后端
  • 运维

【免费下载链接】ajenti

Ajenti Core and stock plugins

项目地址:https://gitcode.com/gh_mirrors/aj/ajenti
点击查看免费下载
上一篇:tsParticles 实战教程:三步为网站搭建动态粒子动画背景
下一篇:Dolphin-2.9.2-Phi-3-Medium编程能力实战:10个代码生成与调试案例详解

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/27 7:11:24

滨州装修开工前必做的 7 件事,少一件都容易耽误工期

很多滨州的朋友拿到新房钥匙&#xff0c;急着赶紧开工装修&#xff0c;恨不能当天就砸墙&#xff0c;结果开工没两天就因为手续不全被叫停&#xff0c;要么就是没准备好耽误工期。其实装修开工前的准备工作特别重要&#xff0c;准备做足了&#xff0c;后面才能顺顺利利&#xf…

作者头像 李华
网站建设 2026/9/27 7:09:14

图片裁剪框为什么会越界?四个方向的边界计算方法

在维护我自己的图片处理项目图片猫&#xff08;PicCat&#xff09;时&#xff0c;我重新检查了快速编辑器的裁剪逻辑。一个容易混淆的现象是&#xff1a;裁剪框已经被限制在图片内&#xff0c;拖到边缘时却仍会“滑走”。原因可能不是少写了一个判断&#xff0c;而是把移动整个…

作者头像 李华
网站建设 2026/9/27 7:08:08

AI论文写作工具怎么选?开题报告适用的8款工具对比

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/27 6:59:20

yolo下载地址

一、先搞清楚&#xff1a;你想下载哪个 YOLO&#xff1f; YOLO 系列并不是由一个团队统一维护的&#xff0c;不同版本分属不同作者/公司。下载前先认准"官方仓库"&#xff0c;避免下到第三方修改版&#xff1a; 版本 维护方 状态 YOLOv1~v3 Joseph Redmon (dark…

作者头像 李华