feat(student-bff): 完整实现 student-bff 聚合层

包含 src 全部实现、Dockerfile、shared-ts/bff 包等
This commit is contained in:
SpecialX
2026-07-10 19:10:51 +08:00
parent e5ca4c6c7b
commit f585080e70
55 changed files with 7141 additions and 252 deletions

View File

@@ -0,0 +1,693 @@
/**
* DownstreamClient - BFF 模式 v2 标准下游调用抽象.
*
* 仲裁依据: coord-final-decisions §2 B8 (回写 teacher-bff, 3 个 BFF 统一使用)
* coord-final-decisions §2 B2 (首次实现即 gRPC, 禁止 HTTP fetch)
*
* 核心能力:
* - gRPC 调用封装 (@grpc/grpc-js + @grpc/proto-loader)
* - mock 模式 (env.MOCK_UPSTREAM=true 时返回固定数据, 上游就绪后移除)
* - 超时/重试/traceId 透传
* - 错误归一化 (下游 gRPC 错误 → BFF_STUDENT_BAD_GATEWAY)
* - 熔断器集成点 (P6 接入 opossum)
*
* 使用方式:
* const client = new DownstreamClient(env);
* const userInfo = await client.call('iam', 'GetUserInfo', { userId }, { metadata: { 'x-user-id': userId } });
*/
import { promises } from 'node:fs';
import path from 'node:path';
import * as grpc from '@grpc/grpc-js';
import * as protoLoader from '@grpc/proto-loader';
import type { PackageDefinition, ServiceClientConstructor } from '@grpc/grpc-js';
import { logger } from './logger.js';
/**
* 下游调用配置.
*/
export interface CallOptions {
/** 超时时间(毫秒), 默认 5000ms */
timeoutMs?: number;
/** 重试次数, 默认 2 */
retryCount?: number;
/** 重试退避基数(毫秒), 默认 100ms, 指数退避 */
retryBackoffMs?: number;
/** 全链路追踪 ID (从 x-request-id header 获取) */
traceId?: string;
/** gRPC metadata (含 x-user-id / x-user-roles / x-dataScope) */
metadata?: Record<string, string>;
}
/**
* 下游调用结果(成功).
*/
export interface DownstreamResult<T> {
success: true;
data: T;
}
/**
* 下游调用结果(失败).
*/
export interface DownstreamFailure {
success: false;
error: {
code: string;
message: string;
service: string;
method: string;
status?: number;
traceId?: string;
};
}
export type DownstreamResponse<T> = DownstreamResult<T> | DownstreamFailure;
/**
* 下游服务配置.
*/
export interface DownstreamServiceConfig {
/** 服务名 (iam / core-edu / content / data-ana / msg / ai) */
name: string;
/** gRPC 端点 (如 'localhost:50052') */
grpcUrl: string;
/** proto 文件路径 (相对项目根) */
protoPath: string;
/** proto package 名 (如 'next_edu_cloud.iam.v1') */
packageName: string;
/** 是否启用该下游 (按阶段扩展, 见 president §2.4) */
enabled: boolean;
/** 是否为必需依赖 (失败返回 503, 否则软失败) */
required: boolean;
}
/**
* DownstreamClient 配置.
*/
export interface DownstreamClientConfig {
/** 是否启用 mock 模式 (上游未就绪时) */
mockUpstream: boolean;
/** 是否开发模式 (DEV_MODE 放行越权等) */
devMode: boolean;
/** 下游服务配置列表 */
services: DownstreamServiceConfig[];
/** 默认超时(ms) */
defaultTimeoutMs: number;
/** 默认重试次数 */
defaultRetryCount: number;
/** 默认重试退避基数(ms) */
defaultRetryBackoffMs: number;
}
/**
* Mock 数据提供器接口.
* 各 BFF 自行实现, 提供各 RPC 的 mock 数据.
*/
export type MockDataProvider = (
service: string,
method: string,
request: unknown,
) => unknown | undefined;
/**
* DownstreamClient - BFF 模式 v2 标准下游调用抽象.
*
* 3 个 BFF (teacher-bff / student-bff / parent-bff) 统一使用 (B8 裁决).
*/
export class DownstreamClient {
private readonly clients = new Map<string, grpc.Client>();
private readonly serviceDefs = new Map<string, DownstreamServiceConfig>();
private readonly packageDefs = new Map<string, PackageDefinition>();
private mockProvider: MockDataProvider | null = null;
constructor(private readonly config: DownstreamClientConfig) {
for (const svc of config.services) {
this.serviceDefs.set(svc.name, svc);
}
}
/**
* 设置 mock 数据提供器.
* 各 BFF 在初始化时注入自己的 mock 数据实现.
*/
setMockProvider(provider: MockDataProvider): void {
this.mockProvider = provider;
}
/**
* gRPC 调用封装.
*
* @param service 下游服务名 (如 'iam' / 'core-edu')
* @param method RPC 方法名 (如 'GetUserInfo')
* @param request 请求 message
* @param options 调用配置
* @returns 响应数据, 失败时抛出 DownstreamError
*/
async call<TRequest, TResponse>(
service: string,
method: string,
request: TRequest,
options?: CallOptions,
): Promise<TResponse> {
const svc = this.serviceDefs.get(service);
if (!svc) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_UNKNOWN_SERVICE',
message: `Unknown downstream service: ${service}`,
service,
method,
});
}
if (!svc.enabled) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_DISABLED',
message: `Downstream service ${service} is not enabled in current stage`,
service,
method,
});
}
// mock 模式: 上游未就绪时返回固定数据
if (this.config.mockUpstream && this.mockProvider) {
const mockData = this.mockProvider(service, method, request);
if (mockData !== undefined) {
logger.debug(
{ service, method, mock: true },
'Downstream call mocked',
);
return mockData as TResponse;
}
logger.warn(
{ service, method },
'No mock data provider for downstream call, falling through to gRPC',
);
}
const timeoutMs = options?.timeoutMs ?? this.config.defaultTimeoutMs;
const retryCount = options?.retryCount ?? this.config.defaultRetryCount;
const retryBackoffMs =
options?.retryBackoffMs ?? this.config.defaultRetryBackoffMs;
let lastError: unknown = null;
for (let attempt = 0; attempt <= retryCount; attempt++) {
try {
const data = await this.invokeGrpc<TRequest, TResponse>(
svc,
method,
request,
timeoutMs,
options?.metadata,
);
return data;
} catch (err) {
lastError = err;
if (attempt < retryCount) {
const backoff = retryBackoffMs * Math.pow(2, attempt);
logger.warn(
{ service, method, attempt: attempt + 1, retryCount, backoff, err },
'Downstream call failed, retrying',
);
await sleep(backoff);
}
}
}
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_BAD_GATEWAY',
message: `Downstream ${service}.${method} failed after ${retryCount + 1} attempts`,
service,
method,
traceId: options?.traceId,
cause: lastError,
});
}
/**
* gRPC server-streaming 调用封装.
*
* 用于流式 RPC (如 ai.StreamChat), 返回 AsyncIterable 逐块产出.
* mock 模式下产出单个 mock chunk 后结束.
*
* @param service 下游服务名
* @param method RPC 方法名 (必须是 server-streaming)
* @param request 请求 message
* @param options 调用配置 (timeoutMs / metadata / traceId)
* @returns AsyncIterable<TResponse>, 每个元素为一个流块
*/
async *callStream<TRequest, TResponse>(
service: string,
method: string,
request: TRequest,
options?: CallOptions,
): AsyncIterable<TResponse> {
const svc = this.serviceDefs.get(service);
if (!svc) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_UNKNOWN_SERVICE',
message: `Unknown downstream service: ${service}`,
service,
method,
});
}
if (!svc.enabled) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_DISABLED',
message: `Downstream service ${service} is not enabled in current stage`,
service,
method,
});
}
// mock 模式: 产出单个 mock chunk 后结束
if (this.config.mockUpstream && this.mockProvider) {
const mockData = this.mockProvider(service, method, request);
if (mockData !== undefined) {
logger.debug(
{ service, method, mock: true },
'Downstream stream call mocked',
);
yield mockData as TResponse;
return;
}
}
const timeoutMs = options?.timeoutMs ?? this.config.defaultTimeoutMs;
const client = await this.getOrCreateClient(svc);
const meta = new grpc.Metadata();
if (options?.metadata) {
for (const [key, value] of Object.entries(options.metadata)) {
meta.set(key, value);
}
}
const deadline = Date.now() + timeoutMs;
const callFn = (client as unknown as Record<string, Function>)[method];
if (typeof callFn !== 'function') {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_METHOD_NOT_FOUND',
message: `Method ${method} not found on service ${svc.name}`,
service: svc.name,
method,
});
}
// 发起 server-streaming 调用
const stream = callFn.call(client, request, meta, { deadline });
// 将 Node ReadableStream 转换为 AsyncIterable
try {
let streamDone = false;
let streamError: Error | null = null;
const chunkQueue: TResponse[] = [];
let resolveWait: ((v: { done: true } | { done: false; value: TResponse }) => void) | null = null;
stream.on('data', (chunk: TResponse) => {
if (resolveWait) {
const r = resolveWait;
resolveWait = null;
r({ done: false, value: chunk });
} else {
chunkQueue.push(chunk);
}
});
stream.on('end', () => {
streamDone = true;
if (resolveWait) {
const r = resolveWait;
resolveWait = null;
r({ done: true });
}
});
stream.on('error', (err: Error) => {
streamError = err;
streamDone = true;
if (resolveWait) {
const r = resolveWait;
resolveWait = null;
r({ done: true });
}
});
while (!streamDone || chunkQueue.length > 0) {
if (chunkQueue.length > 0) {
yield chunkQueue.shift()!;
continue;
}
if (streamDone) break;
const result = await new Promise<{ done: true } | { done: false; value: TResponse }>(
(resolve) => {
resolveWait = resolve;
},
);
if (result.done) break;
yield result.value;
}
if (streamError) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_STREAM_ERROR',
message: streamError.message,
service: svc.name,
method,
traceId: options?.traceId,
cause: streamError,
});
}
} finally {
// 确保流被销毁
stream.destroy?.();
}
}
/**
* 并行调用多个下游服务, 部分失败容错 (Promise.allSettled).
* 用于 Dashboard 类聚合场景 (president §2.6 方案 B 降级).
*
* @returns 每个调用的结果 (success 或 failure), 不抛出异常
*/
async callAll<T extends readonly DownstreamCallSpec[]>(
calls: T,
): Promise<{ -readonly [K in keyof T]: DownstreamResponse<unknown> }> {
const results = await Promise.allSettled(
calls.map((c) =>
this.call(c.service, c.method, c.request, c.options).then(
(data): DownstreamResponse<unknown> => ({ success: true, data }),
),
),
);
return results.map((r, idx) => {
if (r.status === 'fulfilled') {
return r.value;
}
const spec = calls[idx];
const err =
r.reason instanceof DownstreamError
? r.reason
: new DownstreamError({
code: 'BFF_DOWNSTREAM_UNKNOWN_ERROR',
message: String(r.reason),
service: spec.service,
method: spec.method,
});
return {
success: false,
error: {
code: err.code,
message: err.message,
service: err.service,
method: err.method,
traceId: spec.options?.traceId,
},
};
}) as { -readonly [K in keyof T]: DownstreamResponse<unknown> };
}
/**
* 关闭所有 gRPC 连接 (优雅关闭时调用).
*/
async close(): Promise<void> {
for (const [name, client] of this.clients) {
client.close();
logger.debug({ service: name }, 'gRPC client closed');
}
this.clients.clear();
}
/**
* 检查下游服务可达性 (用于 /readyz 探针, §2.4).
*/
async checkHealth(service: string): Promise<boolean> {
const svc = this.serviceDefs.get(service);
if (!svc || !svc.enabled) {
return false;
}
if (this.config.mockUpstream) {
return true;
}
try {
const client = await this.getOrCreateClient(svc);
return new Promise<boolean>((resolve) => {
client.waitForReady(
Date.now() + 2000,
(err) => resolve(!err),
);
});
} catch {
return false;
}
}
/**
* 获取所有已启用的必需下游服务名 (用于 /readyz 探针).
*/
getRequiredServices(): string[] {
return Array.from(this.serviceDefs.values())
.filter((s) => s.enabled && s.required)
.map((s) => s.name);
}
/**
* 获取所有已启用的下游服务名.
*/
getEnabledServices(): string[] {
return Array.from(this.serviceDefs.values())
.filter((s) => s.enabled)
.map((s) => s.name);
}
/**
* 内部: 执行 gRPC 调用.
*/
private async invokeGrpc<TRequest, TResponse>(
svc: DownstreamServiceConfig,
method: string,
request: TRequest,
timeoutMs: number,
metadata?: Record<string, string>,
): Promise<TResponse> {
const client = await this.getOrCreateClient(svc);
const meta = new grpc.Metadata();
if (metadata) {
for (const [key, value] of Object.entries(metadata)) {
meta.set(key, value);
}
}
return new Promise<TResponse>((resolve, reject) => {
const deadline = Date.now() + timeoutMs;
const call = (client as unknown as Record<string, Function>)[method];
if (typeof call !== 'function') {
reject(
new DownstreamError({
code: 'BFF_DOWNSTREAM_METHOD_NOT_FOUND',
message: `Method ${method} not found on service ${svc.name}`,
service: svc.name,
method,
}),
);
return;
}
call.call(
client,
request,
meta,
{ deadline },
(err: grpc.ServiceError | null, response: TResponse) => {
if (err) {
reject(
new DownstreamError({
code: mapGrpcErrorCode(err.code),
message: err.message,
service: svc.name,
method,
status: err.code,
}),
);
return;
}
resolve(response);
},
);
});
}
/**
* 内部: 获取或创建 gRPC client (channel 复用).
*/
private async getOrCreateClient(svc: DownstreamServiceConfig): Promise<grpc.Client> {
let client = this.clients.get(svc.name);
if (client) {
return client;
}
const packageDef = await this.loadPackageDefinition(svc);
const proto = grpc.loadPackageDefinition(packageDef) as unknown as Record<
string,
unknown
>;
const packageObj = this.getNestedPackage(proto, svc.packageName);
const ServiceCtor = this.findServiceCtor(packageObj, svc);
client = new ServiceCtor(
svc.grpcUrl,
grpc.credentials.createInsecure(),
);
this.clients.set(svc.name, client);
return client;
}
/**
* 内部: 加载 proto package definition (缓存).
*/
private async loadPackageDefinition(
svc: DownstreamServiceConfig,
): Promise<PackageDefinition> {
let pkgDef = this.packageDefs.get(svc.name);
if (pkgDef) {
return pkgDef;
}
const fullPath = path.resolve(process.cwd(), svc.protoPath);
try {
await promises.access(fullPath);
} catch {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_PROTO_NOT_FOUND',
message: `Proto file not found: ${fullPath}`,
service: svc.name,
method: '<init>',
});
}
pkgDef = protoLoader.loadSync(fullPath, {
keepCase: false,
longs: String,
enums: String,
defaults: true,
oneofs: true,
});
this.packageDefs.set(svc.name, pkgDef);
return pkgDef;
}
/**
* 内部: 按 packageName 点号路径获取嵌套 package 对象.
*/
private getNestedPackage(
root: Record<string, unknown>,
packageName: string,
): Record<string, unknown> {
const parts = packageName.split('.');
let current: Record<string, unknown> = root;
for (const part of parts) {
const next = current[part];
if (typeof next !== 'object' || next === null) {
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_PACKAGE_NOT_FOUND',
message: `Package ${packageName} not found in proto (missing part: ${part})`,
service: '',
method: '<init>',
});
}
current = next as Record<string, unknown>;
}
return current;
}
/**
* 内部: 在 package 对象中查找第一个 service 构造器.
* proto-loader 将 service 暴露为 ServiceClientConstructor.
*/
private findServiceCtor(
packageObj: Record<string, unknown>,
svc: DownstreamServiceConfig,
): ServiceClientConstructor {
for (const [key, value] of Object.entries(packageObj)) {
if (
typeof value === 'function' &&
'service' in (value as object)
) {
return value as ServiceClientConstructor;
}
}
throw new DownstreamError({
code: 'BFF_DOWNSTREAM_SERVICE_NOT_FOUND',
message: `No gRPC service found in package ${svc.packageName} for service ${svc.name}`,
service: svc.name,
method: '<init>',
});
}
}
/**
* 并行调用规范 (用于 callAll).
*/
export interface DownstreamCallSpec {
service: string;
method: string;
request: unknown;
options?: CallOptions;
}
/**
* 下游调用错误.
*/
export class DownstreamError extends Error {
readonly code: string;
readonly service: string;
readonly method: string;
readonly status?: number;
readonly traceId?: string;
readonly cause?: unknown;
constructor(params: {
code: string;
message: string;
service: string;
method: string;
status?: number;
traceId?: string;
cause?: unknown;
}) {
super(params.message);
this.name = 'DownstreamError';
this.code = params.code;
this.service = params.service;
this.method = params.method;
this.status = params.status;
this.traceId = params.traceId;
this.cause = params.cause;
}
}
/**
* 将 gRPC status code 映射为 BFF 错误码.
*/
function mapGrpcErrorCode(code: grpc.status | number): string {
switch (code) {
case grpc.status.UNAVAILABLE:
return 'BFF_DOWNSTREAM_UNAVAILABLE';
case grpc.status.DEADLINE_EXCEEDED:
return 'BFF_DOWNSTREAM_TIMEOUT';
case grpc.status.UNAUTHENTICATED:
return 'BFF_DOWNSTREAM_UNAUTHENTICATED';
case grpc.status.PERMISSION_DENIED:
return 'BFF_DOWNSTREAM_PERMISSION_DENIED';
case grpc.status.NOT_FOUND:
return 'BFF_DOWNSTREAM_NOT_FOUND';
case grpc.status.INVALID_ARGUMENT:
return 'BFF_DOWNSTREAM_INVALID_ARGUMENT';
case grpc.status.UNIMPLEMENTED:
return 'BFF_DOWNSTREAM_UNIMPLEMENTED';
case grpc.status.INTERNAL:
return 'BFF_DOWNSTREAM_INTERNAL';
default:
return 'BFF_DOWNSTREAM_BAD_GATEWAY';
}
}
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

