# 从 10 个 Go 开源项目，学会组件设计

给熟悉 Java／Scala／Python／C++ 的工程师。跳过语法入门，直接讨论抽象成本、状态所有权、并发协议和可测试性。

这门课的目标：当需求改变时，你知道应该修改哪个组件；当错误或取消发生时，你知道谁负责收尾；当性能下降时，你能提出可以验证的解释。

## 00 · 样本、排名与读法

“Top”在这里指 GitHub stars，**不代表代码质量排名**。查询使用 `language:Go is:public fork:false`，按 stars 降序；保留 API 原始响应和时间，不把今天的默认分支当作稳定发布版。

查询快照：**2026-09-05 14:28:13 UTC**（洛杉矶 07:28:13 PDT）。API 报告 incomplete_results=false。

| 原始排名 | 项目 | Stars | 本课处理 |
|---|---|---:|---|
| 1 | [avelino/awesome-go](https://github.com/avelino/awesome-go) | 183,219 | 资源目录，保留榜单 |
| 2 | [ollama/ollama](https://github.com/ollama/ollama) | 180,216 | 纳入源码研习 |
| 3 | [golang/go](https://github.com/golang/go) | 137,518 | 纳入源码研习 |
| 4 | [kubernetes/kubernetes](https://github.com/kubernetes/kubernetes) | 126,370 | 纳入源码研习 |
| 5 | [microsoft/TypeScript](https://github.com/microsoft/TypeScript) | 110,897 | 纳入源码研习 |
| 6 | [fatedier/frp](https://github.com/fatedier/frp) | 109,224 | 纳入源码研习 |
| 7 | [JuliusBrussee/caveman](https://github.com/JuliusBrussee/caveman) | 103,745 | Go 核心 BSL，保留榜单 |
| 8 | [infiniflow/ragflow](https://github.com/infiniflow/ragflow) | 90,086 | 纳入源码研习 |
| 9 | [gohugoio/hugo](https://github.com/gohugoio/hugo) | 89,704 | 纳入源码研习 |
| 10 | [gin-gonic/gin](https://github.com/gin-gonic/gin) | 89,172 | 纳入源码研习 |
| 11 | [syncthing/syncthing](https://github.com/syncthing/syncthing) | 88,317 | 纳入源码研习 |
| 12 | [junegunn/fzf](https://github.com/junegunn/fzf) | 82,826 | 纳入源码研习 |

[GitHub 查询入口](https://github.com/search?q=language%3AGo+is%3Apublic+fork%3Afalse&type=repositories&s=stars&o=desc)会随时间变化；[ranking.json](ranking.json) 保留本次原始响应。


原始榜单里的 `awesome-go` 是资源目录，不作为大型实现的教学样本；`caveman` 的固定版本采用分目录许可，Go engine 等目录属于 BSL-1.1，因此没有纳入这份 OSS 实现课程。两者都保留在原始表中，向下补入 Syncthing 和 fzf。筛选依据见 [awesome-go 仓库说明][awesome-readme]、[Caveman 分目录许可][caveman-license]。TypeScript 和 RAGFlow 是这次查询真实出现的结果，下面分析的是它们的 Go 实现。

**证据边界。** 每个源码链接固定到完整 commit 和行号。我们精读了指定接口、调用链和部分测试，未审计整个仓库，也没有构建这些上游工程。文中的“源码观察”是可定位的事实；“迁移建议”和原创示例是教学判断。下载过的文件不等于全部精读过。

不要从 K8s 的 `main` 开始逐层追踪数千个依赖。先读一个契约，再找一个调用方、一个实现和一个失败测试，最后才扩大范围。

推荐顺序：**Go 标准库 → Gin → Hugo → frp → RAGFlow → TypeScript → Syncthing → fzf → Ollama → K8s**。这个顺序逐步增加状态和时间维度，与 stars 顺序不同。

每章练习都可以按同一流程完成：先预测行为，读对应函数，关掉源码复述不变量，再写一个反例测试。能说出“这个模式在什么条件下会失效”，才算学会。

## 01 · 从你熟悉的语言迁移

优秀的 Go 组件通常很直接：具体类型保存状态，方法完成动作，小接口描述某个调用点需要的能力，普通函数把它们组装起来。抽象的价值要体现在修改范围和契约上。

| 你已有的经验 | 在 Go 中保留什么 | 需要调整什么 |
|---|---|---|
| Java 的 interface、DI、包封装 | 替换边界、构造时注入、隐藏实现 | 不必给每个 struct 配 `IService`；接口通常由使用者定义；组装常是一段普通代码 |
| Scala 的函数组合、trait、泛型 | 用函数表达策略，用类型表达约束 | embedding 没有 trait 的虚方法覆盖语义；不用给每个流程建立一套高阶抽象 |
| Python 的 duck typing、context manager | 关注对象能做什么，资源获取后就安排释放 | 接口满足关系在编译期检查；`defer` 到函数返回才执行，不是到代码块结束 |
| C++ 的值语义、RAII、move、const | 认真分析复制成本和资源所有者 | 没有通用 RAII 析构／move-only 保证；slice 的复制不复制底层数组；只读性常靠 API 契约 |
| coroutine／Future | 取消传播、并发上限、等待完成 | 一条 `go f()` 不建立父子生命周期；`cancel()` 也不等于 `join()` |

官方 [Code Review Comments](https://go.dev/wiki/CodeReviewComments#interfaces) 倾向在使用方定义接口、实现方返回具体类型，并建议在真实用途出现之后再抽象。这是默认选择：后面会看到标准库和插件系统返回接口的合理例外。

### 三个 Go 特有的边界问题

**1. 隐式实现不等于没有契约。** 方法签名吻合只证明“能调用”。是否并发安全、返回 slice 能不能改、取消后多久返回、错误怎样分类，都必须另行定义。

**2. `T` 和 `*T` 的 method set 不同。** 普通命名类型 `T` 的方法集中不含接收者为 `*T` 的方法；`*T` 包含两者。变量可取地址时的调用便利，不能替代接口赋值规则。有锁或可变状态的对象一般通过指针使用，而且使用后不能复制。参见 [Go 规范：Method sets](https://go.dev/ref/spec#Method_sets)。

**3. 接口中的 typed nil 不是 nil 接口。** `var p *MyError = nil; var err error = p` 可以让 `err != nil`。成功路径明确 `return nil`，不要返回装有 nil 指针的 error。依赖注入也可能遭遇同样的问题。参见 [Go FAQ](https://go.dev/doc/faq#nil_error)。

<details><summary>自测：接口越小，耦合就一定越小吗？</summary>

不一定。即使接口只有一个方法，参数如果是巨大的 `AppContext`、`*gorm.DB` 或框架 Context，调用方仍然依赖这些类型和生命周期。方法数量、参数类型、错误协议、状态所有权要一起看。

</details>


## 02 · Go 标准库：用能力组合，避免类型层级

### 需求：复制来自任意来源的字节

如果按 Java 类层次思考，容易先设计 `AbstractInputSource`，再让文件、网络和内存对象继承它。Go 的 `io.Reader` 只描述读取能力。文件和网络连接不必因为字节复制这个需求就被归入同一个继承树。

**源码路线：** [Reader 契约][go-io] → [copyBuffer 的选择顺序][go-copy] → [MultiReader 的组合][go-multi]。

`io.Copy` 先询问来源有没有 `WriterTo`，再询问目标有没有 `ReaderFrom`；都没有才运行通用读写循环。这是“最小基础契约 + 可选能力”的扩展方式。基础调用方仍然只需要 Reader／Writer。

`MultiReader` 持有多个 Reader，自己也提供读取能力。调用方无需知道它内部是单个来源还是组合来源。它还复制传入的接口 slice，避免调用者改变列表；这**没有深复制每个 Reader 对象**。

### 迁移到自己的代码

以下是原创用法片段，省略 imports；它让导出逻辑只关心字节能力：

```go
func Export(dst io.Writer, header string, body io.Reader) error {
    input := io.MultiReader(strings.NewReader(header), body)
    if _, err := io.Copy(dst, input); err != nil {
        return fmt.Errorf("export document: %w", err)
    }
    return nil
}
```

调用时可以传文件、`bytes.Buffer` 或压缩 writer。这个函数不擅自关闭它们：资源由调用方创建，关闭责任也留在调用方。若函数自行打开文件，则应在函数内安排关闭。

### 真正值得读的是边界条件

Reader 允许同一次调用返回 `n > 0` 和非 nil error。正确循环先消费这些字节，再处理错误；不能写成 `if err != nil { return err }` 后才处理 `buf[:n]`。`io.Copy` 还处理短写，且不会把正常 EOF 作为复制失败返回。[对应循环][go-copy]

**适用条件：** 调用方需要稳定的小能力，多种实现共享语义。**不宜照搬：** 为每个私有辅助函数都创建一个只有单实现、没有独立契约的接口；或用不断增长的可选 type assertion 代替清楚的能力分组。

<details><summary>练习：给 Export 加审计计数，应改所有 Reader 吗？</summary>

不需要。最简单是使用 `io.Copy` 的字节计数返回值；只有确实需要逐次观察 Read 时才增加 Reader wrapper。包装可能隐藏底层的 `WriterTo` 能力，所以还要评估是否改变性能路径。不要为了“Decorator 模式”忽略现成返回值。

</details>

## 03 · Gin：middleware 是控制流，不是注解魔法

### 需求：一个请求共享鉴权、追踪和业务处理

**源码路线：** [Context.Next／Abort][gin-flow] → [请求对象池][gin-pool] → [Abort 的上游测试][gin-tests]。

Gin 用 handler slice 和一个推进位置组织执行。`Next` 在当前 middleware 中执行后续 handler，返回之后继续当前函数的尾部代码。`Abort` 把位置移动到终止区间，**不会替你从当前函数返回**。因此拒绝请求之后通常还需要显式 `return`。

这和你熟悉的 around advice 接近，但执行顺序可以直接沿普通函数调用展开。不要把“没有调用 `Next`”简单理解为一定阻断 Gin 后续 handler；外层的推进循环仍可能继续。需要终止后续链时使用 Abort。

执行顺序示意：outer.before → Require → inner.before → handler → inner.after → outer.after。拒绝时跳过 inner 与 handler。

我们用标准库写一个可直接运行测试的版本。`http.HandlerFunc` 是命名函数类型带方法的适配器，不需要创建一个只有 `ServeHTTP` 的类。[标准库适配器源码][go-handler]


```go
// Package middleware demonstrates composition through http.Handler.
package middleware

import "net/http"

// Trace marks normal execution before and after the next handler.
// mark must be safe for concurrent requests. This is a teaching trace, not a
// panic recovery layer or an HTTP status/latency metrics implementation.
func Trace(name string, mark func(string), next http.Handler) http.Handler {
	return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
		mark(name + ":before")
		next.ServeHTTP(w, r)
		mark(name + ":after")
	})
}

// Require stops the handler chain if allowed rejects the request.
// It demonstrates control flow; callers supply the actual access policy.
func Require(allowed func(*http.Request) bool, next http.Handler) http.Handler {
	return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
		if !allowed(r) {
			http.Error(w, "forbidden", http.StatusForbidden)
			return
		}
		next.ServeHTTP(w, r)
	})
}
```


注意 `Require` 的 `return`，以及拒绝后外层 Trace 的 after 仍然执行。顺序错误会改变日志、鉴权和恢复行为，因此测试要检查整个 trace，而不只是最终 HTTP status。

### 对象复用使生命周期成为契约

Gin 在请求结束后将 Context 放回 pool。若 goroutine 保留原始 `*gin.Context`，可能在下一次请求复用时读到错误状态。`Copy` 对部分容器做复制，Request 指针仍是共享的，map 中的任意值也不因此深复制。[Context.Copy][gin-copy]

迁移建议：后台任务尽量只接收提取后的 ID、不可变 payload 和明确的生命周期。`c.Copy()` 不能延长 HTTP 连接的写入窗口，也不能让请求取消自动变成持久任务系统。

**不要照搬：** 未经测量就给业务对象加 `sync.Pool`。Pool 允许缓存项被回收，不能充当可靠资源池；对象归还后也不应继续持有可变引用。[sync.Pool 文档](https://pkg.go.dev/sync#Pool)

<details><summary>练习：为什么成功的日志链是 outer-before → inner-before → handler → inner-after → outer-after？</summary>

每一层在调用 next 前记录 before，next 返回后记录 after。因此进入次序和退出次序相反。如果把 after 放进 defer，panic 路径行为又会改变，需要明确记录语义。练习项目的 Trace 只记录正常返回路径，测试同时覆盖允许和拒绝请求。

</details>

## 04 · Hugo：配置、依赖与可选能力各司其职

### 需求：渲染内容，但不同渲染器能力不一样

**源码路线：** [converter 的小接口族][hugo-converter] → [Deps 的组装信息][hugo-deps] → [文件系统 wrapper][hugo-fs]。

Hugo 的基础 Converter 只要求转换。可选的 ParseRenderer 把解析和渲染分开，使支持该能力的实现可以先提取目录，再渲染。Provider 负责按文档上下文创建 converter；`NewProvider` 把一个创建函数适配为有名字、有 New 方法的对象。

这给出了两种独立的变化轴：**怎样创建对象**和**对象能做什么**。不要因为需要一个 factory，就连带引入依赖容器、生命周期框架和几十个注册表。

### 对组件输入做一次分类

| 输入类别 | 例子 | 建议位置 |
|---|---|---|
| 必需依赖 | 文档来源、输出存储、渲染器 | 构造函数的明确参数 |
| 固定配置 | 并发数、超时、输出格式 | 有名字的 Config／Options |
| 本次调用数据 | key、文档内容、请求上下文 | 方法参数 |
| 运行时状态 | 缓存、正在处理的任务 | 私有字段，由组件管理 |

Hugo 的 Deps 很大，是大型站点构建过程的组装信息载体。这是源码事实，并不意味着你应该给每个组件都传 `*App` 或 `*Deps`。若一个 formatter 只需要时钟和 writer，就只传这两个；这样它的真实依赖在签名中可见。

### Functional options 什么时候值得用

Gin 的构造入口确实接收 `OptionFunc`。[源码][gin-options] 但我们的 mini reconciler 只有两个必需依赖，因此直接 `New(source, dest)` 最清楚。公开库有很多独立、可扩展的可选设置时再考虑 `WithTimeout` 一类函数。

迁移建议：必需依赖保持显式；配置先填默认值，再应用覆盖，最后验证。不要在 option 执行过程中启动 goroutine 或建立半初始化资源。若 `0` 是合法业务值，不要同时把它定义为“没填写”；用 pointer、显式标记或不同构造入口区分。

**不要照搬：** 每个配置字段都变成 option；用一个巨大 dependency bag 隐藏十几个依赖；把配置对象当成可以随时无锁修改的共享 map。

<details><summary>练习：产品增加“生成目录”功能，要向所有 Converter 加一个方法吗？</summary>

先明确目录生成是否是所有渲染器的必需能力。如果只有部分实现支持，可以保留基础转换接口，额外定义解析／目录能力。调用方必须明确不支持时的行为，而不是让每个实现返回一个无意义的空结果。Hugo 的接口拆分就是可参照的例子。

</details>

## 05 · frp：组合、策略和插件注册的边界

### 需求：不同代理协议共享部分行为

**源码路线：** [Proxy 契约与工厂][frp-base] → [GeneralTCPProxy 的 embedding][frp-tcp] → [Plugin 的创建和生命周期][frp-plugin]。

frp 的 GeneralTCPProxy 嵌入 `*BaseProxy` 复用方法。插件部分使用名字到创建函数的映射，建立对象后通过 `Name／Handle／Close` 工作。这对应你熟悉的 Strategy 和 Factory，但关键是注册时机、配置类型和连接生命周期。

### embedding 不是虚方法继承

下面是原创语义片段，省略 package 和 main；注释给出在 main 中调用后的结果：

```go
type Base struct{}
func (Base) Name() string { return "base" }
func (b Base) Describe() string { return b.Name() }

type Child struct{ Base }
func (Child) Name() string { return "child" }

// Child{}.Name()     -> "child"
// Child{}.Describe() -> "base"
```

调用提升后的 Describe 时，方法中的接收者仍是 Base；它不会通过“实际子类”重新分派 Name。若共享算法需要调用变化行为，显式传函数／接口，或者把对象放到独立字段里委托。不要写依赖“子类覆盖”的模板方法后期待 Go 自动实现它。

### 插件的本质是可替换契约

frp 的注册表是包级 map，重复注册会 panic，具体实现常在 `init` 注册。这种设计依赖初始化期完成注册等使用前提；从这几行代码不能推出运行期可以随意并发增删插件。

迁移建议：只有两三个固定策略时，构造函数里的 switch 通常足够。确实要让插件组合随构建变化时，才考虑显式 registry。动态配置首先验证，再创建实例，最后进入运行期。不要在热路径依赖反射去弥补不清晰的类型边界。

<details><summary>练习：只有两个后端，是否需要设计 Registry、Factory、Provider 三层？</summary>

通常不需要。先在程序入口选出一个实现，再传给使用者。只有“后端发现”“创建多个实例”“不同实例生命周期”成为独立需求时，才拆出相应职责。结构应该解释当前变化点，而不是提前模拟一个生态系统。

</details>

## 06 · RAGFlow：接口的一致性比方法数量更重要

### 需求：同一功能使用不同对象存储

**源码路线：** [Storage 契约][rag-contract] → [MemoryStorage 的防御性复制][rag-memory] → [Factory 的选择逻辑][rag-factory]。

MemoryStorage 在写入时复制字节，读取时再次复制，并用 RWMutex 保护 map。你可以把它理解为对可变内存做所有权隔离：锁保护内部操作，复制防止调用结束后外部继续修改内部数组。

它是值得学习的内存适配器。与此同时，Storage 接口有较多能力，Factory 使用 singleton 和全局配置。这些选择并不能自动让所有使用者都低耦合。

### 读源码要检查承诺是否吻合

这次版本中，Storage.Get 的注释说缺失时可返回 nil；MemoryStorage.Get 对缺失返回包装过的 `ErrMemoryNotFound`。只看接口名，无法得出跨后端一致的 missing 语义。[契约][rag-contract]、[实现][rag-memory]

我们不在这里推断整个系统有缺陷；需要继续读适配调用方才能判断影响。但你写自己的 API 时应该主动消除这类不确定性：区分“缺失”“内容为空”“权限失败”“暂时不可用”，并让各实现服从同一组契约测试。

另一个具体例子：RAGFlow 的 RetrievalService 只有一个 Search 方法，参数中却有 `*gorm.DB`。[源码][rag-search] 它仍然暴露了 ORM 依赖。因此“只有一个方法”不能等同于“纯领域边界”。

### 本课程采用的契约

```go
type Reader interface {
    Read(ctx context.Context, key string) ([]byte, error)
}
```

签名之外还要写清楚：缺失返回可被 `errors.Is` 识别的 ErrNotFound；成功返回调用方独占的 bytes；支持并发调用；检查取消；空 bytes 是合法值。我们的 `Store` 在此基础上只增加写入能力。

**不要照搬：** 把一个通用 Storage 接口传遍系统；用 `(nil, nil)` 同时表示多种情况；以为用了 Mutex 就允许调用方任意修改返回的 slice。

<details><summary>练习：如果业务只读取对象，为什么参数不直接用完整 Store？</summary>

接收 Reader 可以让只读实现、缓存读取器和最小 fake 都满足依赖；也明确表达组件没有写权限。返回新接口并不是目的，限制所需能力才是目的。若确实只传具体类型且没有变化边界，直接传具体类型也可以。

</details>

## 07 · TypeScript：泛型用于共享结构，函数用于变化行为

### 需求：编译器需要稳定遍历次序，也需要改写 AST

本章读取的是 `microsoft/TypeScript/tsc/internal` 下的 Go 实现。**源码路线：** [OrderedMap][ts-map] → [公开 API 的测试][ts-tests] → [NodeVisitor][ts-visitor] → [有变化才复制][ts-clone]。

OrderedMap 用一个 map 负责查找，一个 key slice 保存插入顺序。`K comparable` 表达键必须可比较；`V any` 不限制值。这里泛型消除的是容器结构的重复，没有吞掉业务语义。

注意几个实际设计细节：零值在 Set 时懒初始化；更新已有 key 不把它再追加到列表；迭代用 `iter.Seq`／`Seq2`；`yield` 返回 false 时停止；删除还要维护 key slice，因此不能假定所有操作都是 O(1)。`noCopy` 提供给 vet 的复制检查线索，**不是运行时锁，也不使容器线程安全**。

### Iterator 与你熟悉的 generator 有什么关系

`iter.Seq[T]` 可以理解为把 yield 函数交给遍历逻辑，调用方的 break 会反馈为 false。它本身不表示启动了新的 goroutine。这个 OrderedMap 还特意用动态长度循环，让遍历能看到迭代期间追加的新项目；这是具体容器的契约，不应泛化为所有 Go iterator 的行为。

### Visitor 不需要几十个子类

NodeVisitor 中保存 Visit 函数和 hook 函数。VisitSlice 先扫描，发现节点被替换／删除等变化时才复制已访问的前缀，未改变时返回原 slice。函数处理变化策略，结构体保存遍历依赖，避免仅为覆盖一个行为建立类层次。

这让你能提出清楚的不变量：未变的输入保持结构共享；新结果不能意外修改旧视图。它不是深度 immutable 的自动保证，回调是否原地修改 Node 仍需契约约束。

### 了解性能技巧，但晚一点使用

同一仓库还有 [Arena[T]][ts-arena]，用批量分配减少细碎分配。它仍使用 Go 管理的数组，旧块是否能回收取决于引用；不能类比为“调用 arena.free 就手动释放全部对象”。先从 benchmark／profile 证明分配是瓶颈，再评估对象存活期和别名风险。

<details><summary>练习：为什么不把业务 Service 统一写成 Service[T, R, E]？</summary>

容器的算法和不变量在不同 T 上通常稳定；业务服务的权限、事务、错误和生命周期却可能完全不同。如果类型参数越来越多、内部到处 type switch，说明你正在把不相关行为压进一个模板。泛型应该减少结构重复，而不是掩盖差异。

</details>

## 08 · Syncthing：把状态修改变成受控协议

### 需求：运行中的多个组件一起响应配置变化

**源码路线：** [Modify 与 Serve][sync-config] → [验证、发布、通知][sync-commit] → [Committer 的承诺][sync-contract]。

Syncthing 的调用方提交一个修改函数。服务按顺序处理修改：取得当前配置副本，让函数修改副本，准备并验证，发布新配置，通知订阅者，并等待本轮处理。配置读取还使用 mutex；这不是“有一个事件循环就永远不需要锁”。

这个模式可以迁移到热更新 routing table、服务发现缓存和小型管理控制面：**修改入口集中；读取结果有约定；修改顺序可解释。** 比允许所有模块拿到共享 `map` 后各自加锁更容易检查。

### 三个不能混淆的状态

1. 新配置通过验证、写入内存。
2. 所有组件完成配置回调；有组件可能要求重启。
3. 配置持久化成功。

源码的 CommitConfiguration 返回 false 会标记需要重启，不会神奇回滚已经响应的所有组件。Save 另有路径。因此不要把“收到 Waiter”或“回调结束”描述成分布式原子提交。

### 迁移时的工程约束

修改函数保持短小、同步，不在里面做任意网络调用，不把传入的配置指针存到外部。发布给订阅者的快照按只读使用；复制 struct 不等于所有嵌套引用都自动隔离。队列满时应该有明确反馈，停止接收更新和等待回调退出也需要完整协议。

**适用条件：** 多个组件共享一份需要有序更新的状态。**不宜照搬：** 一个只有两个字段的局部对象，为了“Actor 模式”专门开启 goroutine 和消息队列。

<details><summary>练习：热配置验证通过，但一个组件不能在线应用，应该返回成功吗？</summary>

先区分“配置被接受”和“所有运行行为已经生效”。可以接受配置并返回 requiresRestart，也可以在前置验证阶段拒绝；选择取决于产品契约。不能悄悄把部分生效说成全部生效，更不能只靠一个 bool 同时表达所有状态。

</details>

## 09 · fzf：交互事件可以合并，业务命令未必可以

### 需求：用户快速输入，不值得算完每个过时查询

**源码路线：** [EventBox][fzf-box] → [Matcher.Loop][fzf-loop] → [分片 worker][fzf-workers] → [取消与等待][fzf-cancel]。

EventBox 用 `map[EventType]any` 存放事件；同一类型的新 Set 会覆盖旧值。这种按类型合并很适合 UI 最新状态。它不是保证每一条事件都投递的 FIFO，也不承诺不同类型间的全局时间顺序。

Matcher 将搜索拆成有限个 worker，按 CPU／配置和工作量限定并发；每个 worker 使用自己的临时工作区。取消旧扫描时，还需要等待相关计算结束，才能安全修改它们可能正在读取的数据。

### 从机制反推产品语义

| 事件 | 可以丢掉中间状态吗？ | 处理策略 |
|---|---|---|
| 搜索框从 `g` 到 `go` 到 `golang` | 通常可以 | 合并／取消过时查询，最终展示最新结果 |
| 期望副本数从 2 到 3 到 4 | 控制器通常可以重新读取最新目标 | key 去重、重新协调当前状态 |
| “给账户加 10”连续三次 | 不可以 | 每个命令有身份，按事务／幂等键处理 |
| 进度从 31% 到 32% 到 33% | 通常可以 | 保留最新进度 |

这才是队列选型的前提。你已经懂 coroutine，下一步是明确**哪些工作可以被取消、合并、重试或重复执行**。

### 一个值得检查的锁边界

EventBox.Wait 在持锁时调用 callback。这个 callback 若反过来调用同一 EventBox 的 Set，就可能尝试重复获取同一把锁。源码中的使用方式约束了它；你写公共 callback API 时，应明确“是否在锁内回调”，或取出副本后再解锁调用。[Wait／Set][fzf-box]

<details><summary>练习：把所有搜索请求放进一个无限队列，再启动更多 goroutine，能解决响应延迟吗？</summary>

未必。你可能把 CPU 花在用户已经不关心的查询上，排队时间继续增长。先定义最新结果语义，丢弃过时工作，再设并发和缓存边界。并发数不是吞吐量或延迟的无限旋钮。

</details>

## 10 · Ollama：接纳请求、占用资源、结束工作要分开

### 需求：模型加载昂贵，多个请求共享有限资源

**源码路线：** [Scheduler 字段与初始化][ollama-sched] → [请求接纳][ollama-admit] → [请求完成／模型到期][ollama-complete] → [函数注入测试][ollama-test]。

Scheduler 使用容量受配置限制的请求 channel，满了会返回 ErrMaxQueue；已加载且可复用的 runner 有快捷路径。加载资源、并行使用已加载资源、请求完成和到期回收，是不同动作。相应状态既有 channel 事件，也有 mutex 保护，不能把整个对象简单叫成无锁 actor。

`loadFn／newServerFn／getGpuFn` 等函数字段提供了替换边界。测试把真实启动依赖替换成函数，就能制造加载失败。你不必为一个调用点创造 `AbstractModelLoaderFactoryProvider`。

### 缓冲队列是接纳策略的一部分

有界 queue 必须回答满了怎么办：阻塞并支持取消、立即拒绝，还是按明确规则覆盖。只写 `make(chan Job, 100)` 没有完成设计。超时预算也应包括排队时间，而不是出队后重新给完整超时。

源码为单次请求的 success／error 通道设置了容量 1，能让一次结果发送不依赖接收者恰好同时就绪。但容量 1 并不能普遍解决多次发送、永不停止的生产者和共享对象所有权问题。

### 取消、退出和资源复用是三件事

此版本的 Scheduler.Run 启动 goroutine 后立即返回，没有从这个方法提供 join 保证。请求对象保留自己的 context 以跨队列携带该请求生命周期；这是具有明确含义的异步工作项，不能据此把 Context 当作普通 Service 的永久字段。[Run 与接纳路径][ollama-admit]、[请求使用 runner][ollama-request]

迁移建议：先明确构造、运行和停止 API。简单组件优先让 `Run(ctx) error` 阻塞到内部工作全部退出；如果必须异步启动，另给 Wait／Done。取消是合作信号，不能强杀一个不检查 ctx 的函数。[context 文档](https://pkg.go.dev/context)

<details><summary>练习：一次模型请求超时，是否应该立刻释放该模型的全部资源？</summary>

不能直接推出。其他请求可能仍在使用相同 runner；单个请求的寿命、共享模型的引用状态和闲置回收策略必须分开。先终止该请求的工作，确认引用关系，再按资源策略决定是否卸载。

</details>

## 11 · Kubernetes：重新协调当前状态，而不是重放每个事件

### 需求：事件会重复，状态会变化，写入会失败

**源码路线：** [DeploymentController 的依赖][k8s-controller] → [processNextWorkItem][k8s-worker] → [syncDeployment][k8s-sync] → [workqueue 状态机][k8s-queue]。

DeploymentController 把 key 放进队列。worker 取出 key，再从 lister 读取当前 Deployment；对象已经删除时可能正常结束，修改前先 DeepCopy，避免改坏共享 informer cache。这是典型的 **reconcile loop**：对照当前期望与观察状态，执行必要动作。

事件相当于“这个对象值得再检查一次”，不是必须逐条重放的交易日志。缓存也可能落后于服务端，所以设计仍需处理冲突、重新读取和重试；不能因为当前从 cache 没看到变化，就推断远端没有变化。

### 五行调用顺序，背后是两个不同协议

原创伪代码，突出职责而非复刻上游实现：

```go
key := queue.Get()
defer queue.Done(key)       // 结束本轮 processing 状态
err := reconcile(ctx, key)
if err == nil { queue.Forget(key) } // 清除重试历史
if err != nil { queue.AddRateLimited(key) }
```

真实实现还要处理 shutdown、错误分类和重试上限。DeploymentController 成功／特定情形调用 Forget；可重试错误受次数上限控制，达到上限也会结束本轮重试。[上游 worker][k8s-worker]

**Done 不等于 Forget。** Done 改变“谁正在处理”及是否需要再入队；Forget 清除 rate limiter 的历史，不会替你结束 processing，也不是删除所有待处理工作。[rate limiter wrapper][k8s-rate]

队列练习：依次执行 Add(A)、Get(A)、Add(A)、Done(A)，观察 dirty、processing 和可获取队列。互动版可以逐步操作。

### 队列为什么不是一个普通 channel

它至少管理三个状态集合：

| 状态 | 意义 |
|---|---|
| queue | 可以被 worker 获取的 key 顺序 |
| dirty | 需要处理的 key；可能已在处理但又收到新变化 |
| processing | 已取走、尚未 Done 的 key |

正常操作中的关键不变量是：queue 中的 key 属于 dirty，且不属于 processing。Get 把 key 加入 processing，并清除本轮 dirty；若处理期间再次 Add，相同 key 被标 dirty，但不会同时分配给另一个 worker；Done 发现它仍 dirty，就重新入队。[Add／Get／Done][k8s-queue]

同一个队列内的 key 协调不等于整个分布式系统 exactly-once。跨进程／重启／外部写入还需要幂等操作、版本检查或事务。尤其不要把“加一”这样的不可幂等副作用原样放进可重试 reconcile。

### 生命周期也是组件契约

DeploymentController.Run 等待 cache 首次同步，启动 worker；退出时关闭队列并等待 worker。这里的 shutdown 和等待是明确可定位的代码。[Run][k8s-lifetime] 这比只给 goroutine 一个 ctx 更完整。

测试也在验证这个协议：例如 [TestAddWhileProcessing][k8s-test] 主动让处理中的项目再次加入队列。学它的行为场景，比背“Controller／Observer／Factory”这些名字更有用。

<details><summary>练习：Add(A) → Get(A) → Add(A) → Add(A) → Done(A)，队列中有几个 A？</summary>

一个。第一次 Get 已清除本轮 dirty。处理中第一次 Add 把它重新标 dirty，第二次被合并；Done 才重新放回队列。再 Get 并 Done、且期间没有新 Add 后，三个集合里都不再有 A。若只 Forget 不 Done，processing 状态并没有完成。

</details>

## 12 · 把这些观察变成自己的组件边界

### 从一个具体需求开始

假设需求是：“读取源端文档，把目标端更新到相同内容；重复运行不做多余写入；可以取消，同时最多执行 N 个不同 key。”

先写顺序版本：读源 → 读目标 → 比较 → 必要时写。不要先创建 `domain／application／infra／factory／manager` 一整套空目录。第二个真实后端出现，或需要独立验证错误路径时，再把存储依赖提取为使用方的 Reader／Store。并发调度成为独立需求时，再移入 batch。

编译依赖：cmd → batch／reconcile／memstore；memstore → reconcile 的契约。运行调用：batch → work 函数 → Reconciler → 存储适配器。

### 用变化原因决定 package

| Package | 它知道什么 | 它拥有的状态／职责 | 改动它的理由 |
|---|---|---|---|
| `reconcile` | 读取／写入契约、期望与当前内容 | 一次协调的规则；不拥有后台任务 | 比较规则、缺失语义改变 |
| `memstore` | 如何在内存保存 bytes | map、mutex、复制边界 | 存储机制改变 |
| `batch` | key 和一个工作函数 | worker 数、队列关闭、取消与等待 | 并发调度策略改变 |
| `cmd/syncdemo` | 哪些具体实现组装在一起 | 进程级 context、运行顺序 | 启动配置／部署环境改变 |
| `middleware` | HTTP 请求和 next | 请求链的外围行为 | HTTP 接入需求改变 |

我们把 ErrNotFound 放在使用方定义的存储契约附近，所以 memstore 依赖 reconcile 的契约；reconcile 不 import memstore。这是一个小工程里的具体依赖方向。如果错误／值类型将被多个独立用例共享，才评估提取一个有业务含义的小 package，不必预先建立通用 `types` 大包。

### 比“Clean Architecture 目录”更有效的四个检查

**修改局部性。** 换成 S3 适配器时，应该主要新增适配器并修改 main 的组装，协调规则和 batch 无需变化。

**输入诚实。** 不把 DB、logger、config 和几十个服务全部藏进 ctx.Value 或 `*Application`；签名应让依赖可见。

**封装有内容。** 字段私有是为了维护不变量；如果一层 getter／setter 完整暴露同一个 mutable map，只是增加输入成本，没有保护状态。

**读取路径短。** 一次普通业务调用最好能在少量相邻函数中读完。四层都只有 `return next.Do(...)` 时，说明分层没有解释变化点。

`internal` 是 Go 工具执行的导入边界；`cmd` 是组织多个程序入口的常见方式。`pkg` 不是编译器赋予的“公开 API”标记，也没有要求所有工程采用同一套目录。参见 [官方模块布局](https://go.dev/doc/modules/layout)。

<details><summary>练习：现在加 PostgreSQL，实现了 Read／Write，协调器就自动正确了吗？</summary>

还没有。你要检查缺失错误是否转换、Read 的数据所有权、Write 是否整体替换、事务和并发语义、ctx 是否传到底层、同 key 的竞争，以及测试是否覆盖这些契约。能赋值给接口是编译期条件，不是行为一致性的证明。

</details>

## 13 · 可运行实验：一个小型 reconciler

[下载完整 Go 实验代码](go-oss-lab.zip)，解压得到 `go-oss-lab/`，只有标准库依赖。要求 Go 1.25+；本次实际验证环境为 Go 1.26.2、macOS arm64。当前上游默认分支可能需要更新的工具链，实验不依赖它们的构建环境。

```sh
unzip go-oss-lab.zip
cd go-oss-lab
go run ./cmd/syncdemo
go test ./...
go test -race ./...
go vet ./...
```

预期 demo 输出：

```text
pass 1: changed=2
pass 2: changed=0
after source update: changed=1
```

第一次把两个文档写入目标。第二次重新读取并比较，没有变化就不写。源端改变一个文档后，下一轮只有一个写入。输入里重复的 `guide` 在同一轮 batch 内被去重。

### 13.1 协调规则不需要知道内存存储


```go
// Package reconcile converges destination bytes toward a source snapshot.
// It owns the use case and its dependency contracts, not storage mechanisms.
package reconcile

import (
	"bytes"
	"context"
	"errors"
	"fmt"
)

// ErrNotFound means the requested key has no value. Empty bytes are a value.
// Adapters must translate their backend's missing-object error to this error.
var ErrNotFound = errors.New("object not found")

// Reader returns caller-owned bytes or an error wrapping ErrNotFound.
// Implementations must honor cancellation and support concurrent calls.
type Reader interface {
	Read(ctx context.Context, key string) ([]byte, error)
}

// Store replaces a whole value. Write must not retain the caller's byte slice.
// Callers must not mutate the input while Write is executing.
type Store interface {
	Reader
	Write(ctx context.Context, key string, value []byte) error
}

// Result describes one completed reconciliation, not a global system state.
type Result struct {
	Changed bool
}

// Reconciler reads current state on each call. It has no background goroutines.
// Different keys may be reconciled concurrently. The caller must serialize
// calls for the same key if it needs to prevent concurrent duplicate writes.
type Reconciler struct {
	source Reader
	dest   Store
}

// New wires dependencies without acquiring resources. Dependencies must be
// non-nil, including the concrete values stored in their interfaces.
func New(source Reader, dest Store) *Reconciler {
	return &Reconciler{source: source, dest: dest}
}

// Reconcile copies the current source value only when the destination differs.
// A missing source is an error; this use case never deletes destination data.
// It provides convergence on repeated calls, not a distributed transaction.
func (r *Reconciler) Reconcile(ctx context.Context, key string) (Result, error) {
	if key == "" {
		return Result{}, errors.New("key is required")
	}
	if err := ctx.Err(); err != nil {
		return Result{}, err
	}
	want, err := r.source.Read(ctx, key)
	if err != nil {
		return Result{}, fmt.Errorf("read source %q: %w", key, err)
	}
	have, err := r.dest.Read(ctx, key)
	if err != nil && !errors.Is(err, ErrNotFound) {
		return Result{}, fmt.Errorf("read destination %q: %w", key, err)
	}
	if err == nil && bytes.Equal(want, have) {
		return Result{}, nil
	}
	if err := r.dest.Write(ctx, key, want); err != nil {
		return Result{}, fmt.Errorf("write destination %q: %w", key, err)
	}
	return Result{Changed: true}, nil
}
```


阅读时抓住三处：目标读取失败不能统统按“缺失”处理；空内容与不存在不同；比较一致才返回 Changed=false。每次调用都重新读取源端，因此源端变动会在后续调用中继续收敛。

<details><summary>13.2 内存适配器：Mutex 与 Clone 解决两个不同问题</summary>


```go
// Package memstore implements an in-memory adapter for the reconcile contracts.
package memstore

import (
	"bytes"
	"context"
	"fmt"
	"sync"

	"example.com/go-oss-lab/internal/reconcile"
)

// Store is safe for concurrent use. Its zero value is ready to use.
// A Store must not be copied after first use.
type Store struct {
	mu     sync.RWMutex
	values map[string][]byte
}

var _ reconcile.Store = (*Store)(nil)

// Read returns an independent copy. Cancellation is checked after acquiring
// the lock; the mutex acquisition itself cannot be interrupted by ctx.
func (s *Store) Read(ctx context.Context, key string) ([]byte, error) {
	s.mu.RLock()
	defer s.mu.RUnlock()
	if err := ctx.Err(); err != nil {
		return nil, err
	}
	value, ok := s.values[key]
	if !ok {
		return nil, fmt.Errorf("read %q: %w", key, reconcile.ErrNotFound)
	}
	return bytes.Clone(value), nil
}

// Write atomically replaces one value within this process.
// This in-memory store does not persist data across process restarts.
func (s *Store) Write(ctx context.Context, key string, value []byte) error {
	s.mu.Lock()
	defer s.mu.Unlock()
	if err := ctx.Err(); err != nil {
		return err
	}
	if s.values == nil {
		s.values = make(map[string][]byte)
	}
	s.values[key] = bytes.Clone(value)
	return nil
}
```


锁避免内部 map 并发访问失序，Clone 避免锁释放后外部继续写内部数组。两个问题不能互相替代。Read 检查 ctx，但 mutex 等锁期间本身不可取消；这里临界区短且不做网络 I/O。

</details>

<details><summary>13.3 batch：限制并发，收到错误后取消，并等待所有 worker</summary>


```go
// Package batch owns bounded, finite concurrent work and its lifetime.
package batch

import (
	"context"
	"errors"
	"fmt"
	"sync"
)

// Run processes each distinct key at most once in this invocation.
// The first cancellation cause stops dispatch and is returned after all
// workers exit. Work already started may complete. Results are not rolled back.
// fn must honor ctx, be safe for concurrent calls, and must not panic.
// Run does not retry, serialize keys across invocations, or persist a queue.
func Run(ctx context.Context, keys []string, workers int, fn func(context.Context, string) error) error {
	if workers < 1 {
		return errors.New("workers must be positive")
	}
	if fn == nil {
		return errors.New("work function is required")
	}
	ctx, cancel := context.WithCancelCause(ctx)
	defer cancel(nil)
	jobs := make(chan string)
	var wg sync.WaitGroup
	for range min(workers, len(keys)) {
		wg.Go(func() {
			for {
				select {
				case <-ctx.Done():
					return
				case key, ok := <-jobs:
					if !ok || ctx.Err() != nil {
						return
					}
					if err := fn(ctx, key); err != nil {
						cancel(fmt.Errorf("process %q: %w", key, err))
						return
					}
				}
			}
		})
	}
	seen := make(map[string]struct{}, len(keys))
dispatch:
	for _, key := range keys {
		if _, exists := seen[key]; exists {
			continue
		}
		seen[key] = struct{}{}
		select {
		case <-ctx.Done():
			break dispatch
		case jobs <- key:
		}
	}
	close(jobs) // This function is the only sender and therefore owns closure.
	wg.Wait()   // Cancellation requests exit; Wait observes that exit completed.
	return context.Cause(ctx)
}
```


无缓冲 jobs 提供 dispatch 处的背压，固定 worker 数限制同时执行的函数数。唯一发送方关闭 jobs；`context.WithCancelCause` 保存先发生的取消原因；Wait 把运行期间的所有 worker 收回来。

</details>

### 这个实验明确承诺什么

| 性质 | 是否提供 |
|---|---|
| 稳定输入下重复协调不重复写 | 提供，测试验证 |
| 进程内 map 安全、读写 bytes 不共享内部数组 | 提供，测试及 race 检查 |
| 单次 batch 中同 key 去重、并发上限 | 提供，确定性测试验证 |
| 出错／取消后等待已启动 worker 退出 | 提供，前提是工作函数合作响应 ctx |
| 每个调用都能被强制超时终止 | 不提供；Go 不强杀任意函数 |
| 多次并发 Run 之间的同 key 串行 | 不提供；需要共享协调机制 |
| 动态到达事件、延迟重试和持久队列 | 不提供；这是有限输入 batch |
| 跨存储原子性、exactly-once、自动删除 | 不提供；接口和用例均未作此承诺 |

`Reconcile` 的读－比较－写不是跨后端事务。两个并发调用处理同一 key，可能都判断需要写；源端在一次调用中变化，也可能让目标短暂落后。真实产品按需求增加 key 串行、版本／CAS 检查或事务，并保证后续有再次协调的触发。

“第一次失败取消其他 worker”也不撤销已经成功的写入。错误后可以重跑，因为这个用例的写入是替换为期望值；若改为发邮件、扣款或递增计数，必须重新设计副作用协议。

远端 Write 还可能已经提交，只是成功回执丢失。此时返回 error 或 Changed=false 都不能证明“没有副作用”；需要重新读取、幂等重试或使用操作身份核对。实验中的 Changed 描述本次确认完成的结果，不是外部世界的事务证明。

## 14 · Clean code：写出能够审查的契约

### 错误既是诊断信息，也是 API 的一部分

中间层补充动作和 key：`fmt.Errorf("read source %q: %w", key, err)`。上层按 `errors.Is`／`errors.As` 决策，别比较错误字符串。使用 `%w` 等于允许调用方观察底层错误身份；是否公开某个驱动错误，要在边界处决定。[Go 官方错误设计说明](https://go.dev/blog/go1.13-errors)

通常在最了解操作结果的边界记录错误，中间层返回并补充上下文，避免每层重复打一条相同堆栈。重试循环可以记录重试事件，但应区别于最终失败。用户输入、I/O 失败和资源不足用返回错误处理；panic 不应替代正常业务分支。

只有当调用方能做出不同决定时，才值得新增 sentinel／typed error。多个操作部分成功时，应返回足够结果信息；随便一个 bool 往往不能说明已发生什么。

### 零值可用，与构造函数不矛盾

memstore 的零值可用，因为它能在首次 Write 初始化 map。reconciler 需要外部依赖，因此要求构造时提供有效实现。不要强求每个类型零值都可执行，也别只为了写 `NewX` 而禁止一个自然可用的零值。

构造尽量只组装。确需打开文件、启动监听或建立连接时，失败必须可回滚，成功后的关闭责任必须明确。不要把“对象已创建”“已开始运行”“完全可服务”揉成一个无法判断的状态。

### 引用和生命周期速查

| 容易误判的写法 | 实际要检查什么 |
|---|---|
| `copy := originalStruct` | map／slice／pointer 是否仍引用同一份数据；锁和 once 是否被复制 |
| `snapshot := slices.Clone(items)` | 只复制元素；元素如果是 pointer／slice，内部对象仍可能共享 |
| `defer file.Close()` 写在长循环里 | 直到外层函数返回才关闭；需要时把单次处理提取为函数 |
| `go work(ctx)` 后立即返回 | 谁保留工作寿命、接收错误、等待退出 |
| 从 pool 取对象，Put 后继续用 | 另一个调用者可能已经获得相同可变对象 |
| `append(s, x)` 后认为原 slice 不变 | 容量足够时会复用底层数组；是否允许写入取决于所有权 |

`defer` 的参数在 defer 语句执行时求值；关闭／释放顺序是后进先出。写文件时，Flush／Close 也可能有需要返回的错误；不能因为用了 defer 就把成功视为理所当然。[Effective Go：Defer](https://go.dev/doc/effective_go#defer)

### 并发和性能检查

channel 用来传递工作／信号／所有权；mutex 用来保护共享状态不变量。哪种更好取决于状态归属和操作粒度，不能通过“全换成 channel”消除竞争。GC 管理可达内存，不能替你结束一个阻塞的 goroutine。可见性和 happens-before 关系仍需依据同步操作判断。[Go Memory Model](https://go.dev/ref/mem)

先让函数行为清楚，再测量。容器泛型、手写循环、缓存、pool、arena 都可以有价值，但必须能回答哪个 allocation 或哪段 CPU 时间值得优化。测试项目提供一个 4 KiB 拷贝读取 benchmark，便于观察所有权隔离的代价：

```sh
go test ./internal/memstore -run '^$' -bench BenchmarkReadCopy -benchmem
```

不要从一次结果推断其他机器或生产负载的性能。race detector 也只能发现本次执行路径中的竞争；它不是并发正确性的证明。

### 不要机械复制旧 Go 教程

Go 1.22 起，在相应语言版本下循环内新声明的迭代变量具有每轮独立语义；反复添加 `v := v` 不再是普遍必要修复。外部复用变量、共享指针对象和并发写 map 仍然可能竞争。检查 `go.mod` 和实际共享的数据，而不是背旧反例。[官方说明](https://go.dev/blog/loopvar-preview)

## 15 · 用测试描述组件，而不是描述实现

测试应让更换内部实现后仍能成立。本实验用外部测试包检验可观察行为，并用少量手写 fake 制造错误；不用生成几十个接口 mock 来验证“某函数调用了某函数”。

| 实验测试 | 它防止的回归 |
|---|---|
| `TestConvergesAndDoesNotRewrite` | 每轮无条件写入；把空内容误判为缺失 |
| `TestPreservesErrorsAndAvoidsUnsafeWrite` | 目标读失败被当作不存在；错误身份丢失 |
| `TestOwnsBytes` | 修改调用方输入／返回值，污染存储内部状态 |
| `TestBoundsConcurrencyAndDeduplicates` | 每 key 启动一个 goroutine；重复 key 同轮多次执行 |
| `TestCancelWaitsForAllWorkers` | Run 已返回但 worker 仍活着 |
| `TestFailureCancelsSiblingsAndKeepsCause` | 触发取消的业务错误被 sibling 的 context.Canceled 覆盖 |
| `TestTraceOrderAndRejection` | 拒绝请求后仍调用业务；middleware 顺序错乱 |

### 并发测试尽量使用可控事件

不要 `time.Sleep(100*time.Millisecond)` 之后“希望 worker 已经启动”。测试用 gate channel 控制何时允许继续，用 Go 1.25+ 的 `testing/synctest` 等待测试 bubble 中其他 goroutine 都阻塞，再作断言。它适用于受控的并发逻辑；真实网络、系统调用和外部依赖有其适用限制。[synctest 文档](https://pkg.go.dev/testing/synctest)

<details><summary>查看实际测试：取消后，返回前必须观察到全部 worker 退出</summary>


```go
func TestCancelWaitsForAllWorkers(t *testing.T) {
	synctest.Test(t, func(t *testing.T) {
		ctx, cancel := context.WithCancel(t.Context())
		defer cancel()
		var started, exited atomic.Int64
		done := make(chan error, 1)
		go func() {
			done <- batch.Run(ctx, []string{"a", "b", "c", "d"}, 3, func(ctx context.Context, _ string) error {
				started.Add(1)
				defer exited.Add(1)
				<-ctx.Done()
				return ctx.Err()
			})
		}()
		synctest.Wait()
		if started.Load() != 3 {
			t.Errorf("started=%d", started.Load())
		}
		cancel()
		if err := <-done; !errors.Is(err, context.Canceled) {
			t.Fatal(err)
		}
		if exited.Load() != started.Load() {
			t.Fatalf("returned with active workers: started=%d exited=%d", started.Load(), exited.Load())
		}
	})
}
```


</details>

上游测试可以教你选场景：Gin 检查执行顺序，TypeScript 检查容器外部行为，K8s 检查处理中重新入队。我们阅读过这些选定测试；本次执行的是课程自带 lab 的测试，不把它们混为“上游全部验证通过”。

### Code review 时的十个问题

1. 一个组件究竟负责什么变化？它的名字能否说明责任？
2. import 方向是否让业务规则依赖具体基础设施？
3. 接口是否来自真实调用需求？参数是否偷偷带入巨大依赖？
4. 谁创建／关闭资源？谁能修改返回的引用？
5. error 是否区分缺失、拒绝、取消和暂时失败？
6. 每条 goroutine 在什么条件下退出？谁等待它？
7. 队列和并发有没有上限？排队也计入超时了吗？
8. 重复执行是否安全？“幂等”是否覆盖外部副作用？
9. 测试覆盖了取消、部分成功、别名和失败，而不只是 happy path 吗？
10. 优化是否有测量依据？引入的复杂度能否由收益解释？

## 16 · 六次练习，把知识变成写代码的能力

每次建议 60–90 分钟。先做再看参考答案；练习在课程目录的独立 lab 中完成。

| 次数 | 阅读与动手 | 可验收结果 |
|---|---|---|
| 1 | 标准库、Gin；自己重写 Trace／Require | 两条执行顺序测试通过；能解释拒绝后的外层 after |
| 2 | Hugo、frp、RAGFlow；画依赖图并加一个只读来源 | 不修改协调器即可接入；明确缺失／空值协议 |
| 3 | TypeScript；写一个保持插入顺序的小泛型集合 | 覆盖新增、覆盖已有 key、删除、break、零值行为 |
| 4 | Syncthing、fzf；设计一个只保留最新值的 UI 更新队列 | 写出为何可合并；演示相同机制不适用于三次增量命令 |
| 5 | Ollama；给有限 worker 加“接纳上限／拒绝”实验 | 能区分活跃数、待处理数和拒绝数；取消后全部退出 |
| 6 | K8s；在实验中增加动态 key queue 和失败重试 | 测试处理中再 Add、重复合并、失败延迟、成功清历史、shutdown |

### 最终作业：增加一个文件存储适配器

从只有内存的 demo 演进，但不要一次做完整文件同步产品。约定 key 由程序内部生成，先禁止路径分隔符和路径穿越。实现 Read／Write，转换缺失错误，明确一次整体替换的语义；处理写入、关闭、替换失败以及临时文件清理。

这个阶段涉及真实持久化：临时文件后 rename 可以提供某些同文件系统下的替换性质，但**不自动等于跨平台、断电可恢复的持久事务**。若需要 crash durability，继续研究目标平台的 fsync／目录同步和恢复流程。本课程未实现该适配器，不会把练习建议当成已验证的持久化保证。

验收时至少证明：更换存储不改协调规则；缺失和空文件可区分；失败不会被误报为 Changed=true；source 不存在不会悄悄删除 destination；并发和取消行为与契约一致。

<details><summary>进阶练习的设计提示</summary>

动态 key queue 不能只复用有限 batch 的 seen 集合：处理期间再次到来的变化必须触发后续一轮。需要明确 dirty／processing 状态，并为同 key 的串行、延迟重试和 shutdown 分别写场景。不要在 workqueue.Len()>0 的检查和 Get 之间推断原子性。

文件后端如无法直接满足原契约，应修改或收紧契约，然后重新审查所有调用方；不能只让编译通过。一个好的组件边界允许你发现需求不一致，而不是把不一致藏起来。

</details>

完成后，应当能在两分钟内说明一个组件的职责、依赖、所有权、错误和退出协议，并在代码中指到对应位置。

## 17 · 源码索引与可复现证据

下面是按研习路径整理的固定版本入口。每个项目还可以沿课程中的函数行号链接继续读。`research/ranking.json` 保存原始查询；`repositories.json` 保存排名、commit、许可标识和筛选原因；`reading-map.json` 保存本课程的定位范围；`source-index.json` 保存下载文件的 SHA-256。

| 样本 | 主要研习内容 | 固定 commit | 源码入口 |
|---|---|---|---|
| ollama/ollama | 资源调度、背压、函数注入 | [`83ed7d9965b1`](https://github.com/ollama/ollama/commit/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8) | [定位函数](https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched.go#L60-L104) |
| golang/go | 能力接口、函数适配、组合 | [`c5941983810b`](https://github.com/golang/go/commit/c5941983810b68ba93c30f0ef22c91ad63fb3e5c) | [定位函数](https://github.com/golang/go/blob/c5941983810b68ba93c30f0ef22c91ad63fb3e5c/src/io/io.go#L55-L107) |
| kubernetes/kubernetes | 协调循环、key 队列、幂等与退出 | [`b2ec8b6fefac`](https://github.com/kubernetes/kubernetes/commit/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe) | [定位函数](https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/pkg/controller/deployment/deployment_controller.go#L481-L519) |
| microsoft/TypeScript | 泛型容器、Visitor、按需复制 | [`1f70213d4922`](https://github.com/microsoft/TypeScript/commit/1f70213d4922b434345f639b441681e470c7cfc1) | [定位函数](https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/collections/ordered_map.go#L15-L187) |
| fatedier/frp | embedding、策略、插件生命周期 | [`832df8dff66d`](https://github.com/fatedier/frp/commit/832df8dff66d1b0caa0d776c450418f982756df1) | [定位函数](https://github.com/fatedier/frp/blob/832df8dff66d1b0caa0d776c450418f982756df1/pkg/plugin/client/plugin.go#L36-L71) |
| infiniflow/ragflow | 存储契约、适配器、引用所有权 | [`0c28d59ea1d3`](https://github.com/infiniflow/ragflow/commit/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b) | [定位函数](https://github.com/infiniflow/ragflow/blob/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b/internal/storage/memory.go#L32-L117) |
| gohugoio/hugo | 依赖组装、Provider、可选能力 | [`9c2527f8558e`](https://github.com/gohugoio/hugo/commit/9c2527f8558e3a2253d9c69e30d82d96a62d1cff) | [定位函数](https://github.com/gohugoio/hugo/blob/9c2527f8558e3a2253d9c69e30d82d96a62d1cff/markup/converter/converter.go#L44-L133) |
| gin-gonic/gin | middleware、上下文复用、顺序测试 | [`dcaa4296d111`](https://github.com/gin-gonic/gin/commit/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9) | [定位函数](https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/context.go#L195-L245) |
| syncthing/syncthing | 配置更新协议、状态所有权 | [`9af3c75f377c`](https://github.com/syncthing/syncthing/commit/9af3c75f377c51a7e3f285704025385104d8f428) | [定位函数](https://github.com/syncthing/syncthing/blob/9af3c75f377c51a7e3f285704025385104d8f428/lib/config/wrapper.go#L214-L291) |
| junegunn/fzf | 事件合并、有限 worker、取消 | [`52f4319a72c1`](https://github.com/junegunn/fzf/commit/52f4319a72c17e123396cc3a2e6abf2e96e9d753) | [定位函数](https://github.com/junegunn/fzf/blob/52f4319a72c17e123396cc3a2e6abf2e96e9d753/src/util/eventbox.go#L8-L54) |


本地复查源文件可以运行 `python3 research/fetch_sources.py`。脚本只下载固定 commit 的源码与许可，不执行上游代码。源码缓存和原始大目录树留在本地，未纳入课程代码提交；课程中的永久链接不依赖缓存存在。

[查看实验验证记录](verification.json)：包含 Go 实验的测试命令、环境和已验证的源码提交。上游项目未构建，也未进行全库审计。


[go-io]: https://github.com/golang/go/blob/c5941983810b68ba93c30f0ef22c91ad63fb3e5c/src/io/io.go#L55-L107
[go-copy]: https://github.com/golang/go/blob/c5941983810b68ba93c30f0ef22c91ad63fb3e5c/src/io/io.go#L375-L453
[go-multi]: https://github.com/golang/go/blob/c5941983810b68ba93c30f0ef22c91ad63fb3e5c/src/io/multi.go#L13-L76
[go-handler]: https://github.com/golang/go/blob/c5941983810b68ba93c30f0ef22c91ad63fb3e5c/src/net/http/server.go#L2344-L2352
[gin-flow]: https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/context.go#L195-L245
[gin-copy]: https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/context.go#L103-L156
[gin-pool]: https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/gin.go#L661-L676
[gin-options]: https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/gin.go#L193-L248
[gin-tests]: https://github.com/gin-gonic/gin/blob/dcaa4296d111981ffb31ac3eba90bb63e1eb5ab9/middleware_test.go#L159-L207
[hugo-converter]: https://github.com/gohugoio/hugo/blob/9c2527f8558e3a2253d9c69e30d82d96a62d1cff/markup/converter/converter.go#L44-L133
[hugo-deps]: https://github.com/gohugoio/hugo/blob/9c2527f8558e3a2253d9c69e30d82d96a62d1cff/deps/deps.go#L53-L138
[hugo-fs]: https://github.com/gohugoio/hugo/blob/9c2527f8558e3a2253d9c69e30d82d96a62d1cff/hugofs/nosymlinks_fs.go#L49-L86
[frp-plugin]: https://github.com/fatedier/frp/blob/832df8dff66d1b0caa0d776c450418f982756df1/pkg/plugin/client/plugin.go#L36-L71
[frp-base]: https://github.com/fatedier/frp/blob/832df8dff66d1b0caa0d776c450418f982756df1/client/proxy/proxy.go#L42-L120
[frp-tcp]: https://github.com/fatedier/frp/blob/832df8dff66d1b0caa0d776c450418f982756df1/client/proxy/general_tcp.go#L23-L47
[rag-contract]: https://github.com/infiniflow/ragflow/blob/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b/internal/storage/types.go#L56-L103
[rag-memory]: https://github.com/infiniflow/ragflow/blob/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b/internal/storage/memory.go#L32-L117
[rag-factory]: https://github.com/infiniflow/ragflow/blob/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b/internal/storage/storage_factory.go#L28-L95
[rag-search]: https://github.com/infiniflow/ragflow/blob/0c28d59ea1d362d9b6aa7481eed48c7fd9a95f0b/internal/agent/tool/retrieval_service.go#L80-L114
[ts-map]: https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/collections/ordered_map.go#L15-L187
[ts-tests]: https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/collections/ordered_map_test.go#L13-L104
[ts-visitor]: https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/ast/visitor.go#L9-L32
[ts-clone]: https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/ast/visitor.go#L130-L191
[ts-arena]: https://github.com/microsoft/TypeScript/blob/1f70213d4922b434345f639b441681e470c7cfc1/tsc/internal/core/arena.go#L7-L43
[sync-config]: https://github.com/syncthing/syncthing/blob/9af3c75f377c51a7e3f285704025385104d8f428/lib/config/wrapper.go#L214-L291
[sync-commit]: https://github.com/syncthing/syncthing/blob/9af3c75f377c51a7e3f285704025385104d8f428/lib/config/wrapper.go#L304-L345
[sync-contract]: https://github.com/syncthing/syncthing/blob/9af3c75f377c51a7e3f285704025385104d8f428/lib/config/wrapper.go#L40-L90
[fzf-box]: https://github.com/junegunn/fzf/blob/52f4319a72c17e123396cc3a2e6abf2e96e9d753/src/util/eventbox.go#L8-L54
[fzf-loop]: https://github.com/junegunn/fzf/blob/52f4319a72c17e123396cc3a2e6abf2e96e9d753/src/matcher.go#L81-L147
[fzf-workers]: https://github.com/junegunn/fzf/blob/52f4319a72c17e123396cc3a2e6abf2e96e9d753/src/matcher.go#L171-L241
[fzf-cancel]: https://github.com/junegunn/fzf/blob/52f4319a72c17e123396cc3a2e6abf2e96e9d753/src/matcher.go#L243-L271
[ollama-sched]: https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched.go#L60-L104
[ollama-admit]: https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched.go#L169-L240
[ollama-complete]: https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched.go#L367-L407
[ollama-request]: https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched.go#L475-L491
[ollama-test]: https://github.com/ollama/ollama/blob/83ed7d9965b1ee07e0f0b29fd46e47c31f0fcab8/server/sched_test.go#L60-L93
[k8s-controller]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/pkg/controller/deployment/deployment_controller.go#L67-L100
[k8s-worker]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/pkg/controller/deployment/deployment_controller.go#L481-L519
[k8s-sync]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/pkg/controller/deployment/deployment_controller.go#L574-L615
[k8s-lifetime]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/pkg/controller/deployment/deployment_controller.go#L171-L199
[k8s-queue]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/staging/src/k8s.io/client-go/util/workqueue/queue.go#L190-L302
[k8s-rate]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/staging/src/k8s.io/client-go/util/workqueue/rate_limiting_queue.go#L129-L147
[k8s-test]: https://github.com/kubernetes/kubernetes/blob/b2ec8b6fefac451a2dedafc4dd71f2f16c7a6abe/staging/src/k8s.io/client-go/util/workqueue/queue_test.go#L112-L170
[awesome-readme]: https://github.com/avelino/awesome-go/blob/681aa563a169472af1d00dbe4d1976a291f5f0bd/README.md
[caveman-license]: https://github.com/JuliusBrussee/caveman/blob/5184b3d11ac6a1acb7d44b9bfaa31698157cff97/LICENSING.md
