设计与实现Design

代码不多,边界很清楚:Kafka 的事交给 moonkafka。

There is not much code, and the boundary is sharp: Kafka belongs to moonkafka.

模块边界Module boundary

cmd/main                        daqing/moonkafka                 Apache Kafka
────────────────────────────    ──────────────────────────────   ────────────
argument parsing                broker connection                broker
UTF-8 conversion      ────────► metadata handling       ────────► topic log
result rendering                Kafka protocol
                                Producer::send
                                Consumer::poll

左边这一列就是本项目全部的职责。协议编码、连接管理、元数据刷新、偏移量跟踪都在 moonkafka 里,cmd/main 没有一行自定义的 Kafka 客户端逻辑。

The left column is the whole of this project's responsibility. Protocol encoding, connection management, metadata refresh and offset tracking all live in moonkafka; cmd/main contains no custom Kafka client logic at all.

包导入
cmd/maindaqing/moonkafka、moonbitlang/async、core/encoding/utf8、core/env、core/string
tests/itestasync、async/fs、async/shell、core/argparse、core/string
PackageImports
cmd/maindaqing/moonkafka, moonbitlang/async, core/encoding/utf8, core/env, core/string
tests/itestasync, async/fs, async/shell, core/argparse, core/string

模块依赖是 daqing/moonkafka@0.3.4 与 moonbitlang/async@0.22.4。moon.mod 声明 preferred_target = "native";两个包都限制 supported_targets = "+native",cmd/main 还为 native 目标加了 cc-link-flags = "-lz"。

The module depends on daqing/moonkafka@0.3.4 and moonbitlang/async@0.22.4. moon.mod declares preferred_target = "native"; both packages constrain supported_targets = "+native", and cmd/main adds cc-link-flags = "-lz" for the native target.

produce 的路径The produce path

取参数、建连接、发一条、打印偏移量——全部发生在 @async.with_task_group 里,连接在作用域结束时由 defer 关闭。

Read the arguments, connect, send once, print the offset — all inside an @async.with_task_group, with the connection closed by defer when the scope ends.

async fn run_producer(args : Array[String]) -> Unit {
  guard args.get(2) is Some(topic) && args.get(3) is Some(value) else {
    println("usage: produce <topic> <value> [key] [host] [port]")
    return
  }
  let key = args.get(4).map(text => @utf8.encode(text))
  let host = args.get(5).unwrap_or("127.0.0.1")
  let port = parse_port(args, 6)

  @async.with_task_group(group => {
    let producer = @moonkafka.Producer::connect(
      group~, host~, port~, topic~, acks=-1,
    )
    defer producer.close()

    let offset = producer.send(key?, value=@utf8.encode(value))
    println("sent to topic \"\{topic}\" at offset \{offset}")
  })
}

consume 的路径The consume path

async fn run_consumer(args : Array[String]) -> Unit {
  let (max_messages, args) = take_max_messages(args)
  guard args.get(2) is Some(topic) else { /* usage */ return }
  let host = args.get(3).unwrap_or("127.0.0.1")
  let port = parse_port(args, 4)
  let start_from = parse_start_from(args.get(5))

  @async.with_task_group(group => {
    let consumer = @moonkafka.Consumer::connect(
      group~, host~, port~, topic~, start_from~,
    )
    defer consumer.close()

    let mut received = 0
    for ;; {
      for record in consumer.poll() {
        let key = record.key_utf8().unwrap_or("<null>")
        let value = record.value_utf8().unwrap_or("<tombstone>")
        println("offset=\{record.offset} timestamp=\{record.timestamp} ...")
        received = received + 1
      }
      if max_messages is Some(limit) && received >= limit {
        break
      }
    }
  })
}

外层是无限循环,内层遍历 poll() 返回的一批记录。 --max-messages 的判定放在批处理循环之外:计数达标就跳出, 因此读够 n 条后退出码是 0,不会因为多读了一批而超发。

The outer loop is unbounded; the inner one walks the batch poll() returns. The --max-messages check sits outside the batch loop: once the count is reached the program breaks, so it exits 0 after exactly n records rather than overshooting by a batch.

这是一个「简单消费者」This is a simple consumer

moonkafka 的 Consumer 在注释里写得很清楚:它拉取单个主题的所有分区, 不做 consumer group 协调,偏移量只记在内存里。这对一个演示程序正合适, 但有两件事需要知道:

moonkafka's Consumer says so itself: it fetches all partitions of one topic, without consumer group coordination, and tracks offsets in memory. That suits a demo, with two consequences worth knowing:

