internal/sdk/(目录名沿用历史,里面是各服务的 client,不是对外发布的 SDK)。每个子包把一个远端服务(BlockDB gRPC、链节点 JSON-RPC、Meta gRPC、Iceberg/Glue)包装成 IO 子系统认识的 iocore.BackendAdapter,再由 Worker 或 sync-invoker 在启动时装配进 WorkerIOScope。Builder、Writer Plugin 和 Function Code View 通过 TaskIO.Read/Write 或直接持有 adapter 使用它们。整体位置见 架构总览,调用链的上半段见 IO 访问子系统。
职责与边界
- 每个 SDK 包只负责一件事:把
iocore.ReadReq/iocore.WriteReq翻译成一次远端调用,并把响应放进ReadResult.Data/WriteResult.Data。请求类型由 SDK 包定义,不由internal/io/core定义。 - SDK 不做缓存、不做 singleflight、不做 task 窗口、不做准入控制。这些都在 IO 子系统;SDK 只提供
Registration(...)或assembly.Module帮助调用方把 adapter 包进adaptive.WrapBackend。 - SDK 负责错误分类:实现
iocore.ErrorClassifier.ClassifyError(决定 retryable / non-retryable / param)和adaptive.OutcomeClassifier.ClassifyAdaptiveOutcome(决定 AIMD 是否收缩)。本地 context 取消/超时一律不算下游过载。 - BlockDB SDK 只暴露原子 RPC。BatchWrite 的 Init → PUT → Commit 编排在
internal/plugin/event/batchwrite,不在 SDK(docs/specs/2026-07-30-bundle-write-batch-api.md§2)。 - 面向 Python executor 的 bridge 路径(BlockDB bridge、NodeRPC bridge、capability backend)只透传 protobuf/JSON-RPC 字节,不重建请求体。
- 连接策略属于 adapter 实例:指向 NLB 的 gRPC client 用
grpcclient.NewMultiConn,每个 adapter 内相同地址复用一个ClientConn,没有进程级连接注册表。 internal/duckdb/borrowed只提供 DuckDB 原生绑定之上的连接池与结果解码;SQL 拼装、S3 凭证、过滤下推留在internal/plugin/callbuilder/bundlescan。
代码位置
核心类型与接口
所有 adapter 都实现同一个接口(internal/io/core/model/adapter.go):
blockdb
AdapterConfig(adapter.go):每个 gRPC service 一个地址字段(TableReadAddr、TableWriteAddr、TableScanAddr、BlockReadAddr、BlockWriteAddr、BatchWriteAddr、BlockSubscribeAddr、L2BlockAddr、TimeReadAddr、TimeWriteAddr)加MaxRecvMsgSize(0 表示 32 MiB)。相同地址复用一个ClientConn;未配置的 service 在被调用时返回blockdb:service_not_configured(param 类错误)。TableSubscribeAddr字段存在并从BLOCKDB_TABLE_SUBSCRIBE_ADDR读入,但NewAdapter不会拨号它;它只在internal/worker/app/function_code.go里作为启用 BlockDB Function Code View 的条件之一。Adapter(adapter.go):typed 后端,BackendBlockDB = "blockdb"。Read用类型开关分派ReadOp[*blockdbpb.GetRowRequest]、BatchGetRowsRequest、FilterRowsRequest、GetBlockRowsRequest、GetEventRequest、GetStateRequest、GetValueRequest,把ResultSet转成 JSON 行数组;JSON/LIST/DICT 类逻辑列按ColumnMeta.logical_type还原为对象。Write分派WriteOp[*UpsertRowsRequest]、WriteOp[*DeleteRowsRequest]、WriteOp[*UpsertTimeRowsRequest]、BlockWriteReq、InitWriteJobReq、CommitWriteJobReq、CancelWriteJobReq。ReadOp[T]/WriteOp[T](types.go):泛型包装,把任意 proto 请求变成iocore.ReadReq/WriteReq,附带Key(cache key)、Op(统计名)、RTO/WTO(超时)。Filter、EqFilter、AndFilter、FilterJSON生成参数化 WHERE;Value(any)转*blockdbpb.Value。- 三种表语义(
proto/table.proto、block.proto、time.proto):Table(L1 行表:GetRow/BatchGetRows/FilterRows/UpsertRows/DeleteRows/Scan)、Block(L2 区块作用域:GetEvent/GetState带current_block,UpsertBlock带Block{block_id,height,timestamp},L2BlockReadService.GetBlockRows)、Time(GetValue(row_id,time_at)、UpsertTimeRows(time_at))。BlockWriteReq是多表 L2 写:Data map[table][]*WriteRow,adapter 对每张表发一次UpsertBlock。 - BatchWrite(
batch_write.go、proto/batch.proto):InitWriteJobReq/CommitWriteJobReq/CancelWriteJobReq各包一个BatchWriteService请求,响应以 proto 字节放进WriteResult.Data。上传 presigned URL、校验job_id一致性等编排在internal/plugin/event/batchwrite/writer.go。 BridgeAdapter(bridge_adapter.go):BackendBlockDBBridge = "blockdb_executor",接收 Pythonblockdb_bridge.BridgeChannel经 UDS 二进制 sidecar 发来的原始请求 protobuf,用capabilitywire.DecodeBinaryRequest解出 gRPC method path 后按bridgeReadMethods/bridgeWriteMethods白名单转发;bridgeBlockedWriteMethods拒绝函数代码发起 BatchWrite 生命周期;Scan/Subscribe流式方法不支持。响应字节原样放进ReadResult.Data。ScanClient/Adapter.ScanAll(online.go):TableScanService.Scan的流式读取;Function Code View 复用 typedAdapter.ScanAll。SubscribeClient目前只有构造与连接持有,没有Subscribe方法。- 错误分类(
adapter.go、adaptive.go):classifyBlockDBError解析 gRPC status 里的Error NNNN (SQLSTATE)得到mysql:NNNN/proxysql:NNNNcode,再按 gRPC code 映射;classifyAdaptiveOutcome把ResourceExhausted/DeadlineExceeded/Unavailable、传输错误和一组 MySQL/TiDB/ProxySQL 超时短语判为过载。 Registration(cfg, adaptive.Config, opts...)/BridgeRegistration(...):返回iocore.BackendRegistration;WithMultiConnConfig注入grpcclient.MultiConnConfig。- 测试替身:
DelayedGRPCStub(grpc_stub.go,供cmd/blockdb_grpc_stub挂到真实 gRPC server)、DelayedBridgeStubAdapter(bridge_stub.go,进程内 bridge stub,STUB_BLOCKDB_DELAY_MS触发)。
noderpc
NodeRPCReq(types.go):Go 侧直接调用的 JSON-RPC 2.0 请求,BackendNodeRPC = "rpc",Operation()返回 method,Timeout()固定 10s。Adapter(adapter.go):Read把请求体 POST 到endpoint;bridge 请求(Python executor 发来的rawBridgeRequester)从 JSON 元数据取params.chain,POST 到endpoint/<chain>,body 直接使用 sidecar 字节。Write不支持。每个请求带x-load-deadline/x-load-priority/x-load-retries头,HTTP 2xx 中的{"error":...}解析为jsonRPCError,非 2xx 变成HTTPStatusError。HTTPClientConfig{TransportCount, RefreshIntervalMs}+rotatingHTTPTransport(transport.go):8 个独立http.Transportshard 轮询,每 30s 平滑替换一个;HTTPPoolSizesForAdmission从 admission 初始/上限并发推导 idle 连接预算。Registration(endpoint, adaptive.Config, opts...)。- 语义批处理(
batching_backend.go、batch_core.go、batch_protocol.go、batch_config.go、batch_window*.go):BatchingBackend实现BackendAdapter并包在 admission 之外(assembly.Module.OuterWrap)。BatchConfig{Mode, MaxWait, MaxItems, FallbackCooldown}、BatchMode=off|shadow|on、BatchConfigFromEnv()、NewBatchingBackend(inner, cfg, BatchObserver)。可批方法只有getAddressCode、getStorageAt、getAddressBalance、getAddressNonce,且 block context 必须是固定 hash/height;物理请求是内部方法blockx_stateReadBatch,载荷是 BSRB/1 二进制。batchobserved.Observer()提供 obs 指标钩子。 - 二进制 sidecar:bridge 请求的直通路径在
Adapter.readFromBridge——缺少RawBinaryRequest()直接返回noderpc:invalid_bridge_request,没有 JSON 重建回退。Python 侧在python/blockx_executor/noderpc_bridge.py。
meta 与 logicaltypes
internal/sdk/meta/proto/*.proto的go_package指向 blockdb 仓库;make proto-meta用--go_opt=M...把三个文件重映射到internal/sdk/meta/gen(包名gen,调用方通常别名metapb)。logicaltypes.Adapter(logicaltypes/adapter.go):BackendLogicalTypes = "logicalTypes",ReadReq{TableID, RTimeout}→metapb.TableMetaServiceClient.GetTable→ JSON 编码的TableSchema{Type string, Columns []logical_types.Column}。ClassifyError把 schema/配置类 gRPC code 判为不可重试。LocalStubAdapter按 table ID("*"兜底)返回本地 schema。Registration(host, adaptive.Config, opts...)。meta.Adapter(meta/meta.go):HTTP GET/api/v1/meta/get_table_metadata?id=的旧客户端,BackendMeta = "meta"、MetaReadReq。仓库内没有生产代码 import 它,只有包内单测和//go:build live的meta_live_test.go。Worker 的META_ADDR走的是logicaltypes。
iceberg
Config{Namespace, Region}、NewPlanner(ctx, cfg)(planner.go):用 AWS 默认凭证链(EC2 上即 instance role)构造 Glue catalog,构造时不访问 Glue/S3。DataFilePlanner接口只有PlanDataFiles(ctx, tableName) ([]string, error);Planner.PlanDataFiles加载表、规划当前快照全部 data file,遇到 delete file 直接报错。ResolveDataFilesReq{TableID, PhysicalTableID, RequiredBundle, RTimeout}(io_adapter.go):BackendIcebergResolve = "iceberg",CacheKey()只用逻辑TableID;实现CachedReadValidator.AcceptCachedRead——缓存索引的 high water ≥RequiredBundle才复用,否则触发一次全表刷新(对应提交 “iceberg: refresh indexes above bundle high water”)。IcebergIOAdapter:只读,把路径列表编码成BFI1索引(bundle_file_index.go:公共前缀 + 每 bundle 一条记录,同一 bundle 出现两个 live 文件即报错)。LookupBundleFileIndex(encoded, bundle)、BundleFileIndexHighWater(encoded)供bundlescan解码。- Worker 把它注册为
SystemCache: true、ReadOnly: true的模块,准入是固定并发IO.SystemIOCacheRefreshConcurrency(默认 2)。所需 IAM 见 blockx 仓库docs/deploy.md§3.3(glue:GetTable、s3:GetObject,SSE-KMS 时加kms:Decrypt)。
localtestservice
Service(service.go)直接实现生成的LocalTestServiceServer:HelloWorld返回"hello, <name>";BlockDBGetRow通过注入的BlockDBReader(宿主的 raw blockdb adapter 或 stub)发一次ReadOp[*blockdbpb.GetRowRequest]。Adapter(adapter.go):BackendLocalTestService = "localtestservice",拆 capabilitywire 二进制信封、按 method pathproto.Unmarshal、进程内调Service。Module(blockDB BlockDBReader) assembly.Module:只读、固定准入 4。examples/local_test_service/是 Python 侧端到端示例(blockx-py 的LocalTestService类 → BridgeChannel → worker);新增 capability backend 的逐文件模板见 blockx 仓库docs/capability-backend-guide.md。
router
Adapter(router/adapter.go):BackendRouter = "router",只支持routerpkg.RouterFindFunctionsFullMethodName,业务在依赖github.com/chaintable/router/go的Service.FindFunctions。Module()只读、固定准入 100。
grpcclient
New(target, opts...):DNSround_robin的 insecure ClientConn。NewMultiConn(target, MultiConnConfig, opts...):注册名为blockx_multi_conn的自定义 balancer,在一个ClientConn内维持Connections条 SubConn,每RefreshIntervalMs换一条;DefaultMultiConnConfig()= 8 条 / 10000 ms。grpc.WithDisableServiceConfig()防止 resolver 下发的 service config 覆盖策略。DelimitedWriteRowEncoder(delimited_write_row.go):按logical_types.Column决定Valueoneof kind,直接写 length-delimitedblockdb.v1.WriteRow字节,不构造生成的 proto message;BatchWrite 上传文件用它。
duckdb/borrowed
Pool(pool.go):NewPool(db, maxOpen, maxIdle)、Acquire(ctx)、Release(conn)、Close();WatchContext(ctx, conn)(context.go)用duckdb_interrupt把 context 取消传给 DuckDB;ScanResult(ctx, res, yield)(rows.go)按 chunk 迭代并把 VARCHAR/BLOB 以零拷贝视图交给回调——回调返回后值失效,需要保留必须复制。- 转换逻辑改编自
duckdb-go,直接调用duckdb-go-bindings;LICENSE.duckdb-go、LICENSE.duckdb-go-bindings与README.md记录来源与所有权。
cmd/blockdb_grpc_stub
main.go用blockdb.RegisterDelayedGRPCStub(registrar, readDelay)起一个真实 gRPC server:读 RPC 等待-read-delay后返回固定 token 形状(GetRow、GetState),写 RPC 立即返回;-probe-address模式只探测已运行的 stub 并校验其 delay。用途和 worker 侧环境要求见cmd/blockdb_grpc_stub/README.md。
数据流 / 执行流程
1
装配
internal/worker/app/app.go(sync-invoker 在 cmd/syncinvoker/io.go)把每个后端声明为一个 assembly.Module{Kind, Enabled, Build, Stub, Admission, OuterWrap}。Enabled 为假时用 Stub(如 NODE_RPC_ENDPOINT 为空用 devstub,META_ADDR 为空用 logicaltypes.LocalStubAdapter)。assembly.Build 统一经 adaptiveobserved.WrapBackendRegistration 套上带指标的准入包装(internal/io/assembly/assembly.go)。2
请求进入
Go 调用方构造 SDK 请求类型(如
blockdb.ReadOp[*blockdbpb.GetBlockRowsRequest])交给 TaskIO.Read;Python 调用方的 gRPC/JSON-RPC 字节经 UDS sidecar 到 Worker,被包成带 RawBinaryRequest() 的 IORequest,Backend() 由 wire 上的 backend 名决定。3
SDK 执行
adapter 只看自己认识的请求类型;不认识的返回错误。gRPC 调用前用
clientid.IntoOutgoingGRPC(ctx) 传播 client id,HTTP 用 clientid.IntoHTTP。4
错误回流
adapter 返回的 error 先经
ClassifyAdaptiveOutcome 喂给 AIMD,再经 ClassifyError 变成 IOError 决定 Call/IO 层重试。两处都把本地 context 结束判为 Ignore / timeout,而不是下游过载。状态与生命周期
- gRPC 连接:
NewMultiConn的 balancer 为每个 slot 先建 candidate SubConn,READY 后才提升并排空旧连接;首次刷新在(0, refreshInterval]内随机延迟。Adapter.Close()关闭全部ClientConn。 - NodeRPC transport:
rotatingHTTPTransport每个 refresh 间隔退役一个 shard,退役 shard 等在飞请求结束再关闭;Close()拒绝新请求并等待全部排空。 - 批处理 group:
groupCollecting → groupFlushing → groupDone;同 chain、同固定 block、同 timeout、同 client id 的请求进同一 group,MaxWait(默认 500µs,Linux 上用 timerfd 计时)或MaxItems(默认 16)触发 flush。收到-32601进入 Worker 级 cooldown(默认 60s),到期后只放一个 half-open 探针 batch。Close()让收集中的 group 失败并取消在飞物理请求。 - 幂等:BlockDB 各 RPC 幂等,所以 IO 层可以对 Result 阶段的 Init/Commit/Cancel 做有限重试;NodeRPC 只有 Read 路径。
配置
Worker 的字段在internal/worker/app/config.go(WorkerFullConfig),环境变量覆盖在 internal/worker/app/app.go;默认值从 DefaultWorkerFullConfig() 读。
扩展点
- 新增一个 BlockDB RPC:改
internal/sdk/blockdb/proto/*.proto→make proto-blockdb→ 在Adapter.dispatchRead/dispatchWrite加ReadOp[*NewRequest]分支 → 如需对 Python 开放,把FullMethodName加进bridge_adapter.go的bridgeReadMethods/bridgeWriteMethods→ 在blockdb_test.go(bufconn + gomock server)加用例。 - 新增一个 backend:实现
BackendAdapter(可选ErrorClassifier、OutcomeClassifier),提供Registration或Module,然后在internal/worker/app/app.go的backendModules加一条声明。gRPC 风格的 capability backend 直接照internal/sdk/localtestservice复制,步骤在docs/capability-backend-guide.md。 - 新增可批处理的 NodeRPC 方法:改
batch_core.go的batchMethods与batch_protocol.go的 BSRB kind,同时改 Leafage 侧 handler,并遵循docs/specs/noderpc-semantic-batching.md§7.3 的版本演进规则。 - 调整连接策略:只改对应 adapter 的
grpcclient.MultiConnConfig,不要新增全局连接池。 - 重新生成 proto:需要
protoc、protoc-gen-go、protoc-gen-go-grpc在$(go env GOPATH)/bin(Makefile顶部有安装命令)。make proto会先跑proto-blockdb、proto-meta、proto-localtestservice,再生成api/grpc/*。三个 SDK 目标都带require_unimplemented_servers=false。
测试
blockdb:blockdb_test.go用bufconn起 gRPC server,mock_server_test.go是 mockgen 生成的 service mock;partial_config_test.go覆盖未配置 service 的干净失败;batch_write_test.go校验 proto 与 BlockDB 仓库契约一致;bridge_*_test.go覆盖 bridge 白名单与原始字节路径。noderpc:noderpc_test.go用httptest.Server;batching_backend_test.go覆盖跨 task 合批、错误隔离、-32601回退与 cooldown 探针;batch_window_linux_test.go在 timerfd 不可用时t.Skip。- 需要外部服务的测试都会跳过:
TestAdapter_ReadReal需要NODERPC_REAL_URL(且非-short);meta_live_test.go需要-tags live和META_API_URL。CI 的make test用-short,因此不会碰真实服务。 iceberg:planner_test.go只测 delete-file 拒绝和索引编码,不访问 Glue。duckdb/borrowed:pool_test.go、cells_test.go、rows_test.go需要 cgo 与 DuckDB 绑定,随go test编译。- E2E:
e2e/system/logical_type/用metapb起本地 Meta gRPC 服务测 logical types 与 dbscan;进程级测试各自起 mock gRPC server(cmd/worker/worker_process_blockdb_bridge_test.go实现TableReadServiceServer,cmd/worker/worker_process_function_code_test.go用internal/testutil.StartMockFunctionBlockDB);性能基准用cmd/blockdb_grpc_stub。
相关文档
- blockx 仓库
docs/specs/2026-07-30-bundle-write-batch-api.md:BatchWrite Job 生命周期,SDK 只提供三个原子 RPC。 - blockx 仓库
docs/specs/noderpc-semantic-batching.md:批处理 facade 的准入规则、group 状态机、BSRB/1 协议、配置与回退。 - blockx 仓库
docs/specs/2026-08-06-noderpc-binary-sidecar.md:NodeRPC bridge 二进制直通。 - blockx 仓库
docs/specs/2026-07-23-grpc-client-balancing.md:NewMultiConn的设计与各组件配置。 - blockx 仓库
docs/specs/io-subsystem.md、docs/io-backend-module-design.md、docs/capability-backend-guide.md:backend adapter 契约与新增后端指南。 - blockx 仓库
docs/deploy.md§3.3(Iceberg IAM)、§9.2(Worker 环境变量)。 - 站内:IO 访问子系统、Plugin 系统、Function Code View、Python Executor、Bundle 集群、协议与接口、测试组织与命令。