> ## Documentation Index
> Fetch the complete documentation index at: https://docs.blockx.chaintable.com/llms.txt
> Use this file to discover all available pages before exploring further.

# 可观测性

> BlockX 的结构化日志、Prometheus 指标、OpenTelemetry tracing、usage 采集与 client id 传播

可观测性不是一个独立进程，而是 Worker、Coordinator、Sync Invoker、Bundle 集群和 Python Executor 共享的一组 adapter 层基础设施。它由三个 Go 包组成：`internal/obs`（日志、指标、tracing）、`internal/usage`（计费用量采集）、`internal/common/clientid`（请求归属 ID 传播）。它们在系统中的位置见 [架构总览](/architecture/overview)。

## 职责与边界

* `internal/obs` 负责：安装进程级 `slog` handler、把 ctx 里的 `task_id`/`call_id`/`instance_id` 自动注入日志、声明所有 Prometheus 指标并提供 `Record*` 函数、初始化 OpenTelemetry tracer、启动 `/metrics` HTTP 端点。
* `internal/usage` 负责：按 `client_id` 累计每个 task 收敛时的 executor 有效 CPU 时间，定期冻结成 record 写入 Kafka topic `chaintable-usage`。
* `internal/common/clientid` 负责：在 gRPC/HTTP 边界读写 `client-id` 与兼容 header `x-instance-id`，在进程内用 ctx 携带完整 ID。
* 不负责：日志聚合栈、告警规则、Prometheus 服务端、OTLP collector。这些都在部署侧。
* 不变量：`core/` 包（Sans-IO）不 import `internal/obs`；只有 adapter 和 `cmd/` 代码打日志、记指标、开 span。详见 [设计原则](/architecture/design-principles)。
* 不变量：日志 `instance_id`、chaintable-log `trace_id`、Prometheus `instance_id` label 都用 `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_id` label 的 vec 必须同时登记到 `internal/obs/metrics_reaper.go` 的 `instanceIDVecs`/`functionIDVecs`，否则 `TestReaperVecListsMatchDeclarations` 会失败。

## 代码位置

| 路径                                                                                                                            | 用途                                                                                             |
| ----------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------- |
| `internal/obs/logger.go`                                                                                                      | `Init`、`ContextHandler` 的 ctx 绑定 helper（`WithTaskID` 等）、所有 `Attr*`/`LogType*` 常量               |
| `internal/obs/handler.go`                                                                                                     | `ContextHandler`：把 ctx 中的关联字段注入每条 `slog.Record`                                                |
| `internal/obs/chlog_bridge.go`                                                                                                | chaintable-log Kafka 桥接：把每条 slog 记录同时转发到 `instance-logs` topic                                 |
| `internal/obs/table_dedupe.go`                                                                                                | `table_read`/`table_write` 日志按 `(instance, table, operation)` 时间窗去重                            |
| `internal/obs/metrics.go`                                                                                                     | 绝大多数 `blockx_task_*`/`blockx_worker_*` 指标声明、`Record*` 函数、`StartPrometheusEndpoint`             |
| `internal/obs/metrics_reaper.go`                                                                                              | `instance_id`/`function_id` 序列 TTL 清扫器                                                         |
| `internal/obs/metrics_io_adaptive.go`                                                                                         | `blockx_io_adaptive_*` 自定义 Collector（IO 自适应并发限制）                                               |
| `internal/obs/metrics_noderpc_batch.go`                                                                                       | `blockx_noderpc_semantic_batch_*`                                                              |
| `internal/obs/metrics_syncinvoke.go`                                                                                          | `blockx_syncinvoke_*`                                                                          |
| `internal/obs/tracing.go`                                                                                                     | `InitTracing`、`Tracer`、`AddUDSEvent`、tracing 环境变量常量                                            |
| `internal/usage/collector.go`                                                                                                 | `Collector`、`Record`/`Sink` 接口、active/pending 缓冲与 flush 逻辑                                     |
| `internal/usage/kafka.go`                                                                                                     | `KafkaSink`、`NewKafkaCollector`（启动时探测 topic 分区）                                                |
| `internal/usage/metrics.go`                                                                                                   | `blockx_usage_*` 三个发送健康度指标                                                                     |
| `internal/common/clientid/client_id.go`                                                                                       | header 常量、`FromIncomingGRPC`/`IntoOutgoingGRPC`/`IntoHTTP`/`LogKey`                            |
| `internal/coordinator/adapters/etcd_watcher_metrics.go`                                                                       | `blockx_coordinator_registry_workers`                                                          |
| `internal/io/core/stats/`                                                                                                     | task 级 IO/call 延迟的 CKMS 分位数累积器（喂给 `task_finished` 日志和 `RecordTask*Converged`）                  |
| `internal/worker/app/app.go`、`internal/coordinator/app/app.go`、`cmd/syncinvoker/main.go`                                      | 三类进程的接线：`obs.Init` → `InitTracing` → `StartPrometheusEndpoint` → `InitChlog` → usage collector |
| `python/blockx_executor/logging_setup.py`、`tracing_setup.py`                                                                  | Python executor 侧对应实现，共用同一组环境变量                                                                |
| `analyze_spans.py`（仓库根）                                                                                                       | 把 Go/Python 两种 span JSON 归一并按 trace 打印 span 树                                                  |
| `docs/specs/logging.md`、`blockx-log-types.md`、`tracing.md`、`client-id-propagation.md`、`2026-07-13-blockx-usage-collection.md` | 对应 spec                                                                                        |