Consumer::poll() 接受 max_wait_ms(默认 500)与 max_bytes(默认 1 MiB)两个可选参数,本 demo 都用默认值。

Consumer::poll() takes optional max_wait_ms (500 by default) and max_bytes (1 MiB by default); the demo uses both defaults.

参数解析Argument parsing

--max-messages 需要前置处理:它在解析早期被 take_max_messages 摘出来,返回「限制值 + 剩余位置参数」。

--max-messages needs a pre-pass: take_max_messages pulls it out early and returns the limit together with the remaining positional arguments, in order.

fn take_max_messages(args : Array[String]) -> (Int?, Array[String])
代价:位置参数必须写在前面 因为 flag 与它的取值从序列里消失,写在中间的 --max-messages 会让后面的位置参数前移一格。细节见 命令行参考。
The cost: positionals go first Since the flag and its value vanish from the sequence, a --max-messages placed in the middle shifts later positionals up by one. Details in the CLI reference.

parse_port 与 parse_start_from 在同一层做校验:端口用 @string.parse_int 转换并在 catch 里报错,起始位置用模式匹配 限定为两个取值。

parse_port and parse_start_from validate at the same layer: the port is converted with @string.parse_int and reports errors from a catch block, while the start position is pattern-matched down to its two legal values.

错误处理Error handling

场景做法
缺少必需参数guard … else { println(usage); return }——打印用法后正常结束
取值非法println(…) ; abort("")——先给出人能看懂的说明,再中止
数字解析失败@string.parse_int(v) catch { _ => … },把异常换成一条消息
资源释放defer consumer.close() / defer producer.close(),连接一定被关掉
任务生命周期@async.with_task_group(group => …)——moonkafka 的连接需要借一个 task group
SituationApproach
Missing required argumentguard … else { println(usage); return } — print usage, return normally
Invalid valueprintln(…) ; abort("") — a readable message first, then abort
A number fails to parse@string.parse_int(v) catch { _ => … }, turning the exception into a message
Releasing resourcesdefer consumer.close() / defer producer.close() — the connection is always closed
Task lifetime@async.with_task_group(group => …) — moonkafka's connections borrow a task group

harness 的设计How the harness is built

命令是值,不是副作用Commands are values, not side effects

一条外部命令建模为 Step:用途、程序、参数。它由一批共享的构造函数生成, --dry-run 只是把这些值渲染出来,真实运行才去执行它们。 因为两边同源,计划不可能与真实行为脱节。

An external command is modelled as a Step: a purpose, a program, and its arguments. Shared builders produce those values; --dry-run merely renders them while a real run executes them. Because both sides come from the same source, the plan cannot drift from what actually happens.

struct Step {
  purpose : String
  program : String
  args : Array[String]
}

三种执行语义Three execution semantics

函数用途
run_requiredharness 依赖的前置步骤。失败就中止整个运行——后续阶段的结果没有意义
run_observed结果本身就是断言的一部分。无论成败都记录退出码与耗时
run_until_success对 produce 的有界重试,用来吸收 broker 的启动竞态
FunctionPurpose
run_requiredA step the harness depends on. Failure aborts the whole run — later phases could not mean anything
run_observedThe outcome is itself part of the assertion. Exit code and duration are recorded either way
run_until_successBounded retries around the produce, absorbing the broker's startup race

断言与中止Assertions and aborting

Runner 只记三件事:跑了多少条断言、失败多少、是否已经中止。 check 记录一条结论;fail 额外把运行标记为中止,后续阶段直接跳过。

Runner tracks three things: how many checks ran, how many failed, and whether the run has already aborted. check records a verdict; fail additionally marks the run aborted so later phases are skipped.

进程的失败通过 suberror ItestError::Aborted(String) 表达,并且刻意放在 main 的最后一行抛出——那时容器已经停掉、资源已经释放, 进程才带着一行可读的消息以非零码退出。

Failure is expressed as suberror ItestError::Aborted(String), raised deliberately on the last line of main: by then the container is stopped and resources are released, and the process exits non-zero with one readable message.

外部命令不会挂住 harnessExternal commands cannot hang the harness

每条命令都经 @shell.Cmd(...).output(timeout_ms~) 执行,带硬超时。 连「程序根本起不来」这种情况也算作一次普通的失败步骤(exit_code = -1 加一条说明), 而不是抛异常——比如一个没装的 compose provider。

Every command runs through @shell.Cmd(...).output(timeout_ms~) under a hard timeout. Even "the program could not be spawned at all" counts as an ordinary failed step (exit_code = -1 plus an explanation) rather than an exception — a compose provider that is not installed, say.