练习 25:steer——往正在跑的轮次里插一句话

上一章的进程常驻了,但它一干活就聋。轮次跑起来之后没人读键盘,你打的字 攒在终端缓冲里,等这一轮全部跑完才被当成下一句话读走。上一章的常见问题 里我说过这件事:队是终端替你排的,不是你的代码排的。

这一章把键盘接管过来。轮次跑着的时候你照样能说话,而且说的话有两种去处:

  • 插话:直接打。它会在模型下一次发请求之前塞进这一轮,模型立刻看到, 当场改方向——不用等这一轮结束,也不用重来一遍。
  • 排队/q 开头。它不掺进这一轮,等这一轮收工,单独跑一轮。

顺带清一笔旧账。练习 20 做并发扇出的时候留下过一个已知的洞:几个子 agent 同时要你批准,几个 confirm 一起读同一个 os.Stdin,提示交错着打、回答 落到谁头上全看调度。当时写成了常见问题,说"这一章不修"。现在必须修了—— 键盘马上要归输入 goroutine 独占,confirm 再自己伸手去读,坏的就不只是 观感了。这一章结束时你会拿竞态检测器亲眼看到那个洞,以及它被堵上。

敲进去

在练习 24 的代码上继续写。

第一件事,把读键盘收进一个 goroutine:

// ---- 输入层:全程序唯一的键盘读者 ----

// startInputReader 把"读键盘"这件事收进一个 goroutine,一行一行往外送。
// 上一章读键盘是在主循环里同步做的,轮次跑起来就没人读了——你打的字只能
// 攒在终端缓冲里等这一轮结束。收进 goroutine 之后,轮次跑着的时候键盘
// 照样有人听,这一章要的插话才成立。
//
// 全程序只有这一个地方碰 os.Stdin。这不是洁癖:练习 20 的并发扇出翻车
// 就翻在"好几个地方同时读同一个 os.Stdin"上。
func startInputReader() <-chan string {
	lines := make(chan string)
	go func() {
		defer close(lines) // Ctrl+D:关掉 channel,收到的人自己知道该收摊
		for {
			text, err := stdin.ReadString('\n')
			if err != nil {
				return
			}
			lines <- strings.TrimSpace(text)
		}
	}()
	return lines
}

第二件事,收件箱。中途进来的话先存这儿:

// queuePrefix 让你明说这句话不要插进当前这一轮。不带前缀的默认是插话。
const queuePrefix = "/q "

// inboxItem 是一条中途进来的消息。standalone 为真表示用户明说了"排队":
// 它不掺进正在跑的这一轮,等这一轮收工,单独跑一轮。蒸馏自 octo 的
// queuedTurn.standalone——网页端的 Cmd+Enter、TUI 的 Ctrl+Q 都是这个意思。
type inboxItem struct {
	text       string
	standalone bool
}

// inbox 是一个能被多个 goroutine 同时写的队列。写它的是输入循环,读它的
// 是正在跑的轮次,两边不在一个 goroutine 上,所以要加锁。
type inbox struct {
	mu    sync.Mutex
	items []inboxItem
}

func (ib *inbox) enqueue(text string, standalone bool) {
	if strings.TrimSpace(text) == "" {
		return
	}
	ib.mu.Lock()
	ib.items = append(ib.items, inboxItem{text: text, standalone: standalone})
	ib.mu.Unlock()
}

// drainSteer 只取能插进当前这一轮的那些,明说要排队的原地不动——它们
// 存在的意义就是单独跑一轮,掺进来就白说了。
func (ib *inbox) drainSteer() []string {
	ib.mu.Lock()
	defer ib.mu.Unlock()
	var out []string
	kept := ib.items[:0]
	for _, it := range ib.items {
		if it.standalone {
			kept = append(kept, it)
			continue
		}
		out = append(out, it.text)
	}
	ib.items = kept
	return out
}

// drainQueued 取走剩下的(排队的那些),一轮结束之后由 repl 逐条跑。
func (ib *inbox) drainQueued() []string {
	ib.mu.Lock()
	defer ib.mu.Unlock()
	var out []string
	for _, it := range ib.items {
		out = append(out, it.text)
	}
	ib.items = nil
	return out
}