## 核心类型与接口

| 标识符                                                                                                              | 文件                                       | 说明                                                                                                        |
| ---------------------------------------------------------------------------------------------------------------- | ---------------------------------------- | --------------------------------------------------------------------------------------------------------- |
| `obs.Init()`                                                                                                     | `internal/obs/logger.go`                 | 读 `BLOCKX_LOG_LEVEL`/`BLOCKX_LOG_FORMAT`，安装 JSON（或 text）handler 并包一层 `ContextHandler` 为 `slog.Default()`  |
| `obs.ContextHandler`                                                                                             | `internal/obs/handler.go`                | 实现完整 `slog.Handler`；`Handle` 时从 ctx 取 `task_id`/`call_id`/`executor_id`/`attempt_seq`/`instance_id` 追加到记录 |
| `obs.WithTaskID` / `WithCallID` / `WithExecutorID` / `WithAttemptSeq` / `WithInstanceID` / `WithExecutorCallIDs` | `internal/obs/logger.go`                 | 把关联字段绑到 ctx；空字符串不绑定                                                                                       |
| `obs.BindProcess(key, value)`                                                                                    | `internal/obs/logger.go`                 | 把进程级常量追加到默认 logger base attrs；启动期调用一次                                                                     |
| `obs.Attr*` / `obs.LogType*`                                                                                     | `internal/obs/logger.go`                 | 结构化字段名与 `type` 字段枚举，调用点必须用常量                                                                              |
| `obs.FailureAttrs(component, code, err, extra...)`                                                               | `internal/obs/metrics.go`                | 错误路径统一 attrs；debug 级别时附带 goroutine stack                                                                  |
| `obs.RecordTaskCallsConverged` / `RecordTaskIOConverged` / `RecordTaskPhase`                                     | `internal/obs/metrics.go`                | task 收敛时一次性写 call/IO/phase 指标                                                                             |
| `obs.PrometheusConfig` / `obs.StartPrometheusEndpoint` / `obs.HTTPHandler`                                       | `internal/obs/metrics.go`                | 抓取端点配置与启动；额外路由（pprof、健康探针）挂同一 HTTP server                                                                 |
| `obs.RegisterIOAdaptiveMetrics(backend, provider)`                                                               | `internal/obs/metrics_io_adaptive.go`    | 每个自适应 IO backend 注册一次快照函数                                                                                 |
| `obs.InitTracing(ctx, service)` / `obs.Tracer(name)` / `obs.TracingEnabled()`                                    | `internal/obs/tracing.go`                | 按环境变量选 exporter；未配置时安装 noop provider                                                                      |
| `obs.ChlogConfig` / `NewChlogLogger` / `InitChlog`                                                               | `internal/obs/chlog_bridge.go`           | chaintable-log Kafka 桥接的配置与安装                                                                             |
| `usage.Recorder` / `usage.Sample` / `usage.Record` / `usage.Sink`                                                | `internal/usage/collector.go`            | Worker 到 collector 的非阻塞入口、单 task 样本、Kafka 消息、发送后端                                                         |
| `usage.Collector` / `NewKafkaCollector`                                                                          | `internal/usage/collector.go`、`kafka.go` | 进程内聚合器及其 Kafka 构造函数                                                                                       |
| `clientid.FromIncomingGRPC` / `WithContext` / `FromContext` / `IntoOutgoingGRPC` / `IntoHTTP` / `LogKey`         | `internal/common/clientid/client_id.go`  | client id 全部读写规则                                                                                          |

