Documentation
¶
Overview ¶
Package kairosflux 是 KairosFlux 的可导入双模式引擎 API(顶层包,发布批次 阶段 A):embedded(纯进程内)与 server(同一套 API 的网络壳)共用同一个 Engine,时态内核五操作 + 审计 + Context/Proposal 访问口在两种模式下逐字节 同语义。跨仓调用方(ChronoBrew/ChronoBrew 的 sample 与 E2E 测试)用 `import kairosflux "github.com/ChronoBrew/KairosFlux"` 即可获得全部能力, 不需要 import 任何 internal 包。
三种构造形态:
- NewEmbedded:embedded 模式。DataDir 不存在会创建,不开网络监听, 适合"把 KairosFlux 当库嵌进另一个进程"(E2E 测试、样例、CI 装置)。
- Open:打开既有数据目录(embedded/server 由 Options.Port 决定,Port<=0 为 embedded)。DataDir 必须已存在——重启恢复(kill -9 后一致性验证、 长期数据复用)用这个形态。
- Serve:server 模式。DataDir 不存在会创建,构造即绑定监听(Port 必填), Addr() 返回监听地址,Close() 优雅关停。server 模式的网络壳与生产装配 Node(service/node.go)是同一个核心接线(v1 Router + v2 RouterV2 + 采集过滤钩子),差别是 Node 额外挂分片路由/网关准入/下游投递/周期指标 这些按配置门控的生产子系统(Engine 的 server 模式不挂,接口面以"真的 会换实现"为准,不造插件注册表)。
与 Node 的关系:Node 仍是生产服务器装配(cmd/kairosflux-server 的薄壳), 读进程级 config.G;Engine 是给"把 KairosFlux 当库/当嵌入式服务"的调用方 用的实例级 API,每个 Engine 持有自己的数据目录与监听配置,互不干扰。 Engine 只支持 standalone 模式(WAL 持久化);Raft 集群是另一个部署形态, 不在本 API 的范围内。
已知边界(诚实标注,不是遗漏):
- Engine 的进程内写路径(PutVersioned)不经过协议层角色强制 (service/router_v2.go 的 handlePutVersioned 对 source 的 agent 身份 校验只覆盖 PUT_VERSIONED 线协议帧)——进程内调用方本身就是可信代码, 角色强制是"线协议上防越权"的机制,不适用于同进程直调。server 模式下 经网络的写入仍受完整协议层强制。
- kairnet.Server 的连接数上限(config.G.MaxConn)仍读进程级全局配置 (kairnet/server.go 的 acceptLoop,改动会牵动既有行为,未在本批次触碰): Engine 的 server 模式沿用进程默认值 1000。
- Engine 不加载 contracts/*.schema.json(service 的 schema.LoadContractsDefault 是生产契约校验的前提;Engine 的采集过滤钩子对未注册前缀走默认放行 回退路径——见 service/ingesthook/filter.go 的 validate 回退),嵌入式 场景不需要仓库内契约文件即可运行。
- 存储引擎无独立 Close(flush/compaction 工作协程随进程退出);Engine.Close 负责排空并关闭 WAL 文件句柄与(server 模式下)网络壳,保证同进程内 Open 复用同一数据目录安全。
Index ¶
- Constants
- type ContextBundle
- type ContextRequest
- type Engine
- func (e *Engine) Addr() string
- func (e *Engine) BuildContext(req ContextRequest, contractsDir, redlinesPath string) (ContextBundle, error)
- func (e *Engine) Close() error
- func (e *Engine) GetAsOf(logical string, asOfNanos int64) (Version, bool, error)
- func (e *Engine) ListVersions(logical string) ([]Version, error)
- func (e *Engine) ListWrites(prefix string, tFromNanos, tToNanos int64, sourceFilter string) (ListWritesResult, error)
- func (e *Engine) PutVersioned(logical string, payload []byte, writeNanos int64, source string, ...) (uint64, error)
- func (e *Engine) ReplayFingerprint(prefix string, asOfNanos int64) (ReplayResult, error)
- func (e *Engine) SubmitProposal(p Proposal) (fingerprint string, seq uint64, err error)
- type ListWritesResult
- type Options
- type Proposal
- type ProposalKind
- type ReplayResult
- type Version
- type WriteEnvelope
Constants ¶
const ( ProposalFactor = aiplane.ProposalFactor ProposalHypothesis = aiplane.ProposalHypothesis ProposalExperiment = aiplane.ProposalExperiment ProposalRecommendation = aiplane.ProposalRecommendation ProposalReview = aiplane.ProposalReview )
ProposalKind 常量(值即 aiplane 的枚举值,跨仓调用方不需要 import internal 包即可引用)。
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ContextBundle ¶
type ContextBundle = aiplane.ContextBundle
ContextBundle 是 BuildContext 的输出(确定性上下文包,字段契约见 contracts/aiplane/context.schema.json)。
type ContextRequest ¶
type ContextRequest = aiplane.ContextRequest
ContextRequest 是 BuildContext 的唯一输入参数(as-of 语义)。
type Engine ¶
type Engine struct {
// contains filtered or unexported fields
}
Engine 是 KairosFlux 的双模式引擎:embedded 与 server 共用同一个时态内核 (PUT_VERSIONED/GET_AS_OF/LIST_VERSIONS/REPLAY_FINGERPRINT/LIST_WRITES), server 模式额外持有网络壳(kairnet.Server,v1+v2 双协议)。Engine 的方法 在两种模式下逐字节同语义——这是"server 模式 = 同一 API 的网络壳"的落点。
func NewEmbedded ¶
NewEmbedded 构造 embedded 模式引擎:DataDir 不存在则创建,不开网络监听。 失败以 error 返回(不 panic、不静默降级),调用方据此决定如何处理。
func Open ¶
Open 打开既有数据目录构造引擎(embedded/server 由 Options.Port 决定): DataDir 必须已存在,否则报错。这是"重启恢复"形态——kill -9 之后用同一 DataDir 重新 Open,WAL 重放 + SSTable 加载使数据回到崩溃前状态 (崩溃安全顺序见 internal/temporal 包文档:先落版本键、再落 :current 指针)。
func Serve ¶
Serve 构造 server 模式引擎:DataDir 不存在则创建,Port 必须 >0;构造完成 即已绑定监听(kairnet.Server.Start 同步绑定,绑定失败以 error 返回), Addr() 可读,调用方后续只需 Close() 优雅关停。网络壳与 Node 共用同一套 核心接线(v1 Router + v2 RouterV2 + 采集过滤钩子),Node 的分片路由/准入/ 投递/指标等按配置门控的生产子系统不挂载(见本文件顶部文档)。
func (*Engine) BuildContext ¶
func (e *Engine) BuildContext(req ContextRequest, contractsDir, redlinesPath string) (ContextBundle, error)
BuildContext 是 Context 访问口:为研究员 agent 组装确定性上下文包 (数据集契约 + 摘要、测过哪些因子、策略状态、风控红线 + 各自摘要), 同一请求(同一 as_of + 同一底层账本/契约/红线文件)两次调用逐字节相同。 contractsDir 通常是 contracts/ 目录,redlinesPath 通常是 riskredlines/redlines.json(委托 internal/aiplane.BuildContext)。
func (*Engine) Close ¶
Close 优雅关停引擎(幂等):server 模式先按 kairnet 三段式排空在途请求并 关闭监听,随后排空并关闭 WAL 文件句柄(service.KVServer.Close 负责)—— 此后同进程内用同一 DataDir 重新 Open 是安全的(WAL 文件唯一写者不变量 恢复)。存储引擎的 flush/compaction 工作协程随进程退出,无需也不存在单独 的关闭入口(见 storage.Engine 文档)。
func (*Engine) GetAsOf ¶
GetAsOf 返回 logical 在 asOfNanos 时刻可见的最新版本(as-of 语义:绝不 返回 asOfNanos 之后写入的版本)。该时刻无可见版本时 found=false。
func (*Engine) ListVersions ¶
ListVersions 返回 logical 的全部版本,按 seq 升序(从未写过返回空切片)。
func (*Engine) ListWrites ¶
func (e *Engine) ListWrites(prefix string, tFromNanos, tToNanos int64, sourceFilter string) (ListWritesResult, error)
ListWrites 是审计查询:扫描 prefix 下的全部版本键,按 [tFromNanos, tToNanos](<=0 表示对应方向无界)与 sourceFilter(""=不过滤)筛出每次 历史写入,返回信封列表(按 LogicalKey,Seq 升序,确定性输出)与按来源 聚合计数。temporal.Store 的幂等/崩溃安全保证使 HashOK 可作"数据在写入 之后是否发生过漂移"的逐条自检。
func (*Engine) PutVersioned ¶
func (e *Engine) PutVersioned(logical string, payload []byte, writeNanos int64, source string, schemaVersion uint32) (uint64, error)
PutVersioned 写入一个不可变版本,返回分配到的 seq。与网络协议 PUT_VERSIONED 同一语义(source/schemaVersion 落进操作元数据信封,供 LIST_WRITES 审计 按来源/契约版本过滤);进程内直调不经过协议层 agent 角色强制(见本文件 顶部"已知边界")。writeNanos 是写入时刻(unix 纳秒),调用方控制——E2E 测试用可控时间戳构造 as-of 定点语义;生产通常传 time.Now().UnixNano()。
func (*Engine) ReplayFingerprint ¶
func (e *Engine) ReplayFingerprint(prefix string, asOfNanos int64) (ReplayResult, error)
ReplayFingerprint 对 prefix(数据集)重放出每个逻辑键的最新状态并计算 确定性指纹(sha256 over (LogicalKey, Seq, Payload) 集合)。asOfNanos<=0 无时间上界并核对 :current 指针(Mismatches 是真实对账结果);asOfNanos>0 按此刻重放(Bounded=true,不做 :current 对账,见 ReplayResult.Bounded 的 文档——调用方必须区分"核对通过"与"没核对")。
func (*Engine) SubmitProposal ¶
SubmitProposal 是 agent 写提议的唯一入口(Proposal 访问口):字段校验 + 角色强制(只接受 Proposal 对象,写落 proposal: 键空间)委托 internal/aiplane.SubmitProposal,返回提议指纹与本次写入分配到的 seq。 与 jobctl.V2Store(网络瘦客户端)是同一个 ReadWriter 语义的不同实现—— 这里是进程内直连 TemporalStore。
type ListWritesResult ¶
type ListWritesResult = service.ListWritesResult
ListWritesResult 是 LIST_WRITES 的结果(Entries + BySource 聚合, 见 service.ListWritesResult)。
type Options ¶
type Options struct {
// DataDir 是数据目录(WAL 与 SSTable 的落盘处)。必填。
DataDir string
// Host 是 server 模式的监听主机(默认 "127.0.0.1")。embedded 模式不使用。
Host string
// Port 是 server 模式的监听端口。<=0 构造 embedded 模式(不开网络监听);
// Serve 要求 >0。
Port int
// WindowSafetyValveN 透传给 RouterV2 的 ack=window 安全阀 N
// (service/router_v2.go 的 windowSafetyValveN 文档);<=0 用生产默认值
// 1000。embedded 模式下无网络连接,不参与 ack 协商,本字段无意义。
WindowSafetyValveN uint32
}
Options 是 Engine 的唯一构造参数(NewEmbedded/Open/Serve 三形态共用)。
type ProposalKind ¶
type ProposalKind = aiplane.ProposalKind
ProposalKind 枚举 agent 能提交的提议种类(factor/hypothesis/ experiment/recommendation/review)。
type ReplayResult ¶
type ReplayResult = service.ReplayResult
ReplayResult 是 REPLAY_FINGERPRINT 的结果(KeyCount/Fingerprint/ Mismatches/Bounded,见 service.ReplayResult)。
type Version ¶
Version 是 PUT_VERSIONED 写入的一条不可变版本记录(LogicalKey/Seq/ WriteNanos/Source/SchemaVer/PersistedHash/Payload,见 internal/temporal.Version)。
type WriteEnvelope ¶
type WriteEnvelope = service.WriteEnvelope
WriteEnvelope 是 LIST_WRITES 返回的单条审计信封 (LogicalKey/Seq/WriteNanos/Source/SchemaVer/PersistedHash/Payload/ HashOK,见 service.WriteEnvelope)。
Directories
¶
| Path | Synopsis |
|---|---|
|
Package client 是 KairosFlux 的 Go 客户端 SDK。
|
Package client 是 KairosFlux 的 Go 客户端 SDK。 |
|
Package cluster 是分片集群的控制面:决定「一个 key 归属哪个分片、哪个物理节点」, 以及节点间的读择优与转发连接复用。
|
Package cluster 是分片集群的控制面:决定「一个 key 归属哪个分片、哪个物理节点」, 以及节点间的读择优与转发连接复用。 |
|
cmd
|
|
|
kairosflux-bench
command
命令 kairosflux-bench 是压测工具集:
|
命令 kairosflux-bench 是压测工具集: |
|
kairosflux-bench-grpc
command
|
|
|
kairosflux-cli
command
命令 kairosflux-cli 是 KairosFlux 的命令行客户端,基于 client SDK。
|
命令 kairosflux-cli 是 KairosFlux 的命令行客户端,基于 client SDK。 |
|
kairosflux-crashsim
command
命令 kairosflux-crashsim 是找茬面的崩溃模拟器(阶段 B):独立进程打开 embedded 引擎连续写入 N 条版本化记录,每 50 条用原子 rename 更新一次 progress 文件(记录"已拿到 fsync 确认的最大下标"),供外部观察者在任意 时刻 SIGKILL(或 rlimit 触发信号)后验证崩溃一致性:
|
命令 kairosflux-crashsim 是找茬面的崩溃模拟器(阶段 B):独立进程打开 embedded 引擎连续写入 N 条版本化记录,每 50 条用原子 rename 更新一次 progress 文件(记录"已拿到 fsync 确认的最大下标"),供外部观察者在任意 时刻 SIGKILL(或 rlimit 触发信号)后验证崩溃一致性: |
|
kairosflux-grpc-server
command
命令 kairosflux-grpc-server 是 Kair(kairnet TLV) 协议的基准测试/协议对照服务端, 不是生产摄入入口——生产摄入用 cmd/kairosflux-server(Kair)。
|
命令 kairosflux-grpc-server 是 Kair(kairnet TLV) 协议的基准测试/协议对照服务端, 不是生产摄入入口——生产摄入用 cmd/kairosflux-server(Kair)。 |
|
kairosflux-ingest
command
命令 ingest 是 M1 的 A1 层压测:进程内直接驱动存储引擎,证明存储引擎 在内存封顶下扛得住高频顺序写。
|
命令 ingest 是 M1 的 A1 层压测:进程内直接驱动存储引擎,证明存储引擎 在内存封顶下扛得住高频顺序写。 |
|
kairosflux-jobctl
command
kairosflux-jobctl 是 M3 声明式 Job 控制面的独立命令行入口(docs/方案- BanDB-时态内核与AI数据平面.md §M3)。
|
kairosflux-jobctl 是 M3 声明式 Job 控制面的独立命令行入口(docs/方案- BanDB-时态内核与AI数据平面.md §M3)。 |
|
kairosflux-sample-audit
command
命令 kairosflux-sample-audit 是"audit 审计导出"能力样例(发布批次阶段 A): embedded 引擎写入多来源数据后,用 LIST_WRITES 导出 append-only JSONL 审计文件(每行一条写入信封:logical_key/seq/write_ts/source/schema_ver/ payload_hash/payload_b64/hash_ok),末尾追加清单行(导出条数 + export_fingerprint——导出内容的确定性摘要,见 exportManifestRecord 的 文档:它是"本次导出文件"的完整性校验值,不是数据集状态指纹,二者不可 互相比较)。
|
命令 kairosflux-sample-audit 是"audit 审计导出"能力样例(发布批次阶段 A): embedded 引擎写入多来源数据后,用 LIST_WRITES 导出 append-only JSONL 审计文件(每行一条写入信封:logical_key/seq/write_ts/source/schema_ver/ payload_hash/payload_b64/hash_ok),末尾追加清单行(导出条数 + export_fingerprint——导出内容的确定性摘要,见 exportManifestRecord 的 文档:它是"本次导出文件"的完整性校验值,不是数据集状态指纹,二者不可 互相比较)。 |
|
kairosflux-sample-embedded
command
命令 kairosflux-sample-embedded 是"embedded 模式进程内全流程"能力样例 (发布批次阶段 A):把 KairosFlux 当库嵌进进程,跑通 合成数据 → PUT_VERSIONED 版本化写入 → GET_AS_OF 定点读取 → LIST_VERSIONS 版本清单 → REPLAY_FINGERPRINT 重放指纹(:current 对账)→ LIST_WRITES 审计 全链路,全部走 kairosflux.go(仓库根)的可导入 API,不开任何网络监听。
|
命令 kairosflux-sample-embedded 是"embedded 模式进程内全流程"能力样例 (发布批次阶段 A):把 KairosFlux 当库嵌进进程,跑通 合成数据 → PUT_VERSIONED 版本化写入 → GET_AS_OF 定点读取 → LIST_VERSIONS 版本清单 → REPLAY_FINGERPRINT 重放指纹(:current 对账)→ LIST_WRITES 审计 全链路,全部走 kairosflux.go(仓库根)的可导入 API,不开任何网络监听。 |
|
kairosflux-sample-server
command
命令 kairosflux-sample-server 是"server 模式 + Python 客户端"能力样例 (发布批次阶段 A):kairosflux.Serve 起一个真实监听端口,同一套 API 作为 网络壳对外服务——先由 Go 侧 v2 瘦客户端(kairnet/codec + negotiate + proto 拼帧,与 kairosflux-cli 同一模式)经真实线协议完成 PUT_VERSIONED → GET_AS_OF 往返,再由 Python 客户端(仓库 client/python/bandb_client.py:v1 直写直读 + v2 ack=window 批量写 + FLUSH 对账 + STAT 累计计数)演示跨语言访问。
|
命令 kairosflux-sample-server 是"server 模式 + Python 客户端"能力样例 (发布批次阶段 A):kairosflux.Serve 起一个真实监听端口,同一套 API 作为 网络壳对外服务——先由 Go 侧 v2 瘦客户端(kairnet/codec + negotiate + proto 拼帧,与 kairosflux-cli 同一模式)经真实线协议完成 PUT_VERSIONED → GET_AS_OF 往返,再由 Python 客户端(仓库 client/python/bandb_client.py:v1 直写直读 + v2 ack=window 批量写 + FLUSH 对账 + STAT 累计计数)演示跨语言访问。 |
|
kairosflux-server
command
与 server_pprof.go(//go:build pprof) 互斥:两者各自定义 main,缺少本约束会使 `go build -tags pprof` 因 main 重复声明而失败,pprof 构建不可用。
|
与 server_pprof.go(//go:build pprof) 互斥:两者各自定义 main,缺少本约束会使 `go build -tags pprof` 因 main 重复声明而失败,pprof 构建不可用。 |
|
Package config 定义 KairosFlux 的全局运行配置及其加载。
|
Package config 定义 KairosFlux 的全局运行配置及其加载。 |
|
internal
|
|
|
admission
Package admission 提供网关入口的「自适应并发限流 / 准入控制」。
|
Package admission 提供网关入口的「自适应并发限流 / 准入控制」。 |
|
aiplane
Package aiplane 实现 M4 AI 数据平面(docs/方案-BanDB-时态内核与AI数据平面.md §M4):让"Agent 只读真相、只写提议、引擎裁决"成为正式接口。
|
Package aiplane 实现 M4 AI 数据平面(docs/方案-BanDB-时态内核与AI数据平面.md §M4):让"Agent 只读真相、只写提议、引擎裁决"成为正式接口。 |
|
credit
Package credit 提供字节级信用池(令牌桶式背压):写入方 Acquire 占用字节信用, 持久化方 Release 归还信用;信用不足时 Acquire 阻塞,从而把未持久化数据的内存占用 限制在预算之内。
|
Package credit 提供字节级信用池(令牌桶式背压):写入方 Acquire 占用字节信用, 持久化方 Release 归还信用;信用不足时 Acquire 阻塞,从而把未持久化数据的内存占用 限制在预算之内。 |
|
identity
Package identity 是"写请求发起方角色"这一条规则的唯一真相来源:把 internal/aiplane.Role(agent 只能写 Proposal、引擎不受限)从 API 层的一道 闸门升级为协议层强制(M4 上报缺口,任务书:"把 M4 的 WriteAsAgent API 层 闸门升级为协议层强制")时,service/router_v2.go 需要在 handlePutVersioned 里按解出的 source 字段做角色校验——但 service 不能 import internal/aiplane: internal/aiplane/integration_test.go 与 internal/jobctl/v2store_integration_test.go 都是同包内部测试(package aiplane / package jobctl,不是 _test 后缀的外部 测试包)且都 import service 起真实服务端做端到端验证,若 service 反过来 import internal/aiplane,会在 `go test ./internal/aiplane/...` 编译测试 二进制时触发"import cycle not allowed in test"(已用最小复现验证过这个 编译期错误,不是理论风险)。
|
Package identity 是"写请求发起方角色"这一条规则的唯一真相来源:把 internal/aiplane.Role(agent 只能写 Proposal、引擎不受限)从 API 层的一道 闸门升级为协议层强制(M4 上报缺口,任务书:"把 M4 的 WriteAsAgent API 层 闸门升级为协议层强制")时,service/router_v2.go 需要在 handlePutVersioned 里按解出的 source 字段做角色校验——但 service 不能 import internal/aiplane: internal/aiplane/integration_test.go 与 internal/jobctl/v2store_integration_test.go 都是同包内部测试(package aiplane / package jobctl,不是 _test 后缀的外部 测试包)且都 import service 起真实服务端做端到端验证,若 service 反过来 import internal/aiplane,会在 `go test ./internal/aiplane/...` 编译测试 二进制时触发"import cycle not allowed in test"(已用最小复现验证过这个 编译期错误,不是理论风险)。 |
|
jobctl
Package jobctl 实现 M3 对象模型与声明式 Job 控制面(docs/方案-BanDB- 时态内核与AI数据平面.md §M3):daily 流水线从 shell 串联升级为"声明式 任务 + 本地 reconcile 循环"。
|
Package jobctl 实现 M3 对象模型与声明式 Job 控制面(docs/方案-BanDB- 时态内核与AI数据平面.md §M3):daily 流水线从 shell 串联升级为"声明式 任务 + 本地 reconcile 循环"。 |
|
kvgrpc
Package kvgrpc 是 KairosFlux 的 gRPC 传输实现,位于 internal/ 之下:它是内部传输, 不是对外契约。
|
Package kvgrpc 是 KairosFlux 的 gRPC 传输实现,位于 internal/ 之下:它是内部传输, 不是对外契约。 |
|
metrics
Package metrics 提供零依赖的进程内可观测性:一组原子计数器 + 仪表回调, 以「周期性 slog 快照」作为暴露出口——headless 边缘设备直接 tail 日志即可观测, 无需开端口、无需 Prometheus/Grafana 等外部基础设施。
|
Package metrics 提供零依赖的进程内可观测性:一组原子计数器 + 仪表回调, 以「周期性 slog 快照」作为暴露出口——headless 边缘设备直接 tail 日志即可观测, 无需开端口、无需 Prometheus/Grafana 等外部基础设施。 |
|
temporal
Package temporal 定义 KairosFlux 时态内核(M0)的版本化记录语义。
|
Package temporal 定义 KairosFlux 时态内核(M0)的版本化记录语义。 |
|
codec
Package codec 是 Kair 帧格式的编解码层:字节 ↔ Message 的转换,只认字节, 不认 msgID 该分派给谁、不认连接生命周期——这是重构 RFC (docs/rfc/bannet-重构.md)第一步迁移的目标包,从根包 kairnet 的 message.go/datapack.go 原样搬入,本步不改变任何字节布局或行为,只搬家。
|
Package codec 是 Kair 帧格式的编解码层:字节 ↔ Message 的转换,只认字节, 不认 msgID 该分派给谁、不认连接生命周期——这是重构 RFC (docs/rfc/bannet-重构.md)第一步迁移的目标包,从根包 kairnet 的 message.go/datapack.go 原样搬入,本步不改变任何字节布局或行为,只搬家。 |
|
dispatch
Package dispatch 是分发层:路由表(msgID → Handler)、worker 池调度、 panic 隔离——见 docs/rfc/bannet-重构.md C.2/C.5,msghandle.go 整体迁入。
|
Package dispatch 是分发层:路由表(msgID → Handler)、worker 池调度、 panic 隔离——见 docs/rfc/bannet-重构.md C.2/C.5,msghandle.go 整体迁入。 |
|
handler
Package handler 是 kairnet 面向业务代码的契约层:用户实际要实现/使用的 唯一一组类型(Handler/Request/HookAction/Conn),边界不因本次重构而 变复杂——这是重构 RFC(docs/rfc/bannet-重构.md)C.2 所说的"业务层"。
|
Package handler 是 kairnet 面向业务代码的契约层:用户实际要实现/使用的 唯一一组类型(Handler/Request/HookAction/Conn),边界不因本次重构而 变复杂——这是重构 RFC(docs/rfc/bannet-重构.md)C.2 所说的"业务层"。 |
|
lifecycle
Package lifecycle 是连接生命周期层:一个显式的状态机(Idle→Active→ Closing→Closed),取代此前散落在 connection.go 里的裸 context.Context+sync.Once 组合——见 docs/rfc/bannet-重构.md B.2/C.1: 重构前没有任何地方能回答"这个连接现在处于什么状态",状态只存在于一堆 副作用的组合里。
|
Package lifecycle 是连接生命周期层:一个显式的状态机(Idle→Active→ Closing→Closed),取代此前散落在 connection.go 里的裸 context.Context+sync.Once 组合——见 docs/rfc/bannet-重构.md B.2/C.1: 重构前没有任何地方能回答"这个连接现在处于什么状态",状态只存在于一堆 副作用的组合里。 |
|
negotiate
Package negotiate 实现 Kair v1/v2 的 HELLO 协商(docs/rfc/Kair-2.md §5/§5.1):v2 客户端连接后先发一个 v1 格式的 HELLO 帧,按是否在超时内 收到 v2 格式的响应判断对端版本,零破坏地与 v1 共存于同一条连接。
|
Package negotiate 实现 Kair v1/v2 的 HELLO 协商(docs/rfc/Kair-2.md §5/§5.1):v2 客户端连接后先发一个 v1 格式的 HELLO 帧,按是否在超时内 收到 v2 格式的响应判断对端版本,零破坏地与 v1 共存于同一条连接。 |
|
transport
Package transport 是传输层:只认原始字节的收发(Reader/Writer 循环、 连接注册表),不认 Kair 帧内容该分派给谁——见 docs/rfc/bannet-重构.md C.2/C.5,是重构第五步(最后一步)的迁移目标。
|
Package transport 是传输层:只认原始字节的收发(Reader/Writer 循环、 连接注册表),不认 Kair 帧内容该分派给谁——见 docs/rfc/bannet-重构.md C.2/C.5,是重构第五步(最后一步)的迁移目标。 |
|
Package predicate 提供一个最小的字段谓词,用于边缘查询的服务端下推: 对 JSON 值取出指定字段,与操作数按算子比较,只让命中的行通过—— 从而「只回传命中切片」而非整段原始流。
|
Package predicate 提供一个最小的字段谓词,用于边缘查询的服务端下推: 对 JSON 值取出指定字段,与操作数按算子比较,只让命中的行通过—— 从而「只回传命中切片」而非整段原始流。 |
|
Package proto 定义客户端/服务端的命名协议常量。
|
Package proto 定义客户端/服务端的命名协议常量。 |
|
本文件是 Raft 的类型、构造、生命周期与只读访问器。
|
本文件是 Raft 的类型、构造、生命周期与只读访问器。 |
|
delivery
Package delivery 是数仓写入前置缓冲的「下游投递层」骨架:把本地缓冲的数据 按批投递到一个或多个下游 sink(ClickHouse / Doris / 湖仓 / 文件)。
|
Package delivery 是数仓写入前置缓冲的「下游投递层」骨架:把本地缓冲的数据 按批投递到一个或多个下游 sink(ClickHouse / Doris / 湖仓 / 文件)。 |
|
delivery/governance
Package governance 是投递层的「治理」子包,借鉴 dubbo-go 的服务治理模型但落在数据面: 把多个下游 sink 当作一组要被治理的后端——熔断(breaker)、健康感知路由(router)、 健康探测(health)、退避重试(retry)。
|
Package governance 是投递层的「治理」子包,借鉴 dubbo-go 的服务治理模型但落在数据面: 把多个下游 sink 当作一组要被治理的后端——熔断(breaker)、健康感知路由(router)、 健康探测(health)、退避重试(retry)。 |
|
delivery/offset
Package offset 是投递层的「强一致 offset」子包:把每个 sink 的投递进度(游标) 持久化为一条 KV,从而在进程崩溃/重启后从已提交位置续投,而非从头重投。
|
Package offset 是投递层的「强一致 offset」子包:把每个 sink 的投递进度(游标) 持久化为一条 KV,从而在进程崩溃/重启后从已提交位置续投,而非从头重投。 |
|
ingesthook
Package ingesthook 提供一个挂在采集入口的真实 PreHandle 过滤钩子示例: 在数据落盘前完成「丢弃畸形帧 + 时间戳单调性校验 + schema 校验 + 字段脱敏」四件事, 把「可编程边缘采集缓冲网关」从挂载点变成有内容的演示。
|
Package ingesthook 提供一个挂在采集入口的真实 PreHandle 过滤钩子示例: 在数据落盘前完成「丢弃畸形帧 + 时间戳单调性校验 + schema 校验 + 字段脱敏」四件事, 把「可编程边缘采集缓冲网关」从挂载点变成有内容的演示。 |
|
ingesthook/schema
Package schema 提供「按数据类型注册校验规则」的最小注册表:每种业务数据类型 (如行情快照)注册一个 Validator,落盘前按 key 前缀分派到对应校验器。
|
Package schema 提供「按数据类型注册校验规则」的最小注册表:每种业务数据类型 (如行情快照)注册一个 Validator,落盘前按 key 前缀分派到对应校验器。 |
|
本文件是 SSTable 的类型与共用定义:磁盘布局常量、块索引、以及 SSTable 本身。
|
本文件是 SSTable 的类型与共用定义:磁盘布局常量、块索引、以及 SSTable 本身。 |