ekfa/CODE_LOGIC_OVERVIEW.md
2026-05-11 22:42:00 +08:00

11 KiB
Raw Blame History

EFKA 当前代码逻辑梳理

本文档记录当前阅读代码后对项目逻辑的理解,主要基于 apps/efka/src/apps/efka/include/config/ 和 protobuf 生成文件。尚未覆盖每个边界条件和全部测试行为。

1. 项目定位

这个项目是一个 Erlang/OTP 应用,核心作用像一个边缘侧 agent

  • 对本机微服务开放 WebSocket 接入,微服务注册、订阅 topic、上报指标。
  • 与上游服务器建立 TLS 长连接,完成鉴权、上报数据、接收远程命令。
  • 通过 Docker API 管理本机容器,包括部署、启动、停止、杀掉、删除、更新配置和列表查询。
  • 用 DETS 做少量本地持久化,包括服务状态和离线缓存包。

依赖主要有:

  • cowboy:本地 WebSocket server。
  • gun:通过 Unix Socket 访问 Docker HTTP API。
  • ssl:连接上游 TLS server。
  • OTP json 模块JSON 编解码。
  • gpbprotobuf 代码生成。

2. 启动流程

应用入口是 efka_app:start/2

启动步骤:

  1. 设置 Erlang IO 编码为 unicode。
  2. 设置 fullsweep_after,试图加快进程内存回收。
  3. 启动 Cowboy WebSocket server。
  4. 启动顶层 supervisor efka_sup

WebSocket server

  • config/sys.config 读取 efka.websocket_server
  • 监听 /ws
  • 请求处理模块是 service_channel

efka_sup 启动的子进程:

  • efka_logger:部署日志落盘。
  • efka_service_sup:动态管理每个已注册微服务对应的 efka_service 进程。
  • cache_modelDETS 离线缓存。
  • service_modelDETS 服务状态表。
  • efka_subscription:本地 topic 订阅中心。
  • efka_client:连接上游 TLS server 的状态机。

docker application 启动的子进程:

  • docker_task_reporter:部署任务事件缓冲和转发。
  • docker_deploy_manager:部署任务管理器。

docker_events 代码存在于 apps/docker,但目前默认不会启动。

3. 两条通信链路

项目里有两条主要通信链路。

3.1 微服务到 EFKAWebSocket

本机微服务连接 ws://<host>:18080/ws

协议处理在 service_channel

  • WebSocket 收到 ping 直接回复 pong
  • 收到 binary 帧时,第一个字节是帧类型。
  • FRAME_REQUEST = 0x01,用 service_pb 解码为 ServiceRequest
  • FRAME_CAST = 0x03,用 service_pb 解码为 ServiceCast
  • 回复时使用 FRAME_REPLY = 0x02

支持的微服务请求:

  • register:注册服务。
  • subscribe:订阅 topic。

支持的微服务 cast

  • metric_data:上报指标数据。

注册流程:

  1. 微服务发送 ServiceRequest.Register,包含 service_id
  2. service_channel 调用 efka_service_sup:start_service(ServiceId)
  3. 如果该服务进程不存在,则动态启动一个 efka_service;如果已存在,复用已有 pid。
  4. service_channel 调用 efka_service:attach_channel(ServicePid, self()) 绑定当前 WebSocket channel。
  5. efka_service 只允许绑定一个存活 channel已有 channel 存活时返回 channel exists
  6. 注册成功后写入 service_model,服务状态设为 running。
  7. channel 退出时,取消全部订阅,并把服务状态改为 stopped。

指标上报流程:

  1. 微服务发送 ServiceCast.MetricData
  2. service_channel 调用 efka_service:metric_data(ServicePid, RouteKey, Metric)
  3. efka_service 转发给 efka_client:metric_data(RouteKey, Metric)
  4. efka_client 如果处于 activated 状态,直接发给上游;否则写入 cache_model 离线缓存。

3.2 EFKA 到上游TLS 长连接

上游连接逻辑在 efka_client,它是一个 gen_statem