关键签名（从代码复制）：

```go theme={null}
// internal/obs/metrics.go
func StartPrometheusEndpoint(ctx context.Context, service string, cfg PrometheusConfig, extra ...HTTPHandler) (func(context.Context) error, error)

// internal/obs/tracing.go
func InitTracing(ctx context.Context, service string) (func(context.Context) error, error)

// internal/usage/collector.go
type Recorder interface {
	Record(clientID string, sample Sample)
}
type Sink interface {
	Publish(ctx context.Context, records []Record) error
	Close(ctx context.Context) error
}
```

## 数据流

一次请求的归属标识从 gRPC metadata 进入，沿 ctx 流到日志、span、指标和 usage 四个出口：

```mermaid theme={null}
flowchart LR
  A["gRPC 请求 metadata: client-id / x-instance-id"] --> B["clientid.FromIncomingGRPC + WithContext"]
  B --> C["ctx 携带完整 ID, obs.WithTaskID / WithCallID"]
  C --> D["slog.*Context: ContextHandler 注入 instance_id (LogKey)"]
  C --> E["otel span: otelgrpc / uds tracer"]
  C --> F["obs.Record*: label instance_id (LogKey)"]
  C --> G["usage.Recorder.Record: 完整 client_id"]
  C --> H["clientid.IntoOutgoingGRPC / IntoHTTP 到下游 BlockDB / NodeRPC / Meta"]
  D --> D1["stderr JSON"]
  D --> D2["chlog bridge 到 Kafka instance-logs"]
  E --> E1["OTLP gRPC 或 BLOCKX_OTEL_SPAN_FILE"]
  F --> F1["/metrics 被 Prometheus 抓取"]
  G --> G1["Kafka chaintable-usage"]
```

<Steps>
  <Step title="边界归一化">
    `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`。
  </Step>

  <Step title="进程内传播">
    Worker 在 task runtime 里用 `clientid.WithContext(ctx, rt.instanceID)` 重新绑定（`internal/worker/adapters/orchestrator_dispatch.go`、`orchestrator_phases.go`），保证异步 goroutine 和 IO 回调不丢 ID。`obs.WithInstanceID` 只是 `clientid.WithContext` 的别名，不存第二份值。
  </Step>

  <Step title="日志与指标出口">
    `ContextHandler.Handle` 用 `clientid.LogKey` 清洗后写 `instance_id`；`metricInstanceID` 对 Prometheus label 做同样清洗并刷新序列的 last-seen。chlog bridge 用 `LogInstanceIDFromCtx` 作为 `trace_id`。
  </Step>

  <Step title="usage 与下游">
    task 收敛时 `orchestrator_phases.go` 调 `o.usageRecorder.Record(taskInstanceID, taskUsage)`，传完整 ID。`internal/sdk/blockdb`、`logicaltypes`、`noderpc`、`meta` adapter 用 `IntoOutgoingGRPC`/`IntoHTTP` 双写两个 header 给下游。
  </Step>