View File

@@ -0,0 +1,29 @@
/**
* BFF 模式 v2 标准抽象 (B8 裁决).
*
* 3 个 BFF (teacher-bff / student-bff / parent-bff) 统一使用本模块:
* - DownstreamClient: gRPC 调用封装 + mock + 重试 + 错误归一化
* - createBffLogger: pino logger 工厂 (各 BFF 注入 service 名)
* - DownstreamError + DownstreamResponse: 错误与响应类型
*
* 仲裁依据:
* - coord-final-decisions §2 B1 (P2 起直接 GraphQL)
* - coord-final-decisions §2 B2 (首次实现即 gRPC, 禁止 HTTP fetch)
* - coord-final-decisions §2 B8 (回写 teacher-bff, 3 BFF 统一)
* - coord-final-decisions §1 G4 (pino 结构化日志)
* - coord-final-decisions §1 G12 (ESM .js 后缀 import)
* - coord-final-decisions §1 G13 (import type)
*/
export { DownstreamClient, DownstreamError } from "./downstream-client.js";
export type {
CallOptions,
DownstreamClientConfig,
DownstreamServiceConfig,
DownstreamResult,
DownstreamFailure,
DownstreamResponse,
DownstreamCallSpec,
MockDataProvider,
} from "./downstream-client.js";
export { createBffLogger, logger } from "./logger.js";
export type { Logger } from "./logger.js";

