internal/obs(日志、指标、tracing)、internal/usage(计费用量采集)、internal/common/clientid(请求归属 ID 传播)。它们在系统中的位置见 架构总览。
职责与边界
internal/obs负责:安装进程级sloghandler、把 ctx 里的task_id/call_id/instance_id自动注入日志、声明所有 Prometheus 指标并提供Record*函数、初始化 OpenTelemetry tracer、启动/metricsHTTP 端点。internal/usage负责:按client_id累计每个 task 收敛时的 executor 有效 CPU 时间,定期冻结成 record 写入 Kafka topicchaintable-usage。internal/common/clientid负责:在 gRPC/HTTP 边界读写client-id与兼容 headerx-instance-id,在进程内用 ctx 携带完整 ID。- 不负责:日志聚合栈、告警规则、Prometheus 服务端、OTLP collector。这些都在部署侧。
- 不变量:
core/包(Sans-IO)不 importinternal/obs;只有 adapter 和cmd/代码打日志、记指标、开 span。详见 设计原则。 - 不变量:日志
instance_id、chaintable-logtrace_id、Prometheusinstance_idlabel 都用clientid.LogKey去掉instance:前缀;usage 消息和所有下游 header 保留完整原值。 - 不变量:Prometheus label 只允许闭合枚举(
status/phase/backend/kind等)加instance_id/function_id两个受 TTL 清扫器保护的高基数维度。不要用 task_id、call_id、URL、错误文本做 label。 - 不变量:Metrics 用
promauto注册到默认 registry;每个带instance_id或function_idlabel 的 vec 必须同时登记到internal/obs/metrics_reaper.go的instanceIDVecs/functionIDVecs,否则TestReaperVecListsMatchDeclarations会失败。
代码位置
核心类型与接口
关键签名(从代码复制):
数据流
一次请求的归属标识从 gRPC metadata 进入,沿 ctx 流到日志、span、指标和 usage 四个出口:1
边界归一化
internal/worker/adapters/grpc_server.go 的 SubmitTask、internal/syncinvoker/adapters/grpc_server.go 的 Invoke/DebugInvoke、internal/coordinator/adapters/grpc_server.go 都以 ctx = clientid.WithContext(ctx, clientid.FromIncomingGRPC(ctx)) 开头。client-id 优先,缺失回退 x-instance-id,都缺失得到 unknown。2
进程内传播
Worker 在 task runtime 里用
clientid.WithContext(ctx, rt.instanceID) 重新绑定(internal/worker/adapters/orchestrator_dispatch.go、orchestrator_phases.go),保证异步 goroutine 和 IO 回调不丢 ID。obs.WithInstanceID 只是 clientid.WithContext 的别名,不存第二份值。3
日志与指标出口
ContextHandler.Handle 用 clientid.LogKey 清洗后写 instance_id;metricInstanceID 对 Prometheus label 做同样清洗并刷新序列的 last-seen。chlog bridge 用 LogInstanceIDFromCtx 作为 trace_id。4
usage 与下游
task 收敛时
orchestrator_phases.go 调 o.usageRecorder.Record(taskInstanceID, taskUsage),传完整 ID。internal/sdk/blockdb、logicaltypes、noderpc、meta adapter 用 IntoOutgoingGRPC/IntoHTTP 双写两个 header 给下游。日志
格式与字段
- 每条日志一行 JSON,写到 stderr(
internal/obs/logger.go)。time是毫秒精度 UTC(2006-01-02T15:04:05.000Z07:00),与 Python executor 对齐,方便同一进程树的 Go/Python 日志混流解析。 - 关联字段由 ctx 自动注入:
task_id、call_id、executor_id、attempt_seq(绑定后即使为 0 也输出)、instance_id。 - 两条规则(来自
docs/specs/logging.md§4.1):有 ctx 的函数用slog.InfoContext(ctx, ...);没有 ctx 的 goroutine 闭包必须接收调用方用slog.Default().With(obs.AttrCallID, id)预绑定的*slog.Logger。 - 字段名一律走
obs.Attr*常量。错误字段是err。 obs.Init()不带参数;obs.BindProcess没有任何cmd/或app调用点,service/worker_addr只出现在个别显式传参的日志里。
结构化日志类型(type 字段)
obs.LogType* 常量是 type 字段的全部合法值。用户视角最重要的几个:
LogTypeTaskPhaseFinished、LogTypeCallFailed、LogTypeIOError 常量仍在但不再单独输出,内容并入 task_finished。
chaintable-log 桥接
设置CHAINTABLE_LOG_BROKERS 后,obs.InitChlog 把 ContextHandler 的内层 handler 换成 chlogBridgeHandler:每条记录仍写 stderr,同时按 type → chlog WithType、其余 attrs → data payload、ctx 里的 instance id → trace_id 转发到 Kafka(默认 topic instance-logs)。没有 instance id 的记录(coordinator 后台日志等)被 SDK 静默丢弃。
指标
命名与注册
- 命名空间
blockx,snake_case,counter 以_total结尾,直方图带单位后缀(_milliseconds、_seconds_total、_bytes)。 - 声明方式统一为包级
promauto.New*Vec(...),通过obs.Record*函数暴露给 adapter;adapter 不直接持有prometheus.*对象。IO 自适应限流是例外,用自定义prometheus.Collector拉快照。 - task 维度指标在 task 收敛时一次性观测:avg/p90/p99 由
internal/io/core/stats/的 CKMS 累积器在进程内算好,用statlabel 各 observe 一次。执行次数加权平均则用*_duration_seconds_total / *_duration_samples_total两个 counter 相除。
指标族与定义文件
internal/worker/adapters/call_metrics_test.go、executor_metrics_test.go、task_barrier_metrics_test.go 是这些记录点的测试,不是定义文件。internal/worker/adapters/resource_metrics.go 是心跳用的 CPU/内存采样器(WorkerResourceSampler),不导出 Prometheus 指标。
最重要的指标
Tracing
obs.InitTracing(ctx, service)始终设置 W3C TraceContext propagator。exporter 选择:BLOCKX_OTEL_SPAN_FILE优先(JSON-lines 追加写文件,多进程共享)→OTEL_EXPORTER_OTLP_ENDPOINT(OTLP/gRPC,insecure)→ 都没设则 noop。obs.TracingEnabled()让热路径跳过tracer.Start。- 采样:默认
AlwaysSample;OTEL_TRACES_SAMPLER_ARG=0.1这类值切到ParentBased(TraceIDRatioBased)。生产必须显式设置。 - 服务名由入口传入:
blockx-worker、blockx-bundle-worker、blockx-coordinator、blockx-bundle-coordinator、blockx-syncinvoker(cmd/*/main.go的Profile.Service),Python 侧是blockx-executor。 - 自动 span:Coordinator→Worker gRPC 用
otelgrpc.NewClientHandler()/NewServerHandler()(internal/coordinator/adapters/worker_rpc.go、internal/worker/app/app.go)。 - 手工 span:Worker↔Executor UDS 用
internal/worker/adapters/executor/send.go的udsTracer(tracer 名github.com/Chaintable/blockx/worker/uds),span 名uds.send.<MsgType>/uds.recv.<MsgType>/worker.io.dispatch;trace context 通过MessageEnvelope.traceContext字段过 UDS。Sync Invoker 用obs.Tracer("blockx-syncinvoker")。 BLOCKX_OTEL_UDS_EVENTS=1打开obs.AddUDSEvent/AddUDSEventAt的细粒度 UDS pipeline event,仅 perf 调试用。analyze_spans.py读 span 文件,把 Gostdouttrace与 PythonConsoleSpanExporter两种 JSON 归一,按 trace 打印 span 树和 event 偏移。典型用法:
docs/specs/tracing.md §4。
Usage 采集
- 采什么:每个 task 收敛时
taskRuntime.snapshotUsage()(internal/worker/adapters/task_runtime.go)把 executor 上报的有效执行时间(callEffectiveSum,毫秒)换成usage.Sample{CPUTimeMicros}。不含 executor 排队、IO 等待、Builder/Result Handler 的 Go 侧 CPU。 - 归属:
Collector.Record(clientID, sample)按完整client_id累计到active;空 ID 记为unknown;CPUTimeMicros<=0丢弃。 - 往哪写:默认每 5s
Flush,把active冻结成[]usage.Record(idUUID、client_id、service、resource_type、usage毫秒、timestamp冻结时刻),经KafkaSink.Publish同步写 topicchaintable-usage(key 为client_id,RequiredAcks=RequireAll,MaxAttempts=1)。失败时 record 留在pending,下轮以同一id重试;消费端按id去重。 - 服务枚举:
usage.ServiceBlockXWorker = "blockx_worker"(cmd/worker)、usage.ServiceBlockXBundleWorker = "blockx_bundle_worker"(cmd/bundle_worker),resource_type固定usage.ResourceTypeCompute = "compute"。由Profile.UsageService绑定,配置不能改。Sync Invoker 不创建 collector。 - fail-closed:
NewKafkaCollector在未显式Disabled时要求 brokers 非空、topic 有分区,探测失败进程拒绝启动。
Client id 传播
状态与生命周期
- usage Collector:
Start启动 ticker goroutine;Close设置 closing barrier(之后Record被拒绝)、停 ticker、在调用方给的 shutdown ctx 内每 200ms 重试最终Flush,然后sink.Close。Close返回错误视为计费关键错误,Worker 以非零退出。 - Prometheus 序列 TTL:
instance_id/function_id的 label 值闲置超过BLOCKX_METRICS_SERIES_TTL_MS(默认 24h)后由清扫器(默认每 10min,BLOCKX_METRICS_SERIES_SWEEP_INTERVAL_MS)用DeletePartialMatch删除,每轴每轮最多 256 个值。显式非正数关闭。清扫器随StartPrometheusEndpoint启动,即使没配抓取地址也会跑。 - 表读写日志去重:
table_read10 分钟窗口、table_write1 分钟窗口,64 分片、每分片最多 2048 个 key。 - tracing / UDS events 开关:进程启动时读一次环境变量,运行中修改无效。
监控端点
Worker 和 Coordinator 当前没有 HTTP 健康探针,也没有 gRPC health service。
配置
全部通过环境变量(Worker 也接受 JSON 配置文件的对应字段,见internal/worker/app/config.go 的 Prometheus/Chlog/Usage):
部署侧的完整变量清单见 blockx 仓库
docs/deploy.md §9,以及站内 部署概览。
性能分析方法
- 方法论:blockx 仓库
docs/specs/perf-methodology.md(建模 → 测量 → 比对 → 归因 → 决策的闭环);复现指引docs/specs/perf-test-how-to-v2.md;结论在perf-test-report-v2.md。 - 用例:
e2e/perf/按 S0–S5 阶段前缀组织(s0_framework_*、s1_builder_*、s2_call_count_*、s3_cpu_*、s4_io_*、s5_plugin_*),索引在e2e/perf/README.md。产物落e2e/perf/_artifacts/<TestName>/spans.jsonl、pyspy/。 - 跑法:
./test.sh -run TestPerf_XXX在 cgroup(默认 4 vCPU / 8 GB)下跑;直接go test ./e2e/perf/...是无限制对照。带with_trace的子测试自动设置BLOCKX_OTEL_SPAN_FILE。 - 线上压测:
go run ./cmd/perf -coordinator host:port -n 100 -c 4 -mode full|meta|reserve,输出 P50/P90/P99。 - 进程级归因:
WORKER_CPU_PROFILE/WORKER_TRACE写文件,或抓/debug/pprof/profile。
扩展点
- 新增一个日志字段或
type:在internal/obs/logger.go加Attr*/LogType*常量;如果是用户可见的 type,同步更新docs/specs/blockx-log-types.md。 - 新增一个 Prometheus 指标:在
internal/obs/metrics.go(或对应子系统的metrics_*.go)用promauto声明并加Record*函数;label 只能是闭合枚举加instance_id/function_id;带这两个 label 的 vec 必须加入metrics_reaper.go的instanceIDVecs/functionIDVecs。instance_id值必须经过metricInstanceID。 - 给新的 IO backend 加自适应指标:通过
internal/io/adaptive/observed注册,它会调用obs.RegisterIOAdaptiveMetrics(backend, provider)。 - 加 span:只在 adapter 层;先考虑能否用 event(
obs.AddUDSEvent)表达;跨 goroutine 用 child span 而不是往已 End 的父 span 加 event。 - 新增一个需要计费的入口:在
cmd/<entry>/main.go的Profile.UsageService传新的usage.Service*常量;service/resource_type不允许从配置覆盖。 - 新增一个下游 adapter:出站前调用
clientid.IntoOutgoingGRPC(ctx)或clientid.IntoHTTP(req),不要自己读 metadata。
测试
internal/obs/handler_test.go:ContextHandler的 slogtest 合规、字段注入、WithAttrs/WithGroup保持注入、LogKey清洗。internal/obs/chlog_bridge_test.go:桥接双写、type不进 data、group 前缀、trace_id清洗。internal/obs/metrics_test.go、metrics_function_test.go、stream_metrics_test.go、metrics_syncinvoke_test.go、metrics_io_adaptive_test.go:各Record*函数写出的 label 与数值。internal/obs/metrics_reaper_test.go、metrics_reaper_consistency_test.go:TTL 清扫器行为;TestReaperVecListsMatchDeclarations解析源码,保证带instance_id/function_id的 vec 都登记到清扫列表。internal/obs/table_dedupe_test.go:去重窗口与分片上限。internal/usage/collector_test.go、metrics_test.go:closing barrier、Kafka topic 探测、发送失败指标。internal/common/clientid/client_id_test.go:header 优先级、双写、LogKey。- 记录点测试在
internal/worker/adapters/call_metrics_test.go、executor_metrics_test.go、task_barrier_metrics_test.go;进程级日志 schema 断言在cmd/worker/worker_process_logging_test.go。
相关文档
blockx 仓库 spec:docs/specs/logging.md— 日志 schema、字段表、Go/Python 调用契约。docs/specs/blockx-log-types.md— 用户视角的type日志清单、task_finished字段定义、从日志推导指标的方法。docs/specs/tracing.md— 启用方式、完整 span/event 清单、analyze_spans.py用法、排障表。docs/specs/client-id-propagation.md— header 优先级、出站双写、日志/监控清洗规则。docs/specs/2026-07-13-blockx-usage-collection.md— usage 消息格式、CPU 口径、flush 与故障语义。docs/specs/perf-methodology.md、docs/specs/perf-test-how-to-v2.md— 性能测试方法与复现。docs/deploy.md§6、§9 — 部署侧可观测性配置与环境变量。