</Steps>

## 日志

### 格式与字段

* 每条日志一行 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` 字段的全部合法值。用户视角最重要的几个：

| `type`                                             | 触发点                                                                                 | 说明                                                                                                                             |
| -------------------------------------------------- | ----------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------ |
| `task_finished`（常量名 `LogTypeTaskConverged`）        | `internal/worker/adapters/orchestrator_phases.go`                                   | task 终态时一条汇总日志：`phases[]`、`call_functions[]`、`call_failed_clusters[]`、`io_error_clusters[]`、`io_backends`、`block_latency_ms` 等 |
| `task_activation_failed`                           | 同上                                                                                  | task 激活失败                                                                                                                      |
| `table_read` / `table_write` / `table_write_debug` | `internal/obs/table_dedupe.go`、`internal/plugin/event/*`                            | 表读写审计，按 `(instance, table, operation)` 在 10 分钟 / 1 分钟窗口内去重                                                                     |
| `audit_reject`                                     | `internal/worker/adapters/audit_gate.go`、`internal/syncinvoker/adapters/service.go` | 代码审计拒绝（enforce 与 dark 模式共用 type，靠 msg 区分）                                                                                      |
| `executor_wedge_reaped`                            | `internal/worker/adapters/orchestrator_wedge.go`                                    | 卡死 executor 被硬杀，值得告警                                                                                                           |
| `function_code_refresh`                            | `internal/functioncode/adapters/redis/syncer.go`                                    | Function Code View 刷新                                                                                                          |
| `usage_publish_failed` / `usage_publish_recovered` | `internal/usage/collector.go`                                                       | usage Kafka 发送失败/恢复，按 record 逐条输出                                                                                              |

`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 累积器在进程内算好，用 `stat` label 各 observe 一次。执行次数加权平均则用 `*_duration_seconds_total / *_duration_samples_total` 两个 counter 相除。

### 指标族与定义文件

| 指标族                                                                                 | 定义文件                                                    | 记录点                                                                                                                |
| ----------------------------------------------------------------------------------- | ------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------ |
| `blockx_task_*`（phase、calls、function、io、cache、latency）                              | `internal/obs/metrics.go`                               | `internal/worker/adapters/orchestrator_phases.go`、`orchestrator_actor.go`                                          |
| `blockx_worker_*`（active tasks、dispatcher、executor、stream、writer、barrier、retention） | `internal/obs/metrics.go`                               | `internal/worker/adapters/orchestrator.go`、`orchestrator_dispatch.go`、`executor/adapter.go`、`executor/handlers.go` |
| `blockx_function_code_*`                                                            | `internal/obs/metrics.go`                               | `internal/functioncode/adapters/redis/syncer.go`                                                                   |
| `blockx_io_adaptive_*`                                                              | `internal/obs/metrics_io_adaptive.go`                   | `internal/io/adaptive/observed/backend.go`                                                                         |
| `blockx_noderpc_semantic_batch_*`                                                   | `internal/obs/metrics_noderpc_batch.go`                 | `internal/sdk/noderpc/batchobserved/observer.go`                                                                   |
| `blockx_syncinvoke_*`                                                               | `internal/obs/metrics_syncinvoke.go`                    | `internal/syncinvoker/adapters/service.go`                                                                         |
| `blockx_usage_*`                                                                    | `internal/usage/metrics.go`                             | `internal/usage/collector.go`                                                                                      |
| `blockx_coordinator_registry_workers`                                               | `internal/coordinator/adapters/etcd_watcher_metrics.go` | 同文件                                                                                                                |

`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 指标。

### 最重要的指标

| 指标名                                                                                                                   | 类型                  | labels                                        | 含义                                              |
| --------------------------------------------------------------------------------------------------------------------- | ------------------- | --------------------------------------------- | ----------------------------------------------- |
| `blockx_task_phase_duration_milliseconds`                                                                             | histogram           | `instance_id, phase, status`                  | 每个 phase 墙钟耗时；`phase="all"` 是 task 总耗时          |
| `blockx_task_calls_total`                                                                                             | counter             | `instance_id, status`                         | task 收敛时累计的 call 成功/失败数                         |
| `blockx_task_function_calls_total`                                                                                    | counter             | `instance_id, function_id, call_kind, status` | 按 function 的终态逻辑 call 数                         |
| `blockx_task_function_effective_duration_seconds_total` / `_whole_duration_seconds_total` / `_duration_samples_total` | counter             | `instance_id, function_id, call_kind`         | 执行次数加权的 CPU / CPU+IO 耗时求和与样本数                   |
| `blockx_task_calls_effective_duration_milliseconds` / `blockx_task_calls_whole_duration_milliseconds`                 | histogram           | `instance_id, function_id, call_kind, stat`   | 每 task avg/p90/p99 各 observe 一次                 |
| `blockx_task_io_ops_total`                                                                                            | counter             | `instance_id, backend, operation, status`     | IO 操作数                                          |
| `blockx_task_io_backend_duration_milliseconds`                                                                        | histogram           | `instance_id, backend, operation, stat`       | 每 task IO backend 延迟 avg/p90/p99                |
| `blockx_task_io_cache_hits_total` / `blockx_task_io_singleflight_hits_total`                                          | counter             | `instance_id, backend`                        | 缓存 / singleflight 命中                            |
| `blockx_task_block_latency_milliseconds` / `blockx_task_upstream_latency_milliseconds`                                | histogram           | `instance_id`                                 | 从区块时间戳 / 上游表完成到 task 收敛的延迟                      |
| `blockx_task_errors_total`                                                                                            | counter             | `instance_id, component, error_code`          | task 级分类错误                                      |
| `blockx_worker_active_tasks`                                                                                          | gauge               | —                                             | WorkerCore 当前保留的 task 数                         |
| `blockx_worker_dispatcher_calls`                                                                                      | gauge               | `state`                                       | dispatcher 中 ready/running/waiting call 数       |
| `blockx_worker_executor_calls`                                                                                        | gauge               | `state`                                       | executor 心跳上报的内部 call 状态计数                      |
| `blockx_worker_executor_cpu_cores` / `blockx_worker_executor_cpu_seconds_total`                                       | gauge / counter     | —                                             | executor 进程 CPU 占用                              |
| `blockx_worker_logical_calls_terminal_total`                                                                          | counter             | `kind, outcome`                               | 逻辑 call 终态计数                                    |
| `blockx_worker_stream_productions_active` / `blockx_worker_stream_credit_waiters`                                     | gauge               | —                                             | 流式 call 生产与背压                                   |
| `blockx_io_adaptive_limit` / `blockx_io_adaptive_pressure`                                                            | gauge               | `backend, kind` / `backend, state`            | 自适应并发限制当前值与压力                                   |
| `blockx_syncinvoke_calls_total` / `blockx_syncinvoke_call_duration_milliseconds`                                      | counter / histogram | `status, failure_code` / `status`             | Sync Invoker 调用量与时延                             |
| `blockx_usage_publish_failures_total` / `blockx_usage_pending_records`                                                | counter / gauge     | `service, resource_type`                      | usage Kafka 发送健康度，建议 `increase(...[1m]) > 0` 告警 |
| `blockx_coordinator_registry_workers`                                                                                 | gauge               | `format`                                      | Coordinator 各 registry 格式发现的 worker 数           |

## 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 文件，把 Go `stdouttrace` 与 Python `ConsoleSpanExporter` 两种 JSON 归一，按 trace 打印 span 树和 event 偏移。典型用法：

```bash theme={null}
: > /tmp/blockx_spans.jsonl
BLOCKX_OTEL_SPAN_FILE=/tmp/blockx_spans.jsonl BLOCKX_OTEL_UDS_EVENTS=1 \
  go test -v -count=1 -run TestWorkerProcess_RpcWait ./cmd/worker/...
python3 analyze_spans.py /tmp/blockx_spans.jsonl
```

完整 span/event 清单见 blockx 仓库 `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`（`id` UUID、`client_id`、`service`、`resource_type`、`usage` 毫秒、`timestamp` 冻结时刻），经 `KafkaSink.Publish` 同步写 topic `chaintable-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 传播

| 位置              | 做什么                                                                                                   |
| --------------- | ----------------------------------------------------------------------------------------------------- |
| 入站 gRPC handler | `clientid.WithContext(ctx, clientid.FromIncomingGRPC(ctx))`：`client-id` > `x-instance-id` > `unknown` |
| 进程内             | `clientid.FromContext(ctx)` 取完整值；`obs.InstanceIDFromCtx` 同义（空时返回 `unknown`）                           |
| 日志/指标/chlog     | `clientid.LogKey(id)`：`instance:<id>` → `<id>`；无前缀原样返回                                                |
| 出站 gRPC         | `clientid.IntoOutgoingGRPC(ctx)`：复制 metadata，覆盖写两个 header                                             |
| 出站 HTTP         | `clientid.IntoHTTP(req)`：同上，从 `req.Context()` 取值                                                      |
| Shadow submit   | `internal/worker/adapters/shadow_forwarder.go` 把 ID 改写为 `<原值>-shadow` 后再 `IntoOutgoingGRPC`           |

## 状态与生命周期

* **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_read` 10 分钟窗口、`table_write` 1 分钟窗口，64 分片、每分片最多 2048 个 key。
* **tracing / UDS events 开关**：进程启动时读一次环境变量，运行中修改无效。

## 监控端点

| 端点                                               | 进程                                                              | 说明                                                                                                        |
| ------------------------------------------------ | --------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------- |
| `GET /metrics`（路径可由 `BLOCKX_PROMETHEUS_PATH` 覆盖） | worker、bundle worker、coordinator、bundle coordinator、syncinvoker | 由 `obs.StartPrometheusEndpoint` 在 `BLOCKX_PROMETHEUS_LISTEN_ADDR` 上服务；地址为空则不监听                            |
| `GET /debug/pprof/`                              | worker、bundle worker                                            | `internal/worker/app/app.go` 的 `workerDebugHTTPHandlers` 挂在同一 server 上；`WORKER_DEBUG_PPROF=false` 关闭，默认开启 |
| `GET /readyz`、`GET /healthz`                     | syncinvoker                                                     | `cmd/syncinvoker/health.go`；`/readyz` 在依赖未就绪或 draining 时返回 503                                            |

Worker 和 Coordinator 当前没有 HTTP 健康探针，也没有 gRPC health service。

## 配置

全部通过环境变量（Worker 也接受 JSON 配置文件的对应字段，见 `internal/worker/app/config.go` 的 `Prometheus`/`Chlog`/`Usage`）：

| 环境变量                                                                                    | 默认                 | 作用                                                                                                | 定义位置                                                           |
| --------------------------------------------------------------------------------------- | ------------------ | ------------------------------------------------------------------------------------------------- | -------------------------------------------------------------- |
| `BLOCKX_LOG_LEVEL`                                                                      | `info`             | `debug`/`info`/`warn`/`error`，非法值回退 `info`；Go 与 Python 共用                                         | `internal/obs/logger.go`                                       |
| `BLOCKX_LOG_FORMAT`                                                                     | `json`             | `json`/`text`                                                                                     | 同上                                                             |
| `BLOCKX_PROMETHEUS_LISTEN_ADDR`                                                         | 空（不监听）             | `/metrics` 监听地址，如 `:9090`                                                                         | `internal/obs/metrics.go`                                      |
| `BLOCKX_PROMETHEUS_PATH`                                                                | `/metrics`         | 抓取路径                                                                                              | 同上                                                             |
| `BLOCKX_METRICS_SERIES_TTL_MS`                                                          | `86400000`         | 高基数序列闲置回收 TTL；非正数关闭                                                                               | `internal/obs/metrics_reaper.go`                               |
| `BLOCKX_METRICS_SERIES_SWEEP_INTERVAL_MS`                                               | `600000`           | 清扫周期                                                                                              | 同上                                                             |
| `OTEL_EXPORTER_OTLP_ENDPOINT`                                                           | 空（noop）            | OTLP/gRPC 地址，可带 `http://` 前缀                                                                      | `internal/obs/tracing.go`                                      |
| `BLOCKX_OTEL_SPAN_FILE`                                                                 | 空                  | JSON-lines span 文件；与 OTLP 互斥，文件优先                                                                 | 同上                                                             |
| `OTEL_TRACES_SAMPLER_ARG`                                                               | 空（全采）              | 头部采样比例 0.0–1.0                                                                                    | 同上                                                             |
| `BLOCKX_OTEL_UDS_EVENTS`                                                                | 关                  | `1` 打开 UDS pipeline events                                                                        | 同上                                                             |
| `CHAINTABLE_LOG_BROKERS`                                                                | 空（关闭）              | chaintable-log Kafka brokers，逗号分隔                                                                 | `internal/obs/chlog_bridge.go`                                 |
| `CHAINTABLE_LOG_TOPIC`                                                                  | `instance-logs`    | chaintable-log topic                                                                              | 同上                                                             |
| `BLOCKX_USAGE_DISABLED`                                                                 | `false`            | 显式关闭计费；关闭时跳过其余 usage 校验                                                                           | `internal/usage/collector.go`                                  |
| `BLOCKX_USAGE_BROKERS`                                                                  | 空（未禁用时启动失败）        | usage Kafka brokers                                                                               | 同上                                                             |
| `BLOCKX_USAGE_TOPIC`                                                                    | `chaintable-usage` | usage topic                                                                                       | 同上                                                             |
| `BLOCKX_USAGE_FLUSH_INTERVAL_MS`                                                        | `5000`             | flush 周期，合法范围 `(0, 5000]`                                                                         | 同上                                                             |
| `WORKER_DEBUG_PPROF`                                                                    | `true`             | 是否在 metrics server 上挂 `/debug/pprof/`                                                             | `internal/worker/app/app.go`                                   |
| `WORKER_CPU_PROFILE` / `WORKER_TRACE` / `WORKER_MUTEX_PROFILE` / `WORKER_BLOCK_PROFILE` | 空                  | 进程生命周期内写 pprof / runtime trace 到文件；coordinator 对应 `COORDINATOR_CPU_PROFILE` / `COORDINATOR_TRACE` | `internal/worker/app/app.go`、`internal/coordinator/app/app.go` |

部署侧的完整变量清单见 blockx 仓库 `docs/deploy.md` §9，以及站内 [部署概览](/development/deployment)。

## 性能分析方法

* 方法论：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`。

更多见 [测试组织与命令](/development/testing)。

## 扩展点

* **新增一个日志字段或 `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。

## 测试

```bash theme={null}
go test ./internal/obs/... ./internal/usage/... ./internal/common/clientid/...
```

* `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 — 部署侧可观测性配置与环境变量。

站内页面：[Worker](/components/worker)、[Call 执行子系统](/components/call-execution)、[IO 访问子系统](/components/io-subsystem)、[Sync Invoker](/components/sync-invoker)、[部署概览](/development/deployment)、[测试组织与命令](/development/testing)。
