拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

Ajenti Core Push 推送服务解析:基于 Socket.IO 的实时消息广播架构 后端运维【免费下载链接】ajentiAjenti 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.pyservice 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.Queueregister()每次都会新建独立的队列因此不同客户端的消息读取互不干扰先读先得天然满足实时推送的“只关心最新消息”语义。弱引用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 路由 IDcomponent(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, })工作方式可以概括为浏览器与/socket命名空间建立 Socket.IO 连接后on_connect触发端点通过self.spawn(self._reader)在 gevent 中派生一个读协程_reader调用Push.get(self.context).register()拿到自己的广播队列进入阻塞读取循环任何插件调用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, pluginNone)通过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 tasks 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 时需要注意的约束单向通道Push 只负责服务端 → 客户端方向的广播。客户端 → 服务端的反向消息走socket.service.send(plugin, data)socket.emit(message, ...)两者方向不同不要混用。消息不持久化BroadcastQueue的语义是“广播给当前在线订阅者”消息发送时未连接或已断开的客户端不会收到补发。需要保证送达的业务应当自行设计持久化或重试机制。弱引用自动清理断开连接的客户端队列会被broadcast()自动清理无需手动注销但这也意味着服务端无法通过 Push 获知“谁没收到”可靠性完全取决于 Socket 连接本身。路由 ID 即客户端命名空间plugin参数在客户端被拼接为push:plugin事件名因此应使用稳定、有语义的 ID如tasks、push并避免在消息负载中混入未序列化的对象。依赖基础组件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 模块。理解这条链路后无论是插件界面刷新、后台任务进度上报还是自定义实时通知都可以在几十行代码内基于既有基础设施落地。赞分享后端运维【免费下载链接】ajentiAjenti Core and stock plugins项目地址https://gitcode.com/gh_mirrors/aj/ajenti点击查看免费下载相关推荐CodeGuide 系列实战Netty 4.1 服务端群发消息——基于 ChannelGroup 实现广播推送CodeGuide 系列实战Netty 4.1 服务端群发消息——基于 ChannelGroup 实现广播推送 在微信、QQ 的群聊场景中用户发出的每一条消文档教程后端uni-app推送服务uni-push跨端消息推送uni app推送服务uni push跨端消息推送 概述 还在为多端消息推送的兼容性问题头疼吗uni app的uni push服务提供了统一的跨平台消息推送前端跨平台移动开发小程序Sails 实时广播全指南掌握 sails.sockets.blast() 实现全服务器消息推送Sails 实时广播全指南掌握 sails.sockets.blast 实现全服务器消息推送 blast 是 Sails 框架 sails.sockets.后端上一篇tsParticles 实战教程三步为网站搭建动态粒子动画背景下一篇Dolphin-2.9.2-Phi-3-Medium编程能力实战10个代码生成与调试案例详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表