From 087018f269189d45b34710204cea259bc5f61575 Mon Sep 17 00:00:00 2001 From: Robert Betts Date: Tue, 11 Aug 2026 22:11:08 +0100 Subject: [PATCH 1/2] Add native Async h2c path and honest Std.Async docs (#6). Ship AsyncByteTransport / *Async serve-dial-unary without .block on the h2c hot path, keep IO APIs as edge adapters, and document blocking TLS FFI for v1.2.0. Co-authored-by: Cursor --- .github/workflows/ci.yml | 3 +- CHANGELOG.md | 16 +++- Examples/Helloworld/Server.lean | 6 +- Grpc.lean | 2 +- Grpc/Channel.lean | 56 +++++++++++- Grpc/Client.lean | 63 ++++++++++++++ Grpc/Metadata.lean | 2 +- Grpc/Server.lean | 10 ++- Grpc/Stream.lean | 12 +-- H2/Client.lean | 118 ++++++++++++++----------- H2/Server.lean | 150 +++++++++++++++++++++++++++----- H2/Transport.lean | 35 ++++++-- README.md | 11 ++- ROADMAP.md | 24 +++-- SECURITY.md | 2 +- Tests/AsyncH2cLoopback.lean | 43 +++++++++ docs/README.md | 3 +- docs/api-reference.md | 19 ++-- docs/architecture.md | 18 ++-- docs/async-io.md | 52 +++++++++++ docs/getting-started.md | 6 +- docs/packaging.md | 8 +- lakefile.lean | 7 +- 23 files changed, 534 insertions(+), 132 deletions(-) create mode 100644 Tests/AsyncH2cLoopback.lean create mode 100644 docs/async-io.md diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6a843a2..b37a057 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -22,7 +22,7 @@ jobs: chmod +x scripts/fetch-openssl-headers.sh ./scripts/fetch-openssl-headers.sh - name: Build libraries and tests - run: lake build Proofs bytesTests hpackTests h2Tests grpcTests securityTests trailersServer trailersLoopback tlsServer tlsLoopback helloworldServer helloworldClient interopServer interopClient protocGenLean4Grpc benchUnary benchSoak opsSmoke h2specServer routeGuideServer routeGuideClient adcSmoke fakeAdsServer + run: lake build Proofs bytesTests hpackTests h2Tests grpcTests securityTests trailersServer trailersLoopback asyncH2cLoopback tlsServer tlsLoopback helloworldServer helloworldClient interopServer interopClient protocGenLean4Grpc benchUnary benchSoak opsSmoke h2specServer routeGuideServer routeGuideClient adcSmoke fakeAdsServer - name: Formal proofs (compile-time) run: lake build Proofs - name: Unit tests @@ -33,6 +33,7 @@ jobs: ./.lake/build/bin/grpcTests ./.lake/build/bin/securityTests ./.lake/build/bin/trailersLoopback + ./.lake/build/bin/asyncH2cLoopback - name: Native ASAN / UBSan security harness run: | chmod +x scripts/security-asan.sh diff --git a/CHANGELOG.md b/CHANGELOG.md index 5236ba5..501bc89 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,20 @@ # Changelog -All notable changes to lean-grpc are documented here. The package version is the Lake/`Grpc.version` semver (currently **1.1.0**). Git tags such as `v1.1.0` are created manually by maintainers when publishing. +All notable changes to lean-grpc are documented here. The package version is the Lake/`Grpc.version` semver (currently **1.2.0**). Git tags such as `v1.2.0` are created manually by maintainers when publishing. + +## [1.2.0] — 2026-08-11 + +Honest Async IO + native Async h2c path ([#6](https://github.com/RileyBetts/lean-grpc/issues/6)). + +- **Docs honesty:** README / architecture / package description no longer imply end-to-end Async for the blocking IO adapters. New [docs/async-io.md](docs/async-io.md) describes UV-loop ownership, sync adapters, and the TLS sync caveat. +- **Async transport:** `H2.AsyncByteTransport`, `tcpTransportAsync` (send/recv **without** `.block`); `ByteTransport` / `tcpTransport` kept as `.block` adapters. +- **Async h2c:** `listenH2cAsync` / `connectH2cAsync` / `awaitResponseAsync` / `serveH2cAsync` / `Channel.unaryAsync` / `Client.unaryCallAsync` — zero `.block` on accept/connect/send/recv for the Async call chain. +- **Compatibility:** existing `serveH2c`, `unary`, `connectH2c` remain; they `.block` the Async core at the edge. +- **TLS:** in-process OpenSSL FFI stays **blocking** in v1.2.0 (documented); no false “fully async TLS” claim. +- **Tests:** `asyncH2cLoopback` — concurrent Async clients against `serveH2cAsync` on one UV loop. +- **Example:** helloworld server honors `LEAN_GRPC_ASYNC=1` → `serveH2cAsync`. + +Migration: additive — bump pin to `v1.2.0`. Prefer `*Async` for new h2c composition; keep IO APIs for lean-compliance and existing code. ## [1.1.0] — 2026-08-04 diff --git a/Examples/Helloworld/Server.lean b/Examples/Helloworld/Server.lean index 05bd4db..c7ca39b 100644 --- a/Examples/Helloworld/Server.lean +++ b/Examples/Helloworld/Server.lean @@ -3,6 +3,7 @@ Copyright © 2026, Riley Betts Ltd (rileybetts.ai) Released under Apache 2.0 license as described in the file LICENSE. -/ import Grpc +import H2 import Examples.Helloworld.Generated /-- Typed helloworld server (stubs from `Generated.lean` / `./scripts/gen-helloworld.sh`). -/ @@ -22,4 +23,7 @@ def main : IO Unit := do s := Grpc.Channelz.register s counters | _ => s := Grpc.Health.registerWithWatch s healthStatus - Grpc.Server.serveH2c s { host := "127.0.0.1", port := 50051 } + let cfg : H2.ServerConfig := { host := "127.0.0.1", port := 50051 } + match ← IO.getEnv "LEAN_GRPC_ASYNC" with + | some "1" => (Grpc.Server.serveH2cAsync s cfg).block + | _ => Grpc.Server.serveH2c s cfg diff --git a/Grpc.lean b/Grpc.lean index adbfd34..91c88b5 100644 --- a/Grpc.lean +++ b/Grpc.lean @@ -39,5 +39,5 @@ import Grpc.Grpclb import Grpc.Interceptor namespace Grpc -def version : String := "1.1.0" +def version : String := "1.2.0" end Grpc diff --git a/Grpc/Channel.lean b/Grpc/Channel.lean index 952752c..6e77610 100644 --- a/Grpc/Channel.lean +++ b/Grpc/Channel.lean @@ -18,6 +18,9 @@ import Grpc.ServiceConfig import Grpc.Retry import Grpc.Xds import Grpc.XdsAds +import Std.Async + +open Std.Async namespace Grpc @@ -109,7 +112,7 @@ def get (ch : Channel) : IO H2.ClientConn := do let st ← c.state.get -- Keepalive enforcement: unanswered PING → GOAWAY and drop connection. if st.pendingPingAtMs != 0 && now - st.pendingPingAtMs ≥ pingTimeoutMs then - H2.sendFrames c.transport (H2.goAwayFrames st) + H2.sendFramesBlocking c.transport (H2.goAwayFrames st) c.state.set { st with wentAway := true, pendingPing := ByteArray.empty, pendingPingAtMs := 0 } ch.conn.set none let addr ← ch.connAddr.get @@ -121,7 +124,7 @@ def get (ch : Channel) : IO H2.ClientConn := do let last ← ch.lastActivityMs.get if now - last ≥ ch.keepaliveMs && st.pendingPingAtMs == 0 then let payload := ByteArray.mk #[1,2,3,4,5,6,7,8] - H2.sendFrames c.transport #[H2.Frame.ping payload] + H2.sendFramesBlocking c.transport #[H2.Frame.ping payload] c.state.set { st with pendingPing := payload, pendingPingAtMs := now } ch.lastActivityMs.set now return c @@ -300,6 +303,53 @@ def unary (ch : Channel) (service method : String) (request : ByteArray) return last return last +/-- Async unary RPC (h2c send/recv without `.block` on the call path). + Hedging/retry still run via the sync `unary` adapter for v1.2.0. -/ +def unaryAsync (ch : Channel) (service method : String) (request : ByteArray) + (metadata : Metadata := {}) (timeout? : Option String := none) + (compress : Compression.Algorithm := .identity) : Async CallResult := do + match ch.serviceConfig.hedging, ch.serviceConfig.retry with + | some _, _ | _, some _ => + liftM (unary ch service method request metadata timeout? compress) + | none, none => + if request.size > ch.sendMsgSize then + return { status := .resourceExhausted "message too large", message := ByteArray.empty, + headers := #[], trailers := #[] } + let deadlineMs? := + match timeout? with + | some t => Metadata.parseTimeoutMs t + | none => ch.serviceConfig.timeoutMs + if let some 0 := deadlineMs? then + return { status := .deadlineExceeded, message := ByteArray.empty, headers := #[], trailers := #[] } + let mut md := metadata + if let some creds := ch.callCreds then + md ← liftM (creds.apply md) + let mut extra := Metadata.toFields md + if let some t := timeout? then + extra := extra.push (Metadata.timeout t) + else if let some ms := deadlineMs? then + extra := extra.push (Metadata.timeout s!"{ms}m") + liftM (pickAndMaybeReconnect ch) + if ← liftM (checkGoAway ch) then + return goAwayResult + let c ← liftM (get ch) + let t0 ← IO.monoMsNow + let useHttps := + match ch.channelCreds with + | .tls _ => true + | .insecure => false + let res ← Client.unaryCallAsync c service method ch.host request extra compress useHttps + let t1 ← IO.monoMsNow + ch.lastActivityMs.set t1 + let res := + match deadlineMs? with + | some ms => + if t1 - t0 ≥ ms || res.status.code == .deadlineExceeded then + { res with status := Status.deadlineExceeded } + else res + | none => res + return enforceRecvSize ch res + /-- Open a bidirectional client stream for `service/method`. -/ def openStream (ch : Channel) (service method : String) (metadata : Metadata := {}) : IO Stream.ClientStream := do @@ -362,7 +412,7 @@ def goAway (ch : Channel) : IO Unit := do | none => pure () | some c => let st ← c.state.get - H2.sendFrames c.transport (H2.goAwayFrames st) + H2.sendFramesBlocking c.transport (H2.goAwayFrames st) ch.conn.set none def close (ch : Channel) : IO Unit := do diff --git a/Grpc/Client.lean b/Grpc/Client.lean index 0eeb568..48bd3c5 100644 --- a/Grpc/Client.lean +++ b/Grpc/Client.lean @@ -9,6 +9,9 @@ import Grpc.Status import Grpc.Message import Grpc.Metadata import Grpc.Compression +import Std.Async + +open Std.Async namespace Grpc @@ -129,5 +132,65 @@ def unaryCall (c : H2.ClientConn) (service method : String) (authority : String) startUnary c service method authority request extraHeaders compress useHttps finishUnary c sid deadlineMs? +/-- Async variant of `startUnary` (no `.block` on send). -/ +def startUnaryAsync (c : H2.ClientConn) (service method : String) (authority : String) + (request : ByteArray) (extraHeaders : Array Hpack.HeaderField := #[]) + (compress : Compression.Algorithm := .identity) + (useHttps : Bool := false) : Async (UInt32 × Option Nat) := do + let scheme := if useHttps then Metadata.schemeHttps else Metadata.schemeHttp + let mut headers : Array Hpack.HeaderField := + #[ + Metadata.methodPost, + scheme, + ⟨Metadata.ascii ":authority", Metadata.ascii authority⟩, + Metadata.path service method, + Metadata.contentTypeGrpc, + Metadata.teTrailers, + Metadata.userAgent, + Metadata.grpcAcceptEncoding "identity,gzip,deflate,snappy" + ] ++ extraHeaders + if compress != .identity then + headers := headers.push (Metadata.grpcEncoding compress.name) + let body ← liftM (Message.encodeIO request compress) + let deadlineMs? : Option Nat := + Id.run do + for h in extraHeaders do + let (n, v) := headerAscii h + if n == "grpc-timeout" then + match Metadata.parseTimeoutMs v with + | some ms => return some ms + | none => pure () + return none + let sid ← H2.Client.startRequestAsync c headers body true + return (sid, deadlineMs?) + +/-- Async await + decode into `CallResult`. -/ +def finishUnaryAsync (c : H2.ClientConn) (streamId : UInt32) (deadlineMs? : Option Nat) : + Async CallResult := do + let resp ← H2.Client.awaitUnaryAsync c streamId deadlineMs? + let status := statusFromHeaders resp.headers resp.trailers + let respAlg := + Id.run do + for h in resp.headers do + let (n, v) := headerAscii h + if n == "grpc-encoding" then + if let some a := Compression.Algorithm.parse? v then return a + return Compression.Algorithm.gzip + let payloads ← + match ← liftM (Message.decodeAllIO (Bytes.Slice.ofByteArray resp.data) true respAlg + Message.defaultMaxMsgSize) with + | .ok ps => pure ps + | .error _ => pure #[] + let message := payloads.getD 0 ByteArray.empty + return { status, message, headers := resp.headers, trailers := resp.trailers } + +def unaryCallAsync (c : H2.ClientConn) (service method : String) (authority : String) + (request : ByteArray) (extraHeaders : Array Hpack.HeaderField := #[]) + (compress : Compression.Algorithm := .identity) + (useHttps : Bool := false) : Async CallResult := do + let (sid, deadlineMs?) ← + startUnaryAsync c service method authority request extraHeaders compress useHttps + finishUnaryAsync c sid deadlineMs? + end Client end Grpc diff --git a/Grpc/Metadata.lean b/Grpc/Metadata.lean index c18cf2c..08b591f 100644 --- a/Grpc/Metadata.lean +++ b/Grpc/Metadata.lean @@ -159,7 +159,7 @@ def methodGet : Hpack.HeaderField := ⟨ascii ":method", ascii "GET"⟩ def schemeHttp : Hpack.HeaderField := ⟨ascii ":scheme", ascii "http"⟩ def schemeHttps : Hpack.HeaderField := ⟨ascii ":scheme", ascii "https"⟩ def teTrailers : Hpack.HeaderField := ⟨ascii "te", ascii "trailers"⟩ -def userAgent (version : String := "1.1.0") : Hpack.HeaderField := +def userAgent (version : String := "1.2.0") : Hpack.HeaderField := ⟨ascii "user-agent", ascii s!"grpc-lean/{version}"⟩ def http415 : Array Hpack.HeaderField := diff --git a/Grpc/Server.lean b/Grpc/Server.lean index a6a6799..4d3d5d1 100644 --- a/Grpc/Server.lean +++ b/Grpc/Server.lean @@ -14,6 +14,9 @@ import Grpc.PeerIdentity import Grpc.Stream import Grpc.Compression import Grpc.Tls +import Std.Async + +open Std.Async namespace Grpc @@ -305,9 +308,14 @@ def handlerFor (s : Server) (peerIdentity : Option PeerIdentity := none) def serveH2c (s : Server) (cfg : H2.ServerConfig := {}) : IO Unit := H2.Server.listen cfg (handlerFor s none false) +/-- Async h2c serve (accept/send/recv without `.block` on the hot path). -/ +def serveH2cAsync (s : Server) (cfg : H2.ServerConfig := {}) : Async Unit := + H2.Server.listenAsync cfg (handlerFor s none false) + /-- Serve with in-process TLS+ALPN `h2` (see `Grpc.Tls.serveH2` for the plaintext/mTLS decision based on `tlsCfg`). Peer identity is extracted per - accepted connection and threaded into unary context handlers. -/ + accepted connection and threaded into unary context handlers. + TLS byte IO remains blocking OpenSSL FFI in v1.2.0 — see docs/async-io.md. -/ def serveTls (s : Server) (tlsCfg : Tls.Config) (cfg : H2.ServerConfig := {}) : IO Unit := let mtlsRequired := tlsCfg.clientCaPath.isSome Tls.serveH2 tlsCfg cfg fun peerId => handlerFor s peerId mtlsRequired diff --git a/Grpc/Stream.lean b/Grpc/Stream.lean index 9abfcf3..be9608c 100644 --- a/Grpc/Stream.lean +++ b/Grpc/Stream.lean @@ -66,13 +66,13 @@ def send (w : StreamWriter) (msg : ByteArray) : IO Unit := do while st.sendConnWindow ≤ 0 && spins < 50 do spins := spins + 1 IO.sleep 1 - H2.sendFrames w.conn.transport (H2.Frame.dataFragmented w.streamId (Message.encodeId msg) false + H2.sendFramesBlocking w.conn.transport (H2.Frame.dataFragmented w.streamId (Message.encodeId msg) false st.ourSettings.maxFrameSize.toNat) def halfClose (w : StreamWriter) : IO Unit := do if ← w.closed.get then return w.closed.set true - H2.sendFrames w.conn.transport #[H2.Frame.data w.streamId ByteArray.empty true] + H2.sendFramesBlocking w.conn.transport #[H2.Frame.data w.streamId ByteArray.empty true] end StreamWriter @@ -88,7 +88,7 @@ partial def ensureHeaders (r : StreamReader) (fuel : Nat := 200) : IO (Array Hpa if s.endHeaders then r.headers.set (some s.requestHeaders) return s.requestHeaders - match (← r.conn.transport.recv? 65536) with + match (← (r.conn.transport.recv? 65536).block) with | none => throw (IO.userError "EOF") | some chunk => let mut buf ← r.conn.readBuf.get @@ -101,7 +101,7 @@ partial def ensureHeaders (r : StreamReader) (fuel : Nat := 200) : IO (Array Hpa for f in frames do let (st', outs) ← IO.ofExcept (H2.handleFrame st f) st := st' - H2.sendFrames r.conn.transport outs + H2.sendFramesBlocking r.conn.transport outs r.conn.state.set st ensureHeaders r (fuel - 1) @@ -142,7 +142,7 @@ partial def recv? (r : StreamReader) (fuel : Nat := 200) : IO (Option ByteArray) r.done.set true return none if fuel == 0 then throw (IO.userError "timeout waiting message") - match (← r.conn.transport.recv? 65536) with + match (← (r.conn.transport.recv? 65536).block) with | none => r.done.set true return none @@ -157,7 +157,7 @@ partial def recv? (r : StreamReader) (fuel : Nat := 200) : IO (Option ByteArray) for f in frames do let (st', outs) ← IO.ofExcept (H2.handleFrame st f) st := st' - H2.sendFrames r.conn.transport outs + H2.sendFramesBlocking r.conn.transport outs r.conn.state.set st recv? r (fuel - 1) diff --git a/H2/Client.lean b/H2/Client.lean index 88e829b..6b60382 100644 --- a/H2/Client.lean +++ b/H2/Client.lean @@ -18,8 +18,9 @@ open Std.Net namespace H2 +/-- Client connection. `transport` is async underneath; sync APIs `.block` at the edge. -/ structure ClientConn where - transport : ByteTransport + transport : AsyncByteTransport state : IO.Ref ConnState readBuf : IO.Ref ByteArray @@ -55,11 +56,11 @@ def rstToTrailers (code : UInt32) : Array Hpack.HeaderField := ⟨"grpc-message".toUTF8, msg.toUTF8⟩ ] -/-- Finish HTTP/2 client preface + settings exchange on an already-connected transport. -/ -def connectTransport (t : ByteTransport) : IO ClientConn := do +/-- Finish HTTP/2 client preface + settings on an async transport (no `.block`). -/ +def connectTransportAsync (t : AsyncByteTransport) : Async ClientConn := do t.send clientPreface let st := ConnState.create (isServer := false) - sendFrames t #[Frame.settings #[ + sendFramesAsync t #[Frame.settings #[ (.initialWindowSize, st.ourSettings.initialWindowSize), (.maxFrameSize, st.ourSettings.maxFrameSize) ]] @@ -73,16 +74,14 @@ def connectTransport (t : ByteTransport) : IO ClientConn := do | some chunk => buf := Bytes.Pool.pushBytes buf chunk let st0 ← state.get - let (frames, consumed) ← IO.ofExcept - (decodeFrames (Bytes.Slice.ofByteArray buf) st0.ourSettings.maxFrameSize.toNat) + let (frames, consumed) ← liftM (IO.ofExcept + (decodeFrames (Bytes.Slice.ofByteArray buf) st0.ourSettings.maxFrameSize.toNat)) buf := buf.extract consumed buf.size let mut st := st0 for f in frames do - let (st', outs) ← IO.ofExcept (handleFrame st f) + let (st', outs) ← liftM (IO.ofExcept (handleFrame st f)) st := st' - sendFrames t outs - -- SETTINGS ACK is already emitted by `handleFrame`; do not send a second ACK - -- (grpc C-core / Python rejects unexpected SETTINGS ACKs with GOAWAY). + sendFramesAsync t outs if f.type == .settings && !(Flags.has f.flags Flags.ack) then gotSettings := true state.set st @@ -90,19 +89,27 @@ def connectTransport (t : ByteTransport) : IO ClientConn := do readBuf.set buf return { transport := t, state, readBuf } -def connectH2c (host : String) (port : UInt16) : IO ClientConn := do +/-- Sync/TLS entry: wrap blocking `ByteTransport` then run async preface exchange. -/ +def connectTransport (t : ByteTransport) : IO ClientConn := + (connectTransportAsync (.ofBlocking t)).block + +/-- Dial h2c without `.block` on connect/send/recv. -/ +def connectH2cAsync (host : String) (port : UInt16) : Async ClientConn := do let sock ← TCP.Socket.Client.mk let addr ← match IPv4Addr.ofString host with | some a => pure (SocketAddress.v4 { addr := a, port }) | none => pure (SocketAddress.v4 { addr := IPv4Addr.ofParts 127 0 0 1, port }) - (sock.connect addr).block - connectTransport (tcpTransport sock) + sock.connect addr + connectTransportAsync (tcpTransportAsync sock) -/-- Start a new request stream; returns stream id. - When `endStream = false`, headers are sent and the stream stays open for DATA. -/ -def startRequest (c : ClientConn) (headers : Array Hpack.HeaderField) (body : ByteArray) - (endStream : Bool := true) : IO UInt32 := do +/-- Sync adapter over `connectH2cAsync`. -/ +def connectH2c (host : String) (port : UInt16) : IO ClientConn := + (connectH2cAsync host port).block + +/-- Start a new request stream (Async). -/ +def startRequestAsync (c : ClientConn) (headers : Array Hpack.HeaderField) (body : ByteArray) + (endStream : Bool := true) : Async UInt32 := do let mut st ← c.state.get let sid := st.nextClientStreamId st := { st with nextClientStreamId := sid + 2 } @@ -111,29 +118,30 @@ def startRequest (c : ClientConn) (headers : Array Hpack.HeaderField) (body : By st := st.upsertStream s c.state.set st let block := Hpack.encodeHeadersIndexed headers - -- Peers such as grpc-go FullDuplex empty_stream expect END_STREAM on DATA, - -- not on the initial HEADERS, when the request body is empty. if body.size == 0 && endStream then - sendFrames c.transport #[ + sendFramesAsync c.transport #[ Frame.headers sid block false true, Frame.data sid ByteArray.empty true ] else if body.size == 0 && !endStream then - sendFrames c.transport #[Frame.headers sid block false true] + sendFramesAsync c.transport #[Frame.headers sid block false true] else - sendFrames c.transport #[Frame.headers sid block false true] - sendFrames c.transport (Frame.dataFragmented sid body endStream st.ourSettings.maxFrameSize.toNat) + sendFramesAsync c.transport #[Frame.headers sid block false true] + sendFramesAsync c.transport (Frame.dataFragmented sid body endStream st.ourSettings.maxFrameSize.toNat) return sid -/-- Wait for response headers + data + optional trailers. - When `deadlineMs?` is set, RST_STREAM + throw after wall-clock expiry. -/ -partial def awaitResponse (c : ClientConn) (streamId : UInt32) (fuel : Nat := 200) - (deadlineMs? : Option Nat := none) (t0 : Nat := 0) : IO Response := do +def startRequest (c : ClientConn) (headers : Array Hpack.HeaderField) (body : ByteArray) + (endStream : Bool := true) : IO UInt32 := + (startRequestAsync c headers body endStream).block + +/-- Wait for response headers + data + optional trailers (Async — no `.block`). -/ +partial def awaitResponseAsync (c : ClientConn) (streamId : UInt32) (fuel : Nat := 200) + (deadlineMs? : Option Nat := none) (t0 : Nat := 0) : Async Response := do if fuel == 0 then throw (IO.userError "timeout waiting response") if let some ms := deadlineMs? then let now ← IO.monoMsNow if now - t0 ≥ ms then - sendFrames c.transport #[Frame.rstStream streamId 0x8] + sendFramesAsync c.transport #[Frame.rstStream streamId 0x8] throw (IO.userError "DEADLINE_EXCEEDED") let st ← c.state.get if let some s := st.getStream streamId then @@ -158,16 +166,15 @@ partial def awaitResponse (c : ClientConn) (streamId : UInt32) (fuel : Nat := 20 let mut buf ← c.readBuf.get buf := Bytes.Pool.pushBytes buf chunk let st0 ← c.state.get - let (frames, consumed) ← IO.ofExcept - (decodeFrames (Bytes.Slice.ofByteArray buf) st0.ourSettings.maxFrameSize.toNat) + let (frames, consumed) ← liftM (IO.ofExcept + (decodeFrames (Bytes.Slice.ofByteArray buf) st0.ourSettings.maxFrameSize.toNat)) c.readBuf.set (buf.extract consumed buf.size) let mut st := st0 for f in frames do - let (st', outs) ← IO.ofExcept (handleFrame st f) + let (st', outs) ← liftM (IO.ofExcept (handleFrame st f)) st := st' - sendFrames c.transport outs + sendFramesAsync c.transport outs c.state.set st - -- Re-check RST after processing frames (may have just closed the stream). if let some s := st.getStream streamId then if let some code := s.rstErrorCode then return { @@ -176,18 +183,17 @@ partial def awaitResponse (c : ClientConn) (streamId : UInt32) (fuel : Nat := 20 data := s.dataBuf rstErrorCode := some code } - awaitResponse c streamId (fuel - 1) deadlineMs? t0 + awaitResponseAsync c streamId (fuel - 1) deadlineMs? t0 -/-- Await a previously-started unary request's response (headers/data/trailers), - converting a wall-clock deadline timeout into synthesized DEADLINE_EXCEEDED - trailers. Split out from `unary` so callers that need the stream id early - (e.g. parallel hedging, which must be able to `resetStream` a losing - attempt while it is still in flight) can call `startRequest` themselves. -/ -def awaitUnary (c : ClientConn) (streamId : UInt32) (deadlineMs? : Option Nat := none) : - IO Response := do +def awaitResponse (c : ClientConn) (streamId : UInt32) (fuel : Nat := 200) + (deadlineMs? : Option Nat := none) (t0 : Nat := 0) : IO Response := + (awaitResponseAsync c streamId fuel deadlineMs? t0).block + +def awaitUnaryAsync (c : ClientConn) (streamId : UInt32) (deadlineMs? : Option Nat := none) : + Async Response := do let t0 ← IO.monoMsNow try - awaitResponse c streamId 200 deadlineMs? t0 + awaitResponseAsync c streamId 200 deadlineMs? t0 catch e => if (toString e).contains "DEADLINE_EXCEEDED" then return { @@ -200,22 +206,28 @@ def awaitUnary (c : ClientConn) (streamId : UInt32) (deadlineMs? : Option Nat := } else throw e +def awaitUnary (c : ClientConn) (streamId : UInt32) (deadlineMs? : Option Nat := none) : + IO Response := + (awaitUnaryAsync c streamId deadlineMs?).block + +def unaryAsync (c : ClientConn) (headers : Array Hpack.HeaderField) (body : ByteArray) + (deadlineMs? : Option Nat := none) : Async Response := do + let sid ← startRequestAsync c headers body true + awaitUnaryAsync c sid deadlineMs? + def unary (c : ClientConn) (headers : Array Hpack.HeaderField) (body : ByteArray) - (deadlineMs? : Option Nat := none) : IO Response := do - let sid ← startRequest c headers body true - awaitUnary c sid deadlineMs? - -/-- Send RST_STREAM to cancel a stream. - Also marks the stream closed locally with the same error code, so a - client-initiated cancel (e.g. interop `cancel_after_begin`) is reflected - immediately by `awaitResponse` (via `rstToTrailers`) without depending on - anything coming back from the peer. -/ -def resetStream (c : ClientConn) (streamId : UInt32) (errorCode : UInt32 := 0x8) : IO Unit := do - sendFrames c.transport #[Frame.rstStream streamId errorCode] + (deadlineMs? : Option Nat := none) : IO Response := + (unaryAsync c headers body deadlineMs?).block + +def resetStreamAsync (c : ClientConn) (streamId : UInt32) (errorCode : UInt32 := 0x8) : Async Unit := do + sendFramesAsync c.transport #[Frame.rstStream streamId errorCode] let st ← c.state.get match st.getStream streamId with | some s => c.state.set (st.upsertStream { s with state := .closed, rstErrorCode := some errorCode }) | none => pure () +def resetStream (c : ClientConn) (streamId : UInt32) (errorCode : UInt32 := 0x8) : IO Unit := + (resetStreamAsync c streamId errorCode).block + end Client end H2 diff --git a/H2/Server.lean b/H2/Server.lean index 27bfa18..4b21d36 100644 --- a/H2/Server.lean +++ b/H2/Server.lean @@ -44,11 +44,37 @@ private def parseAddr (host : String) (port : UInt16) : IO SocketAddress := do | none => return .v4 { addr := IPv4Addr.ofParts 127 0 0 1, port } +/-- Build response HEADERS frames only; DATA is queued via `queueSend` for flow control. -/ +def responseHeaderFrames (streamId : UInt32) (headers : Array Hpack.HeaderField) + (bodyEmpty : Bool) (trailersEmpty : Bool) (headersAlreadySent : Bool) (finished : Bool) : + Array Frame := + Id.run do + if headersAlreadySent || headers.isEmpty then return #[] + let endOnHeaders := finished && bodyEmpty && trailersEmpty + let hdrBlock := Hpack.encodeHeadersIndexed headers + return #[Frame.headers streamId hdrBlock endOnHeaders true] + +def sendFramesAsync (t : AsyncByteTransport) (frames : Array Frame) : Async Unit := do + for f in frames do + t.send (Frame.encode f) + +/-- Sync helper over `AsyncByteTransport` (`.block` at the edge). -/ +def sendFramesBlocking (t : AsyncByteTransport) (frames : Array Frame) : IO Unit := + (sendFramesAsync t frames).block + def sendFrames (t : ByteTransport) (frames : Array Frame) : IO Unit := do for f in frames do t.send (Frame.encode f) /-- Read until we have at least `n` bytes; returns exactly `n` bytes and any leftover. -/ +partial def recvExactAsync (t : AsyncByteTransport) (need : Nat) + (acc : ByteArray := ByteArray.empty) : Async (ByteArray × ByteArray) := do + if acc.size ≥ need then + return (acc.extract 0 need, acc.extract need acc.size) + match (← t.recv? 65536) with + | none => throw (IO.userError "EOF") + | some chunk => recvExactAsync t need (Bytes.Pool.pushBytes acc chunk) + partial def recvExact (t : ByteTransport) (need : Nat) (acc : ByteArray := ByteArray.empty) : IO (ByteArray × ByteArray) := do if acc.size ≥ need then @@ -57,24 +83,75 @@ partial def recvExact (t : ByteTransport) (need : Nat) (acc : ByteArray := ByteA | none => throw (IO.userError "EOF") | some chunk => recvExact t need (Bytes.Pool.pushBytes acc chunk) -/-- Build response HEADERS frames only; DATA is queued via `queueSend` for flow control. -/ -def responseHeaderFrames (streamId : UInt32) (headers : Array Hpack.HeaderField) - (bodyEmpty : Bool) (trailersEmpty : Bool) (headersAlreadySent : Bool) (finished : Bool) : - Array Frame := - Id.run do - if headersAlreadySent || headers.isEmpty then return #[] - let endOnHeaders := finished && bodyEmpty && trailersEmpty - let hdrBlock := Hpack.encodeHeadersIndexed headers - return #[Frame.headers streamId hdrBlock endOnHeaders true] +/-- Drain buffered reads into frames and handle them (Async transport). -/ +partial def processBufferAsync (t : AsyncByteTransport) (st : ConnState) (buf : ByteArray) + (handler : StreamHandler) : Async (ConnState × ByteArray) := do + let (frames, consumed) ← + match decodeFrames (Bytes.Slice.ofByteArray buf) st.ourSettings.maxFrameSize.toNat with + | .error e => + if e.startsWith "frame too large" then + sendFramesAsync t #[Frame.goAway st.lastPeerStreamId 0x6] + return ({ st with wentAway := true }, ByteArray.empty) + else throw (IO.userError e) + | .ok x => pure x + let mut st := st + let mut replies : Array Frame := #[] + for f in frames do + let (st', outs) ← liftM (IO.ofExcept (handleFrame st f)) + st := st' + replies := replies ++ outs + if let some s := st.getStream f.streamId then + if s.endHeaders && s.state != .closed then + let shouldInvoke := + !s.handlerFinished && (s.endStreamRemote || (s.dataBuf.size > s.dataConsumed)) + if shouldInvoke then + let hdrs := s.requestHeaders + let fresh := s.dataBuf.extract s.dataConsumed s.dataBuf.size + let resp ← liftM (handler s.id hdrs fresh s.endStreamRemote s.responseHeadersSent) + let emitFrames := !resp.headers.isEmpty || resp.body.size > 0 || + (resp.finished && (s.responseHeadersSent || !resp.headers.isEmpty || !resp.trailers.isEmpty)) + if emitFrames then + replies := replies ++ responseHeaderFrames s.id resp.headers + (resp.body.size == 0) resp.trailers.isEmpty s.responseHeadersSent resp.finished + let headersNowSent := s.responseHeadersSent || !resp.headers.isEmpty + let endOnHeaders := resp.finished && resp.body.size == 0 && resp.trailers.isEmpty + && !s.responseHeadersSent && !resp.headers.isEmpty + let s := { s with + dataConsumed := s.dataBuf.size + responseHeadersSent := headersNowSent + handlerFinished := resp.finished || s.handlerFinished + headersBuf := if resp.finished then ByteArray.empty else s.headersBuf + dataBuf := if resp.finished then ByteArray.empty else s.dataBuf + trailersBuf := if resp.finished then ByteArray.empty else s.trailersBuf } + st := st.upsertStream s + if !endOnHeaders && (resp.body.size > 0 || !resp.trailers.isEmpty || + (resp.finished && headersNowSent)) then + let (st', dataFrames) := queueSend st s.id resp.body resp.trailers resp.finished + st := st' + replies := replies ++ dataFrames + else + let advanced := resp.finished || resp.body.size > 0 || !resp.headers.isEmpty + let s := { s with + dataConsumed := if advanced then s.dataBuf.size else s.dataConsumed + responseHeadersSent := s.responseHeadersSent || !resp.headers.isEmpty + handlerFinished := resp.finished || s.handlerFinished + state := if resp.finished then .closed else s.state + headersBuf := if resp.finished then ByteArray.empty else s.headersBuf + dataBuf := if resp.finished then ByteArray.empty else s.dataBuf + trailersBuf := if resp.finished then ByteArray.empty else s.trailersBuf } + st := st.upsertStream s + sendFramesAsync t replies + let rest := buf.extract consumed buf.size + return (st, rest) -/-- Drain buffered reads into frames and handle them. -/ +/-- Drain buffered reads into frames and handle them (blocking `ByteTransport`). -/ partial def processBuffer (t : ByteTransport) (st : ConnState) (buf : ByteArray) (handler : StreamHandler) : IO (ConnState × ByteArray) := do let (frames, consumed) ← match decodeFrames (Bytes.Slice.ofByteArray buf) st.ourSettings.maxFrameSize.toNat with | .error e => if e.startsWith "frame too large" then - sendFrames t #[Frame.goAway st.lastPeerStreamId 0x6] -- FRAME_SIZE_ERROR + sendFrames t #[Frame.goAway st.lastPeerStreamId 0x6] return ({ st with wentAway := true }, ByteArray.empty) else throw (IO.userError e) | .ok x => pure x @@ -86,8 +163,6 @@ partial def processBuffer (t : ByteTransport) (st : ConnState) (buf : ByteArray) replies := replies ++ outs if let some s := st.getStream f.streamId then if s.endHeaders && s.state != .closed then - -- After `finished := true`, only flushPending may run (no re-invoke). - -- Incremental duplex keeps invoking until it returns finished (trailers). let shouldInvoke := !s.handlerFinished && (s.endStreamRemote || (s.dataBuf.size > s.dataConsumed)) if shouldInvoke then @@ -102,7 +177,6 @@ partial def processBuffer (t : ByteTransport) (st : ConnState) (buf : ByteArray) let headersNowSent := s.responseHeadersSent || !resp.headers.isEmpty let endOnHeaders := resp.finished && resp.body.size == 0 && resp.trailers.isEmpty && !s.responseHeadersSent && !resp.headers.isEmpty - -- Keep stream open until flushPending finishes (flow-control WINDOW_UPDATE). let s := { s with dataConsumed := s.dataBuf.size responseHeadersSent := headersNowSent @@ -131,6 +205,32 @@ partial def processBuffer (t : ByteTransport) (st : ConnState) (buf : ByteArray) let rest := buf.extract consumed buf.size return (st, rest) +/-- Serve one accepted client (Async TCP — no `.block` on send/recv). -/ +partial def serveConnAsync (t : AsyncByteTransport) (handler : StreamHandler) : Async Unit := do + let (preface, leftover) ← recvExactAsync t clientPreface.size + if preface != clientPreface then + throw (IO.userError s!"bad client preface got={Bytes.BE.hexDump (Bytes.Slice.ofByteArray preface)}") + let mut st := ConnState.create + sendFramesAsync t #[Frame.settings #[ + (.initialWindowSize, st.ourSettings.initialWindowSize), + (.maxConcurrentStreams, st.ourSettings.maxConcurrentStreams), + (.maxFrameSize, st.ourSettings.maxFrameSize), + (.maxHeaderListSize, st.ourSettings.maxHeaderListSize) + ]] + let mut buf := leftover + if buf.size > 0 then + let (st', rest) ← processBufferAsync t st buf handler + st := st' + buf := rest + while !st.wentAway do + match (← t.recv? 65536) with + | none => break + | some chunk => + buf := Bytes.Pool.pushBytes buf chunk + let (st', rest) ← processBufferAsync t st buf handler + st := st' + buf := rest + /-- Serve one accepted client as h2c prior-knowledge (or TLS already terminated). -/ partial def serveConn (t : ByteTransport) (handler : StreamHandler) : IO Unit := do let (preface, leftover) ← recvExact t clientPreface.size @@ -157,27 +257,37 @@ partial def serveConn (t : ByteTransport) (handler : StreamHandler) : IO Unit := st := st' buf := rest -/-- Listen for h2c connections; each accepted client is served concurrently. -/ -partial def listenH2c (cfg : ServerConfig) (handler : StreamHandler) : IO Unit := do +/-- Listen for h2c connections under `Std.Async` (accept/send/recv without `.block`). + Each connection is scheduled with `background` on the same UV loop. -/ +partial def listenH2cAsync (cfg : ServerConfig) (handler : StreamHandler) : Async Unit := do let server ← TCP.Socket.Server.mk let addr ← parseAddr cfg.host cfg.port server.bind addr server.listen 128 - IO.println s!"H2 h2c listening on {cfg.host}:{cfg.port}" + IO.println s!"H2 h2c listening on {cfg.host}:{cfg.port} (Async)" while true do - let client ← server.accept.block - discard <| IO.asTask (prio := .dedicated) do + let client ← server.accept + background (prio := Task.Priority.dedicated) do try - serveConn (tcpTransport client) handler + serveConnAsync (tcpTransportAsync client) handler catch e => IO.eprintln s!"conn error: {e}" +/-- Sync adapter: runs `listenH2cAsync` to completion via `.block` at the edge. -/ +def listenH2c (cfg : ServerConfig) (handler : StreamHandler) : IO Unit := + (listenH2cAsync cfg handler).block + /-- Serve using a pre-built transport (e.g. in-process TLS accept). -/ def serveTransport (t : ByteTransport) (handler : StreamHandler) : IO Unit := serveConn t handler +/-- Async serve using an async transport. -/ +def serveTransportAsync (t : AsyncByteTransport) (handler : StreamHandler) : Async Unit := + serveConnAsync t handler + namespace Server def listen := listenH2c +def listenAsync := listenH2cAsync end Server end H2 diff --git a/H2/Transport.lean b/H2/Transport.lean index e486b62..d6b23ee 100644 --- a/H2/Transport.lean +++ b/H2/Transport.lean @@ -8,16 +8,41 @@ open Std.Async namespace H2 -/-- Byte pipe under HTTP/2 (plain TCP or TLS). -/ +/-- Async byte pipe under HTTP/2 (plain TCP). Prefer this over `ByteTransport` + when composing under `Std.Async` so send/recv do not call `.block`. -/ +structure AsyncByteTransport where + send : ByteArray → Async Unit + recv? : Nat → Async (Option ByteArray) + close : Async Unit := pure () + +/-- Synchronous byte pipe (legacy / TLS FFI). Each op may block the caller. -/ structure ByteTransport where send : ByteArray → IO Unit recv? : Nat → IO (Option ByteArray) close : IO Unit := pure () -/-- Wrap a connected `Std.Async` TCP client. -/ -def tcpTransport (sock : TCP.Socket.Client) : ByteTransport where - send := fun b => (sock.send b).block - recv? := fun n => (sock.recv? n.toUInt64).block +/-- Wrap a connected `Std.Async` TCP client without `.block`. -/ +def tcpTransportAsync (sock : TCP.Socket.Client) : AsyncByteTransport where + send := fun b => sock.send b + recv? := fun n => sock.recv? n.toUInt64 close := pure () +/-- Lift blocking `IO` send/recv into `Async` (still blocks the UV thread while + the IO runs — used for OpenSSL FFI transports in v1.2.0). -/ +def AsyncByteTransport.ofBlocking (t : ByteTransport) : AsyncByteTransport where + send := fun b => liftM (t.send b) + recv? := fun n => liftM (t.recv? n) + close := liftM t.close + +/-- Sync facade: each op `.block`s the underlying async transport. -/ +def ByteTransport.ofAsync (t : AsyncByteTransport) : ByteTransport where + send := fun b => (t.send b).block + recv? := fun n => (t.recv? n).block + close := t.close.block + +/-- Wrap a connected `Std.Async` TCP client as a **blocking** `ByteTransport` + (compatibility adapter; prefer `tcpTransportAsync` + `*Async` APIs). -/ +def tcpTransport (sock : TCP.Socket.Client) : ByteTransport := + .ofAsync (tcpTransportAsync sock) + end H2 diff --git a/README.md b/README.md index d4cc1c0..6ff150f 100644 --- a/README.md +++ b/README.md @@ -4,9 +4,11 @@ [![Lean](https://img.shields.io/badge/Lean-4.32-purple.svg)](lean-toolchain) [![Docs](https://img.shields.io/badge/docs-rileybetts.ai-0B3D2E.svg)](https://rileybetts.ai/oss/lean-grpc) -General-purpose **Lean 4 gRPC library**: HPACK + HTTP/2 + gRPC framing on `Std.Async.TCP`. +General-purpose **Lean 4 gRPC library**: HPACK + HTTP/2 + gRPC framing over `Std.Async.TCP`. -Standalone Lake package (**1.1.0**). Consumers depend via git tag or, after indexing, [Reservoir](https://reservoir.lean-lang.org/). +Standalone Lake package (**1.2.0**). Consumers depend via git tag or, after indexing, [Reservoir](https://reservoir.lean-lang.org/). + +**Async model:** sockets are `Std.Async.TCP`. Through v1.1.x the public API was blocking `IO` via `.block`. **v1.2.0** adds a native Async h2c path (`serveH2cAsync` / `unaryAsync`) with **zero** `.block` on accept/connect/send/recv; existing IO APIs remain as explicit sync adapters. In-process TLS is still blocking OpenSSL FFI — see [docs/async-io.md](docs/async-io.md). **Docs:** [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc) (curated) · [docs/](docs/README.md) (full in-repo index) @@ -20,7 +22,7 @@ In your `lakefile.lean`: ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.1.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" ``` Then `import Grpc`. After Reservoir lists the package you can use `require «lean-grpc»` without a git URL. Packaging details and the maintainer release checklist: [docs/packaging.md](docs/packaging.md). @@ -36,13 +38,14 @@ Then `import Grpc`. After Reservoir lists the package you can use `require «lea | [Packaging](docs/packaging.md) | Lake/Reservoir layout, consumer contract, release checklist | | [Provenance](docs/provenance.md) | Independent protocol implementation; third-party interop protos | | [Architecture](docs/architecture.md) | Layering and data flow | +| [Async IO](docs/async-io.md) | Std.Async model, sync adapters, TLS caveat | | [API reference](docs/api-reference.md) | Module catalogue | | [Protocol mapping](docs/protocol-mapping.md) | gRPC-over-HTTP/2 mapping for this stack | | [Conformance](docs/conformance.md) | Scorecard, interop matrix, allowlists | | [Formal proofs](docs/proofs.md) | Compile-time theorems for pure codecs | | [TLS / Envoy](docs/tls-envoy.md) | In-process OpenSSL and sidecars | | [CHANGELOG](CHANGELOG.md) | Version history | -| [ROADMAP](ROADMAP.md) | What v1.1.0 shipped vs open proof/hardening follow-ups | +| [ROADMAP](ROADMAP.md) | What v1.2.0 shipped vs open proof/hardening follow-ups | | [CONTRIBUTING](CONTRIBUTING.md) | Dev setup and PR expectations | | [SECURITY](SECURITY.md) | Vulnerability reporting (`security@rileybetts.ai`) | | [Code of Conduct](CODE_OF_CONDUCT.md) | Community standards | diff --git a/ROADMAP.md b/ROADMAP.md index d108469..2346135 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -1,6 +1,6 @@ # Roadmap -lean-grpc **v1.1.0** is the current package tip (Lake / `Grpc.version`): an interop-tested Lean 4 gRPC stack with a CI-gated **`Proofs`** library for selected pure codecs, plus additive mTLS peer-identity / request-context APIs for enterprise AuthN. It is **not** a machine-checked end-to-end PROTOCOL-HTTP2 / TLS / session proof. +lean-grpc **v1.2.0** is the current package tip (Lake / `Grpc.version`): an interop-tested Lean 4 gRPC stack with native Async h2c APIs (honest Std.Async model), a CI-gated **`Proofs`** library for selected pure codecs, plus additive mTLS peer-identity / request-context APIs for enterprise AuthN. It is **not** a machine-checked end-to-end PROTOCOL-HTTP2 / TLS / session proof. This document records what shipped, what is still open, and the next proof/hardening tranches. @@ -25,6 +25,16 @@ First public packaging baseline (interop + Lake/Reservoir layout). Formal Lean p | Dial / LB / retry, health, reflection, channelz | Present; ops demos gated | | ADC / xDS ADS | Mock / FakeAds CI; live Google paths allowlisted | +## Shipped — v1.2.0 (Async honesty) + +Native Async h2c + docs that match the implementation ([#6](https://github.com/RileyBetts/lean-grpc/issues/6), [docs/async-io.md](docs/async-io.md)). + +| Included | Deferred | +|---|---| +| `AsyncByteTransport` / `*Async` h2c serve/dial/unary | Nonblocking TLS BIO / off-loop OpenSSL | +| Sync IO APIs as `.block` adapters | Streaming `*Async` surface beyond unary | +| Concurrent `asyncH2cLoopback` CI | Full Concurrent retry/hedge under Async | + ## Shipped — v1.1.0 (IAM) Additive server AuthN plumbing for verified mTLS peer identity (see [feature/lean-grpc-iam-requirements.md](feature/lean-grpc-iam-requirements.md), [docs/cookbook-interceptors.md](docs/cookbook-interceptors.md)). @@ -51,7 +61,7 @@ Product/packaging release with **selected** compile-time proofs. See [docs/proof **Explicit non-goals (unchanged):** ALTS / GCE channel credentials, HTTP CONNECT proxying, full `cacheable_unary` proxy infrastructure, end-to-end session proofs. -## Toward v1.2.x / later +## Toward v1.3.x / later Work is grouped so each tranche can ship as a minor release without waiting for a full ConnState + Huffman proof stack. @@ -99,11 +109,11 @@ Close remaining audit follow-ups so later minors are not “proved but soft”: ## Milestone sketch ```text -v0.5.0 ──► v1.0.0 ──► v1.1.0 ──► v1.2.x - shipped interop + mTLS peer identity + ConnState + - selected Proofs ServerCallContext general frame/msg - (CI) + provenance (unary IAM) roundtrips / - Huffman trie +v0.5.0 ──► v1.0.0 ──► v1.1.0 ──► v1.2.0 ──► v1.3.x + shipped interop + mTLS peer identity + Async h2c honesty ConnState + + selected Proofs ServerCallContext (#6) + sync adapters general frame/msg + (CI) + provenance (unary IAM) roundtrips / + Huffman trie ``` Dates are intentionally omitted; order matters more than calendar. diff --git a/SECURITY.md b/SECURITY.md index 420cea9..ed1c42b 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -4,7 +4,7 @@ Security fixes are applied to the published tip (`main` / `development` as released). There is no long-term LTS train yet. -**Supported version:** the latest **tagged** release (e.g. `v1.1.0`). Maintainers create tags manually — see [docs/packaging.md](docs/packaging.md). Prefer GitHub Security Advisories for private reports when enabled. +**Supported version:** the latest **tagged** release (e.g. `v1.2.0`). Maintainers create tags manually — see [docs/packaging.md](docs/packaging.md). Prefer GitHub Security Advisories for private reports when enabled. ## Reporting a vulnerability diff --git a/Tests/AsyncH2cLoopback.lean b/Tests/AsyncH2cLoopback.lean new file mode 100644 index 0000000..b8ff412 --- /dev/null +++ b/Tests/AsyncH2cLoopback.lean @@ -0,0 +1,43 @@ +/- +Copyright © 2026, Riley Betts Ltd (rileybetts.ai) +Released under Apache 2.0 license as described in the file LICENSE. +-/ +import Grpc +import H2 +import Proto +import Std.Async +import Std.Async.Timer + +open Std.Async + +/-- Concurrent Async h2c loopback: `serveH2cAsync` + multiple `connectH2cAsync` / + `unaryCallAsync` clients on one UV loop (no `.block` on the Async hot path). -/ +private def oneClient (port : UInt16) (name : String) : Async String := do + let c ← H2.Client.connectH2cAsync "127.0.0.1" port + let req := Proto.HelloRequest.encode { name } + let res ← Grpc.Client.unaryCallAsync c "helloworld.Greeter" "SayHello" "127.0.0.1" req + if res.status.code != .ok then + throw (IO.userError s!"status {res.status.code.toUInt32}") + let reply ← liftM (IO.ofExcept (Proto.HelloReply.decode res.message)) + return reply.message + +def main : IO Unit := do + let port : UInt16 := 50061 + let mut s := Grpc.Server.empty + s := Grpc.Server.register s "helloworld.Greeter" "SayHello" fun reqBytes => do + let req ← IO.ofExcept (Proto.HelloRequest.decode reqBytes) + let reply : Proto.HelloReply := { message := s!"Hello, {req.name}" } + return (Proto.HelloReply.encode reply, Grpc.Status.ok) + (do + background (prio := Task.Priority.dedicated) do + Grpc.Server.serveH2cAsync s { host := "127.0.0.1", port } + sleep (Std.Time.Millisecond.Offset.ofNat 400) + let names : Array String := #["A", "B", "C", "D", "E", "F", "G", "H"] + let msgs : Array String ← Async.concurrentlyAll (names.map fun n => oneClient port n) + if msgs.size != names.size then + throw (IO.userError s!"expected {names.size} replies, got {msgs.size}") + for msg in msgs do + if !("Hello".isPrefixOf msg) then + throw (IO.userError msg) + IO.println "asyncH2cLoopback OK" + ).block diff --git a/docs/README.md b/docs/README.md index 5292623..798533e 100644 --- a/docs/README.md +++ b/docs/README.md @@ -15,13 +15,14 @@ User and contributor docs for the Lean 4 gRPC stack. | [Website sync](website-sync.md) | How curated docs stay aligned with rileybetts.ai | | [Provenance / IP diligence](provenance.md) | Independent protocol implementation note; third-party interop protos | | [Architecture](architecture.md) | Layering, data flow, process model | +| [Async IO](async-io.md) | Std.Async honesty, `*Async` APIs, sync adapters, TLS caveat | | [API reference](api-reference.md) | Primary `Grpc` / `H2` / `Proto` surfaces | | [Protocol mapping](protocol-mapping.md) | How lean-grpc maps onto gRPC-over-HTTP/2 | | [Conformance](conformance.md) | Scorecard, compliance estimates, interop matrix, allowlists | | [Formal proofs](proofs.md) | Compile-time `Proofs` library for pure codecs / maps | | [TLS / Envoy](tls-envoy.md) | In-process OpenSSL and optional sidecars | | [CHANGELOG](../CHANGELOG.md) | Package version history | -| [ROADMAP](../ROADMAP.md) | What v1.1.0 shipped vs open proof/hardening follow-ups | +| [ROADMAP](../ROADMAP.md) | What v1.2.0 shipped vs open proof/hardening follow-ups | | [CONTRIBUTING](../CONTRIBUTING.md) | Dev setup, tests, PR expectations | | [SECURITY](../SECURITY.md) | Vulnerability reporting | | [Code of Conduct](../CODE_OF_CONDUCT.md) | Community standards | diff --git a/docs/api-reference.md b/docs/api-reference.md index a9e5525..88a1f2c 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -1,6 +1,6 @@ # API reference -Lean module catalogue for consumers. Signatures are summarized; see source under `Grpc/`, `H2/`, `Proto/` for full definitions. Version string: `Grpc.version` (currently `1.1.0`). +Lean module catalogue for consumers. Signatures are summarized; see source under `Grpc/`, `H2/`, `Proto/` for full definitions. Version string: `Grpc.version` (currently `1.2.0`). Async vs sync IO: [async-io.md](async-io.md). Import umbrella: `import Grpc` (pulls status, channel, server, credentials, TLS, xDS, ops, etc.). Add `import Proto` for bundled message codecs. @@ -59,7 +59,7 @@ Length-prefixed gRPC frames (`Compressed-Flag` + 4-byte length + payload). `status`, `message` (decoded payload bytes), `headers`, `trailers`. -### `Grpc.Client.unaryCall` +### `Grpc.Client.unaryCall` / `unaryCallAsync` Low-level unary on an `H2.ClientConn` (scheme http/https, user-agent, compression). @@ -69,7 +69,7 @@ Low-level unary on an `H2.ClientConn` (scheme http/https, user-agent, compressio |---|---| | `connectH2c host port` | Plain h2c channel | | `dial target opts svc?` | Resolve + balanc + connect (`dns:///`, `host:port`, `xds:///`) | -| `unary ch service method request metadata? timeout? compress?` | Unary RPC | +| `unary` / `unaryAsync` | Unary RPC (Async = no `.block` on h2c call path; hedge/retry via sync adapter) | | `openStream ch service method metadata?` | Interactive bidi client stream | | `serverStream` / `clientStream` / `bidiStream` | Batch streaming helpers (messages + status) | | `get` / `goAway` / `close` | Connection pool / drain | @@ -121,8 +121,8 @@ Verified mTLS peer certificate identity (OpenSSL; subject DN is **RFC 2253**): | `registerTyped` / `registerTypedWithContext` | Typed unary adapters | | `registerServerStream` / `registerClientStream` / `registerBidi` | Streaming (raw bytes; context deferred) | | `registerServerStreamTyped` / `registerClientStreamTyped` / `registerBidiTyped` | Streaming with decode/encode adapters | -| `serveH2c` | Listen h2c (`peerIdentity = none`) | -| `serveTls` | Listen TLS+ALPN; per-connection peer identity → context handlers | +| `serveH2c` / `serveH2cAsync` | Listen h2c (`peerIdentity = none`); Async path has no `.block` on accept/send/recv | +| `serveTls` | Listen TLS+ALPN; per-connection peer identity → context handlers (**blocking** OpenSSL FFI in v1.2.0) | | `maxMsgSize` | Inbound limit | Bad `content-type` → HTTP **415**. Unknown method / zero timeout → trailers-only gRPC status. @@ -222,11 +222,12 @@ Consumers usually stay in `Grpc.*`. Useful lower APIs: | API | Purpose | |---|---| -| `H2.Client.connectH2c` / `connectTransport` | Client connection | -| `H2.Client.startRequest` / `awaitResponse` / `unary` / `resetStream` | Streams | +| `H2.Client.connectH2c` / `connectH2cAsync` / `connectTransport` | Client connection | +| `H2.Client.startRequest` / `awaitResponse` / `unary` (+ `*Async`) / `resetStream` | Streams | | `H2.Client.rstToTrailers` | RST → synthetic gRPC trailers | -| `H2.Server.listen` / `serveConn` | Server accept loop | -| `H2.ByteTransport` | Pluggable send/recv | +| `H2.Server.listen` / `listenAsync` / `serveConn` / `serveConnAsync` | Server accept loop | +| `H2.AsyncByteTransport` / `tcpTransportAsync` | Native Async send/recv (h2c) | +| `H2.ByteTransport` / `tcpTransport` | Sync facade (`.block` adapters) | | `H2.handleFrame` / `ConnState` | State machine (tested via h2spec) | --- diff --git a/docs/architecture.md b/docs/architecture.md index 199d2c9..81f479c 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -11,7 +11,7 @@ flowchart TB Hpack[Hpack: encode decode] Bytes[Bytes: Slice Pool BE] Native[Native: tls_ffi] - TCP[Std.Async.TCP / ByteTransport] + TCP[Std.Async.TCP / AsyncByteTransport] App --> Grpc App --> Proto @@ -39,8 +39,8 @@ Executables (examples, interop, benches, codegen plugin) sit outside the librari ## Client path -1. **`Channel.dial` / `connectH2c`** — parse target (`host:port`, `dns:///…`, `xds:///…`), resolve addresses, pick via balancer, open `H2.ClientConn` (h2c or TLS transport). -2. **`Channel.unary`** — apply call credentials → build request headers (`:method POST`, `:path`, `content-type`, `te: trailers`, `user-agent`, optional `grpc-timeout` / `grpc-encoding`) → length-prefix encode body → `H2.Client.startRequest` / `awaitResponse`. +1. **`Channel.dial` / `connectH2c`** — parse target (`host:port`, `dns:///…`, `xds:///…`), resolve addresses, pick via balancer, open `H2.ClientConn` (h2c or TLS transport). Async: `connectH2cAsync` / `Channel.unaryAsync`. +2. **`Channel.unary`** — apply call credentials → build request headers (`:method POST`, `:path`, `content-type`, `te: trailers`, `user-agent`, optional `grpc-timeout` / `grpc-encoding`) → length-prefix encode body → `H2.Client.startRequest` / `awaitResponse` (sync adapters `.block` the Async core). 3. **Response** — decode headers/trailers into `Status` (including HTTP non-200 / RST_STREAM mapping), inflate message if compressed, enforce max receive size. 4. **Retry / hedge** — service config may retry status codes with backoff, or launch parallel hedges and RST losers. 5. **Keepalive** — idle PING; unanswered PING → GOAWAY and reconnect. @@ -49,7 +49,7 @@ Streaming uses `Channel.openStream` → `Grpc.Stream` writer/reader over the sam ## Server path -1. **`Server.serveH2c` / `serveTls`** — accept TCP (or TLS listener), run HTTP/2 preface + SETTINGS. +1. **`Server.serveH2c` / `serveH2cAsync` / `serveTls`** — accept TCP (or TLS listener), run HTTP/2 preface + SETTINGS. Async h2c uses `background` per connection (no `.block` on accept/send/recv). 2. **`handlerFor`** — on each stream: parse path / content-type / timeout / encodings; reject bad content-type with **HTTP 415**; trailers-only for unimplemented / zero deadline. 3. Dispatch unary or streaming handler; encode response messages; send headers + DATA + trailers (`grpc-status` / `grpc-message` / optional `grpc-status-details-bin`). @@ -57,10 +57,11 @@ Health, reflection, and channelz are ordinary registered methods (`Grpc.Health`, ## Transport abstraction -`H2.ByteTransport` is a small send/recv interface. Implementations: +`H2.AsyncByteTransport` is the native Async send/recv interface (h2c hot path). `H2.ByteTransport` is the sync facade (`.block` adapters + TLS FFI). Implementations: -- Plain TCP (`H2.Client.connectH2c`, `H2.Server.listenH2c`) -- In-process OpenSSL ALPN `h2` (`Grpc.Native.Tls` → `Tls.connectH2` / `serveH2`) +- Plain TCP Async (`H2.Client.connectH2cAsync`, `H2.Server.listenH2cAsync`) — **no** `.block` on the Async chain +- Plain TCP sync adapters (`connectH2c`, `listenH2c`) — `.block` at the edge +- In-process OpenSSL ALPN `h2` (`Grpc.Native.Tls` → `Tls.connectH2` / `serveH2`) — **blocking FFI** in v1.2.0; see [async-io.md](async-io.md) - Optional sidecar via `LEAN_GRPC_TLS_PROXY` (client dials h2c to a local terminator) ## Credentials @@ -103,7 +104,8 @@ These are intentionally lightweight compared to full OpenTelemetry SDKs. ## Process / concurrency model - Networking uses Lean 4 `Std.Async` (libuv-backed). Prefer one event-loop owner per process for listen/accept. -- Interop loopbacks that need a second peer often **spawn a child executable** (e.g. `trailersLoopback`) instead of sharing one UV loop across conflicting `.block` tasks. +- **v1.2.0+:** `serveH2cAsync` schedules connections with `Std.Async.background` so concurrent h2c clients can share one UV loop. See [async-io.md](async-io.md). +- Interop loopbacks that need a second peer may still **spawn a child executable** (e.g. `trailersLoopback`) when mixing sync `.block` adapters with tasks. - CI runs multi-language peers as separate processes (Go / Python / h2spec). ## Native code diff --git a/docs/async-io.md b/docs/async-io.md new file mode 100644 index 0000000..9994f25 --- /dev/null +++ b/docs/async-io.md @@ -0,0 +1,52 @@ +# Async IO model (v1.2.0) + +lean-grpc uses Lean 4 `Std.Async.TCP` (libuv-backed) for sockets. Through +**v1.1.x** the public API collapsed every accept/connect/send/recv with +`.block`, so the library behaved as **synchronous `IO`** despite Async socket +constructors. That overstated “runs on Std.Async” — see +[GitHub #6](https://github.com/RileyBetts/lean-grpc/issues/6). + +## v1.2.0 model + +| Layer | Behavior | +|---|---| +| `H2.AsyncByteTransport` / `tcpTransportAsync` | Native `Async` send/recv — **no** `.block` | +| `H2.listenH2cAsync` / `connectH2cAsync` / `awaitResponseAsync` | Async accept/dial/RPC wait | +| `Grpc.Server.serveH2cAsync` / `Channel.unaryAsync` | Additive Async gRPC entry points | +| `ByteTransport` / `serveH2c` / `unary` / `connectH2c` | **Sync adapters** — `.block` at the edge only | +| In-process TLS (`SSL_read` / `SSL_write`) | Still **blocking FFI**; wrapped via `AsyncByteTransport.ofBlocking` | + +```text + App (Async) ──► serveH2cAsync / unaryAsync + │ + ▼ + AsyncByteTransport (TCP) ← zero .block on h2c hot path + │ + ▼ + Std.Async.TCP + UV loop + + App (IO) ──► serveH2c / unary ──► (asyncPath).block ← explicit adapter +``` + +## UV loop ownership + +Prefer **one** event-loop owner per process. Nested `.block` inside +`IO.asTask` workers on the same loop is fragile (historical +`TrailersLoopback` note). The Async listen path schedules connections with +`Std.Async.background` instead. + +## TLS caveat (v1.2.0) + +OpenSSL in `native/tls_ffi.c` is blocking. `serveTls` / `connectH2` remain +correct and tested, but they are **not** end-to-end async. Do not claim TLS is +async until a nonblocking BIO or off-loop worker pool lands (follow-up issue). + +## Migration + +| Need | Use | +|---|---| +| Compose without blocking the loop (h2c) | `*Async` APIs | +| Existing code / lean-compliance | Keep `IO` APIs (adapters) | +| mTLS AuthN | `serveTls` + `registerWithContext` (sync FFI underneath) | + +Pin: `@ "v1.2.0"`. diff --git a/docs/getting-started.md b/docs/getting-started.md index 3c65df1..9a8f760 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -4,7 +4,7 @@ This guide gets a Lean 4 unary gRPC server and client running on h2c with **type **Requirements:** Lean 4.32+ (see `lean-toolchain`), OpenSSL headers for the `Grpc` library (`libssl-dev` or `./scripts/fetch-openssl-headers.sh`). Peer gzip optionally uses a zlib helper (`./scripts/build_native.sh`) or system `gzip`. -Package version: **1.1.0** (`lakefile.lean` / `Grpc.version`). +Package version: **1.2.0** (`lakefile.lean` / `Grpc.version`). ## Build the repo @@ -151,12 +151,12 @@ In your `lakefile.lean`: ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.1.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" ``` Then `import Grpc` (and `Proto` if you use the bundled codecs). Public libs: `Bytes`, `Hpack`, `H2`, `Proto`, `Grpc`. Linking OpenSSL (`-lssl -lcrypto`) is pulled in via the `Grpc` Lake library (not via Bytes/Hpack/H2 alone). For peer gzip against foreign stacks, set `LEAN_GRPC_ZLIB_HELPER` after `./scripts/build_native.sh`. Full packaging notes: [packaging.md](packaging.md). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). -After [Reservoir](https://reservoir.lean-lang.org/) indexes the package you can use `require «lean-grpc»` without a git URL. Pin a commit SHA if the `v1.1.0` tag is not yet on the remote you use. +After [Reservoir](https://reservoir.lean-lang.org/) indexes the package you can use `require «lean-grpc»` without a git URL. Pin a commit SHA if the `v1.2.0` tag is not yet on the remote you use. ## Next steps diff --git a/docs/packaging.md b/docs/packaging.md index 82a60b9..510c0ec 100644 --- a/docs/packaging.md +++ b/docs/packaging.md @@ -65,7 +65,7 @@ lake build ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.1.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" ``` Then `import Grpc` (and `Proto` if using bundled codecs). @@ -76,7 +76,7 @@ Then `import Grpc` (and `Proto` if using bundled codecs). require «lean-grpc» -- once listed on reservoir.lean-lang.org ``` -Package version in-tree is **1.1.0** (`lakefile.lean`, `Grpc.version`). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). Creating the git tag on origin is a separate maintainer step (below). +Package version in-tree is **1.2.0** (`lakefile.lean`, `Grpc.version`). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). Creating the git tag on origin is a separate maintainer step (below). ## Maintainer release checklist @@ -89,8 +89,8 @@ Documented only — do **not** automate tagging from CI or agent runs. 5. **Manual version tagging (maintainer only):** on `main`, when you choose to publish: ```bash - git tag -a v1.1.0 -m "lean-grpc 1.1.0" - git push origin v1.1.0 + git tag -a v1.2.0 -m "lean-grpc 1.2.0" + git push origin v1.2.0 ``` Agents and CI must not create or push tags. diff --git a/lakefile.lean b/lakefile.lean index 766c683..107f53a 100644 --- a/lakefile.lean +++ b/lakefile.lean @@ -7,9 +7,9 @@ open Lake DSL open System package «lean-grpc» where - version := v!"1.1.0" + version := v!"1.2.0" keywords := #["grpc", "http2", "hpack", "protobuf", "networking"] - description := "Pure Lean 4 gRPC stack (HTTP/2 + HPACK + gRPC) on Std.Async" + description := "Pure Lean 4 gRPC stack (HTTP/2 + HPACK + gRPC); Std.Async.TCP with native Async h2c APIs" homepage := "https://rileybetts.ai/oss/lean-grpc" license := "Apache-2.0" licenseFiles := #["LICENSE", "NOTICE"] @@ -67,6 +67,9 @@ lean_exe trailersServer where lean_exe trailersLoopback where root := `Tests.TrailersLoopback +lean_exe asyncH2cLoopback where + root := `Tests.AsyncH2cLoopback + lean_exe tlsServer where root := `Tests.TlsServer From 686e90ae369980c8780383fb98b37cf0ff353f36 Mon Sep 17 00:00:00 2001 From: Robert Betts Date: Tue, 11 Aug 2026 23:06:48 +0100 Subject: [PATCH 2/2] Add off-loop TLS so OpenSSL does not stall the UV loop (#10). Expose connectH2Async/serveTlsAsync with dedicated-thread SSL I/O, keep sync IO paths for safe composition, and ship asyncTlsLoopback + v1.3.0 docs. Co-authored-by: Cursor --- .github/workflows/ci.yml | 3 +- CHANGELOG.md | 14 +++- Grpc.lean | 2 +- Grpc/Metadata.lean | 2 +- Grpc/Server.lean | 7 +- Grpc/Tls.lean | 77 +++++++++++++++++--- H2/Client.lean | 10 ++- H2/Transport.lean | 15 +++- README.md | 6 +- ROADMAP.md | 24 +++++-- Tests/AsyncTlsLoopback.lean | 139 ++++++++++++++++++++++++++++++++++++ docs/README.md | 2 +- docs/api-reference.md | 4 +- docs/architecture.md | 5 +- docs/async-io.md | 49 +++++++------ docs/getting-started.md | 6 +- docs/packaging.md | 8 +-- lakefile.lean | 7 +- 18 files changed, 320 insertions(+), 60 deletions(-) create mode 100644 Tests/AsyncTlsLoopback.lean diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b37a057..d639ba4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -22,7 +22,7 @@ jobs: chmod +x scripts/fetch-openssl-headers.sh ./scripts/fetch-openssl-headers.sh - name: Build libraries and tests - run: lake build Proofs bytesTests hpackTests h2Tests grpcTests securityTests trailersServer trailersLoopback asyncH2cLoopback tlsServer tlsLoopback helloworldServer helloworldClient interopServer interopClient protocGenLean4Grpc benchUnary benchSoak opsSmoke h2specServer routeGuideServer routeGuideClient adcSmoke fakeAdsServer + run: lake build Proofs bytesTests hpackTests h2Tests grpcTests securityTests trailersServer trailersLoopback asyncH2cLoopback asyncTlsLoopback tlsServer tlsLoopback helloworldServer helloworldClient interopServer interopClient protocGenLean4Grpc benchUnary benchSoak opsSmoke h2specServer routeGuideServer routeGuideClient adcSmoke fakeAdsServer - name: Formal proofs (compile-time) run: lake build Proofs - name: Unit tests @@ -34,6 +34,7 @@ jobs: ./.lake/build/bin/securityTests ./.lake/build/bin/trailersLoopback ./.lake/build/bin/asyncH2cLoopback + ./.lake/build/bin/asyncTlsLoopback - name: Native ASAN / UBSan security harness run: | chmod +x scripts/security-asan.sh diff --git a/CHANGELOG.md b/CHANGELOG.md index 501bc89..4edd73e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,18 @@ # Changelog -All notable changes to lean-grpc are documented here. The package version is the Lake/`Grpc.version` semver (currently **1.2.0**). Git tags such as `v1.2.0` are created manually by maintainers when publishing. +All notable changes to lean-grpc are documented here. The package version is the Lake/`Grpc.version` semver (currently **1.3.0**). Git tags such as `v1.3.0` are created manually by maintainers when publishing. + +## [1.3.0] — 2026-08-11 + +Off-loop TLS so OpenSSL does not stall the UV loop ([#10](https://github.com/RileyBetts/lean-grpc/issues/10)). + +- **`H2.runOffLoop` / `AsyncByteTransport.ofBlockingOffLoop`:** blocking `IO` (including `SSL_*`) runs on `IO.asTask .dedicated`, then resumes Async — UV stays free. +- **Async TLS APIs:** `Tls.connectH2Async`, `Tls.serveH2Async`, `Server.serveTlsAsync`. Sync `connectH2` / `serveH2` / `serveTls` remain pure blocking `IO` (safe under `IO.asTask`; no nested `Async.block`). +- **Honesty:** still **blocking** OpenSSL FFI — documented as off-loop, not nonblocking BIO. +- **Tests:** `asyncTlsLoopback` — concurrent off-loop TLS clients + Async h2c on one UV loop; in-process `serveTlsAsync` smoke. +- **Docs:** [docs/async-io.md](docs/async-io.md) updated for the v1.3.0 TLS model. + +Migration: additive — bump pin to `v1.3.0`. Prefer `*Async` TLS under Async composition; keep IO APIs for lean-compliance. ## [1.2.0] — 2026-08-11 diff --git a/Grpc.lean b/Grpc.lean index 91c88b5..a7dd9d8 100644 --- a/Grpc.lean +++ b/Grpc.lean @@ -39,5 +39,5 @@ import Grpc.Grpclb import Grpc.Interceptor namespace Grpc -def version : String := "1.2.0" +def version : String := "1.3.0" end Grpc diff --git a/Grpc/Metadata.lean b/Grpc/Metadata.lean index 08b591f..f57cf4e 100644 --- a/Grpc/Metadata.lean +++ b/Grpc/Metadata.lean @@ -159,7 +159,7 @@ def methodGet : Hpack.HeaderField := ⟨ascii ":method", ascii "GET"⟩ def schemeHttp : Hpack.HeaderField := ⟨ascii ":scheme", ascii "http"⟩ def schemeHttps : Hpack.HeaderField := ⟨ascii ":scheme", ascii "https"⟩ def teTrailers : Hpack.HeaderField := ⟨ascii "te", ascii "trailers"⟩ -def userAgent (version : String := "1.2.0") : Hpack.HeaderField := +def userAgent (version : String := "1.3.0") : Hpack.HeaderField := ⟨ascii "user-agent", ascii s!"grpc-lean/{version}"⟩ def http415 : Array Hpack.HeaderField := diff --git a/Grpc/Server.lean b/Grpc/Server.lean index 4d3d5d1..7f11168 100644 --- a/Grpc/Server.lean +++ b/Grpc/Server.lean @@ -315,10 +315,15 @@ def serveH2cAsync (s : Server) (cfg : H2.ServerConfig := {}) : Async Unit := /-- Serve with in-process TLS+ALPN `h2` (see `Grpc.Tls.serveH2` for the plaintext/mTLS decision based on `tlsCfg`). Peer identity is extracted per accepted connection and threaded into unary context handlers. - TLS byte IO remains blocking OpenSSL FFI in v1.2.0 — see docs/async-io.md. -/ + OpenSSL stays blocking FFI but runs **off** the UV loop in v1.3.0 — see docs/async-io.md. -/ def serveTls (s : Server) (tlsCfg : Tls.Config) (cfg : H2.ServerConfig := {}) : IO Unit := let mtlsRequired := tlsCfg.clientCaPath.isSome Tls.serveH2 tlsCfg cfg fun peerId => handlerFor s peerId mtlsRequired +/-- Async TLS serve: accept/handshake/`SSL_*` off-loop; connections on `background`. -/ +def serveTlsAsync (s : Server) (tlsCfg : Tls.Config) (cfg : H2.ServerConfig := {}) : Async Unit := + let mtlsRequired := tlsCfg.clientCaPath.isSome + Tls.serveH2Async tlsCfg cfg fun peerId => handlerFor s peerId mtlsRequired + end Server end Grpc diff --git a/Grpc/Tls.lean b/Grpc/Tls.lean index ad40a6c..c6979da 100644 --- a/Grpc/Tls.lean +++ b/Grpc/Tls.lean @@ -6,6 +6,9 @@ import H2 import Grpc.Resolver import Grpc.Native.Tls import Grpc.PeerIdentity +import Std.Async + +open Std.Async namespace Grpc.Tls @@ -33,10 +36,21 @@ structure Config where def envoySidecarNotes : String := "Optional: terminate TLS at Envoy/Caddy with alpn_protocols: [h2], cluster to 127.0.0.1:." -/-- Connect with TLS+ALPN `h2`. +/-- In-process OpenSSL dial (blocking). Used by sync `connectH2`. -/ +private def dialInProcess (host : String) (port : UInt16) (cfg : Config) : IO H2.ClientConn := do + if cfg.insecureSkipVerify then + IO.eprintln "WARN: Tls.Config.insecureSkipVerify=true → peer certificate not verified" + let ca := (cfg.caPath.map (·.toString)).getD "" + let sni := cfg.serverName.getD host + let clientCert := (cfg.certPath.map (·.toString)).getD "" + let clientKey := (cfg.keyPath.map (·.toString)).getD "" + let conn ← Grpc.Native.Tls.dial host port ca sni clientCert clientKey cfg.insecureSkipVerify + H2.Client.connectTransport (Grpc.Native.Tls.transport conn) + +/-- Connect with TLS+ALPN `h2` (blocking `IO`; safe under `IO.asTask`). Priority: 1. `LEAN_GRPC_TLS_PROXY` — h2c to local sidecar/proxy (legacy). - 2. In-process OpenSSL (system CA / `caPath` + hostname verify unless `insecureSkipVerify`). + 2. In-process OpenSSL. 3. `LEAN_GRPC_TLS_INSECURE_FALLBACK=1` — plain h2c (dev only). -/ def connectH2 (host : String) (port : UInt16) (cfg : Config := {}) : IO H2.ClientConn := do if let some proxy := ← IO.getEnv "LEAN_GRPC_TLS_PROXY" then @@ -48,21 +62,66 @@ def connectH2 (host : String) (port : UInt16) (cfg : Config := {}) : IO H2.Clien IO.eprintln "WARN: LEAN_GRPC_TLS_INSECURE_FALLBACK=1 → h2c (no TLS)" H2.Client.connectH2c host port | _ => + dialInProcess host port cfg + +/-- Connect with TLS+ALPN `h2` under `Std.Async` (off-loop OpenSSL). + Blocking dial/handshake/preface run on dedicated threads — not nonblocking BIO. -/ +def connectH2Async (host : String) (port : UInt16) (cfg : Config := {}) : Async H2.ClientConn := do + if let some proxy := ← IO.getEnv "LEAN_GRPC_TLS_PROXY" then + let addr ← IO.ofExcept (Resolver.parseTarget proxy) + IO.eprintln s!"Tls.connectH2Async via LEAN_GRPC_TLS_PROXY={proxy} (ALPN terminated externally)" + return (← H2.Client.connectH2cAsync addr.host addr.port) + match ← IO.getEnv "LEAN_GRPC_TLS_INSECURE_FALLBACK" with + | some "1" => + IO.eprintln "WARN: LEAN_GRPC_TLS_INSECURE_FALLBACK=1 → h2c (no TLS)" + H2.Client.connectH2cAsync host port + | _ => do if cfg.insecureSkipVerify then IO.eprintln "WARN: Tls.Config.insecureSkipVerify=true → peer certificate not verified" let ca := (cfg.caPath.map (·.toString)).getD "" let sni := cfg.serverName.getD host let clientCert := (cfg.certPath.map (·.toString)).getD "" let clientKey := (cfg.keyPath.map (·.toString)).getD "" - let conn ← Grpc.Native.Tls.dial host port ca sni clientCert clientKey cfg.insecureSkipVerify - H2.Client.connectTransport (Grpc.Native.Tls.transport conn) + let conn ← H2.runOffLoop + (Grpc.Native.Tls.dial host port ca sni clientCert clientKey cfg.insecureSkipVerify) + H2.Client.connectTransportOffLoopAsync (Grpc.Native.Tls.transport conn) -/-- Serve with in-process TLS+ALPN when cert/key are set; otherwise h2c (+ optional sidecar). - TLS listen binds loopback only (see native `INADDR_LOOPBACK`). +/-- Serve with in-process TLS+ALPN under `Std.Async`. Accept/handshake run off-loop; + each connection is a dedicated background fiber with blocking SSL I/O. -/ +partial def serveH2Async (cfg : Config) (h2cfg : H2.ServerConfig) + (mkHandler : Option PeerIdentity → H2.StreamHandler) : Async Unit := do + match cfg.certPath, cfg.keyPath with + | some cert, some key => + let clientCa := (cfg.clientCaPath.map (·.toString)).getD "" + let listener ← H2.runOffLoop + (Grpc.Native.Tls.listen h2cfg.port cert.toString key.toString clientCa) + let mtlsNote := if cfg.clientCaPath.isSome then " (mTLS: client cert required)" else "" + IO.println s!"H2 TLS+ALPN listening on 127.0.0.1:{h2cfg.port}{mtlsNote} (Async off-loop)" + while true do + let acceptResult ← + try + pure (Sum.inl (← H2.runOffLoop (Grpc.Native.Tls.accept listener))) + catch e => + pure (Sum.inr e) + match acceptResult with + | .inr e => + IO.eprintln s!"tls accept error: {e}" + | .inl conn => + let peerId ← H2.runOffLoop (Grpc.Native.Tls.peerIdentity? conn) + background (prio := Task.Priority.dedicated) do + try + liftM (H2.serveTransport (Grpc.Native.Tls.transport conn) (mkHandler peerId)) + catch e => + IO.eprintln s!"tls conn error: {e}" + | _, _ => + match ← IO.getEnv "LEAN_GRPC_TLS_INSECURE_FALLBACK" with + | some "1" => H2.Server.listenAsync h2cfg (mkHandler none) + | _ => + IO.eprintln s!"Tls.serveH2Async: serving h2c on {h2cfg.host}:{h2cfg.port}; set certPath/keyPath for in-process TLS, or use sidecar. {envoySidecarNotes}" + H2.Server.listenAsync h2cfg (mkHandler none) - `mkHandler` receives the verified peer identity for each accepted connection - (`none` when the client presented no certificate). Failed accepts/handshakes - are logged and the listen loop continues. -/ +/-- Serve with in-process TLS+ALPN (blocking `IO` accept loop; per-conn `IO.asTask`). + Does not go through `Async.block` — safe when the process is dedicated to serving. -/ partial def serveH2 (cfg : Config) (h2cfg : H2.ServerConfig) (mkHandler : Option PeerIdentity → H2.StreamHandler) : IO Unit := do match cfg.certPath, cfg.keyPath with diff --git a/H2/Client.lean b/H2/Client.lean index 6b60382..4b099a4 100644 --- a/H2/Client.lean +++ b/H2/Client.lean @@ -89,10 +89,18 @@ def connectTransportAsync (t : AsyncByteTransport) : Async ClientConn := do readBuf.set buf return { transport := t, state, readBuf } -/-- Sync/TLS entry: wrap blocking `ByteTransport` then run async preface exchange. -/ +/-- Sync/TLS entry: wrap blocking `ByteTransport` then run async preface via `.block`. + Prefer calling this from an off-loop dedicated task (`H2.runOffLoop`) so OpenSSL + does not stall the UV thread. -/ def connectTransport (t : ByteTransport) : IO ClientConn := (connectTransportAsync (.ofBlocking t)).block +/-- Dial/preface on a dedicated thread (blocking OpenSSL), then return a `ClientConn` + whose later Async send/recv use `ofBlockingOffLoop` (issue #10). -/ +def connectTransportOffLoopAsync (t : ByteTransport) : Async ClientConn := do + let c ← runOffLoop (connectTransport t) + pure { c with transport := .ofBlockingOffLoop t } + /-- Dial h2c without `.block` on connect/send/recv. -/ def connectH2cAsync (host : String) (port : UInt16) : Async ClientConn := do let sock ← TCP.Socket.Client.mk diff --git a/H2/Transport.lean b/H2/Transport.lean index d6b23ee..3747fdd 100644 --- a/H2/Transport.lean +++ b/H2/Transport.lean @@ -28,12 +28,25 @@ def tcpTransportAsync (sock : TCP.Socket.Client) : AsyncByteTransport where close := pure () /-- Lift blocking `IO` send/recv into `Async` (still blocks the UV thread while - the IO runs — used for OpenSSL FFI transports in v1.2.0). -/ + the IO runs). Prefer `ofBlockingOffLoop` for OpenSSL / other blocking FFI. -/ def AsyncByteTransport.ofBlocking (t : ByteTransport) : AsyncByteTransport where send := fun b => liftM (t.send b) recv? := fun n => liftM (t.recv? n) close := liftM t.close +/-- Run blocking `IO` on a dedicated thread, then resume the Async waiter. + Keeps `SSL_read` / `SSL_write` / similar FFI off the UV loop (issue #10). -/ +def runOffLoop (act : IO α) (prio := Task.Priority.dedicated) : Async α := do + let t ← IO.asTask act prio + Async.ofAsyncTask t + +/-- Lift blocking `IO` send/recv into `Async` without stalling the UV loop: + each op runs on a dedicated task (off-loop OpenSSL model for v1.3.0). -/ +def AsyncByteTransport.ofBlockingOffLoop (t : ByteTransport) : AsyncByteTransport where + send := fun b => runOffLoop (t.send b) + recv? := fun n => runOffLoop (t.recv? n) + close := runOffLoop t.close + /-- Sync facade: each op `.block`s the underlying async transport. -/ def ByteTransport.ofAsync (t : AsyncByteTransport) : ByteTransport where send := fun b => (t.send b).block diff --git a/README.md b/README.md index 6ff150f..985574c 100644 --- a/README.md +++ b/README.md @@ -6,9 +6,9 @@ General-purpose **Lean 4 gRPC library**: HPACK + HTTP/2 + gRPC framing over `Std.Async.TCP`. -Standalone Lake package (**1.2.0**). Consumers depend via git tag or, after indexing, [Reservoir](https://reservoir.lean-lang.org/). +Standalone Lake package (**1.3.0**). Consumers depend via git tag or, after indexing, [Reservoir](https://reservoir.lean-lang.org/). -**Async model:** sockets are `Std.Async.TCP`. Through v1.1.x the public API was blocking `IO` via `.block`. **v1.2.0** adds a native Async h2c path (`serveH2cAsync` / `unaryAsync`) with **zero** `.block` on accept/connect/send/recv; existing IO APIs remain as explicit sync adapters. In-process TLS is still blocking OpenSSL FFI — see [docs/async-io.md](docs/async-io.md). +**Async model:** sockets are `Std.Async.TCP`. **v1.2.0** added a native Async h2c path (`serveH2cAsync` / `unaryAsync`) with **zero** `.block` on accept/connect/send/recv. **v1.3.0** runs blocking OpenSSL on dedicated threads (`connectH2Async` / `serveTlsAsync` / `ofBlockingOffLoop`) so TLS does not stall the UV loop — still not nonblocking BIO; see [docs/async-io.md](docs/async-io.md). **Docs:** [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc) (curated) · [docs/](docs/README.md) (full in-repo index) @@ -22,7 +22,7 @@ In your `lakefile.lean`: ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.3.0" ``` Then `import Grpc`. After Reservoir lists the package you can use `require «lean-grpc»` without a git URL. Packaging details and the maintainer release checklist: [docs/packaging.md](docs/packaging.md). diff --git a/ROADMAP.md b/ROADMAP.md index 2346135..63f7582 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -1,6 +1,6 @@ # Roadmap -lean-grpc **v1.2.0** is the current package tip (Lake / `Grpc.version`): an interop-tested Lean 4 gRPC stack with native Async h2c APIs (honest Std.Async model), a CI-gated **`Proofs`** library for selected pure codecs, plus additive mTLS peer-identity / request-context APIs for enterprise AuthN. It is **not** a machine-checked end-to-end PROTOCOL-HTTP2 / TLS / session proof. +lean-grpc **v1.3.0** is the current package tip (Lake / `Grpc.version`): an interop-tested Lean 4 gRPC stack with native Async h2c APIs, **off-loop** in-process TLS (blocking OpenSSL on dedicated threads), a CI-gated **`Proofs`** library for selected pure codecs, plus additive mTLS peer-identity / request-context APIs for enterprise AuthN. It is **not** a machine-checked end-to-end PROTOCOL-HTTP2 / TLS / session proof. This document records what shipped, what is still open, and the next proof/hardening tranches. @@ -25,13 +25,23 @@ First public packaging baseline (interop + Lake/Reservoir layout). Formal Lean p | Dial / LB / retry, health, reflection, channelz | Present; ops demos gated | | ADC / xDS ADS | Mock / FakeAds CI; live Google paths allowlisted | +## Shipped — v1.3.0 (TLS off-loop) + +Blocking OpenSSL runs off the UV loop ([#10](https://github.com/RileyBetts/lean-grpc/issues/10), [docs/async-io.md](docs/async-io.md)). + +| Included | Deferred | +|---|---| +| `runOffLoop` / `ofBlockingOffLoop` | Nonblocking BIO / `SSL_ERROR_WANT_*` + UV readiness | +| `connectH2Async` / `serveTlsAsync` | Streaming `*Async` surface beyond unary | +| `asyncTlsLoopback` CI (TLS + h2c concurrent) | Full Concurrent retry/hedge under Async TLS | + ## Shipped — v1.2.0 (Async honesty) Native Async h2c + docs that match the implementation ([#6](https://github.com/RileyBetts/lean-grpc/issues/6), [docs/async-io.md](docs/async-io.md)). | Included | Deferred | |---|---| -| `AsyncByteTransport` / `*Async` h2c serve/dial/unary | Nonblocking TLS BIO / off-loop OpenSSL | +| `AsyncByteTransport` / `*Async` h2c serve/dial/unary | Streaming `*Async` beyond unary; Concurrent retry/hedge under Async | | Sync IO APIs as `.block` adapters | Streaming `*Async` surface beyond unary | | Concurrent `asyncH2cLoopback` CI | Full Concurrent retry/hedge under Async | @@ -109,11 +119,11 @@ Close remaining audit follow-ups so later minors are not “proved but soft”: ## Milestone sketch ```text -v0.5.0 ──► v1.0.0 ──► v1.1.0 ──► v1.2.0 ──► v1.3.x - shipped interop + mTLS peer identity + Async h2c honesty ConnState + - selected Proofs ServerCallContext (#6) + sync adapters general frame/msg - (CI) + provenance (unary IAM) roundtrips / - Huffman trie +v0.5.0 ──► v1.0.0 ──► v1.1.0 ──► v1.2.0 ──► v1.3.0 ──► later 1.x + shipped interop + mTLS peer identity + Async h2c honesty TLS off-loop (#10) ConnState + + selected Proofs ServerCallContext (#6) + sync adapters + serveTlsAsync general frame/msg + (CI) + provenance (unary IAM) roundtrips / + Huffman trie ``` Dates are intentionally omitted; order matters more than calendar. diff --git a/Tests/AsyncTlsLoopback.lean b/Tests/AsyncTlsLoopback.lean new file mode 100644 index 0000000..f8fb733 --- /dev/null +++ b/Tests/AsyncTlsLoopback.lean @@ -0,0 +1,139 @@ +/- +Copyright © 2026, Riley Betts Ltd (rileybetts.ai) +Released under Apache 2.0 license as described in the file LICENSE. +-/ +import Grpc +import H2 +import Proto +import Std.Async +import Std.Async.Timer + +open Std.Async + +/-- Concurrent off-loop TLS clients + Async h2c on one UV loop (issue #10). + TLS server runs in a subprocess (`tlsServer`); clients use `H2.runOffLoop` so + OpenSSL cannot stall the UV loop while h2c clients progress. + Skips if `openssl` is unavailable. -/ + +private def run (cmd : String) (args : Array String) : IO Unit := do + let out ← IO.Process.output { cmd, args } + if out.exitCode != 0 then + throw (IO.userError s!"{cmd} {args} failed: {out.stderr}") + +private def haveOpenssl : IO Bool := do + try + let out ← IO.Process.output { cmd := "openssl", args := #["version"] } + return out.exitCode == 0 + catch _ => return false + +private def genCerts (dir : System.FilePath) : IO Unit := do + IO.FS.createDirAll dir + let d := dir.toString + run "openssl" #["genrsa", "-out", s!"{d}/ca.key", "2048"] + run "openssl" #["req", "-x509", "-new", "-nodes", "-key", s!"{d}/ca.key", "-sha256", + "-days", "3650", "-out", s!"{d}/ca.crt", "-subj", "/CN=LeanGrpcAsyncTlsCA"] + run "openssl" #["genrsa", "-out", s!"{d}/server.key", "2048"] + run "openssl" #["req", "-new", "-key", s!"{d}/server.key", "-out", s!"{d}/server.csr", + "-subj", "/CN=127.0.0.1"] + run "openssl" #["x509", "-req", "-in", s!"{d}/server.csr", "-CA", s!"{d}/ca.crt", + "-CAkey", s!"{d}/ca.key", "-CAcreateserial", "-out", s!"{d}/server.crt", + "-days", "825", "-sha256"] + +private def spawnTlsServer (binDir dir : System.FilePath) (port : UInt16) : + IO (IO.Process.Child {}) := do + let d := dir.toString + IO.Process.spawn { + cmd := (binDir / "tlsServer").toString + env := #[ + ("GRPC_PORT", some (toString port.toNat)), + ("TLS_CERT", some s!"{d}/server.crt"), + ("TLS_KEY", some s!"{d}/server.key") + ] + } + +/-- Blocking EmptyCall with connect retries (server spawn race). -/ +private partial def tlsEmptyIO (port : UInt16) (ca : System.FilePath) (tag : String) + (attempts : Nat := 12) : IO String := do + let cfg : Grpc.Tls.Config := { + caPath := some ca + serverName := some "127.0.0.1" + } + try + let c ← Grpc.Tls.connectH2 "127.0.0.1" port cfg + let res ← Grpc.Client.unaryCall c "grpc.testing.TestService" "EmptyCall" "127.0.0.1" + ByteArray.empty (useHttps := true) + if res.status.code != .ok then + throw (IO.userError s!"tls empty status {res.status.code.toUInt32} tag={tag}") + match String.fromUTF8? res.message with + | some s => pure s!"{tag}:{s}" + | none => throw (IO.userError "non-utf8") + catch e => + if attempts ≤ 1 then throw e + IO.sleep 250 + tlsEmptyIO port ca tag (attempts - 1) + +/-- Off-loop TLS EmptyCall client. -/ +private def tlsEmptyClient (port : UInt16) (ca : System.FilePath) (tag : String) : Async String := + H2.runOffLoop (tlsEmptyIO port ca tag) + +private def h2cClient (port : UInt16) (name : String) : Async String := do + let c ← H2.Client.connectH2cAsync "127.0.0.1" port + let req := Proto.HelloRequest.encode { name } + let res ← Grpc.Client.unaryCallAsync c "helloworld.Greeter" "SayHello" "127.0.0.1" req + if res.status.code != .ok then + throw (IO.userError s!"h2c status {res.status.code.toUInt32}") + let reply ← liftM (IO.ofExcept (Proto.HelloReply.decode res.message)) + return reply.message + +def main : IO Unit := do + if !(← haveOpenssl) then + IO.println "asyncTlsLoopback SKIPPED (no openssl)" + return + let dir := (← IO.currentDir) / ".lake" / "build" / "async-tls-loopback-certs" + genCerts dir + let binDir ← IO.appDir + let tlsPort : UInt16 := 50071 + let h2cPort : UInt16 := 50072 + let mut s := Grpc.Server.empty + s := Grpc.Server.register s "helloworld.Greeter" "SayHello" fun reqBytes => do + let req ← IO.ofExcept (Proto.HelloRequest.decode reqBytes) + let reply : Proto.HelloReply := { message := s!"Hello, {req.name}" } + return (Proto.HelloReply.encode reply, Grpc.Status.ok) + let srv ← spawnTlsServer binDir dir tlsPort + try + (do + background (prio := Task.Priority.dedicated) do + Grpc.Server.serveH2cAsync s { host := "127.0.0.1", port := h2cPort } + sleep (Std.Time.Millisecond.Offset.ofNat 600) + let tlsTags : Array String := #["T1", "T2", "T3", "T4"] + let h2cNames : Array String := #["H1", "H2", "H3", "H4"] + let tlsJobs := tlsTags.map (fun t => tlsEmptyClient tlsPort (dir / "ca.crt") t) + let h2cJobs := h2cNames.map (fun n => h2cClient h2cPort n) + let msgs ← Async.concurrentlyAll (tlsJobs ++ h2cJobs) + if msgs.size != tlsTags.size + h2cNames.size then + throw (IO.userError s!"expected {tlsTags.size + h2cNames.size} replies, got {msgs.size}") + for msg in msgs do + if !("T".isPrefixOf msg || "Hello".isPrefixOf msg) then + throw (IO.userError msg) + -- Phase 2: in-process `serveTlsAsync` + off-loop client (same UV loop as h2c). + let asyncTlsPort : UInt16 := 50073 + let tlsCfg : Grpc.Tls.Config := { + certPath := some (dir / "server.crt") + keyPath := some (dir / "server.key") + } + let mut s2 := Grpc.Server.empty + s2 := Grpc.Server.registerWithContext s2 "grpc.testing.TestService" "EmptyCall" + fun _ _ => pure (ByteArray.empty, Grpc.Status.ok) + background (prio := Task.Priority.dedicated) do + Grpc.Server.serveTlsAsync s2 tlsCfg { host := "127.0.0.1", port := asyncTlsPort } + sleep (Std.Time.Millisecond.Offset.ofNat 400) + let body ← tlsEmptyClient asyncTlsPort (dir / "ca.crt") "ASYNC" + if !("ASYNC:".isPrefixOf body) then + throw (IO.userError body) + IO.println "asyncTlsLoopback OK" + -- Background accept loops (h2c + serveTlsAsync) keep dedicated threads alive; + -- force a clean exit so Lake/CI do not hang in task_manager finalization. + IO.Process.exit 0 + ).block + finally + srv.kill diff --git a/docs/README.md b/docs/README.md index 798533e..8770e10 100644 --- a/docs/README.md +++ b/docs/README.md @@ -22,7 +22,7 @@ User and contributor docs for the Lean 4 gRPC stack. | [Formal proofs](proofs.md) | Compile-time `Proofs` library for pure codecs / maps | | [TLS / Envoy](tls-envoy.md) | In-process OpenSSL and optional sidecars | | [CHANGELOG](../CHANGELOG.md) | Package version history | -| [ROADMAP](../ROADMAP.md) | What v1.2.0 shipped vs open proof/hardening follow-ups | +| [ROADMAP](../ROADMAP.md) | What v1.3.0 shipped vs open proof/hardening follow-ups | | [CONTRIBUTING](../CONTRIBUTING.md) | Dev setup, tests, PR expectations | | [SECURITY](../SECURITY.md) | Vulnerability reporting | | [Code of Conduct](../CODE_OF_CONDUCT.md) | Community standards | diff --git a/docs/api-reference.md b/docs/api-reference.md index 88a1f2c..c8603c4 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -1,6 +1,6 @@ # API reference -Lean module catalogue for consumers. Signatures are summarized; see source under `Grpc/`, `H2/`, `Proto/` for full definitions. Version string: `Grpc.version` (currently `1.2.0`). Async vs sync IO: [async-io.md](async-io.md). +Lean module catalogue for consumers. Signatures are summarized; see source under `Grpc/`, `H2/`, `Proto/` for full definitions. Version string: `Grpc.version` (currently `1.3.0`). Async vs sync IO: [async-io.md](async-io.md). Import umbrella: `import Grpc` (pulls status, channel, server, credentials, TLS, xDS, ops, etc.). Add `import Proto` for bundled message codecs. @@ -122,7 +122,7 @@ Verified mTLS peer certificate identity (OpenSSL; subject DN is **RFC 2253**): | `registerServerStream` / `registerClientStream` / `registerBidi` | Streaming (raw bytes; context deferred) | | `registerServerStreamTyped` / `registerClientStreamTyped` / `registerBidiTyped` | Streaming with decode/encode adapters | | `serveH2c` / `serveH2cAsync` | Listen h2c (`peerIdentity = none`); Async path has no `.block` on accept/send/recv | -| `serveTls` | Listen TLS+ALPN; per-connection peer identity → context handlers (**blocking** OpenSSL FFI in v1.2.0) | +| `serveTls` / `serveTlsAsync` | Listen TLS+ALPN; per-connection peer identity → context handlers (OpenSSL **off-loop** under Async in v1.3.0) | | `maxMsgSize` | Inbound limit | Bad `content-type` → HTTP **415**. Unknown method / zero timeout → trailers-only gRPC status. diff --git a/docs/architecture.md b/docs/architecture.md index 81f479c..1eef5f8 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -61,7 +61,7 @@ Health, reflection, and channelz are ordinary registered methods (`Grpc.Health`, - Plain TCP Async (`H2.Client.connectH2cAsync`, `H2.Server.listenH2cAsync`) — **no** `.block` on the Async chain - Plain TCP sync adapters (`connectH2c`, `listenH2c`) — `.block` at the edge -- In-process OpenSSL ALPN `h2` (`Grpc.Native.Tls` → `Tls.connectH2` / `serveH2`) — **blocking FFI** in v1.2.0; see [async-io.md](async-io.md) +- In-process OpenSSL ALPN `h2` (`Grpc.Native.Tls` → `Tls.connectH2` / `serveH2` / `*Async`) — **blocking FFI**, **off-loop** under Async in v1.3.0; see [async-io.md](async-io.md) - Optional sidecar via `LEAN_GRPC_TLS_PROXY` (client dials h2c to a local terminator) ## Credentials @@ -104,7 +104,8 @@ These are intentionally lightweight compared to full OpenTelemetry SDKs. ## Process / concurrency model - Networking uses Lean 4 `Std.Async` (libuv-backed). Prefer one event-loop owner per process for listen/accept. -- **v1.2.0+:** `serveH2cAsync` schedules connections with `Std.Async.background` so concurrent h2c clients can share one UV loop. See [async-io.md](async-io.md). +- **v1.2.0+:** `serveH2cAsync` schedules connections with `Std.Async.background` so concurrent h2c clients can share one UV loop. +- **v1.3.0+:** TLS accept/handshake/`SSL_*` run via `H2.runOffLoop` / dedicated fibers so OpenSSL does not stall that loop. See [async-io.md](async-io.md). - Interop loopbacks that need a second peer may still **spawn a child executable** (e.g. `trailersLoopback`) when mixing sync `.block` adapters with tasks. - CI runs multi-language peers as separate processes (Go / Python / h2spec). diff --git a/docs/async-io.md b/docs/async-io.md index 9994f25..1a87360 100644 --- a/docs/async-io.md +++ b/docs/async-io.md @@ -1,4 +1,4 @@ -# Async IO model (v1.2.0) +# Async IO model (v1.3.0) lean-grpc uses Lean 4 `Std.Async.TCP` (libuv-backed) for sockets. Through **v1.1.x** the public API collapsed every accept/connect/send/recv with @@ -6,7 +6,7 @@ lean-grpc uses Lean 4 `Std.Async.TCP` (libuv-backed) for sockets. Through constructors. That overstated “runs on Std.Async” — see [GitHub #6](https://github.com/RileyBetts/lean-grpc/issues/6). -## v1.2.0 model +## v1.2.0+ h2c model | Layer | Behavior | |---|---| @@ -14,39 +14,48 @@ constructors. That overstated “runs on Std.Async” — see | `H2.listenH2cAsync` / `connectH2cAsync` / `awaitResponseAsync` | Async accept/dial/RPC wait | | `Grpc.Server.serveH2cAsync` / `Channel.unaryAsync` | Additive Async gRPC entry points | | `ByteTransport` / `serveH2c` / `unary` / `connectH2c` | **Sync adapters** — `.block` at the edge only | -| In-process TLS (`SSL_read` / `SSL_write`) | Still **blocking FFI**; wrapped via `AsyncByteTransport.ofBlocking` | + +## v1.3.0 TLS off-loop model + +In-process OpenSSL (`SSL_read` / `SSL_write` / handshake) remains **blocking FFI**. +**v1.3.0** runs that FFI on **dedicated threads** so the UV loop is not stalled +([#10](https://github.com/RileyBetts/lean-grpc/issues/10)): + +| Layer | Behavior | +|---|---| +| `H2.runOffLoop` / `AsyncByteTransport.ofBlockingOffLoop` | `IO.asTask .dedicated` + await — SSL off UV | +| `Tls.connectH2Async` / `serveH2Async` / `Server.serveTlsAsync` | Accept/handshake off-loop; conns on `background` | +| `Tls.connectH2` / `serveH2` / `serveTls` | Blocking `IO` (safe under `IO.asTask`; no nested `Async.block`) | + +This is **off-loop blocking TLS**, not nonblocking BIO. Do not claim “fully async +OpenSSL” until a BIO/`WANT_READ` integration lands (follow-up issue). ```text - App (Async) ──► serveH2cAsync / unaryAsync + App (Async) ──► serveH2cAsync / unaryAsync (h2c, zero .block) + └─► serveTlsAsync / connectH2Async (TLS off-loop) │ ▼ - AsyncByteTransport (TCP) ← zero .block on h2c hot path + dedicated Task ──► SSL_* (blocking FFI) │ - ▼ - Std.Async.TCP + UV loop + UV loop stays free for Std.Async.TCP - App (IO) ──► serveH2c / unary ──► (asyncPath).block ← explicit adapter + App (IO) ──► serveTls / connectH2 (plain blocking IO accept/dial) ``` ## UV loop ownership Prefer **one** event-loop owner per process. Nested `.block` inside -`IO.asTask` workers on the same loop is fragile (historical -`TrailersLoopback` note). The Async listen path schedules connections with -`Std.Async.background` instead. - -## TLS caveat (v1.2.0) - -OpenSSL in `native/tls_ffi.c` is blocking. `serveTls` / `connectH2` remain -correct and tested, but they are **not** end-to-end async. Do not claim TLS is -async until a nonblocking BIO or off-loop worker pool lands (follow-up issue). +`IO.asTask` / `runOffLoop` workers is fragile (historical `TrailersLoopback` +note). Use pure `IO` for sync TLS helpers when already on a dedicated thread; +use `Std.Async.background` for connection scheduling under Async. ## Migration | Need | Use | |---|---| | Compose without blocking the loop (h2c) | `*Async` APIs | -| Existing code / lean-compliance | Keep `IO` APIs (adapters) | -| mTLS AuthN | `serveTls` + `registerWithContext` (sync FFI underneath) | +| Compose under Async without freezing UV (TLS) | `connectH2Async` / `serveTlsAsync` / `runOffLoop` | +| Existing code / lean-compliance | Keep `IO` APIs | +| mTLS AuthN | `serveTls` or `serveTlsAsync` + `registerWithContext` | -Pin: `@ "v1.2.0"`. +Pin: `@ "v1.3.0"`. diff --git a/docs/getting-started.md b/docs/getting-started.md index 9a8f760..bec329b 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -4,7 +4,7 @@ This guide gets a Lean 4 unary gRPC server and client running on h2c with **type **Requirements:** Lean 4.32+ (see `lean-toolchain`), OpenSSL headers for the `Grpc` library (`libssl-dev` or `./scripts/fetch-openssl-headers.sh`). Peer gzip optionally uses a zlib helper (`./scripts/build_native.sh`) or system `gzip`. -Package version: **1.2.0** (`lakefile.lean` / `Grpc.version`). +Package version: **1.3.0** (`lakefile.lean` / `Grpc.version`). ## Build the repo @@ -151,12 +151,12 @@ In your `lakefile.lean`: ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.3.0" ``` Then `import Grpc` (and `Proto` if you use the bundled codecs). Public libs: `Bytes`, `Hpack`, `H2`, `Proto`, `Grpc`. Linking OpenSSL (`-lssl -lcrypto`) is pulled in via the `Grpc` Lake library (not via Bytes/Hpack/H2 alone). For peer gzip against foreign stacks, set `LEAN_GRPC_ZLIB_HELPER` after `./scripts/build_native.sh`. Full packaging notes: [packaging.md](packaging.md). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). -After [Reservoir](https://reservoir.lean-lang.org/) indexes the package you can use `require «lean-grpc»` without a git URL. Pin a commit SHA if the `v1.2.0` tag is not yet on the remote you use. +After [Reservoir](https://reservoir.lean-lang.org/) indexes the package you can use `require «lean-grpc»` without a git URL. Pin a commit SHA if the `v1.3.0` tag is not yet on the remote you use. ## Next steps diff --git a/docs/packaging.md b/docs/packaging.md index 510c0ec..c25b9da 100644 --- a/docs/packaging.md +++ b/docs/packaging.md @@ -65,7 +65,7 @@ lake build ```lean require «lean-grpc» from git - "https://github.com/RileyBetts/lean-grpc.git" @ "v1.2.0" + "https://github.com/RileyBetts/lean-grpc.git" @ "v1.3.0" ``` Then `import Grpc` (and `Proto` if using bundled codecs). @@ -76,7 +76,7 @@ Then `import Grpc` (and `Proto` if using bundled codecs). require «lean-grpc» -- once listed on reservoir.lean-lang.org ``` -Package version in-tree is **1.2.0** (`lakefile.lean`, `Grpc.version`). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). Creating the git tag on origin is a separate maintainer step (below). +Package version in-tree is **1.3.0** (`lakefile.lean`, `Grpc.version`). Hosted docs: [rileybetts.ai/oss/lean-grpc](https://rileybetts.ai/oss/lean-grpc). Creating the git tag on origin is a separate maintainer step (below). ## Maintainer release checklist @@ -89,8 +89,8 @@ Documented only — do **not** automate tagging from CI or agent runs. 5. **Manual version tagging (maintainer only):** on `main`, when you choose to publish: ```bash - git tag -a v1.2.0 -m "lean-grpc 1.2.0" - git push origin v1.2.0 + git tag -a v1.3.0 -m "lean-grpc 1.3.0" + git push origin v1.3.0 ``` Agents and CI must not create or push tags. diff --git a/lakefile.lean b/lakefile.lean index 107f53a..8b5ef3c 100644 --- a/lakefile.lean +++ b/lakefile.lean @@ -7,9 +7,9 @@ open Lake DSL open System package «lean-grpc» where - version := v!"1.2.0" + version := v!"1.3.0" keywords := #["grpc", "http2", "hpack", "protobuf", "networking"] - description := "Pure Lean 4 gRPC stack (HTTP/2 + HPACK + gRPC); Std.Async.TCP with native Async h2c APIs" + description := "Pure Lean 4 gRPC stack (HTTP/2 + HPACK + gRPC); Async h2c + off-loop TLS" homepage := "https://rileybetts.ai/oss/lean-grpc" license := "Apache-2.0" licenseFiles := #["LICENSE", "NOTICE"] @@ -70,6 +70,9 @@ lean_exe trailersLoopback where lean_exe asyncH2cLoopback where root := `Tests.AsyncH2cLoopback +lean_exe asyncTlsLoopback where + root := `Tests.AsyncTlsLoopback + lean_exe tlsServer where root := `Tests.TlsServer