• 简体中文
  • RPC

    规划 API: 本页描述 router、procedure、typed client、streaming、metadata 和 resolver 的最终用户用法;用户心智参考 vextjs-rpc 的成熟模型。commflow@0.0.2 尚未发布 RPC runtime,当前可运行导出见 当前版本快速开始

    适用场景

    RPC 适合以“调用服务方法”的方式组织服务间通信:

    场景是否适合
    服务间过程调用适合
    希望共享输入/输出契约并获得类型推导适合
    不想在业务代码里关心 URL 和 HTTP method适合
    需要 unary、服务端流、客户端流或双向流适合
    需要统一透传 requestId、traceparent、tenantId适合
    需要原生 HTTP Response 语义使用 Request
    长连接双向消息使用 Socket

    核心概念

    概念用户理解
    router把多个业务模块组织成一棵 RPC 契约树。
    procedure一个远端方法,可以是 querymutation 或 streaming procedure。
    client调用方通过 typed client 访问远端 procedure。
    handler / adapter服务提供方把 router 暴露到 VextJS、Express、Koa、Fastify 或 Node HTTP。
    resolver决定某个服务名或 procedure 调到哪个地址。
    loadBalancer在多个 endpoint 中选择一个实例。
    metadata随调用透传 requestId、traceparent、tenantId 等上下文。

    定义共享契约

    RPC 的推荐心智不是手写字符串路径,而是先定义共享 router,再让服务端和客户端共用这份契约。

    import {
      createCommflowRpcRouter,
      mutation,
      query
    } from 'commflow';
    
    export const userRpc = createCommflowRpcRouter({
      user: {
        getById: query({
          input: { id: 'string!' },
          output: { id: 'string!', name: 'string!' },
          resolve: async ({ input, ctx }) => {
            return ctx.services.user.findById(input.id);
          }
        }),
    
        create: mutation({
          input: { name: 'string:1-64!', email: 'email!' },
          output: { id: 'string!' },
          resolve: async ({ input, ctx }) => {
            return ctx.services.user.create(input);
          }
        })
      }
    });
    
    export type UserRpc = typeof userRpc;

    上面的 schema 写法表达用户契约:输入、输出、上下文和 procedure 类型在同一份 router 中定义。具体 schema adapter 可以替换,但不改变这个用户模型。

    调用 unary procedure

    import {
      createCommflowRpcClient,
      staticResolver
    } from 'commflow';
    import { userRpc, type UserRpc } from './contracts/user-rpc';
    
    const rpc = createCommflowRpcClient<UserRpc>({
      router: userRpc,
      resolver: staticResolver({
        user: ['https://user-service.internal/rpc']
      }),
      timeout: 5000,
      retry: 1
    });
    
    const user = await rpc.user.getById({ id: '42' });

    推荐把 procedure 按业务域分组:

    命名方式示例
    service.methoduser.getById
    resource.operationorder.cancel
    domain.actionbilling.createInvoice

    流式 RPC

    RPC 不只适合 unary 调用。参考 vextjs-rpc 的用户模型,commflow RPC 区分三类 streaming procedure:

    类型适用场景
    serverStream服务端持续返回日志、进度或订阅结果。
    clientStream客户端持续上传 chunk,服务端最终返回汇总结果。
    bidiStream双向持续收发,例如交互式任务、协作或会话。
    import {
      bidiStream,
      createCommflowRpcRouter,
      serverStream
    } from 'commflow';
    
    export const streamRpc = createCommflowRpcRouter({
      log: {
        tail: serverStream({
          input: { service: 'string!' },
          chunk: { line: 'string!' },
          resolve: async function* ({ input, ctx, signal }) {
            for await (const line of ctx.logs.tail(input.service, { signal })) {
              yield { line };
            }
          }
        })
      },
    
      chat: {
        echo: bidiStream({
          input: { roomId: 'string!' },
          chunkIn: { text: 'string!' },
          chunkOut: { text: 'string!' },
          resolve: async function* ({ stream }) {
            for await (const message of stream) {
              yield { text: `echo:${message.text}` };
            }
          }
        })
      }
    });

    客户端消费:

    for await (const chunk of rpc.log.tail({ service: 'user' })) {
      console.log(chunk.line);
    }
    
    async function* messages() {
      yield { text: 'hello' };
    }
    
    for await (const reply of rpc.chat.echo({ roomId: 'room-1' }, messages())) {
      console.log(reply.text);
    }

    VextJS 接入

    VextJS 是 commflow 的首个重点消费者。宿主 route 负责 RPC 路径;commflow RPC 只提供 handler、context bridge 和 typed client,不能由插件隐藏注册服务端入口。

    规划 adapter API: 下列模块展示 runtime 发布后的 VextJS 接入形态,不能复制为 commflow@0.0.2 的当前导入。关键规则不变:宿主 route 负责 /rpc 路径、鉴权、中间件、OpenAPI 和热重载;commflow 只提供 handler、context bridge 与 client。

    先把服务端入口放在 route 文件中。src/routes/rpc.ts 的文件前缀就是 /rpc,因此子路径使用 /

    // src/routes/rpc.ts
    import { defineRoutes } from 'vextjs';
    import { appRpcHandler } from '../lib/commflow-rpc-handler';
    
    export default defineRoutes((app) => {
      app.post('/', async (request, response) => {
        return appRpcHandler.handleVext(request, response);
      });
    });

    appRpcHandler 应在普通应用模块(例如 src/lib/commflow-rpc-handler.ts)中通过未来的 createCommflowRpcHandler() 创建;它不是插件副作用。

    再把进程内 RPC client 单独放进插件文件。插件只挂载应用能力,不注册服务端入口:

    // src/plugins/commflow-rpc.ts
    import { createCommflowRpcClient, staticResolver } from 'commflow';
    import { definePlugin } from 'vextjs';
    import { appRpc } from '../contracts/rpc';
    
    export default definePlugin({
      name: 'commflow-rpc',
      setup(app) {
        app.extend('rpc', createCommflowRpcClient({
          router: appRpc,
          resolver: staticResolver({
            user: ['http://user-service:3000/rpc'],
            order: ['http://order-service:3000/rpc']
          }),
          metadata: () => ({ source: 'vext' })
        }));
      }
    });

    如果只需要 VextJS 进程内客户端,也可以只挂载 client;这同样是规划 API:

    export default definePlugin({
      name: 'commflow-rpc-client',
      setup(app) {
        app.extend('rpc', createCommflowRpcClient({
          router: appRpc,
          metadata: () => ({
            source: 'order-service'
          }),
          resolver: staticResolver({
            user: ['http://user-service:3000/rpc']
          })
        }));
      }
    });

    通用 Node/Express/Koa/Fastify 用户也遵循同一规则:先创建 handler,再由宿主框架自己选择 path、中间件、鉴权与限流。

    const rpcHandler = createCommflowRpcHandler({
      router: appRpc,
      createContext: ({ request }) => ({
        requestId: request.headers.get('x-request-id') ?? undefined
      })
    });

    宿主 route 负责 RPC 路径;commflow 负责 RPC handler。

    const rpc = createCommflowRpcClient({
      router: appRpc,
      resolver: staticResolver({
        user: ['http://user-service:3000/rpc'],
        order: ['http://order-service:3000/rpc']
      })
    });

    业务代码中调用:

    export class OrderService {
      constructor(private app: VextApp) {}
    
      async createOrder(userId: string, items: OrderItem[]) {
        const user = await this.app.rpc.user.getById({ id: userId });
        return this.app.rpc.order.create({ userId: user.id, items });
      }
    }

    默认建议:

    • 默认路径使用 /rpc,需要隔离内部入口时再改成 /internal/rpc
    • 默认复用 VextJS 的 app.servicesrequest.authrequest.requestIdx-request-idtraceparent
    • 默认出站请求复用后续 commflow request core,而不是在 RPC 里重新实现 fetch 编排。
    • 需要单侧能力时,只创建 handler 或只挂载 client。

    通用客户端

    非 VextJS 用户应能直接创建通用 RPC client:

    import {
      createCommflowRpcClient,
      roundRobin,
      staticResolver
    } from 'commflow';
    import { appRpc, type AppRpc } from './contracts/rpc';
    
    export const rpc = createCommflowRpcClient<AppRpc>({
      router: appRpc,
      resolver: staticResolver({
        user: [
          'http://user-1.internal/rpc',
          'http://user-2.internal/rpc'
        ],
        order: ['http://order.internal/rpc']
      }),
      loadBalancer: roundRobin(),
      timeout: 3000,
      retry: 1,
      metadata: {
        source: 'gateway'
      }
    });

    metadata 与上下文透传

    metadata 用来随每次 RPC 调用透传平台上下文:

    const rpc = createCommflowRpcClient<AppRpc>({
      router: appRpc,
      resolver,
      metadata: ({ procedure }) => ({
        requestId: getCurrentRequestId(),
        traceparent: getCurrentTraceparent(),
        tenantId: getCurrentTenantId(),
        source: 'order-service',
        procedure
      })
    });

    边界:

    • requestId / traceparent 通常由网关或入口中间件生成。
    • commflow 负责透传和暴露这些值,不应凭空生成平台 trace。
    • 服务端 procedure 可以从 metactx 中读取这些值。

    服务发现与负载均衡

    生产环境通常由 Kubernetes Service、服务网格或网关提供稳定地址;commflow RPC 不应该强迫用户使用内置注册中心。

    场景建议
    已有服务网格resolver 指向稳定 Service DNS 或 ClusterIP。
    开发环境使用 staticResolver 写死一组地址。
    多实例无网格使用 staticResolver + roundRobin 作为轻量兜底。
    需要注册中心实现 resolver adapter 接入 Nacos、Consul 或内部平台。

    重试边界

    RPC 的 retry 不能只看网络错误,还要看 procedure 类型:

    类型默认重试建议
    query可对短暂 transport error 保守重试。
    mutation默认不自动重试,除非用户提供幂等键和显式 retry policy。
    serverStream只在尚未收到首个 chunk 的建链阶段允许重试。
    clientStream / bidiStream默认不自动重放,避免重复写入或乱序。

    错误处理

    RPC 错误建议分层:

    错误含义
    transport errortimeout、network、aborted。
    protocol errorprocedure 不存在、参数不合法。
    business error服务返回的业务失败。
    contract error类型契约不匹配。

    统一错误结构应至少让用户判断:

    {
      code: 'NOT_FOUND',
      message: 'User not found',
      status: 404,
      retryable: false,
      details: { id: '42' }
    }

    服务端可以抛出结构化业务错误:

    throw new CommflowRpcError('NOT_FOUND', 'User not found', {
      status: 404,
      retryable: false,
      details: { id: input.id }
    });

    更多恢复策略见 错误与重试

    与 request 的区别

    维度requestRPC
    用户心智发 HTTP 请求调用服务过程
    入口method + pathtyped procedure + input
    服务端通常是已有 HTTP API需要暴露 router/handler
    返回Response 或封装响应业务结果、stream chunk 或 RPC error
    类型契约可选核心能力之一
    retry主要看 HTTP method / error还要看 query/mutation/streaming 类型