状态包括:

  • disconnected:未连接。
  • auth:已连接,等待鉴权响应。
  • restricted:鉴权受限,不能正常推送数据,但可以接收部分授权命令。
  • activated:已激活,可正常收发。

连接流程:

  1. 初始化后立即触发 create_transport
  2. 读取 efka.tls_server_address,用 ssl:connect/4 建立 TLS 连接。
  3. TLS socket 使用 {packet, 4}{active, true}
  4. 连接成功后发送 AuthRequest
  5. 收到鉴权成功 reply 后进入 activated,并触发 flush_cache
  6. 连接或鉴权失败时关闭 socket5 秒后重连。

上游协议同样用第一个字节区分帧类型:

  • FRAME_REQUEST = 0x01
  • FRAME_REPLY = 0x02
  • FRAME_CAST = 0x03

但 protobuf 使用的是 message_pb,不是微服务侧的 service_pb

efka_client 上报的内容:

  • metric_data业务指标数据activated 时实时发送,否则进入 DETS 缓存。
  • task_event_streamDocker 部署任务流式日志,只在 activated 时发送。
  • close_task_event_stream:任务结束事件,只在 activated 时发送。

efka_client 接收的内容:

  • RequestFrame.container_request:远程容器管理请求,交给 docker_container_service
  • CastFrame.command:授权命令,COMMAND_AUTH 用于切换鉴权/受限状态。
  • CastFrame.pub:上游发布的 topic 消息,交给 efka_subscription:publish/3

4. 本地 Pub/Sub 逻辑

efka_subscription 是本地订阅中心。

订阅:

  • service_channel 收到微服务 subscribe 后调用 efka_subscription:subscribe(Topic, self())
  • 同一个 pid 对同一个 topic 重复订阅会被忽略。
  • 首次看到某个 subscriber pid 时,会 monitor 该 pidpid 退出后自动移除订阅。

topic 匹配规则:

  • / 分隔 topic components。
  • * 表示单级匹配。
  • + 表示多级匹配,但只能出现在末尾。
  • 完全匹配优先级最高,不过当前代码里的 order 字段只是计算并保存,没有实际排序使用。

发布:

  • 上游发来的 CastFrame.Pub 会调用 efka_subscription:publish(Topic, Qos, Content)
  • 如果匹配到订阅者,向对应 service_channel 发送 {topic_broadcast, Topic, Content}
  • service_channel 再编码为 ServiceCast.TopicEvent 推送给微服务。
  • 如果没有订阅者且 Qos != 0,消息会暂存在 remaining_messages
  • 新订阅建立时,会尝试把匹配的 remaining messages 补发给该订阅者。

5. Docker 容器管理链路

上游发来的容器请求最终进入 docker_container_service:handle_request/1

支持动作:

  • list:列出容器。
  • deploy:部署容器。
  • start:启动容器。
  • stop:停止容器。
  • kill:发送 kill 信号。
  • remove:删除容器。
  • config:更新容器配置文件。

简单操作:

  • start/stop/kill/remove/list 直接调用 docker_commands
  • docker_commands 通过短生命周期普通 docker_client 进程访问 /var/run/docker.sock
  • HTTP 客户端用 gun:open_unix/2
  • 流式请求使用 {docker_client, Ref, Event} 消息返回,事件包括 {response, Status, Headers}{data, Bin}done{error, Reason}

部署操作:

  1. docker_container_service 收到 deploy 后调用 docker_deploy_manager:deploy(TaskId, Params)
  2. docker_deploy_manager 根据 docker.root_dir 和参数里的 container_name/container_dir 确保容器目录存在。
  3. 它启动一个独立 docker_deployer 进程,并 monitor 该部署进程。
  4. docker_deployer 执行实际部署步骤。
  5. 部署过程事件写入 docker_task_reporter
  6. docker_task_reporter 在上游连接 activated 时把事件转给 efka_client;未激活时保留队列并定时重试。

