Skip to content

Latest commit

 

History

History
1011 lines (763 loc) · 28.8 KB

File metadata and controls

1011 lines (763 loc) · 28.8 KB

Betterfly2 API 文档

本文档描述了 Betterfly2 当前 HTTP API 与内部接口。

文档版本: 1.4 最后更新: 2026-07-12

目录

  1. Push Service API
  2. Call Service API
  3. ABTest Service API
  4. Storage HTTP API
  5. Storage Kafka API
  6. Storage Protobuf 定义

Push Service API

PushService 通过现有 /ws Protobuf 连接管理普通 APNs 和 PushKit VoIP token,并在普通消息发送或来电时请求 APNs。协议定义位于 ../proto/push/push_interface.proto。

注册与删除推送 token

客户端登录后发送 df_interface.RequestMessage.push_request.register_voip_token,携带客户端稳定 device_id、PKPushCredentials.token 的小写十六进制字符串,以及真实的 SANDBOX 或 PRODUCTION 环境。服务端返回 df_interface.ResponseMessage.push_event。

同一设备的 token 更新会替换旧记录;同一个 APNs token 登录到新用户时会转移到新用户,避免把来电推给前一个账号。客户端退出账号或 PushKit token 失效时发送 push_request.unregister_voip_token,携带 device_id 和对应环境。

普通远程通知使用 register_apns_token 和 unregister_apns_token,token 来自 didRegisterForRemoteNotificationsWithDeviceToken。普通 APNs token 与 PushKit token 独立保存,不能混用。

PushKit 与 CallKit 恢复流程

VoIP Push payload 示例:

{
  "aps": {"content-available": 1},
  "event": "incoming_call",
  "call_id": "00112233445566778899aabbccddeeff",
  "call_uuid": "00112233-4455-6677-8899-aabbccddeeff",
  "caller_user_id": 1001,
  "call_type": "video",
  "has_video": true,
  "expires_at": "2026-07-10T10:00:45Z"
}

客户端收到 PushKit 回调后必须立即使用 call_uuid 向 CallKit 报告来电。随后恢复登录 WebSocket,并发送 call_request.resume_call(call_id);服务端返回正常的 INCOMING_CALL,其中包含 SDP offer 和 ICE servers。若返回 CALL_NOT_FOUND 或 INVALID_STATE,客户端应结束刚刚报告的 CallKit 通话,因为主叫可能已经取消或响铃已经超时。

为覆盖 DF Pod 突然退出但 Redis 路由租约尚未过期的窗口,每个来电都会同时尝试 VoIP Push;在线客户端可能同时收到 WebSocket INCOMING_CALL 和 PushKit payload,必须使用相同的 call_id/call_uuid 幂等去重,只创建一个 CallKit 通话。

APNs 环境

PushService 根据每条 token 的 environment 自动选择 sandbox 或 production endpoint,不使用全局环境开关。VoIP topic 为 com.Voltline.Betterfly2.voip,普通通知 topic 为 com.Voltline.Betterfly2;私钥通过 APNS_PRIVATE_KEY_BASE64 或 APNS_PRIVATE_KEY_PATH 注入。

PushService 健康检查为 GET /health 和 GET /ready,默认端口 8086,不需要暴露到公网。

普通消息通知与调试后台

普通消息请求通过校验后,DataForwarding Service 会以 best-effort 方式向 push-service topic 发布推送任务。PushService 根据数据库资料生成标题和正文:私聊使用发送者资料,群聊使用群资料;通知同时携带客户端 Notification Service Extension 所需的通信通知元数据。WebSocket 在线不会阻止 APNs 投递,因为同一账号可能还有其他离线设备;前台是否展示横幅由客户端决定。

配置 PUSH_ADMIN_TOKEN 后可访问 GET /push/admin,并通过受保护的管理 API 调试普通通知、VoIP Push 和全量普通通知。未配置令牌时,页面和管理 API 均返回 404。完整说明见 PushService 文档。


Call Service API

CallService 提供一对一 WebRTC 语音和视频通话的控制面。媒体流不经过 CallService;客户端通过 WebRTC 建立点对点连接,无法直连时由 Coturn 中继。

服务健康检查为 GET /health 和 GET /ready,默认端口 8085。客户端通话 API 继续复用已登录的 /ws Protobuf 连接,具体消息见本文“通话协议”章节。

通话协议

