- 后端
- 运维
【免费下载链接】ajenti
Ajenti Core and stock plugins
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, })工作方式可以概括为:
- 浏览器与
/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, 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 时需要注意的约束:
- 单向通道: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 模块。理解这条链路后,无论是插件界面刷新、后台任务进度上报,还是自定义实时通知,都可以在几十行代码内基于既有基础设施落地。
- 后端
- 运维
【免费下载链接】ajenti
Ajenti Core and stock plugins
相关推荐
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.
后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考