docker_deployer 当前部署步骤:

  1. 上报“开始部署容器”。
  2. 调用 ensure_container_absent/2,但当前实现只是上报“开始创建容器”,没有真正删除旧容器。
  3. 规范化镜像名,没有 tag 时补 :latest
  4. 调用 Docker API 拉取镜像,docker_deployerRef 接收 docker_client 流式消息并上报拉取过程。
  5. 创建容器。
  6. 创建空的 service.conf 配置文件。
  7. 写部署摘要到 efka_logger
  8. 关闭任务事件流,状态为 successfail

创建容器参数:

  • docker_container_builder 把 protobuf 的 DockerCreateOptions 转成 Docker API JSON。
  • 会自动给容器注入环境变量 CONTAINER_NAME=<name>
  • 会自动添加容器内配置路径 /usr/local/etc/service.conf
  • 会自动把宿主机 service.conf bind 到容器内 /usr/local/etc/service.conf

配置更新:

  • docker_container_service:update_container_config/2 根据容器名找到目录。
  • 写入该目录下的 service.conf
  • 容器目录映射由 docker_helper 通过 .container_dir 指针文件维护。

6. 本地持久化

6.1 service_model

使用 DETS 表 service

文件位置来自 efka.dets_dir,默认配置是:

"/usr/local/code/tmp/dets/service.dets"

记录结构在 apps/efka/include/efka_tables.hrl

  • service_id
  • container_name
  • meta_data
  • status
  • create_ts
  • update_ts

用途:

  • 注册成功时插入或更新服务。
  • channel 关闭时把服务状态改成 stopped。
  • 支持查询所有服务、运行中服务和单个服务状态。

6.2 cache_model

使用 DETS 表 cache

作用:

  • efka_client 不在 activated 状态时,把待上报 packet 缓存下来。
  • efka_client 激活后循环 fetch_next -> send -> delete 刷缓存。

缓存 id 使用 os:system_time(microsecond) 生成。

7. 日志

有两套日志:

  • OTP loggerconfig/sys.config 里配置到 console 和 log/debug.log
  • efka_logger:自定义部署日志,按日期写到 code:root_dir() ++ "/log/" 下。

efka_logger 主要用于部署任务摘要和失败原因。

8. 配置项

关键配置来自 config/sys.config

  • docker.root_dir:容器相关目录根路径。
  • dets_dirDETS 文件目录。
  • websocket_server:本地 WebSocket 监听配置。
  • tls_server_address:上游 TLS server 地址。
  • auth:上游鉴权信息。

注意:

  • service_modelcache_model 打开 DETS 前假设 dets_dir 已存在,代码里没有显式创建目录。
  • docker_client 固定使用 /var/run/docker.sock,每次请求内部由短生命周期普通进程执行 open/request/close。

9. 当前看到的几个注意点

  • README.md 描述的是 JSON-RPC 风格 WebSocket但当前代码实际处理的是 binary protobuf 帧README 可能已经过期。
  • service_channelregister 只使用 service_id,没有使用 README 中提到的 meta_data/container_name
  • efka_client:send_result_reply/3send_error_reply/3 编码 ReplyFrame 后没有加 FRAME_REPLY 前缀;接收侧是否期望裸 protobuf 需要确认。
  • docker_deployer:ensure_container_absent/2 当前没有真正确保旧容器不存在,只是上报日志。
  • efka_subscription 计算了 topic order,但匹配广播时没有使用优先级排序。
  • cache_model 使用 DETS bag但 id 由微秒时间生成,理论上极端并发下可能碰撞。
  • docker_commands 部分错误响应解码没有统一使用 [return_maps],有些分支可能匹配不到 map。
  • docker_events 存在但未启动,且使用 shell 命令 docker events,与其他 Docker API 访问方式不同。

10. 一句话主流程

微服务通过 WebSocket 注册到 EFKA本地 channel 把指标交给对应 efka_service,再由 efka_client 通过 TLS 上报给上游;上游通过同一条 TLS 连接下发 pub/sub 消息和容器管理请求pub/sub 再广播回本地微服务,容器请求则通过 Docker Unix Socket 在本机执行。