协议定义位于 ../proto/call/call_interface.proto。客户端将 call_interface.ClientRequest 放入 df_interface.RequestMessage.call_request,并携带当前登录 JWT。服务端返回 df_interface.ResponseMessage.call_event。

客户端命令:

  • get_config: 获取 STUN/TURN 地址和短期 TURN 凭证。
  • initiate: 发起一对一语音或视频通话,必须携带 callee_user_id、AUDIO/VIDEO 和 SDP offer。
  • accept: 被叫接听,必须携带 call_id 和 SDP answer。
  • reject: 被叫拒绝响铃中的通话。
  • hangup: 任一参与者取消或结束通话。
  • ice_candidate: 将 trickle ICE candidate 转发给另一参与者。
  • resume_call: PushKit 唤醒并重新登录后,使用 call_id 恢复待接来电的 SDP/ICE。

服务端事件:

  • CALL_CONFIG: ICE server 列表。
  • OUTGOING_CALL: 服务端已创建通话,主叫获得唯一 call_id。
  • INCOMING_CALL: 被叫收到来电、主叫 ID、通话类型及 SDP offer。
  • CALL_ACCEPTED: 主叫收到 SDP answer,双方进入 ACTIVE。
  • CALL_REJECTED: 被叫拒绝。
  • CALL_ENDED: 挂断、取消、断连或响铃超时。
  • ICE_CANDIDATE_RECEIVED: 收到对端 candidate。
  • CALL_ERROR: 离线、忙线、越权、状态冲突或参数错误。

典型时序:

caller -> DF -> call-service: initiate(offer)
call-service -> caller: OUTGOING_CALL(call_id)
call-service -> callee: INCOMING_CALL(call_id, offer)
callee -> DF -> call-service: accept(call_id, answer)
call-service -> caller: CALL_ACCEPTED(answer)
caller <-> call-service <-> callee: ice_candidate
caller <========= WebRTC media / Coturn relay =========> callee
either side -> call-service: hangup(call_id)
call-service -> both sides: CALL_ENDED

身份安全边界:

  • caller_user_id 不由客户端填写,DF 从已认证 WebSocket 会话中注入 InternalRequest.user_id。
  • CallService 会验证接听者必须是被叫方,ICE 与挂断操作者必须是通话参与者。
  • 同一用户同时只能占用一个 RINGING 或 ACTIVE 通话。
  • TURN 使用基于共享密钥生成的短期 HMAC 凭证,TURN_SHARED_SECRET 不会下发给客户端。

部署端口

  • 8085/tcp: CallService 健康检查,仅需内网开放。
  • 3478/udp、3478/tcp: STUN/TURN 监听端口,需要对客户端开放。
  • 49160-49200/udp: Coturn 媒体中继端口范围,需要对客户端开放。

生产环境必须将 TURN_EXTERNAL_IP 设置为服务器公网 IP,并将 TURN_PUBLIC_HOST 设置为客户端可访问的域名或公网 IP;也可以通过 CALL_STUN_URLS、CALL_TURN_URLS 显式覆盖自动生成的地址。TURN_SHARED_SECRET 必须在 CallService 与 Coturn 中保持一致。

一对一支持 PushKit 离线唤醒;群语音/群视频可选接入自托管 LiveKit SFU,复用 call_request/call_event, 不使用一对一 SDP/ICE 信令。新增群 API、字段号、邀请与权限见群通话。

ABTest Service API

ABTestService 提供统一实验配置与稳定分流能力。当前主要用于客户端实验,接口设计已预留服务端实验入口。

基础URL: http://localhost:8082

客户端获取实验配置

接口: GET /abtest/v1/client/config

Query参数:

  • device_id 必需,客户端设备唯一ID
  • platform 可选,例如 ios、android
  • app_version 可选,例如 1.2.0
  • os 可选,例如 iOS
  • system_version 可选,例如 17.4

示例:

GET /abtest/v1/client/config?device_id=device-001&platform=ios&app_version=1.2.0&system_version=17.4

响应:

{
  "server_time": "2026-04-29T10:00:00Z",
  "merged_config": {
    "enable_new_chat_ui": true
  },
  "experiments": [
    {
      "experiment_id": 1,
      "experiment_key": "new_chat_ui",
      "experiment_type": "client",
      "group_key": "variant",
      "version": 3,
      "start_time": "2026-04-29T10:00:00Z",
      "end_time": "2026-05-06T10:00:00Z",
      "duration_seconds": 604800,
      "config": {
        "enable_new_chat_ui": true
      }
    }
  ]
}

