
Envoy AWS-EventStream-Parser 过滤器实战从 AWS Bedrock 流式响应提取元数据【免费下载链接】envoyCloud-native high-performance edge/middle/service proxy项目地址: https://gitcode.com/GitHub_Trending/en/envoyAWS-EventStream-Parser 是 Envoy 中的一个 HTTP 响应过滤器专门解析application/vnd.amazon.eventstream二进制流式协议把 AWS EventStream 消息的载荷与头部中的值提取出来并写入 dynamic metadata供访问日志、自定义过滤器、指标导出与链路追踪使用。本篇文章以 官方配置文档 为主线结合仓库内 proto 定义、过滤器实现与协议解析器源码完整讲解其配置方式、工作流程、统计指标与底层原理帮助你在 AWS Bedrock 流式 API 等场景下落地可观测性与成本追踪方案。过滤器概述该过滤器以类型 URLtype.googleapis.com/envoy.extensions.filters.http.aws_eventstream_parser.v3.AwsEventstreamParser进行配置v3 API 定义见 aws_eventstream_parser.proto。它目前只处理响应体通过 encoder 路径生效核心价值在于从 AWS EventStream 流式响应如 AWS Bedrock 的流式 API中提取 token 用量等信息写入 dynamic metadata 后供日志、指标、自定义过滤器、追踪系统消费采用类型化扩展架构typed extension过滤器自身负责 EventStream 二进制协议解析而消息载荷的解析交由可插拔的 content parser 完成JSON、XML、protobuf 等格式。在仓库中的注册信息可从 extensions_build_config.bzl 与 extensions_metadata.yaml 查到过滤器名称为envoy.filters.http.aws_eventstream_parser对应 Bazel 目标//source/extensions/filters/http/aws_eventstream_parser:config。该特性于近期版本引入变更记录见 changelogs/current/new_features/aws_eventstream_parser__http_filter_added.rst。配置一个过滤器需要三部分content parser指定如何解析消息载荷并提取值例如 JSON parsercontent parser 内的 rules定义选择器路径selector paths与元数据动作metadata actions可选的header rules直接从 EventStream 消息头部提取值写入元数据。当某条规则命中时提取到的值会被写入指定的 metadata namespace 与 key随后可被访问日志、自定义过滤器、指标系统或追踪链路消费。典型使用场景AWS Bedrock 可观测性与成本追踪AWS Bedrock 流式 API 在流式响应结束时会通过 EventStream 协议返回 token 用量信息。该过滤器可以把 token 计数等元数据提取出来用于日志、指标和可观测性http_filters: - name: envoy.filters.http.aws_eventstream_parser typed_config: type: type.googleapis.com/envoy.extensions.filters.http.aws_eventstream_parser.v3.AwsEventstreamParser response_rules: content_parser: name: envoy.content_parsers.json typed_config: type: type.googleapis.com/envoy.extensions.content_parsers.json.v3.JsonContentParser rules: - rule: selectors: - key: amazon-bedrock-invocationMetrics - key: inputTokenCount on_present: metadata_namespace: envoy.lb key: input_tokens type: NUMBER - rule: selectors: - key: amazon-bedrock-invocationMetrics - key: outputTokenCount on_present: metadata_namespace: envoy.lb key: output_tokens type: NUMBER上例从 EventStream 消息中提取inputTokenCount和outputTokenCount写入envoy.lbnamespace。提取后的元数据有多种消费方式日志访问日志通过%DYNAMIC_METADATA(envoy.lb:input_tokens)%引用动态元数据指标导出自定义 stats sink 可以读取该元数据并产出指标自定义过滤器下游过滤器可以读取并据此做出动作链路追踪元数据可附加到 trace span 上。工作原理以 AWS EventStream JSON content parser 为例过滤器按以下步骤工作Content-Type 检查过滤器检查响应Content-Type头是否为期望的application/vnd.amazon.eventstream仅比较 media type、忽略参数。源码实现位于 filter.cc 的isEventstreamContentType()先用StringUtil::cropRight(content_type, ;)裁掉;之后的参数再做大小写不敏感比较HTTP Content-Type 本身大小写不敏感。若 Content-Type 不匹配或缺失mismatched_content_type_计数器递增并跳过后续处理。EventStream 二进制协议解析按 AWS EventStream 规范解析消息校验 CRC 校验和并正确处理跨多个数据块chunk拆分传输的消息。协议解析器是独立的静态工具类EventstreamParser位于 eventstream_parser.cc。逐消息处理对每条完整的 EventStream 消息先按header rules评估消息头部再取出 payload 字节交给配置的content parser。JSON 解析JSON content parser 把 payload 当作 JSON 解析按配置的选择器在对象中导航取值。按规则写入元数据根据 content parser 中定义的规则写元数据on_present选择器成功从某条消息提取到值后立即执行on_missing延迟到流结束执行仅当on_present从未执行且选择器路径未找到时触发on_error延迟到流结束执行仅当on_present从未执行且发生解析错误时触发。延迟执行的语义on_missing与on_error的延迟执行保证早期消息里没有目标字段时不会阻止后续消息成功提取值。Bedrock 信封自动解包值得注意的一个实现细节processMessage()在把 payload 交给 content parser 之前会先调用unwrapBedrockEnvelope()自动探测 Bedrock 的InvokeModelWithResponseStream信封。Bedrock 会把实际载荷包装为类似{bytes: base64, p: ...}的 JSON该函数解析 JSON 后取bytes字段并做 Base64 解码得到内层 payload 再交给 parser若不是这种信封结构则返回nullopt直接使用原始 payload。相关逻辑见 filter.cc 与头文件中的注释 filter.h。过滤器的编码器链路过滤器继承自Http::PassThroughEncoderFilter见 filter.h在响应路径上实现三个回调encodeHeaders检查 Content-Type 并决定是否进入解析流程encodeData把数据累积进内部 Buffer调用processBuffer()尝试解析完整消息在流结束或提前处理完毕时调用finalizeRules()最后writeMetadata()把待批量的元数据通过streamInfo().setDynamicMetadata()落盘encodeTrailers若尚未结束处理补齐finalizeRules()。processBuffer()采用单趟循环反复调用EventstreamParser::parseMessage()遇解析错误数据损坏、流分帧被破坏则记录eventstream_error_并停止处理遇不完整消息则中断等待更多数据解析出完整消息后先处理头部规则再处理 payload。元数据写入采用按 namespace 批量累积structs_by_namespace_类型为absl::flat_hash_mapstd::string, Protobuf::Struct最后统一调用setDynamicMetadata提交见addMetadata()与writeMetadata()实现。配置详解完整配置示例http_filters: - name: envoy.filters.http.aws_eventstream_parser typed_config: type: type.googleapis.com/envoy.extensions.filters.http.aws_eventstream_parser.v3.AwsEventstreamParser response_rules: content_parser: name: envoy.content_parsers.json typed_config: type: type.googleapis.com/envoy.extensions.content_parsers.json.v3.JsonContentParser rules: - rule: selectors: - key: usage - key: total_tokens on_present: metadata_namespace: envoy.lb key: tokens type: NUMBER header_rules: - header_name: :event-type on_present: metadata_namespace: envoy.lb key: event_type核心配置项说明response_rules处理 EventStream 响应流的规则配置其中包含 content_parser 与 header_rules 两部分。response_rules.content_parser一个 TypedExtensionConfig 类型的扩展指定如何解析并提取消息载荷中的值。当前可用的 parserenvoy.content_parsers.json解析 JSON 内容用类 JSONPath 的选择器提取值配置项见 json_content_parser.proto。在 config.cc 的工厂方法中FilterConfig构造时通过Config::Utility::getAndCheckFactoryContentParser::NamedContentParserConfigFactory()从 TypedExtensionConfig 实例化 parser 工厂并把statsPrefix()拼进统计名。content_parser 在 proto 校验中被标记为 required。JSON Content Parser 配置使用envoy.content_parsers.json时在 typed_config 内配置 rulesrules要应用的规则列表每条规则包含rulejson-to-metadata 规则配置字段包括selectors选择器列表指定如何从 JSON payload 提取值例如usage-total_tokens这样的逐级 key 导航on_present选择器成功提取值时要写入的元数据on_missing选择器路径未找到时要写入的元数据on_error解析出错时要写入的元数据。stop_processing_after_matchesRuleConfig 层控制该规则成功匹配多少次后停止评估0默认对每条内容都评估后匹配覆盖先前的值除非设置preserve_existing_metadata_value等效于提取最后一次出现的值1首次成功匹配后停止评估该规则适合提取流早期就出现的值如模型名N 1预留未来使用当前被校验规则拒绝lte: 1。JSON content parser 的运行时行为可在 json_content_parser_impl.h 中看到每条规则内部维护match_count_、ever_matched_on_present 是否触发过和selector_not_found_状态并在整个流式会话中累计any_parse_error_从而支撑 on_present/on_missing/on_error 的语义。response_rules.header_rules可选的 EventStream 消息头部提取规则列表由过滤器直接评估不经 content parser。每条 header rule 包含header_name要匹配的 EventStream 头部名大小写敏感on_present头部存在时的元数据动作。头部的类型化值会自动转换为 Protobuf ValueBoolTrue/BoolFalse→bool_valueByte/Short/Int32/Int64/Timestamp→number_valueString→string_valueByteArray→string_value十六进制编码UUID→string_value格式化为xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx可选字段value可以覆盖头部的实际值。类型转换实现在headerValueToProtobufValue()见 filter.cc其中 UUID 使用 16 字节数组按 4-2-2-2-6 分段格式化为小写十六进制字符串。on_missing到流结束头部都未在任一消息中出现时的元数据动作此时必须设置value字段因为没有头部值可用。stop_processing_after_matches控制该规则成功匹配多少次后停止评估0默认对所有消息评估后匹配覆盖先前的值等效提取最后一次出现1首次匹配后停止适合提取流早期就出现的头部减少后续消息的处理开销N 1预留当前被拒绝lte: 1。当所有header rules 与所有content parser rules 都设置了匹配上限且均已满足时过滤器会整体停止处理该流processing_complete_ true这一点由allHeaderRulesSatisfied()与processMessage()返回的stop_processing联合判断。如果metadata_namespace为空默认使用envoy.filters.http.aws_eventstream_parser源码中常量FilterConfig::DefaultHeaderNamespace见 filter.h。停止处理的行为细节从 filter.cc 的processBuffer()可以看到停止机制的完整流程每解析出一条完整消息先处理头部规则processMessageHeaders再处理 payloadprocessMessage若 content parser 返回需要停止stop_processing为 true且所有设置了上限的 header rules 均已满足则置processing_complete_并在encodeData/encodeTrailers中触发finalizeRules()此时剩余未消费的 on_missing/on_error 延迟动作会按规则执行。统计指标过滤器在http.stat_prefix.aws_eventstream_parser.resp.parser_prefix*命名空间下输出统计。其中stat_prefix来自所属 HTTP connection managerHttpConnectionManager.stat_prefixparser_prefix来自 content parserJSON parser 为json.统计前缀拼接逻辑见FilterConfig::generateStats()filter.h。名称类型描述resp.parser_prefix.metadata_addedCounter成功写入的元数据条目总数resp.parser_prefix.metadata_from_fallbackCounter使用 on_missing 或 on_error 回退值写入的元数据条目总数resp.parser_prefix.mismatched_content_typeCounterContent-Type 与期望类型不匹配的响应总数resp.parser_prefix.empty_payloadCounterpayload 为空的 EventStream 消息总数resp.parser_prefix.parse_errorCountercontent parser 解析 payload 失败的消息总数resp.parser_prefix.preserved_existing_metadataCounter因 content parser 规则设置了preserve_existing_metadata_value而未写入元数据的次数resp.parser_prefix.eventstream_errorCounterEventStream 协议错误总数CRC 不匹配、格式非法resp.parser_prefix.type_conversion_errorCounter值无法转换为合法 Protobuf Value 类型的次数这些统计在 filter.cc 的相应路径上递增例如processMessage()中对空 payload 递增empty_payload_、对解析失败递增parse_error_addMetadata()中当 Protobuf Value 类型未设置时递增type_conversion_error_。preserved_existing_metadata会在 pending 批次或已提交的 dynamic metadata 中已存在同名 key 时递增filter.cc。AWS EventStream 协议实现细节过滤器按 AWS EventStream 规范实现了二进制协议解析核心代码在 eventstream_parser.cc 与 eventstream_parser.h二进制协议解析包含 prelude、headers、payload、trailer 的二进制消息格式。其中 prelude 固定 12 字节total_length4 字节 headers_length4 字节 prelude_crc4 字节trailer 为末尾 4 字节message_crc。CRC 校验同时校验 prelude CRC覆盖前 8 字节与消息 CRC覆盖除最后 4 字节外的全部内容使用 zlib 的crc32()实现与 AWS EventStream 规范一致。任一不匹配都会返回DataLossErrorPrelude CRC mismatch/Message CRC mismatch。消息分帧正确处理消息边界与不完整消息。解析器先检查缓冲区是否够 12 字节的 prelude随后校验total_length的合理性缓冲区不足整条消息时返回不完整结果bytes_consumed 0由过滤器等待更多数据。分块传输处理跨多个 TCP 包/HTTP chunk 拆分的消息。过滤器的processBuffer()会把数据累积进内部 Buffer线性化后循环解析并在每次解析后buffer_.drain()已消费的字节filter.cc。此外解析器内置了多道防御性上限eventstream_parser.h防止恶意或损坏的流导致无界缓冲与资源耗尽MIN_MESSAGE_SIZE 16 字节prelude trailer 的最小消息MAX_PAYLOAD_SIZE 24 MBMAX_HEADERS_SIZE 128 KBMAX_TOTAL_LENGTH prelude 最大头部 最大 payload trailer 的合计上限MAX_HEADER_STRING_LENGTH 327672^15 - 1。头部值支持 10 种类型BoolTrue、BoolFalse、Byte、Short、Int32、Int64、ByteArray、String、Timestamp、Uuid使用absl::variant承载对应类型bool / int8 / int16 / int32 / int64 / string / 16 字节数组。parseHeaders()按规范逐条解析头部1 字节名字长度、名字字节、1 字节类型字节随后按类型读取定长或带 2 字节长度前缀的变长值对截断、空字符串、超长、未知类型等情况均返回明确的错误状态。安全注意事项CRC 校验对每条消息的 prelude 与正文都做 CRC32 校验可检测损坏或被篡改的消息。需要说明的是CRC 是完整性校验而非消息认证码MAC正如 eventstream_parser.h 注释所指出的攻击者完全可以构造 prelude CRC 合法但total_length荒谬的消息因此解析器同时依赖MAX_TOTAL_LENGTH等上限来防止无界缓冲这是纵深防御的一部分资源上限payload、headers、单条字符串头部均设有最大长度防止超大消息造成内存压力解析失败即停止EventStream 流分帧一旦损坏便无法恢复过滤器会记录eventstream_error_并停止处理后续数据避免在损坏的流上继续产生错误的元数据错误降级解析错误不会导致响应被中断过滤器始终返回Continue元数据提取失败不影响正常代理转发。测试与验证仓库为该过滤器提供了完整的单元测试与集成测试可作为学习与验证的参考filter_test.cc过滤器单元测试覆盖 Content-Type 匹配、Bedrock 信封解包、header rules、停止处理、延迟动作等路径integration_test.cc端到端集成测试验证真实 HTTP 流场景下的行为config_test.cc配置解析与校验测试eventstream_parser_test.ccEventStream 二进制协议解析器的单元测试eventstream_parser_fuzz_test.cc针对协议解析器的模糊测试保障对畸形输入的健壮性。在部署前建议结合实际 Bedrock 响应样本在测试环境先行验证选择器路径与元数据 namespace/key 的映射是否符合预期再接入访问日志格式%DYNAMIC_METADATA(...)%或指标消费链路。【免费下载链接】envoyCloud-native high-performance edge/middle/service proxy项目地址: https://gitcode.com/GitHub_Trending/en/envoy创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考