import { BaseModule } from 'rcc-basemodule'; import { ModuleInfo } from 'rcc-basemodule'; /** * 流水线处理接口 */ export interface IPipelineProcessor { /** * 处理输入数据,转换为流水线格式 * @param input - 原始输入数据 * @param context - 处理上下文 * @returns 流水线格式的数据 */ processInput(input: any, context?: Record): Promise; /** * 处理输出数据,从流水线格式转换 * @param output - 流水线格式的输出数据 * @param context - 处理上下文 * @returns 标准化的输出数据 */ processOutput(output: PipelineData, context?: Record): Promise; /** * 验证流水线数据格式 * @param data - 待验证的数据 * @returns 验证结果 */ validatePipelineData(data: any): ValidationResult; /** * 获取支持的流水线格式 * @returns 支持的格式列表 */ getSupportedFormats(): string[]; } /** * 流水线数据结构 */ export interface PipelineData { /** * 数据唯一标识符 */ id: string; /** * 数据类型 */ type: string; /** * 数据版本 */ version: string; /** * 时间戳 */ timestamp: number; /** * 源模块信息 */ source: { moduleId: string; moduleName: string; operation?: string; }; /** * 目标模块信息 */ target?: { moduleId: string; moduleName: string; }; /** * 数据负载 */ payload: any; /** * 元数据 */ metadata: { format: string; encoding?: string; compression?: string; schema?: string; [key: string]: any; }; /** * 流水线阶段信息 */ pipeline: { stage: string; step: number; totalSteps?: number; previousStage?: string; nextStage?: string; }; /** * 处理跟踪信息 */ tracking?: { traceId: string; spanId: string; parentSpanId?: string; correlationId?: string; }; } /** * 验证结果 */ export interface ValidationResult { /** * 是否有效 */ isValid: boolean; /** * 错误信息 */ errors: string[]; /** * 警告信息 */ warnings: string[]; /** * 验证后的数据 */ data?: any; } /** * 流水线处理器配置 */ export interface PipelineProcessorConfig { /** * 默认数据格式 */ defaultFormat: string; /** * 支持的格式列表 */ supportedFormats: string[]; /** * 是否启用数据验证 */ enableValidation: boolean; /** * 是否启用数据跟踪 */ enableTracking: boolean; /** * 是否启用性能监控 */ enablePerformanceMonitoring: boolean; /** * 自定义验证器 */ validators?: Record ValidationResult>; /** * 自定义转换器 */ transformers?: Record Promise>; /** * 错误处理策略 */ errorHandling: { /** * 遇到错误时是否继续 */ continueOnError: boolean; /** * 最大重试次数 */ maxRetries: number; /** * 重试延迟(毫秒) */ retryDelay: number; }; } /** * 默认流水线处理器实现 */ export declare class PipelineProcessor extends BaseModule implements IPipelineProcessor { protected config: PipelineProcessorConfig; private processingStats; constructor(info: ModuleInfo, config?: Partial); /** * 处理输入数据,转换为流水线格式 */ processInput(input: any, context?: Record): Promise; /** * 处理输出数据,从流水线格式转换 */ processOutput(output: PipelineData, context?: Record): Promise; /** * 验证流水线数据格式 */ validatePipelineData(data: any): ValidationResult; /** * 获取支持的流水线格式 */ getSupportedFormats(): string[]; /** * 获取处理统计信息 */ getProcessingStats(): { totalProcessed: number; successful: number; failed: number; averageProcessingTime: number; }; /** * 重置统计信息 */ resetStats(): void; /** * 更新配置 */ updateConfig(config: Partial): void; /** * 初始化处理器 */ initialize(): Promise; /** * 清理资源 */ destroy(): Promise; private generateTraceId; private generateSpanId; private getDataType; private updateStats; /** * 处理消息(扩展BaseModule的消息处理) */ handleMessage(message: any): Promise; } //# sourceMappingURL=PipelineProcessor.d.ts.map