#sse —— SSE 事件解析

    纯逻辑包:吃字节、吐事件,不碰网络、不碰 async。解析规则按 WHATWG EventSource 规范实现,用同步测试逐条钉死。任何字节来源(HTTP 响应体、文件、WebSocket)都能复用同一个解析器。

    import { "q2316367743/moonhttp/sse", }

    #为什么不能用 read_until("\n\n") 切事件

    规范允许行尾是 CRLF / LF / CR 三种,而 CRLF 流上事件边界的字节是 0D 0A 0D 0A——里面没有连续两个 0A,按 "\n\n" 找分隔符永远匹配不到,会一直累积到连接关闭。所以这里按字节扫描,自己识别三种行尾。

    #SseEvent

    字段说明
    event事件类型,来自 event: 字段;流里没给过就是规范默认的 "message"
    data事件数据,来自 data: 字段;同一事件的多条 data: 行用 \n 连接
    id最近一次 id: 的值;None 表示流里还没出现过
    retry最近一次 retry: 的毫秒值;None 表示流里还没出现过

    id 与 retry 是跨事件持久的解析器状态(对齐 EventSource 的内部状态):一旦出现过,就跟着后面每个事件一起交出来。断线重连要用这两个值,本包不自动重连,把它们留给调用方。

    #SseParser

    方法说明
    SseParser::new()干净状态:没有事件,也没有 ID / 重连间隔
    push(chunk)喂一块字节;完整的事件进内部队列
    next()取一个已解析完的事件,没有就 None——None 不代表流结束,可能只是数据还没到齐
    finish()告诉解析器「流到此结束」(读到 EOF 时调用一次),处理末尾那个卡住的裸 CR;未以空行收尾的半行按规范丢弃。重复调用安全

    #用法

    let parser = @sse.SseParser::new()
    parser.push(chunk1)
    parser.push(chunk2) // 一块含 0 个、1 个或多个事件都行;事件也可能横跨多块
    while parser.next() is Some(event) {
    println(event.event + ": " + event.data)
    }
    parser.finish()

    #几个边界

    • 任意切分安全:行尾、字段、甚至一个 UTF-8 字符被切在中间都没问题。缓冲区末尾的裸 CR 会留到下一块再判定——立刻消费就会把 CRLF 当成两次换行,凭空多切出一个空事件。
    • 三种行尾都认:CRLF / LF / CR。流开头的 UTF-8 BOM 会被丢掉。
    • data 缓冲为空的块不发事件(只有 id: / retry: 或注释的块属于这种);data: 写了空值是「data 为空串的事件」,照发。
    • retry: 的值必须全是 ASCII 数字才生效,其它情况整体忽略、不改动已生效的值。

    完整规则与实现取舍见 docs/06-sse.md。

    SseEvent

    pub(all) struct SseEvent {
    event : String
    data : String
    id : String?
    retry : Int?
    } derive(Eq,
    Debug
    )

    一个 SSE 事件,对应规范里的 MessageEvent。

    SseEvent::equal

    fn SseEvent::equal(SseEvent, SseEvent) -> Bool

    SseEvent::not_equal

    fn SseEvent::not_equal(x : SseEvent, y : SseEvent) -> Bool

    SseEvent::to_repr

    SseParser

    pub struct SseParser {
    // private fields
    }

    增量式 SSE 解析器:把任意切分的字节块喂进来,取出完整事件。

    「任意切分」是重点:一次 push 可能包含 0 个、1 个或多个事件, 一个事件也可能横跨多次 push——行尾、字段、甚至一个 UTF-8 字符 都可能被切在中间。调用方不需要关心这些边界。

    id 与 retry 是跨事件持久的解析器状态,对齐规范里 EventSource 的内部状态:一旦出现过,就跟着后面每个事件一起交出去。断线重连要用 这两个值,所以本项目不自动重连,但把它们交给调用方(见 docs/06-sse.md)。

    EOF 时未完成的事件块丢弃(规范行为):没有等到空行的事件不算数, 解析器也不会在流结束时补发。
    impl Debug for SseParser

    SseParser::finish

    fn SseParser::finish(self : SseParser) -> Unit

    告诉解析器「流到此结束」,读到 EOF 时调用一次。

    需要它只有一个原因:末尾那个卡住的 CR。push 遇到缓冲区最后一个字节是 CR 时会先不消费,因为它可能是 CRLF 的前半;流结束时不会再有 LF 了, 这个 CR 就是行尾。不调用它,data: x\r\r 这类以裸 CR 收尾的流会少一个事件。

    未以换行收尾的半行不会被补发:规范要求流结束时丢弃不完整的事件。 重复调用是安全的。

    SseParser::new

    fn SseParser::new() -> SseParser

    创建解析器:干净状态,没有任何事件、也没有 ID / 重连间隔。

    SseParser::next

    fn SseParser::next(self : SseParser) -> SseEvent?

    取一个已经解析完的事件;没有就返回 None。

    None 不代表流结束——也可能只是数据还没到齐,需要继续 push。 「流是否结束」由读取侧判定(见根包的 SseStream::next_event)。

    SseParser::push

    fn SseParser::push(self : SseParser, chunk : Bytes) -> Unit

    喂入一块字节;完整的事件会进内部队列,用 next() 取。

    返回时 pending 已经缩到「最多一行的一部分」,所以长连场景下 内存不会随运行时间增长——除非对端一直不发换行。

    SseParser::to_repr