通用实验求值接口

接口: POST /abtest/v1/evaluate

该接口用于未来服务端实验。subject_type 可以是 device、user、server 等,context 用于传入版本、平台、区域等扩展条件。

{
  "subject_type": "server",
  "subject_id": "dataForwardingService",
  "context": {
    "region": "sg",
    "feature": "message_sync"
  }
}

管理面板与管理API

管理面板:

GET /abtest/admin

管理页面和管理 API 只有在配置 ABTEST_ADMIN_TOKEN 后才启用;未配置时统一返回 404。API 请求使用 Authorization: Bearer <ABTEST_ADMIN_TOKEN> 或 X-Admin-Token。

常用接口:

  • GET /abtest/admin/api/experiments
  • POST /abtest/admin/api/experiments
  • GET /abtest/admin/api/experiments/{id}
  • PUT /abtest/admin/api/experiments/{id}
  • POST /abtest/admin/api/experiments/{id}/start
  • POST /abtest/admin/api/experiments/{id}/pause
  • POST /abtest/admin/api/experiments/{id}/stop
  • POST /abtest/admin/api/experiments/{id}/withdraw
  • POST /abtest/admin/api/experiments/{id}/groups
  • PUT /abtest/admin/api/experiments/{id}/groups/{group_id}
  • DELETE /abtest/admin/api/experiments/{id}/groups/{group_id}
  • POST /abtest/admin/api/experiments/{id}/groups/{group_id}/push_full
  • POST /abtest/admin/api/experiments/{id}/overrides
  • DELETE /abtest/admin/api/experiments/{id}/overrides/{override_id}

运行中和暂停中的实验均可编辑分组 JSON、调整比例和删除分组,管理页分组面板 提供保存/删除按钮,例外面板提供删除按钮。编辑不改变实验状态、时间或分组Key。

分组 PUT 接收 optional traffic_basis_points 和 config(JSON对象),至少提供一项。 未提供的字段保留;config:{} 清空配置;显式比例0可将分组设为只供例外使用。 创建、添加或调整分组时总比例最多10000(100%);JSON单独修改不重写原比例。 删除不自动将剩余组推全,未覆盖的比例继续不参与实验。

删除分组会在同一事务删除指向该组的 force_group 例外,其他例外保留;网页先确认。 至少保留一个分组。当前 rolled_out 的推全组不能删除或修改其100%比例, 但可以修改JSON;需先推全到其他组或撤回推全,再删除它。 分组、例外ID必须属于路径中的实验;跨实验/不存在的ID返回404,状态冲突409, 非法编辑400,数据库故障500,不泄漏数据库细节。

成功修改会原子增加实验版本,并使当前副本的实验与例外缓存失效;其他副本 沿用既有最多5秒TTL收敛,客户端下次获取时生效。没有客户端协议或数据库迁移变化。

PUT /abtest/admin/api/experiments/12/groups/9
Authorization: Bearer <ABTEST_ADMIN_TOKEN>
Content-Type: application/json

{"config":{"enable_new_chat_ui":true}}

分组PUT返回更新后的Group;删除分组或例外返回更新后的Experiment。

推全语义

push_full 不只是把目标分组流量调成 10000。推全后实验会进入 rolled_out 状态,并记录 rollout_group_key。客户端或服务端再次获取该 experiment_key 时,会直接返回该分组配置,不再检查实验开始/结束时间,也不再走流量分桶。撤回推全使用 withdraw,撤回后状态变为 stopped,客户端下一次拉配置将不再收到该 key。

创建实验示例:

{
  "experiment_key": "new_chat_ui",
  "name": "新聊天页",
  "experiment_type": "client",
  "start_time": "2026-04-29T10:00:00Z",
  "duration_seconds": 604800,
  "targeting": {
    "platforms": ["ios"],
    "min_app_version": "1.2.0",
    "min_system_version": "17.0"
  },
  "groups": [
    {
      "group_key": "control",
      "traffic_basis_points": 5000,
      "config": {"enable_new_chat_ui": false}
    },
    {
      "group_key": "variant",
      "traffic_basis_points": 5000,
      "config": {"enable_new_chat_ui": true}
    }
  ]
}

Targeting 规则当前支持:

  • platforms
  • app_versions
  • os、oses 或 operating_systems
  • min_app_version
  • max_app_version
  • system_versions
  • min_system_version
  • max_system_version
  • include: 任意 context 字段白名单
  • exclude: 任意 context 字段黑名单

