fix 语法告警

This commit is contained in:
anlicheng 2026-05-09 16:44:53 +08:00
parent 1bd5221fc6
commit 9195d32ba6
4 changed files with 300 additions and 17 deletions

288
CODE_LOGIC_OVERVIEW.md Normal file
View File

@ -0,0 +1,288 @@
# EFKA 当前代码逻辑梳理
> 本文档记录当前阅读代码后对项目逻辑的理解,主要基于 `src/``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 编解码。
- `gpb`protobuf 代码生成。
## 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_model`DETS 离线缓存。
- `service_model`DETS 服务状态表。
- `efka_subscription`:本地 topic 订阅中心。
- `efka_client`:连接上游 TLS server 的状态机。
- `docker_task_reporter`:部署任务事件缓冲和转发。
- `docker_deploy_manager`:部署任务管理器。
`docker_events` 代码存在,但在 `efka_sup` 中被注释掉了,目前默认不会启动。
## 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_stream`Docker 部署任务流式日志,只在 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` 根据 `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_deployer``Ref` 接收 `docker_client` 流式消息并上报拉取过程。
5. 创建容器。
6. 创建空的 `service.conf` 配置文件。
7. 写部署摘要到 `efka_logger`
8. 关闭任务事件流,状态为 `success``fail`
创建容器参数:
- `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`,默认配置是:
```erlang
"/usr/local/code/tmp/dets/service.dets"
```
记录结构在 `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 logger`config/sys.config` 里配置到 console 和 `log/debug.log`
- `efka_logger`:自定义部署日志,按日期写到 `code:root_dir() ++ "/log/"` 下。
`efka_logger` 主要用于部署任务摘要和失败原因。
## 8. 配置项
关键配置来自 `config/sys.config`
- `root_dir`:容器相关目录根路径。
- `dets_dir`DETS 文件目录。
- `websocket_server`:本地 WebSocket 监听配置。
- `tls_server_address`:上游 TLS server 地址。
- `auth`:上游鉴权信息。
注意:
- `service_model``cache_model` 打开 DETS 前假设 `dets_dir` 已存在,代码里没有显式创建目录。
- `docker_client` 固定使用 `/var/run/docker.sock`,每次请求内部由短生命周期普通进程执行 open/request/close。
## 9. 当前看到的几个注意点
- `README.md` 描述的是 JSON-RPC 风格 WebSocket但当前代码实际处理的是 binary protobuf 帧README 可能已经过期。
- `service_channel``register` 只使用 `service_id`,没有使用 README 中提到的 `meta_data/container_name`
- `efka_client:send_result_reply/3``send_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 在本机执行。

View File

@ -207,6 +207,14 @@ ok
}}}
```
`task_event` 本身不携带 `uuid``iot` 在接收该消息时使用当前已鉴权 `ssl_channel` 绑定的 host UUID把事件路由到内部任务进程 `{UUID, TaskId}`。HTTP 页面通过 SSE 订阅时也必须使用同一组参数:
```http
GET /event_stream?uuid=<host_uuid>&task_id=<task_id>
```
`iot` 会为每个 `{UUID, TaskId}` 维护一个独立的任务进程,用于缓存最近的部署日志、支持多个 SSE listener并在收到 close 事件后结束事件流。
### iot -> efka: pub
```erlang

View File

@ -37,9 +37,8 @@ check_image_exist(Image) when is_binary(Image) ->
{ok, ContainerId :: binary()} | {error, Reason :: any()}.
create_container(ContainerDir, #{container_name := ContainerName0, create := CreateOpts})
when is_list(ContainerDir), is_binary(ContainerName0), is_map(CreateOpts) ->
ContainerName = to_binary(ContainerName0),
Url = lists:flatten(io_lib:format("/containers/create?name=~s", [binary_to_list(ContainerName)])),
Options = docker_container_builder:build_options(ContainerName, ContainerDir, CreateOpts),
Url = lists:flatten(io_lib:format("/containers/create?name=~s", [binary_to_list(ContainerName0)])),
Options = docker_container_builder:build_options(ContainerName0, ContainerDir, CreateOpts),
display_options(Options),
Body = iolist_to_binary(json:encode(Options)),
@ -294,9 +293,3 @@ boolean_to_query_value(true) ->
"true";
boolean_to_query_value(false) ->
"false".
-spec to_binary(unicode:chardata()) -> binary().
to_binary(Value) when is_binary(Value) ->
Value;
to_binary(Value) when is_list(Value) ->
unicode:characters_to_binary(Value).

View File

@ -200,18 +200,12 @@ short_container_id(ContainerId) when is_binary(ContainerId) ->
-spec deploy_container_name(map()) -> binary().
deploy_container_name(#{container_name := ContainerName}) when is_binary(ContainerName) ->
to_binary(ContainerName);
ContainerName;
deploy_container_name(_) ->
throw({deploy_error, <<"invalid deploy params: container_name missing">>}).
-spec deploy_image(map()) -> binary().
deploy_image(#{create := #{config := #{image := Image}}}) when is_binary(Image) ->
to_binary(Image);
Image;
deploy_image(_) ->
throw({deploy_error, <<"invalid deploy params: image missing">>}).
-spec to_binary(unicode:chardata()) -> binary().
to_binary(Value) when is_binary(Value) ->
Value;
to_binary(Value) when is_list(Value) ->
unicode:characters_to_binary(Value).