View File

@@ -0,0 +1,57 @@
/**
* BFF 共享 logger (pino 实例).
*
* 仲裁依据: coord-final-decisions §2 B8 (回写 teacher-bff, 3 个 BFF 统一使用)
* coord-final-decisions §1 G4 (首次实现即结构化日志, 禁止 console.log)
*
* 各 BFF 在自身 logger.ts 中调用 createBffLogger('student-bff' | 'teacher-bff' | 'parent-bff')
* 获取带 service 标签的 pino 实例.
*
* 该 logger 也作为 DownstreamClient 默认 logger, 避免循环依赖.
*/
import pino, { type Logger as PinoLogger, type LoggerOptions } from "pino";
/**
* 默认日志级别 (可通过环境变量 LOG_LEVEL 覆盖).
*/
const DEFAULT_LEVEL = process.env.LOG_LEVEL ?? "info";
/**
* 是否开发模式 (启用 pino-pretty 美化输出).
*/
const IS_DEV = process.env.NODE_ENV === "development";
/**
* 创建 BFF 共享 logger.
*
* @param serviceName 服务名 (student-bff / teacher-bff / parent-bff)
* @param options 额外 pino 配置 (可选)
*/
export function createBffLogger(
serviceName: string,
options?: LoggerOptions,
): PinoLogger {
const opts: LoggerOptions = {
level: DEFAULT_LEVEL,
base: {
service: serviceName,
version: "0.1.0",
},
transport: IS_DEV
? {
target: "pino-pretty",
options: { colorize: true },
}
: undefined,
...options,
};
return pino(opts);
}
/**
* DownstreamClient 默认 logger (无 service 标签, 由调用方覆盖).
* 各 BFF 不应直接使用此 logger, 应使用自身 createBffLogger() 产出实例.
*/
export const logger: PinoLogger = createBffLogger("bff-shared");
export type Logger = PinoLogger;