HTTP API(对外接口)

存储服务提供独立的 HTTP 服务,客户端可以直接通过 HTTP 进行文件的上传和下载操作。

基础URL: http://localhost:8081/storage_service (开发环境)

认证方式: 文件控制面接口都需要在请求头中携带 JWT Token 和用户ID;/health 与 /ready 用于探针检查,可匿名访问

请求头设置:

Authorization: Bearer <JWT_TOKEN>
X-User-ID: <USER_ID>
Content-Type: application/json

或者通过 Query 参数传递 user_id:

?user_id=<USER_ID>

Postman 使用说明

1. 设置 Headers

在 Postman 的 Headers 标签页中添加以下请求头:

Key Value 说明
Authorization Bearer <你的JWT_TOKEN> JWT Token,注意Bearer后面有空格
X-User-ID <你的用户ID> 用户ID,数字
Content-Type application/json 请求体格式

示例:

Authorization: Bearer eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...
X-User-ID: 123
Content-Type: application/json

2. 使用 Postman Environment 变量(推荐)

为了便于管理和切换环境,建议使用 Postman 的 Environment 功能:

  1. 创建 Environment:

    • 点击右上角的 Environment 下拉菜单
    • 选择 "Manage Environments"
    • 点击 "Add" 创建新环境
    • 添加以下变量:
      • base_url: http://localhost:8081
      • jwt_token: <你的JWT Token>
      • user_id: <你的用户ID>
  2. 在请求中使用变量:

    • URL: {{base_url}}/storage_service/upload
    • Headers:
      • Authorization: Bearer {{jwt_token}}
      • X-User-ID: {{user_id}}
  3. 获取 JWT Token:

    • 通过登录接口获取(数据转发服务或认证服务)
    • 登录成功后,从响应中提取 JWT token
    • 将 token 保存到 Environment 变量中

3. 测试步骤

  1. 获取 JWT Token:

    • 先调用登录接口获取有效的 JWT token
    • 将 token 保存到 Environment 变量
  2. 测试上传接口:

    • Method: POST
    • URL: {{base_url}}/storage_service/upload
    • Headers: 使用上述变量
    • Body (raw JSON):
      {
        "file_hash": "<128位十六进制SHA-512>",
        "file_size": 1024
      }
  3. 测试下载接口:

    • Method: GET
    • URL: {{base_url}}/storage_service/download?file_hash=<文件哈希>
    • Headers: 使用上述变量

1. 文件上传(第一阶段:获取上传URL)

接口: POST /storage_service/upload

描述: 客户端请求上传文件,服务端返回预签名上传URL或文件已存在标识。

请求头:

Authorization: Bearer <JWT_TOKEN>
X-User-ID: <USER_ID>
Content-Type: application/json

请求体 (JSON):

{
  "file_hash": "<128位十六进制SHA-512>",
  "file_size": 12345
}

响应 (JSON):

如果文件已存在:

{
  "exists": true
}

如果文件不存在,返回上传URL:

{
  "exists": false,
  "upload_url": "https://rustfs-endpoint/presigned-upload-url",
  "expires_in": 3600
}

错误响应:

  • 400 Bad Request: 请求参数错误
  • 401 Unauthorized: JWT验证失败
  • 500 Internal Server Error: 服务器内部错误

2. 文件上传验证(第二阶段:验证上传的文件)

接口: POST /storage_service/upload/verify

描述: 客户端上传文件完成后,调用此接口验证文件哈希值并保存元数据。

请求头:

Authorization: Bearer <JWT_TOKEN>
X-User-ID: <USER_ID>
Content-Type: application/json

请求体 (JSON):

{
  "file_hash": "<128位十六进制SHA-512>"
}

响应 (JSON):

验证成功:

{
  "success": true
}

验证失败:

{
  "success": false,
  "error_message": "File hash mismatch"
}

错误响应:

  • 400 Bad Request: 请求参数错误
  • 401 Unauthorized: JWT验证失败
  • 500 Internal Server Error: 服务器内部错误

3. 文件下载

接口: GET /storage_service/download?file_hash=<FILE_HASH>

描述: 客户端请求下载文件,服务端返回预签名下载URL。

请求头:

Authorization: Bearer <JWT_TOKEN>
X-User-ID: <USER_ID>

