Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 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
Expand All @@ -33,6 +33,8 @@ jobs:
./.lake/build/bin/grpcTests
./.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
Expand Down
28 changes: 27 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,32 @@
# 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.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

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

Expand Down
6 changes: 5 additions & 1 deletion Examples/Helloworld/Server.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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`). -/
Expand All @@ -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
2 changes: 1 addition & 1 deletion Grpc.lean
Original file line number Diff line number Diff line change
Expand Up @@ -39,5 +39,5 @@ import Grpc.Grpclb
import Grpc.Interceptor

namespace Grpc
def version : String := "1.1.0"
def version : String := "1.3.0"
end Grpc
56 changes: 53 additions & 3 deletions Grpc/Channel.lean
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,9 @@ import Grpc.ServiceConfig
import Grpc.Retry
import Grpc.Xds
import Grpc.XdsAds
import Std.Async

open Std.Async

namespace Grpc

Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
63 changes: 63 additions & 0 deletions Grpc/Client.lean
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,9 @@ import Grpc.Status
import Grpc.Message
import Grpc.Metadata
import Grpc.Compression
import Std.Async

open Std.Async

namespace Grpc

Expand Down Expand Up @@ -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
2 changes: 1 addition & 1 deletion Grpc/Metadata.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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.3.0") : Hpack.HeaderField :=
⟨ascii "user-agent", ascii s!"grpc-lean/{version}"⟩

def http415 : Array Hpack.HeaderField :=
Expand Down
15 changes: 14 additions & 1 deletion Grpc/Server.lean
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ import Grpc.PeerIdentity
import Grpc.Stream
import Grpc.Compression
import Grpc.Tls
import Std.Async

open Std.Async

namespace Grpc

Expand Down Expand Up @@ -305,12 +308,22 @@ 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.
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
12 changes: 6 additions & 6 deletions Grpc/Stream.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Expand All @@ -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)

Expand Down Expand Up @@ -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
Expand All @@ -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)

Expand Down
Loading
Loading