第三件事,也是这一章最要紧的一行:在哪儿把收件箱倒进历史。加在 runTurn 那个 for 循环的最开头,发请求之前:

	for round := 1; round <= maxRounds; round++ {
		// 取用收件箱:位置很关键——在发请求之前,在上一轮工具结果已经
		// 落进历史之后。插话因此是一条独立的用户消息,不是被塞进某个
		// 工具结果里的一段字,模型下一次请求就能原样看到它。octo 的
		// runLoop 把这件事放在同一个位置,理由也是同一句。
		if steers := box.drainSteer(); len(steers) > 0 {
			for _, text := range steers {
				sess.History = append(sess.History, message{Role: "user", Content: text})
			}
			fmt.Fprintf(os.Stderr, "[插话进入这一轮:%d 条,模型这就看到]\n", len(steers))
		}

		r, err := send(ctx, base, apiKey, model, sess.History, reg.definitions())

位置不是随便挑的。往前挪一点,插话就会掉进"模型说要调工具"和"工具结果 回来"这两条消息中间,历史当场不合法;往后挪一点,它就得等下一轮才生效。 只有这个位置,插话既是一条独立的用户消息,又能被这一轮的下一次请求看到。

第四件事,权限确认不再自己读键盘:

// askRequest 是一次"要问人"的请求:一句话,和一个等回答的 channel。
// 工具跑在自己的 goroutine 里,它不去读键盘,而是把这个请求交给正在
// 守着输入的那个循环,然后停在 resp 上等。
type askRequest struct {
	prompt string
	resp   chan bool
}

// askCh 是工具和输入循环之间唯一的通道。没有缓冲:轮次没跑起来的时候
// 没人收,工具就该停在那儿——而不是自作主张放行。
var askCh = make(chan askRequest)

func confirm(ctx context.Context, prompt string) bool {
	resp := make(chan bool, 1)
	select {
	case askCh <- askRequest{prompt: prompt, resp: resp}:
	case <-ctx.Done():
		return false
	}
	select {
	case ok := <-resp:
		return ok
	case <-ctx.Done():
		return false
	}
}

两个 select 都得盯着 ctx.Done(),少一个都会在打断时卡死:轮次被打断, 主循环等轮次收摊,轮次等工具返回,工具等一个再也不会有人回答的问题。

第五件事,runInterruptible 从"守着两个 channel"长成真正的事件循环:

	var pending *askRequest // 正等着回答的那一问
	var askQueue []askRequest
	for {
		select {
		case err := <-done:
			// …… 和上一章一样:报错、收拾历史、返回 ……
			return

		case <-sig:
			cancel()
			<-done
			heal(sess, "[这一轮被用户打断]")
			fmt.Fprintln(os.Stderr, "\n[已打断这一轮。对话还在,接着说]")
			return

		case line, ok := <-lines:
			if !ok {
				lines = nil // Ctrl+D:这个 case 从此不再触发,别空转
				// 键盘从此没人了,悬着的和排队的批准不能永远等下去——
				// 没人能说 y,答案就是 N(fail closed,练习 23 的老规矩)。
				if pending != nil {
					fmt.Fprintln(os.Stderr, "[输入已关闭,没人能批准——按 N 处理]")
					pending.resp <- false
					pending = nil
				}
				for _, q := range askQueue {
					q.resp <- false
				}
				askQueue = nil
				continue
			}
			if pending != nil {
				// 有一问悬着,这一行就是答复,不是插话。
				answer := strings.ToLower(strings.TrimSpace(line))
				pending.resp <- answer == "y" || answer == "yes"
				pending = nil
				if len(askQueue) > 0 {
					pending, askQueue = &askQueue[0], askQueue[1:]
					printAsk(pending.prompt)
				}
				continue
			}
			if standalone := strings.HasPrefix(line, queuePrefix); standalone {
				text := strings.TrimSpace(strings.TrimPrefix(line, queuePrefix))
				box.enqueue(text, true)
				fmt.Fprintln(os.Stderr, "[已排队:这一轮跑完再单独跑它]")
			} else {
				box.enqueue(line, false)
				fmt.Fprintln(os.Stderr, "[已收下:下一次发请求前塞进这一轮]")
			}

		case req := <-askCh:
			// 键盘已经关了(Ctrl+D / 管道读完),这个问题永远等不到 y——
			// 立刻按 N 答复,别让工具吊在一个没人会回答的问题上。
			if lines == nil {
				fmt.Fprintln(os.Stderr, "[输入已关闭,没人能批准——按 N 处理]")
				req.resp <- false
				continue
			}
			// 并发的子 agent 可能同时要批准。一次只问一个,其余排队——
			// 蒸馏自 octo 的模态队列:直接覆盖会把前一个问题的等待方
			// 永远晾在那儿。
			if pending != nil {
				askQueue = append(askQueue, req)
				continue
			}
			r := req
			pending = &r
			printAsk(r.prompt)
		}
	}

最后,一轮跑完,收件箱里可能还有东西:排队的那些,以及赶在收工那一瞬间 才进来、没赶上被取用的插话。

// runFollowUps 把一轮结束后还留在收件箱里的东西跑掉。
//
// 赶不上的插话不能默默丢掉。用户打字的时候模型还在干活,他有理由认为这句
// 话进去了;等他发现没进去,中间已经隔了一轮。赶不上就折成一次跟进的
// 对话——octo 也是这么处理的。
func runFollowUps(base, apiKey, model string, reg *registry, sess *session, window int, lines <-chan string, box *inbox) {
	for {
		if late := box.drainSteer(); len(late) > 0 {
			fmt.Fprintf(os.Stderr, "[插话来晚了:这一轮已经收工,把 %d 条折成一次跟进的对话]\n", len(late))
			runInterruptible(base, apiKey, model, reg, sess, window, strings.Join(late, "\n\n"), lines, box)
			continue
		}
		queued := box.drainQueued()
		if len(queued) == 0 {
			return
		}
		for _, q := range queued {
			fmt.Fprintln(os.Stderr, "[开始排队的那一句]")
			runInterruptible(base, apiKey, model, reg, sess, window, q, lines, box)
		}
	}
}

repl 里开一次输入 goroutine、建一个收件箱,把两者一路传下去,每轮跑完 调一次 runFollowUps

跑起来

cd exercises/ex25
go build -o ex25 .
./ex25

给它一件要跑好几轮的活,然后在它还没跑完的时候接着打字

> 依次创建 5 个文件 a1.txt 到 a5.txt,每个文件里写一行「原始内容」。一次只用一个 write_file,写完一个再写下一个,不要用 bash 验证。
改主意了,剩下还没写的那几个,里面改成写「改过了」

第二行是在第一行还在跑的时候打的。

你应该看到什么

实验一:插话在下一次请求前进这一轮

[round 1] write_file({"path": "a1.txt", "content": "原始内容\n"})
[round 2] write_file({"content": "原始内容\n", "path": "a2.txt"})
[round 3] write_file({"content": "原始内容\n", "path": "a3.txt"})
[已收下:下一次发请求前塞进这一轮]
[round 4] write_file({"content": "原始内容\n", "path": "a4.txt"})
[插话进入这一轮:1 条,模型这就看到]
[round 5] write_file({"content": "改过了\n", "path": "a5.txt"})
完成。5 个文件都已创建:

- a1.txt ~ a4.txt:内容「原始内容」
- a5.txt:内容「改过了」(按你改的主意)

[本轮 6 次请求 · 最后一次输入 2470 tokens(命中缓存 2432)· finish_reason=stop]

磁盘上的东西对得上:

a1.txt: 原始内容
a2.txt: 原始内容
a3.txt: 原始内容
a4.txt: 原始内容
a5.txt: 改过了

只有 a5 改了。a4 是在插话被取用之前就写完的,插话管不着已经发生的事—— 它改的是模型接下来要做什么。

这一轮从头到尾没有重启:[本轮 6 次请求] 是一次连续的对话,插话是其中 第 5 次请求之前多出来的一条用户消息。把会话文件打开,顺序一清二楚:

 9 assistant   write_file(a4.txt)
10 tool        已写入 a4.txt(13 字节)
11 user        改主意了,剩下还没写的那几个,里面改成写「改过了」
12 assistant   write_file({"content": "改过了\n", "path": "a5.txt"})

第 11 条就是插话落地的位置:在上一个工具结果之后,在下一次请求之前。

实验二:插话赶不上这一轮,也不会被吃掉

本机 Ollama(qwen3:4b-instruct)跑同一组,撞出了另一条路。它没听"一次 只用一个 write_file",第一轮就把 5 个文件全写了,于是插话还没来得及被 取用,这一轮已经收工:

[round 1] write_file({"content":"原始内容","path":"c1.txt"})
[round 1] write_file({"content":"原始内容","path":"c2.txt"})
[round 1] write_file({"content":"原始内容","path":"c3.txt"})
[round 1] write_file({"content":"原始内容","path":"c4.txt"})
[round 1] write_file({"content":"原始内容","path":"c5.txt"})
[已收下:下一次发请求前塞进这一轮]
已依次创建并写入文件 c1.txt ~ c5.txt,每个文件内容均为「原始内容」。
[本轮 2 次请求 · …… · finish_reason=stop]

[插话来晚了:这一轮已经收工,把 1 条折成一次跟进的对话]
[round 1] edit_file({"new_string":"改过了","old_string":"原始内容","path":"c4.txt"})
[round 1] edit_file({"new_string":"改过了","old_string":"原始内容","path":"c5.txt"})
已将 c4.txt 和 c5.txt 中的「原始内容」成功修改为「改过了」。

这就是 runFollowUps 那段兜底存在的理由。少了它,用户明明打了字,程序 收下了([已收下] 都打出来了),结果这句话原地蒸发——收下了却不办, 比一开始就不收更糟

实验三:/q 是排队,不是插话

同样的五个文件,中途打 /q 刚才那批文件一共几个?只回答数字

[本轮 6 次请求 · 最后一次输入 2486 tokens(命中缓存 2432)· finish_reason=stop]
5 个文件已依次创建完成,每个文件内容为一行「原始内容」:
- b1.txt …… b5.txt

[开始排队的那一句]
5

[本轮 1 次请求 · 最后一次输入 2354 tokens(命中缓存 1792)· finish_reason=stop]

正在跑的那一轮一点没受影响,五个文件照写照答。排队的那句话在它收工之后 单独跑了一轮,[本轮 1 次请求]——是独立的一轮,不是接在前面那轮尾巴上。

实验四:几个分身同时要你批准,提问一个一个来

这是练习 20 留下的那个洞。让两个子 agent 并发跑,各自都要执行一条 ask 档 命令,两个都需要你点头。左边是从进程启动算起的秒数:

  3.68  ⚠️  模型想执行: sudo -n true
 16.03  允许吗?(y/N) ===== 第一次回答 y =====
 16.03  ⚠️  模型想执行: sudo -n id
 24.05  ===== 第二次回答 y =====
 25.12  [子 agent "执行 sudo -n id" 结束]

第一个问题 3.68 秒就摆在那儿了,第二个问题一直压着不出声,直到 16.03 秒 第一个被回答,它才接上。屏幕上永远只有一个问题,两个 y 各自落在正确的 那一问上。

拿上一章的代码跑同一件事,两个问题当场一起糊在屏幕上。更要命的是用竞态 检测器编译(go build -race)之后,它直接报出根因:

⚠️  模型想执行: sudo -n true
允许吗?(y/N)
⚠️  模型想执行: sudo -n id
允许吗?(y/N) ==================
WARNING: DATA RACE
Read at 0x00c000090208 by goroutine 19:
  bufio.(*Reader).ReadSlice()
  bufio.(*Reader).ReadString()
Previous write at 0x00c000090208 by goroutine 18:
  bufio.(*Reader).fill()
  bufio.(*Reader).ReadSlice()
  bufio.(*Reader).ReadString()

两个 goroutine 同时在同一个 bufio.Reader 上读——这不是"观感不好",是 两个线程在抢同一块内存。这一章的代码跑同一个场景,竞态检测器报 0 次。

实验五:提问悬着的时候按 Ctrl+C

不回答,直接打断:

  1.34  [round 1] bash({"command": "sudo -n true"})
  1.34  ⚠️  模型想执行: sudo -n true
 12.04  允许吗?(y/N) ===== 不回答,直接 Ctrl+C =====
 12.04  [已打断这一轮。对话还在,接着说]
 16.86  没有执行——权限被拒,命令没有跑。

12.04 秒按下,12.04 秒回到提示符。confirm 里那两个 ctx.Done() 就是 为这一刻写的:没有它们,主循环会永远等一个没人回答的问题。

发生了什么

一个循环,四个来源。 上一章的 select 守着两件事:轮次结束、信号。 这一章多了两件:键盘、"要问人"。代码的形状一点没变,只是多了两个 case。 这就是常驻循环这个结构的全部价值——后面几章还要往里加定时唤醒、后台 任务的完成通知,加的都是 case,不是新结构。octo 的常驻循环里同时往里送 消息的来源有十来个,长的还是这个样子。

为什么取用的位置比取用本身重要。 插话最自然的写法,是把它拼进正在 执行的那个工具的结果里——反正模型都要读工具结果。这样写省事,但模型看到 的就不是"用户说了一句话",而是"某个工具的输出里混了一段人话"。octo 的 注释直接点了这件事:放在发请求之前单独倒进历史,是为了让模型看到的中途 插话是一条独立的用户消息,而不是折进工具输出的一段字。这不是洁癖, 是模型分不分得清"谁在说话"的问题。

插话和排队为什么必须分开。 两者都是"跑着的时候进来的话",但意图相反。 插话是"你正在做的这件事,改一下";排队是"这件事你做完,再做另一件"。 把排队的当插话塞进去,模型会以为你在改当前任务的要求;把插话的当排队 留到最后,你眼睁睁看着它把你已经不要的东西做完。octo 用一个 standalone 布尔值区分这两者,网页端是 Cmd+Enter,终端是 Ctrl+Q;我们 这个按行读的界面没有组合键可用,就用一个前缀。

为什么权限确认必须搬家。 一旦键盘归输入 goroutine 独占,confirm 自己去读就是两个读者抢一份字节流——实验四的竞态报告是这句话的硬证据。 搬家之后还白捡两个好处:并发的提问天然排成一列(谁先到谁先问,一次只 问一个),以及打断能穿透到"正在等人回答"的工具。练习 20 那个洞,修的 不是它的表现,是它的成因。

这一章仍然不是一个工具。 和上一章一样,插话、排队、收件箱,模型在 工具列表里一个都看不见。它看得见的只有结果:历史里多出来一条用户消息。 Part 7 的前两章都在搭骨架,从下一章起,加的东西又变回工具了。

常见问题

Q:有提问悬着的时候我打的字算什么? 算答复,不算插话。屏幕上摆着 允许吗?(y/N) 的时候,你打什么都是在回答 它——打一句别的话,它会被当成"不是 y",也就是拒绝。这是有意的:一个 悬而未决的授权问题比一句插话优先。想插话,先把问题回答掉。

Q:插话会不会打断模型正在执行的工具? 不会。插话只在两次请求之间生效,正在跑的那条命令该跑多久还是跑多久。要 让正在跑的东西立刻停,那是 Ctrl+C 的活。这两件事故意分开:插话是"改方向", 打断是"别做了"。

Q:我一次插好几句会怎样? 按顺序全部进历史,每句一条用户消息。drainSteer 一次取完,倒进去的顺序 就是你打的顺序。

Q:为什么排队的那句话跑完,第一轮的上下文还在? 因为 sess.History 从头到尾就一份。排队的那句是新的一轮,但接在同一份 历史后面——实验三里它答得出"5",靠的正是上一轮就在历史里。

Q:/q 这个前缀会不会和用户真想说的话撞上? 会。真要以 /q 开头说话,现在没辙。真实产品用组合键而不是前缀,就是 为了绕开这类冲突。这本书的界面是一行一行读的,没有组合键可用,取舍在 这里说清楚了。

加分练习

  1. 让插话能撤回。 打完回车又后悔,是很常见的事。给收件箱加一个 "把最后一条还没被取用的插话撤掉"的操作,/undo 触发。想清楚一件事: 如果这条插话已经被取用进历史了,撤回该怎么回应——octo 的做法是, 撤不掉就明确告诉用户它已经生效了,而不是假装撤掉了。

  2. 把插话折成一条。 现在一次取用多条插话,会往历史里塞多条用户消息。 改成把它们合并成一条(中间空行隔开)。跑之前先想:合并之后,模型还 分得清这是两句先后说的话吗?两种做法各有代价,写下你选哪个、为什么。

  3. 给排队的那些加个查看和取消。 /queue 列出排着的,/drop N 扔掉 第 N 条。你会发现 drainQueued 这个"取走就清空"的接口不够用了—— 这正是真实产品里队列要能被观察、被修改的原因。

  4. 把"要问人"这条路也用起来。 askCh 现在只有权限确认在用。让 sub_agent 也能反过来问人一个问题(它的任务描述不清时),用同一条 路把问题送到输入循环。你会发现不用改输入循环一行代码——这是这个结构 给的好处,也是 octo 里 ask_user_question 这个工具的由来。