Query参数:

  • file_hash (必需): 文件的 128 位十六进制 SHA-512;服务端会规范化为小写

响应 (JSON):

如果文件存在:

{
  "exists": true,
  "download_url": "https://rustfs-endpoint/presigned-download-url",
  "expires_in": 3600,
  "file_size": 12345
}

如果文件不存在:

{
  "exists": false,
  "error_message": "File not found"
}

错误响应:

  • 400 Bad Request: 请求参数错误(缺少file_hash)
  • 401 Unauthorized: JWT验证失败
  • 500 Internal Server Error: 服务器内部错误

4. 健康检查

接口: GET /health

描述: 检查服务是否正常运行(不需要认证)

响应:

OK

5. 就绪检查

接口: GET /ready

描述: 检查文件控制面的关键依赖是否可用。当前会检查 PostgreSQL 连接和 RustFS bucket。

认证要求: 无

成功响应 (JSON):

{
  "ready": true
}

失败响应 (JSON):

{
  "ready": false,
  "error_message": "database not ready: ..."
}

Kafka MQ API(对内接口)

存储服务通过 Kafka 消息队列接收来自其他服务(主要是数据转发服务)的查询请求。

Topic: storage-service

消息格式: Protobuf 序列化的 RequestMessage,封装在 Envelope 中

响应: 通过请求中的 from_kafka_topic 字段指定的 topic 返回响应


1. 查询文件是否存在

请求消息: QueryFileExists

message QueryFileExists {
  string file_hash = 1;  // 文件SHA512哈希值
}

响应消息: FileExistsRsp

message FileExistsRsp {
  bool exists = 1;
  int64 file_size = 2;
  string storage_path = 3;
}

使用场景: 其他服务需要查询文件是否存在时,通过 Kafka 发送查询请求。


2. 存储新消息

请求消息: StoreNewMessage

message StoreNewMessage {
  int64 from_user_id = 1;
  int64 to_user_id = 2;
  string content = 3;
  string message_type = 4; // text, image, gif, file, audio, video, link
  bool is_group = 5;
  string real_file_name = 6; // 文件消息对应的原始文件名,非文件消息为空
  string client_message_id = 7;
  string client_timestamp = 8;
}

响应消息: StoreMsgRsp

message StoreMsgRsp {
  int64 message_id = 1;
  string client_message_id = 2;
  bool created = 3;
  int64 from_user_id = 4;
  int64 to_user_id = 5;
  string content = 6;
  string message_type = 7;
  bool is_group = 8;
  string real_file_name = 9;
  string client_timestamp = 10;
  string server_timestamp = 11;
}

数据转发服务收到 StoreMsgRsp 后,会向发送方客户端返回 PostAckRsp:

message PostAckRsp {
  int64 message_id = 1;
  string client_message_id = 2;
  string timestamp = 3; // 服务端入库时间
}

3. 查询消息

请求消息: QueryMessage

message QueryMessage {
  int64 message_id = 1;
}

响应消息: MessageRsp

message MessageRsp {
  int64 from_user_id = 1;
  int64 to_user_id = 2;
  string content = 3;
  string timestamp = 4;
  string msg_type = 5; // text, image, gif, file, audio, video, link
  bool is_group = 6;
  string real_file_name = 7;
  int64 message_id = 8;
  bool is_recalled = 9;
  string recalled_at = 10;
  int64 recalled_by = 11;
}

4. 同步消息

请求消息: QuerySyncMessages

message QuerySyncMessages {
  int64 to_user_id = 1;
  string timestamp = 2;
  int32 page_size = 3;
  string cursor_timestamp = 4;
  int64 cursor_message_id = 5;
  bool include_recalled_changes = 6;
  string recall_cursor_timestamp = 7;
  int64 recall_cursor_message_id = 8;
}

响应消息: SyncMessagesRsp

message SyncMessagesRsp {
  repeated MessageRsp msgs = 1;
  bool has_more = 2;
  string next_cursor_timestamp = 3;
  int64 next_cursor_message_id = 4;
  repeated MessageRsp recalled_msgs = 5;
  bool recalls_has_more = 6;
  string next_recall_cursor_timestamp = 7;
  int64 next_recall_cursor_message_id = 8;
}

同步包含自己发送/接收的单聊和当前群成员入群后的群消息。只有认证用户本人能同步。实时 Post 尾部新增 message_id = 9,时间与 ACK、同步均以数据库入库时间为准。


5. 查询用户信息

请求消息: QueryUser

message QueryUser {
  int64 user_id = 1;
}

响应消息: UserInfoRsp

message UserInfoRsp {
  int64 user_id = 1;
  string account = 2;
  string name = 3;
  string avatar = 4;
  string update_time = 5;
}

6. 更新用户名

请求消息: UpdateUserName

message UpdateUserName {
  int64 user_id = 1;
  string new_user_name = 2;
}

响应消息: ResponseMessage (无payload,仅result字段)


7. 更新用户头像

请求消息: UpdateUserAvatar

message UpdateUserAvatar {
  int64 user_id = 1;
  string new_avatar_url = 2;
}

响应消息: ResponseMessage (无payload,仅result字段)


Protobuf 消息定义

通用消息结构

RequestMessage

所有对内接口的请求都封装在 RequestMessage 中:

message RequestMessage {
  string from_kafka_topic = 1;  // 响应发送到的topic
  int64 target_user_id = 2;     // 目标用户ID
  oneof payload {
    StoreNewMessage store_new_message = 3;
    QueryMessage query_message = 4;
    QuerySyncMessages query_sync_messages = 5;
    UpdateUserName update_user_name = 6;
    UpdateUserAvatar update_user_avatar = 7;
    QueryUser query_user = 8;
    QueryFileExists query_file_exists = 9;
  }
}

ResponseMessage

所有对内接口的响应都封装在 ResponseMessage 中:

message ResponseMessage {
  StorageResult result = 1;
  int64 target_user_id = 2;
  oneof payload {
    StoreMsgRsp store_msg_rsp = 3;
    MessageRsp msg_rsp = 4;
    SyncMessagesRsp sync_msgs_rsp = 5;
    UserInfoRsp user_info_rsp = 6;
    FileExistsRsp file_exists_rsp = 7;
  }
}

StorageResult

enum StorageResult {
  OK = 0;
  SERVICE_ERROR = 255;
  RECORD_NOT_EXIST = 1;
}

HTTP 服务消息定义

UploadFileRequest

message UploadFileRequest {
  string file_hash = 1;  // 文件SHA512哈希值
  int64 file_size = 2;   // 文件大小(字节)
}

UploadFileResponse

message UploadFileResponse {
  bool exists = 1;              // 文件是否已存在
  string upload_url = 2;        // 预签名上传URL(如果文件不存在)
  int64 expires_in = 3;         // URL过期时间(秒)
  string error_message = 4;     // 错误信息(如果有)
}

VerifyUploadRequest

message VerifyUploadRequest {
  string file_hash = 1;  // 文件SHA512哈希值
}

VerifyUploadResponse

message VerifyUploadResponse {
  bool success = 1;           // 验证是否成功
  string error_message = 2;  // 错误信息(如果有)
}

DownloadFileRequest

message DownloadFileRequest {
  string file_hash = 1;  // 文件SHA512哈希值
}

DownloadFileResponse

message DownloadFileResponse {
  bool exists = 1;              // 文件是否存在
  string download_url = 2;      // 预签名下载URL(如果文件存在)
  int64 expires_in = 3;         // URL过期时间(秒)
  int64 file_size = 4;          // 文件大小(字节)
  string error_message = 5;     // 错误信息(如果有)
}

文件上传流程

  1. 客户端请求上传

    • 客户端计算文件SHA512哈希值
    • 客户端发送 POST /storage_service/upload 请求,包含文件哈希和大小
    • 服务端验证JWT,检查文件是否已存在
    • 如果文件已存在,返回 exists: true
    • 如果文件不存在,生成预签名上传URL并返回
  2. 客户端上传文件

    • 客户端使用返回的预签名URL直接上传文件到RustFS
    • 上传过程不经过存储服务,直接与RustFS交互
  3. 客户端验证上传

    • 客户端上传完成后,发送 POST /storage_service/upload/verify 请求
    • 服务端从RustFS下载文件并验证哈希值
    • 如果哈希匹配,将文件元数据从待验证状态更新为已验证状态
    • 如果哈希不匹配,删除文件并返回错误
  4. 客户端发送消息

    • 客户端确认文件上传成功后,通过数据转发服务发送消息
    • 文件消息仍然使用普通 Post 报文发送
    • 约定 msg_type = "file"
    • 约定 msg = file_hash
    • 约定 real_file_name = 客户端原始文件名
    • 数据转发服务只转发元数据,不处理文件内容

文件下载流程

  1. 客户端请求下载

    • 客户端从消息中获取文件哈希值
    • 客户端发送 GET /storage_service/download?file_hash=<hash> 请求
    • 服务端验证JWT,检查文件是否存在
    • 如果文件存在,生成预签名下载URL并返回
  2. 客户端下载文件

    • 客户端使用返回的预签名URL直接下载文件
    • 下载过程不经过存储服务,直接与RustFS交互
    • 客户端根据消息中的 real_file_name 恢复文件真实名称

说明:只有完成 upload/verify 校验的文件才会被视为可用文件,未完成校验的待验证文件不会对下载接口或内部文件存在性查询暴露。


环境变量配置

存储服务环境变量

  • HTTP_PORT: HTTP服务端口(进程默认 8080;当前 Docker Compose 默认设置为 8081)
  • PGSQL_DSN: PostgreSQL数据库连接字符串
  • REDIS_ADDR: Redis地址(默认: localhost:6379)
  • KAFKA_BROKER: Kafka broker地址(逗号分隔)
  • KAFKA_STORAGE_TOPIC: Kafka存储服务topic(默认: storage-service)
  • KAFKA_CONSUMER_GROUP: Kafka消费者组(默认: storage-service-group)
  • AUTH_RPC_ADDR: 认证服务gRPC地址(默认: localhost:50051)

RustFS环境变量

  • RUSTFS_REGION: RustFS区域(必需)
  • RUSTFS_ACCESS_KEY_ID: RustFS访问密钥ID(必需)
  • RUSTFS_SECRET_ACCESS_KEY: RustFS秘密访问密钥(必需)
  • RUSTFS_ENDPOINT_URL: RustFS端点URL(必需)
  • RUSTFS_EXTERNAL_ENDPOINT_URL: 客户端可直接访问的RustFS外部地址(可选,推荐在生产环境显式配置)
  • RUSTFS_EXTERNAL_SCHEME: 未设置完整外部地址时使用的协议,例如 https
  • RUSTFS_EXTERNAL_HOST: 未设置完整外部地址时使用的客户端可访问主机;缺省时从 HTTP 请求推导
  • RUSTFS_EXTERNAL_PORT: 未显式配置外部地址时,基于当前HTTP请求推导RustFS地址所使用的端口(默认: 9000)
  • RUSTFS_BUCKET: RustFS存储桶名称(默认: betterfly-files)

错误码说明

HTTP状态码

  • 200 OK: 请求成功
  • 400 Bad Request: 请求参数错误
  • 401 Unauthorized: JWT验证失败
  • 404 Not Found: 资源不存在
  • 500 Internal Server Error: 服务器内部错误

StorageResult枚举

  • OK (0): 操作成功
  • RECORD_NOT_EXIST (1): 记录不存在
  • SERVICE_ERROR (255): 服务内部错误

注意事项

  1. 文件哈希: 所有文件操作都基于SHA512哈希值,客户端必须在上传前计算文件哈希
  2. JWT验证: 所有HTTP接口都需要JWT验证,通过gRPC调用认证服务进行验证
  3. 预签名URL: 上传和下载URL都有有效期(默认1小时),客户端需要在有效期内使用
  4. 文件存储: 文件存储在RustFS中,路径格式为 {hash前2位}/{完整hash}
  5. 元数据存储: 文件元数据存储在PostgreSQL数据库中,包括哈希、大小、存储路径等
  6. 真实文件名: 数据库中不存储文件真实名称,真实名称在数据转发服务的消息中传递

示例代码

上传文件(Go示例)

// 1. 计算文件哈希
hash := calculateSHA512(fileData)

// 2. 请求上传URL
req := UploadFileRequest{
    FileHash: hash,
    FileSize: int64(len(fileData)),
}
resp := requestUploadURL(req, jwt, userID)

// 3. 使用预签名URL上传
if !resp.Exists {
    uploadToRustFS(resp.UploadUrl, fileData)
    
    // 4. 验证上传
    verifyReq := VerifyUploadRequest{FileHash: hash}
    verifyResp := verifyUpload(verifyReq, jwt, userID)
}

下载文件(Go示例)

// 1. 请求下载URL
downloadResp := requestDownloadURL(fileHash, jwt, userID)

// 2. 使用预签名URL下载
if downloadResp.Exists {
    fileData := downloadFromRustFS(downloadResp.DownloadUrl)
    // 使用消息中的real_file_name恢复文件名
}