From cd82ab927578b82555ecb5e1c7f0204613f8bbdc Mon Sep 17 00:00:00 2001 From: Ankur Date: Fri, 31 Jul 2026 22:53:13 +0530 Subject: [PATCH] Spring gRPC on Spring Boot 4: four call types and a symptom-first troubleshooting suite Companion code for the ankurm.com guide. Verified on Spring Boot 4.1.0, spring-grpc 1.1.0, grpc-java 1.80.0, protobuf-java 4.34.2, JDK 25.0.3. results-full.txt is unedited mvn test output: 8 tests, 0 failures. No installs needed - protoc and the gRPC codegen plugin resolve as Maven artifacts, and every test uses the in-process transport. _01_basics all four call types, and how their failure modes differ: iterator semantics on server streaming, the half-close that client streaming hangs without, and the independence of the two streams in bidi. _02_troubleshooting four failures reproduced then fixed - the 4 MB message limit and which side enforces it, the absent default deadline, cancellation that never interrupts a thread, and errors that arrive as UNKNOWN with no description. Findings worth the commit message - The in-process transport CANNOT enforce message size limits: it passes messages by reference and never serialises them. A 4 MB + 1 KB message goes through cleanly in tests and fails in production with RESOURCE_EXHAUSTED. The test asserts this rather than pretending otherwise. Same blind spot covers compression, TLS, keepalive and LB. - Two property names cost real time while writing this: spring.grpc.server.inprocess.name (not in-process) and spring.grpc.client.channel..target (singular channel, and target not address). The second failure surfaces as UnknownHostException on the channel NAME. - gRPC still has no default deadline, and cancellation only sets a Context flag. --- .gitignore | 11 + LICENSE | 21 ++ README.md | 73 +++++ docs/01-troubleshooting.md | 188 +++++++++++++ pom.xml | 108 +++++++ results-full.txt | Bin 0 -> 33032 bytes .../java/com/ankurm/grpc/Application.java | 19 ++ .../ankurm/grpc/orders/OrderServiceImpl.java | 216 ++++++++++++++ src/main/proto/orders.proto | 80 ++++++ .../grpc/_01_basics/FourCallTypesTest.java | 147 ++++++++++ .../HardToDiagnoseTest.java | 265 ++++++++++++++++++ .../com/ankurm/grpc/support/GrpcTestBase.java | 95 +++++++ .../java/com/ankurm/grpc/support/Report.java | 38 +++ 13 files changed, 1261 insertions(+) create mode 100644 .gitignore create mode 100644 LICENSE create mode 100644 README.md create mode 100644 docs/01-troubleshooting.md create mode 100644 pom.xml create mode 100644 results-full.txt create mode 100644 src/main/java/com/ankurm/grpc/Application.java create mode 100644 src/main/java/com/ankurm/grpc/orders/OrderServiceImpl.java create mode 100644 src/main/proto/orders.proto create mode 100644 src/test/java/com/ankurm/grpc/_01_basics/FourCallTypesTest.java create mode 100644 src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java create mode 100644 src/test/java/com/ankurm/grpc/support/GrpcTestBase.java create mode 100644 src/test/java/com/ankurm/grpc/support/Report.java diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..35344dd --- /dev/null +++ b/.gitignore @@ -0,0 +1,11 @@ +target/ +*.class +t1.txt +t2.txt +gen.txt +.idea/ +*.iml +.vscode/ +.DS_Store + +# results-full.txt IS committed on purpose -- it is the evidence for the blog post diff --git a/LICENSE b/LICENSE new file mode 100644 index 0000000..aa5473f --- /dev/null +++ b/LICENSE @@ -0,0 +1,21 @@ +MIT License + +Copyright (c) 2026 Ankur Mhatre + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/README.md b/README.md new file mode 100644 index 0000000..44a7dec --- /dev/null +++ b/README.md @@ -0,0 +1,73 @@ +# Spring gRPC with Spring Boot 4 + +Companion repository for **[Spring gRPC with Spring Boot 4](https://ankurm.com/spring-grpc-spring-boot-4/)** +on [ankurm.com](https://ankurm.com). + +The post covers the concepts. This repository adds the part that is hard to find written down: a +symptom-first troubleshooting guide, with the failures reproduced as tests wherever a test can +honestly reproduce them. + +Verified against **Spring Boot 4.1.0**, **spring-grpc 1.1.0**, **grpc-java 1.80.0**, +**protobuf-java 4.34.2**, **JDK 25.0.3**. Full output in [`results-full.txt`](results-full.txt). + +--- + +## Quick start + +Nothing to install — `protoc` and the gRPC codegen plugin are resolved as Maven artifacts, and every +test runs over the in-process transport, so no port is bound. + +```console +git clone https://ankurm.com/git.app/asmhatre/spring-grpc-boot4.git +cd spring-grpc-boot4 +mvn test +``` + +8 tests, a few seconds. + +--- + +## What is here + +| Test | Covers | +|---|---| +| [`FourCallTypesTest`](src/test/java/com/ankurm/grpc/_01_basics/FourCallTypesTest.java) | Unary, server streaming, client streaming, bidirectional — and how their failure modes differ | +| [`HardToDiagnoseTest`](src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java) | Four failures reproduced and fixed: message size limits, missing deadlines, unobserved cancellation, useless error statuses | + +| Document | Covers | +|---|---| +| [Troubleshooting](docs/01-troubleshooting.md) | Symptom-first index: every entry starts from the error message you actually see | + +The service itself ([`OrderServiceImpl`](src/main/java/com/ankurm/grpc/orders/OrderServiceImpl.java)) +is heavily commented and is where the correct patterns live — cancellation checks, half-close +handling, `StreamObserver` thread-safety, structured errors. + +--- + +## Five findings + +1. **gRPC is a first-class Boot starter now.** `spring-boot-starter-grpc-server` and + `-grpc-client`, version-managed by the Boot BOM. A `@Service` extending the generated `ImplBase` + is registered automatically — no `@GrpcService`, no registration code. +2. **The in-process transport cannot enforce message size limits.** It passes messages by reference + and never serialises them, so a 4 MB + 1 KB message sails through — and fails in production with + `RESOURCE_EXHAUSTED`. Asserted in the test suite. The same blind spot covers compression, TLS, + keepalive and load balancing. +3. **Two property names that cost real time.** `spring.grpc.server.inprocess.name` (not + `in-process`) and `spring.grpc.client.channel..target` (singular `channel`, and `target` not + `address`). Getting the second wrong yields `UnknownHostException` on the channel *name*, which + sends you to look at DNS. +4. **gRPC has no default deadline.** A call without one waits forever. This is the single most + common cause of a gRPC service that silently stops responding. +5. **Cancellation does not interrupt your thread.** It sets a `Context` flag. A server that never + checks it keeps working for a client that left ten minutes ago. + +--- + +## Reference machine + +AMD Ryzen 5 5600U, Windows 11, `java 25.0.3+9-LTS-195`, Maven 3.9.9. + +## Licence + +MIT. See [LICENSE](LICENSE). diff --git a/docs/01-troubleshooting.md b/docs/01-troubleshooting.md new file mode 100644 index 0000000..7bf9c78 --- /dev/null +++ b/docs/01-troubleshooting.md @@ -0,0 +1,188 @@ +# Troubleshooting Spring gRPC on Boot 4 + +Symptom-first. Each entry is something whose error message points somewhere other than its cause. +Entries marked **[tested]** are reproduced in +[`HardToDiagnoseTest`](../src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java); +entries marked **[documented]** need a real network and are stated rather than asserted, because a +test claiming to prove them over in-process transport would be lying. + +--- + +## `UNAVAILABLE: Unable to resolve host ` — **[tested]** + +**Cause, 90% of the time: a typo in a property name**, not a networking problem. + +An unmatched channel name is passed straight to the DNS resolver as a target. Boot does not warn +that the channel is unconfigured. The exact names, both of which are easy to get wrong: + +```properties +spring.grpc.server.inprocess.name=orders-test # NOT "in-process" +spring.grpc.client.channel.orders.target=static://... # SINGULAR "channel"; "target", NOT "address" +``` + +Also check `spring.grpc.client.inprocess.enabled=true` if you are using the in-process transport +from the client side. + +## `IllegalStateException: No grpc channel factory found that supports target : ` + +`GrpcChannelFactory` is a composite that asks each registered factory `supports(target)` **before** +named-channel indirection is applied. Passing a bare logical name to `createChannel()` can therefore +fail even though the name is correctly configured. In application code, inject the stub (via +`@ImportGrpcClients`) and let Boot resolve the target; in tests, pass the full target. + +## `RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 4194304` — **[tested]** + +The default limit is **4 MB, per message, on the receiver**, and the two sides are configured +separately. + +```properties +# large RESPONSE -> the CLIENT is the receiver +spring.grpc.client.channel..inbound.message.max-size=16MB +# large REQUEST -> the SERVER is the receiver +spring.grpc.server.inbound.message.max-size=16MB +``` + +Per call: `stub.withMaxInboundMessageSize(16 * 1024 * 1024)`. + +Half of all "I raised the limit and it still fails" reports are the asymmetry: raising it on the +server does nothing for a large response. + +> **The testing trap.** The in-process transport passes messages **by reference and never +> serialises them**, so it *cannot* enforce size limits. A payload-size regression passes every +> in-process test you have and fails in production. The same applies to compression and anything +> else depending on the wire format. This is demonstrated — the test asserts that a 4 MB + 1 KB +> message goes through in-process without complaint. + +Prefer streaming to large unary messages: the limit is per message, so 1,000 small messages are +fine. Do not raise the limit globally — a large limit turns a malformed request into an OOM. + +## Calls that hang forever — **[tested]** + +**gRPC has no default deadline.** A call without one waits indefinitely. Nothing warns you, and in +testing the server always responds quickly, so it surfaces the day something downstream is slow — +as a thread pool that fills up with no errors in the log. + +```java +stub.withDeadlineAfter(2, TimeUnit.SECONDS).getOrder(request); +``` + +Per channel: `spring.grpc.client.channel..default.deadline=2s`, or register a +`DefaultDeadlineSetupClientInterceptor`. + +Deadlines are **absolute and propagate**: if A calls B with 2s remaining, B sees 2s, not a fresh 2s. +Set the deadline at the edge; do not re-set it at every hop, or a deep chain multiplies its budget. + +## The server keeps working after the client gave up — **[tested]** + +gRPC **does not interrupt your thread** on cancellation. It sets a flag on the `Context`. + +```java +for (Item item : items) { + if (Context.current().isCancelled()) return; // no onCompleted/onError -- stream is closed + responseObserver.onNext(convert(item)); +} +``` + +For blocking work handed to another thread, propagate with `Context.current().wrap(runnable)` — +otherwise the flag is invisible there. Calling `onNext`/`onCompleted` after cancellation throws +`IllegalStateException`. + +## Every error is `UNKNOWN` with no message — **[tested]** + +Any exception escaping a handler becomes `UNKNOWN` with a **null description** — deliberately, since +leaking exception text across a service boundary is an information-disclosure risk. + +```java +Metadata trailers = new Metadata(); +trailers.put(REASON_KEY, "ORDER_LOCKED"); +responseObserver.onError(new StatusRuntimeException( + Status.FAILED_PRECONDITION.withDescription("order is locked"), trailers)); +``` + +Choose the code carefully — it is the contract that retry policies, circuit breakers and dashboards +key off: + +| Do **not** retry | Retry may help | +|---|---| +| `INVALID_ARGUMENT`, `NOT_FOUND`, `ALREADY_EXISTS`, `FAILED_PRECONDITION`, `PERMISSION_DENIED`, `UNAUTHENTICATED` | `UNAVAILABLE`, `DEADLINE_EXCEEDED`, `RESOURCE_EXHAUSTED`, `ABORTED` | + +Getting this wrong makes non-retryable failures look retryable and turns one bad request into a +storm. + +## `IllegalStateException: call already closed` / `half-closed` + +You called `onNext` or `onCompleted` after the stream ended. Usual causes: + +- responding after cancellation or deadline expiry (check `Context.current().isCancelled()` first); +- calling `onCompleted()` twice; +- calling `onError()` and then `onCompleted()` — `onError` is terminal; +- two threads writing to one `StreamObserver`. **`StreamObserver` is not thread-safe.** Concurrent + `onNext` corrupts the stream, and the symptom is usually a deserialization error on the *far* + side, pointing nowhere near the bug. + +## A client-streaming call never completes + +You forgot the half-close. `requestObserver.onCompleted()` is what tells the server "no more +requests"; without it the server's `onCompleted` never fires and the call hangs until the deadline — +or forever, if you did not set one. + +## Retries do nothing — **[documented]** + +gRPC retries are configured through the **service config**, not a client API, and are off by +default: + +```properties +spring.grpc.client.channel.orders.service-config.methodConfig[0].retryPolicy.maxAttempts=4 +spring.grpc.client.channel.orders.service-config.methodConfig[0].retryPolicy.initialBackoff=0.5s +spring.grpc.client.channel.orders.service-config.methodConfig[0].retryPolicy.retryableStatusCodes[0]=UNAVAILABLE +``` + +Three reasons they silently do nothing: + +1. The status you are failing with is not in `retryableStatusCodes`. +2. The response was already partially delivered — gRPC will not retry a committed stream. +3. Retry is disabled on the channel (`enableRetry()` / the corresponding property). + +Only make idempotent methods retryable. gRPC has no idea whether your RPC is safe to repeat. + +## Connections drop every N minutes through a proxy — **[documented]** + +Load balancers and ingress controllers close idle HTTP/2 connections. gRPC reuses one long-lived +connection, so an idle stream is not an idle *connection* — the proxy disagrees. + +```properties +spring.grpc.client.channel..keepalive.time=30s +spring.grpc.client.channel..keepalive.timeout=5s +spring.grpc.server.keepalive.connection.max-idle-time=... +spring.grpc.server.keepalive.connection.max-age=... +``` + +Keep the client's keepalive interval **shorter** than the proxy's idle timeout. Note the server can +reject keepalives that are too frequent (`ENHANCE_YOUR_CALM` / `too_many_pings`), which looks like a +random disconnect — so both ends need to agree. + +## All traffic goes to one backend — **[documented]** + +The default load-balancing policy is `pick_first`, not `round_robin`. With a DNS target resolving to +several addresses, a gRPC client picks one connection and keeps it. This is correct behaviour and +almost never what people expect behind a headless Kubernetes service. + +Set `round_robin` via service config, and remember that DNS re-resolution happens on connection +failure, not on a timer — so scaling up does not redistribute existing connections. + +--- + +## Test transport limitations, collected + +The in-process transport is excellent for testing protocol behaviour and misleading for anything +wire-related. It **does** exercise interceptors, metadata, status codes, deadlines, cancellation and +flow control. It does **not** exercise: + +- message size limits (no serialization — proven in the test suite), +- compression, +- TLS and credentials negotiation, +- keepalive, GOAWAY, idle timeouts, +- name resolution and load balancing, +- HTTP/2 framing and header limits. + +Anything in that second list needs a real port, and belongs in a smaller, slower test tier. diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..22a2413 --- /dev/null +++ b/pom.xml @@ -0,0 +1,108 @@ + + + + 4.0.0 + + + org.springframework.boot + spring-boot-starter-parent + 4.1.0 + + + + com.ankurm.grpc + spring-grpc-boot4 + 1.0.0 + Spring gRPC with Spring Boot 4 + + + 25 + UTF-8 + + + + + + org.springframework.boot + spring-boot-starter-grpc-server + + + org.springframework.boot + spring-boot-starter-grpc-client + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.boot + spring-boot-starter-grpc-server-test + test + + + org.springframework.boot + spring-boot-starter-grpc-client-test + test + + + org.awaitility + awaitility + test + + + + + + + + io.github.ascopes + protobuf-maven-plugin + + ${protobuf-java.version} + + + io.grpc + protoc-gen-grpc-java + ${grpc-java.version} + + + + + + + generate + + + + + + + org.apache.maven.plugins + maven-surefire-plugin + + false + false + + + + + diff --git a/results-full.txt b/results-full.txt new file mode 100644 index 0000000000000000000000000000000000000000..17a9a435ba6c0a4886eb2258b596ffcb5f51438d GIT binary patch literal 33032 zcmeI5ds7@omc{Gujo9zd;hk9F-KEjnwqoQMBqUq5gfyZxo(W>weqAVS z{$0oZRqsEjPi8S_=9?S(gf=K0H2>CosZUe>lX^SQ@%iRi!}o1bbXiq%2BRZeyd$g(l!uw^NI}o>C>Y8=U{+ZtJgmZHzNM@)R4Nu_h zHT7*7B9~9q>*MD8=0EiJp~iVpM&E31i7)H={ZehWbZn>aC}?F~w~u&u*!)@FkX_Px zOEcZmRePFk=+tuI!|diRLfLdqHitbp1j#!};c-&Ry2~g5L0YPrA`B3DbwI&6mE{ z3N4HWZ_pRUoG*01$$h=9YbJN~n%C@>^;&H1N$yMH0@TB=pY?aEaA{4ye<>8~6mCA( zZ!8AxCQ87c;5PWPReqzt!Qp=D(emjiM!&zPKM=wQ#V5M8E?HQT#yKcOcMgh#kSR!m zG(k6hqXGPe2ZI)WnDe?`anAH3tKg?`uzwcdvGwsq;gPWuo!(H7FAE*;iZ+u_dcRON zl$J2|y5#;yqwnjgXS!+@jP6?L^gglCUGefq^^XU+t&!o^Msr_dKGWFvl9NUb>xDDx z>hFap2VHmd39s>G#uspFSmc8=eSAk}v2i4YO%J7uPc<@_h(w?&CHjXJPw+|fjQ{(( zP!39?m(cQIZTY#Vw%5hhv34x_f-LUGGS`!=`a$7c=q(&H_MuPA5PrSX(LM11?Z+z9 z?rg>UMsvU&Bc}h-1Ny%!8Dqml1-6LOBg1|>dpr6(jwbLDqyi_g&e_le3Llr4?6&OX zdtJ98UgAT$I0*vm9^2WzG=W#Qhwc$jefsN|v(#nDMPRkyftM8%-cfyX;v z-rcc?c%Btr61AfLpv^0NZuf9=mibI;^H^ip_YNfAX)R|v=lQfz8ohZStpV*EyW76c z6%HXA{M=upnb@pjzjM0A{t18fH@$b|A&IxcFTWH&ev+K6$#q{(b^Mv;QpX~7Ya)9` z9hapy&owUbKGfm6hTHc($2s|lc;YGp zaeCT)#wYMBGd^M*d#A)TN0Kka4f*A|x*-HnuAxcWBG_xrkXNuSh!MkU5JJ@(0r?0zFE zhn@``9qR8Lp>I8I`teQ0F|IN>OJwHs`?ILs)GU0Jxk}tdv`R&eJPYoH?m}Fm=Pc4TW>Ab#o z&yyjwk#N4J6?Oh!8zb&3^TU&duN#QgS2{w>81dpbeM>sSG3sh~2`ZFSxj#fixv%lp zHF8v^R(tr=^-sUgS$Fvme!3kku@;AZUO3j!_W3^TY`1tDn&~w&)&S}Vs67tr6&cZG z4HA26M|cqzMuhizam=^{j9W0V78dicW9faon%w$r@a?yv+rg#TEc(rLJI1*wEX#~x zn%vqg--gW4n_9C|=cD$XD_~x|qDJiX&8)z)K8mfl{#<(;=R; zSY$)4!Nv;iy4pW3K8wtVbp+Nd+x19n5=$bhi#3Uh!WgQT&f41bfso9+coVA&*jZQB z-6Mf))cP3{>iL_V_o}|f2-BqfHss6dYH$|0dA821>_;s#S3AhnI7ejy|H^81y#~zc zUer5>>bX5b)`Y0&A|dJ;WVFuesB(W(*1kALrL;cxKqHY4#@Jn+c8{3L@9VQWJ`v# z3Kn$bf=13aeY#&(h;P8n*CBapEI9R4>7WRfOj!19%DnS zB!8F;%W|Jx=2s=>vm$dc46+3BSnqF&dUV7(#6VO(v0&D_z}@iUpfh`iB6gvg-Oky^ zdqdAY8Ef-Nv8!45OGd^Xkaxjfs>=A7v?=DjSE|E%I=)=Ea3GKKm5#C-V5K^L`-q28 z3&5_JD^YE|{$%~mmHZe34@YDcRr|@Fkr5mjLGAtHUqOv~0;xW7)ErHTItnKCw!V8G z)z3m<_VK2)k=irqdEM@HIeOt8V8Ch7(c~z(8JIvOG?cBtNp_IMvDe|+mEt?!K5{14 zY~kL~i);_|M#{rpo4zZYW5-~&tG#8~aIHC@W$nvAO~&IWJ3sLoc+FU1M1P zV?{ELg>ToUyua7kFLa)_cR0uXHF}6Mugm$%?^IvmeA?GgQNWI}FS4!JDLNafo)t_* zmw(F~O$2yfnELzS=-bBBcaeddEljOrCiXeFR?x-HdHJu`#mz6?0dBtQ{LWG6EZH6Y zid`zK)Ug)gY=G(w`&X!D&M6kIah&=lwYT(juDp>G4wdbGu3nr`-7DzN+M8>E-oeq{ z#lBru*6TQT6tf!_s;|qtMTSK3%hA4-FZj=d+;x>)%whX}PH@T#2=1PQKp#j`y=L z#y9R2JJ49$&A%u{B4=aSXBApEj54Z5%v+f%kHX(&@6|rh}bo5Jwpc=ksPqD^>J;!d38*qY{M&FH&}CGJo2>^7vAI~w2jlel8Q`Z(jv zov=zyUH_u0m*@mMqqY<)?5pb4yC$g*f#R&Q5J4a6ia5sZm^cbv1l_9z#V!bcSk*m0 z!xeY~+QtS8SDD$s%81#UhcxXj<3M=9|Cj zbxnTfR^c=EaBT`#YOXIH&pYno3;i_Kf}TjOG26r zAqWVwBhI4niBsI^iUp8Gu)<_{_wC{1YTF%^)IhiXS+;qZ))_oBuBUE+_31f4Tk2qN z)!NTp1-1=WHsO5k9S^RcInz7}wDN9O!cW&P=(|~TJLY%xyb!`se)s9c-C!k$lT%nQL34|nA?}f zcXihs8J79S3N^OwUgwI>vh7BE0oprK!WV#_R5!Tu9rWRDcJBbThl{INN%vi$v3`z5 z6UAXCz9Y!A#T95U{wiWCuU{~8?`x)V67h*`f%qd(0IKk?X#99jAfTzMX|d(9hu1M1 z+&T-|@JMq8RS%2S*Ls~2X2$iIv#9ZSlk1x4Q|UeN3kZ+C6Ho4yGqKK#)QJMEU9}}# zQ0hFohOV&5LKWj48faaJq_ULiT=omG19$$G)@*YPDtDm_hYTi!);y8LIG z8f1n)L6TrG_jRHz=FO(Kgie60J{~BJ?V;geKZ|A6HTCnMu1fq#?W5nCLv4=Z6TmWT z1Y2B^U50I9p`+qyM$$`K)Y@l%kLK1RIP~C^MoK((WmuVDQIrMpCz5vJ)_i9xPKh;u?cjJ$ZHkxiZhP zanIJHoL7yMiE@Vp!R>7Vz0rI>y~SDhT7}L;9Olke(Uz^|@A{jxfSOQ=Zb56(Ld?!l zvQm0MBe)Wc=Q!Vt#>O0e2PcVQ`OY}56!@yxj|i{Y)_=NwLElQ5r~ObH8$2f>`#O?I zTS0Qa1ti95i0l)2sfmIRV1+ZnSjj_!(2_l2!1<%RCXZa`9J$yza;bCVvR3@kszycg zhysZ9U1zgy;Q@ZqabuRB<3bK)tOvh|f2_UGiS6GK{f;WgzOZ;_7|0Fm&QWSRMAMTZ zu8YhkY}>mtu^rAp1DPPmL!1;;L~A3a?>z)_vHnBV2Z^)}j_AV*9YNyzY=6$ucZDBN z%3k5aiZ0_{plMuz!~)%+D|ENN|A_VA)pvz=j1_(TB3Xdoph9 zdvf-Zka>Aut2N3xbW0-}o$8}2rN-*Io}a1Z)+6=e8B#%~@(1#{1@-vP6>op0QJK9t zk0qbH3+31B&ul;IFwE(HU`Y9=+v#x{$U_4fL@G5&LXQo+nOytM`3rlPb#pE z+b6bjbgY~ekfw!g|?GhVY`m(9EZhIkLvi@ zh-sVwFI~ws%3@hW+P&$ngb9|4U9K+;2;DMg@z~N<(_Pn&`qWU)ihMxcfqOVtK13D+(_)a_GU&lX5E&bmY zZN&Z$Mdw}px+yJ$>bgJtq>+1GJwMZTWWloy@SpVR)kFHD0v6XkE;GNcj2TYcEB8gX z<{UUXE?s^j9bz(P^7txRgAab1%N@xaZFsJ>+lpeG9rZm85POM{0`ILw#JgBf@~w|n zs3AVG-)%<<22Cmuc6N4b0rv#t)WM@m7v|#l1X%N+_}n=@}v?*(==qgu|2%h(9Mr95v!?KQtXjs15it>Uq*iu$2aV9RSg65r7mu{o z6ks8Xq7XWl_5F$7f0eD=(u;F3AMAWZR`f{o;H=wyt!p3Y_~@vm(JieiGitJkZTa!j<}ns2nwf<@9SXdO#oHp`qAW zEP9B|v~f4J%GedX2X#<`{$m;Snu95;c}4C;rJCoC5`kjneCNJdo?(H#Jk>8W*1G~+ z8!~lhF$e?76K4=zI<`Q!qu#)EAP;Lp)V}GJaU99QZ{m~OEALUk7{&#i$RO;Q!xz(k z*8P%S=qt!(emdu82AdiM{J|S|U%)Tw9kgX-+B!t-mFqkoM)K3Gc*fviq1cqwG$3W3 zwV6B(T}hhS9t6tqHM`+UV&a^&B$E9n>l|A-;ce*6L=3kyABo> zTIWxUa~8M-jphcGhq|JV{F8e;x0iblzAd~&YuZ&zKjS50K{yX4&mYAapzexZcz!sF zE`k_55pq&Q2UYXOn9GQ5Ur!dtdd3y~y`=c-bM+KE6p)W+g>10!i0hfJeIoH+I|>3L zje~g7tTsescVB*pm@MO$U$hFjW)$JoNo>OTh}U!ixBG3xnh5pIp}0Dp`sGLl>WR6X=ld8IU$4)@ zPg{CYWc>?j=pj5w#55UEz;)s-kjm?F=A@}(#}!yPj!Y@B)cc6O$Oym)W0^-z+5(mM z%U-X=)3+sQ(~k><9hzo)q{_w&^TAbuj(Q$vPLh0B%v? zA0Cf)Pc*!eUDO(OuO--BltOO6{rZ>nt^Vb^5E!ekHa1NQzAl^~%a5I) z#>#V|j7SXBMKWlA*lSur+Oko{(T978EZWa)#+n1C?ZIu|^vMX*Xe7d6WgNXHrtuE6 zpe>|i3KD%mk6NdUlUPpphxEjOt)Tt5O{^%l3Y3^$gPb;prGjwbZGv`aG~LuJ$P%oV zcy8}ouc@i|-d+cF4UOHBt|C}oaClP90snDduX`G48r=`XZ!#`00Di)MS82?5;^02L z2;b-HsA!UxE~>tHtEwfT&AjNDDvdR#jc>KKco$8?5zu|Q>LP7wppX0fy%&Z5YT;-( z@uRDgMp(zkv4a(XE@Bk8@(yD2N7$bkQ%4vU7!aXkYD! z5d%fwp?9^A+Y&=Zjxw%;;iach%xLLTz^v~yCYFrzXD3+M-pG-+1d*|K1PTSCo z+oA?vl)I>c-$5Oy=c*}L8o9f3V?_wq0Z^A zcjv^q4NnELpXt`R!e=C9lTm}Gt1w3-0)6m;bt|$l@93d7@>|b7)_BINAdBjUDlOPq zBhloMgPV0!7F68SYVG}!*)9n^9D!haAt?}t>MDq{piy@U3K_*g0_-(s@2J|5s^x)4 zW45#JM}UfNjw6&BZ?Z1Zj-O7~FSx-zv%2n!eZmD+_(7QV+C*TKxenr3--tARDB6hD z^dW?w-XC$D(I&77I!+g7+SEWFZ_jm6Gem=p`)Ru%u4~Ay9kGT@_TA-~_rnqQB_84c$QgqHq^Gr{vLAg7wSFw=-WFp?mIK$UYat=(gcG*@O4RnWGiKGd(2Znu^=UPS`2d4a?28qYLo zdyn);dQZIeNcMmQ@bs?SEw>_Vn2i!NBI<|GgCy#EzMS7cPp+PhQ@5hC%;2Sb0hs0Z z^N#wph7x(kdJP`lQ2^RduljZQW$yHSzK!ZAdX{q7)swIkK1 zpd%&Ljzjw})sYIjyx6aGSK7z!EM{bT@5XN4p=}%p`;Yq+YmA34>;g%o5AHkOK`*KF zz+3Yd-f#uI8nwBCqfsq`U)-<5S$4T{9zP6n7$rTQz^9-duMjJv(K50}%MU)AE@X&y z*WQZg*_3*G-!DX!g{Q>ve->Cd3@yQW>9s9Mjwe_-rg>EXPP?OSVWM0 z5;Y&9Z}#Lxr7gI~c;2;&p4;2{s=#mZ6|x6bjNwdscl$`!!#!r?SQ?ype|!9S`pFg5 z33$S6aux2|3g72W`!;8~yn0x8Wj(=ix-Gg#`uY2Lv4q=24}aY~6x|WKM?RFdoiuKj z!W+UWcBn2%>MQzwQ!m$+yKN`+`k{_@DICe#h*!5Yo4^f66aGKopOt9S-m9&1w+rro cdqhvJVH}j)ah7UgyNcMQ>qBv9(zsCg|9fq1Hvj+t literal 0 HcmV?d00001 diff --git a/src/main/java/com/ankurm/grpc/Application.java b/src/main/java/com/ankurm/grpc/Application.java new file mode 100644 index 0000000..92283c4 --- /dev/null +++ b/src/main/java/com/ankurm/grpc/Application.java @@ -0,0 +1,19 @@ +package com.ankurm.grpc; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +/** + * Companion application for + * Spring gRPC with Spring Boot 4. + * + *

The service lives in {@code orders/}; the interesting material is in {@code src/test}, where + * each failure mode is reproduced and then fixed. + */ +@SpringBootApplication +public class Application { + + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } +} diff --git a/src/main/java/com/ankurm/grpc/orders/OrderServiceImpl.java b/src/main/java/com/ankurm/grpc/orders/OrderServiceImpl.java new file mode 100644 index 0000000..9d638ec --- /dev/null +++ b/src/main/java/com/ankurm/grpc/orders/OrderServiceImpl.java @@ -0,0 +1,216 @@ +package com.ankurm.grpc.orders; + +import com.google.protobuf.ByteString; +import io.grpc.Context; +import io.grpc.Metadata; +import io.grpc.Status; +import io.grpc.StatusRuntimeException; +import io.grpc.stub.StreamObserver; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; + +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; + +/** + * The service implementation. + * + *

On Spring Boot 4 a {@code @Service} bean that extends the generated {@code ImplBase} is + * discovered automatically: it is a {@link io.grpc.BindableService}, and + * {@code GrpcServerServicesAutoConfiguration} registers every such bean with the server. There is + * no {@code @GrpcService} annotation to apply and no registration code to write. + */ +@Service +public class OrderServiceImpl extends OrderServiceGrpc.OrderServiceImplBase { + + private static final Logger log = LoggerFactory.getLogger(OrderServiceImpl.class); + + /** Incremented whenever {@link #slowCall} is entered -- lets tests count server-side attempts. */ + private final AtomicInteger slowCallAttempts = new AtomicInteger(); + + /** Set when a server-streaming call notices that its client has gone away. */ + private final AtomicInteger cancellationsObserved = new AtomicInteger(); + + /** How far ListOrders got before being cancelled; proves the server stopped early. */ + private final AtomicInteger lastListProgress = new AtomicInteger(); + + @Override + public void getOrder(GetOrderRequest request, StreamObserver responseObserver) { + responseObserver.onNext(order(request.getId())); + responseObserver.onCompleted(); + } + + /** + * Server streaming, written the way it should be: checking for cancellation between + * messages. + * + *

If the client goes away -- it cancelled, it timed out, its process died -- gRPC does not + * interrupt your thread. It sets a flag on the {@link Context}. A server that never checks that + * flag keeps computing, keeps querying the database, and keeps calling {@code onNext} into a + * closed stream for the full duration of the work. On a long stream that is a real and + * surprisingly common source of wasted capacity. + */ + @Override + public void listOrders(ListOrdersRequest request, StreamObserver responseObserver) { + int count = request.getCount(); + for (int i = 0; i < count; i++) { + if (Context.current().isCancelled()) { + cancellationsObserved.incrementAndGet(); + lastListProgress.set(i); + log.info("client cancelled after {} of {} messages -- stopping work", i, count); + // Do NOT call onCompleted/onError here: the stream is already closed. Just return. + return; + } + responseObserver.onNext(order("order-" + i)); + if (request.getDelayMillis() > 0) { + sleep(request.getDelayMillis()); + } + } + lastListProgress.set(count); + responseObserver.onCompleted(); + } + + /** + * Client streaming. The returned observer receives the client's messages; the response is sent + * exactly once, from {@code onCompleted}. + */ + @Override + public StreamObserver submitOrders(StreamObserver responseObserver) { + AtomicInteger accepted = new AtomicInteger(); + AtomicLong total = new AtomicLong(); + + return new StreamObserver<>() { + @Override + public void onNext(Order value) { + accepted.incrementAndGet(); + total.addAndGet(value.getAmountCents()); + } + + @Override + public void onError(Throwable t) { + // The CLIENT failed or cancelled. The stream is already dead -- responding here + // throws. Log and release resources, nothing else. + log.info("client-streaming call failed: {}", t.toString()); + } + + @Override + public void onCompleted() { + responseObserver.onNext(SubmitSummary.newBuilder() + .setAccepted(accepted.get()) + .setTotalCents(total.get()) + .build()); + responseObserver.onCompleted(); + } + }; + } + + /** + * Bidirectional streaming. + * + *

The thread-safety rule that is easy to miss: a {@link StreamObserver} is not + * thread-safe. In a bidi call it is legal to call {@code onNext} on the response observer from + * a different thread than the one delivering requests -- but only if you serialise those calls + * yourself. Here every response is emitted from inside {@code onNext}, on the delivering + * thread, which is the simplest correct arrangement. + */ + @Override + public StreamObserver sync(StreamObserver responseObserver) { + return new StreamObserver<>() { + @Override + public void onNext(Order value) { + responseObserver.onNext(value.toBuilder().setStatus(Order.Status.PAID).build()); + } + + @Override + public void onError(Throwable t) { + log.info("bidi call failed: {}", t.toString()); + } + + @Override + public void onCompleted() { + responseObserver.onCompleted(); + } + }; + } + + /** Returns a payload of exactly the requested size, for exercising message size limits. */ + @Override + public void getLargePayload(SizeRequest request, StreamObserver responseObserver) { + byte[] data = new byte[Math.max(0, request.getSizeBytes())]; + responseObserver.onNext(Payload.newBuilder().setData(ByteString.copyFrom(data)).build()); + responseObserver.onCompleted(); + } + + /** Sleeps, so deadlines and cancellation can be exercised deterministically. */ + @Override + public void slowCall(SlowRequest request, StreamObserver responseObserver) { + slowCallAttempts.incrementAndGet(); + sleep(request.getSleepMillis()); + if (Context.current().isCancelled()) { + // The deadline already expired. The client has long since given up; sending now throws + // and adds nothing but noise to the logs. + log.info("slowCall finished but the context was already cancelled -- not responding"); + return; + } + responseObserver.onNext(order("slow")); + responseObserver.onCompleted(); + } + + /** + * Always fails, with a machine-readable reason attached as trailing metadata. + * + *

Returning a bare {@code Status.INTERNAL} tells a caller nothing they can act on. Attaching + * trailers -- or, better, a {@code google.rpc.Status} with typed detail messages -- is how gRPC + * carries structured errors. This uses plain metadata to keep the proto dependency-free. + */ + @Override + public void alwaysFails(GetOrderRequest request, StreamObserver responseObserver) { + Metadata trailers = new Metadata(); + trailers.put(REASON_KEY, "ORDER_LOCKED"); + trailers.put(RETRY_AFTER_KEY, "30"); + responseObserver.onError(new StatusRuntimeException( + Status.FAILED_PRECONDITION.withDescription("order " + request.getId() + " is locked"), + trailers)); + } + + public static final Metadata.Key REASON_KEY = + Metadata.Key.of("x-failure-reason", Metadata.ASCII_STRING_MARSHALLER); + public static final Metadata.Key RETRY_AFTER_KEY = + Metadata.Key.of("x-retry-after-seconds", Metadata.ASCII_STRING_MARSHALLER); + + private static Order order(String id) { + return Order.newBuilder() + .setId(id) + .setCustomer("ankur") + .setAmountCents(1_999) + .setStatus(Order.Status.PENDING) + .build(); + } + + private static void sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + public int slowCallAttempts() { + return slowCallAttempts.get(); + } + + public int cancellationsObserved() { + return cancellationsObserved.get(); + } + + public int lastListProgress() { + return lastListProgress.get(); + } + + public void reset() { + slowCallAttempts.set(0); + cancellationsObserved.set(0); + lastListProgress.set(0); + } +} diff --git a/src/main/proto/orders.proto b/src/main/proto/orders.proto new file mode 100644 index 0000000..2aec69b --- /dev/null +++ b/src/main/proto/orders.proto @@ -0,0 +1,80 @@ +syntax = "proto3"; + +package com.ankurm.grpc.orders; + +option java_multiple_files = true; +option java_package = "com.ankurm.grpc.orders"; +option java_outer_classname = "OrdersProto"; + +// A deliberately small service that still exercises all four gRPC call types, because the +// interesting failure modes differ sharply between them. +service OrderService { + + // Unary: one request, one response. The 90% case, and the only one most tutorials cover. + rpc GetOrder (GetOrderRequest) returns (Order); + + // Server streaming: one request, many responses. Where deadlines and cancellation start to + // matter, and where "the client went away" becomes something you have to handle. + rpc ListOrders (ListOrdersRequest) returns (stream Order); + + // Client streaming: many requests, one response. Where flow control and half-close appear. + rpc SubmitOrders (stream Order) returns (SubmitSummary); + + // Bidirectional streaming: the one that will teach you about StreamObserver thread safety. + rpc Sync (stream Order) returns (stream Order); + + // Used by the troubleshooting suite to return payloads of a requested size, so that message + // size limits can be demonstrated rather than described. + rpc GetLargePayload (SizeRequest) returns (Payload); + + // Sleeps for a requested duration so that deadlines, cancellation and keepalive behaviour can + // be exercised deterministically. + rpc SlowCall (SlowRequest) returns (Order); + + // Always fails, with rich error details attached. + rpc AlwaysFails (GetOrderRequest) returns (Order); +} + +message GetOrderRequest { + string id = 1; +} + +message ListOrdersRequest { + int32 count = 1; + // Milliseconds to wait between emitted messages; lets tests exercise slow producers. + int32 delay_millis = 2; +} + +message Order { + string id = 1; + string customer = 2; + int64 amount_cents = 3; + Status status = 4; + + enum Status { + // proto3 requires the zero value to be the "unset" case. Naming it UNSPECIFIED rather than + // reusing a real state is a convention worth keeping: it makes "field absent" distinguishable + // from "field genuinely has the first value", which is otherwise impossible in proto3. + STATUS_UNSPECIFIED = 0; + PENDING = 1; + PAID = 2; + CANCELLED = 3; + } +} + +message SubmitSummary { + int32 accepted = 1; + int64 total_cents = 2; +} + +message SizeRequest { + int32 size_bytes = 1; +} + +message Payload { + bytes data = 1; +} + +message SlowRequest { + int32 sleep_millis = 1; +} diff --git a/src/test/java/com/ankurm/grpc/_01_basics/FourCallTypesTest.java b/src/test/java/com/ankurm/grpc/_01_basics/FourCallTypesTest.java new file mode 100644 index 0000000..0b74cb1 --- /dev/null +++ b/src/test/java/com/ankurm/grpc/_01_basics/FourCallTypesTest.java @@ -0,0 +1,147 @@ +package com.ankurm.grpc._01_basics; + +import com.ankurm.grpc.orders.GetOrderRequest; +import com.ankurm.grpc.orders.ListOrdersRequest; +import com.ankurm.grpc.orders.Order; +import com.ankurm.grpc.orders.OrderServiceGrpc; +import com.ankurm.grpc.orders.SubmitSummary; +import com.ankurm.grpc.support.GrpcTestBase; +import com.ankurm.grpc.support.Report; +import io.grpc.stub.StreamObserver; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * All four gRPC call types against a Spring Boot 4 server, with the differences that matter. + * + *

Most introductions show unary and stop. The three streaming forms are where the interesting + * failure modes live, and they behave differently enough that experience with one does not transfer + * to the others. + */ +class FourCallTypesTest extends GrpcTestBase { + + @Test + void unary() { + Report.title("Unary: one request, one response"); + + Order order = blockingStub().getOrder(GetOrderRequest.newBuilder().setId("abc").build()); + + Report.bullet("id=%s customer=%s amountCents=%d status=%s", + order.getId(), order.getCustomer(), order.getAmountCents(), order.getStatus()); + assertThat(order.getId()).isEqualTo("abc"); + + Report.takeaway("A @Service extending the generated ImplBase is registered automatically:"); + Report.takeaway("it is a BindableService, and Boot 4 wires every such bean into the server."); + } + + @Test + void serverStreaming() { + Report.title("Server streaming: one request, many responses"); + + Iterator it = blockingStub() + .listOrders(ListOrdersRequest.newBuilder().setCount(5).build()); + + List ids = new ArrayList<>(); + it.forEachRemaining(o -> ids.add(o.getId())); + + Report.bullet("received %d messages: %s", ids.size(), ids); + assertThat(ids).hasSize(5); + + Report.takeaway("The blocking stub returns an Iterator. Each next() may block, and an"); + Report.takeaway("exception surfaces mid-iteration -- so a try/catch around the loop body"); + Report.takeaway("is not the same as one around the call. Wrap the whole iteration."); + } + + @Test + void clientStreaming() throws Exception { + Report.title("Client streaming: many requests, one response"); + + AtomicReference summary = new AtomicReference<>(); + CountDownLatch done = new CountDownLatch(1); + + StreamObserver requests = asyncStub().submitOrders(new StreamObserver<>() { + @Override + public void onNext(SubmitSummary value) { + summary.set(value); + } + + @Override + public void onError(Throwable t) { + done.countDown(); + } + + @Override + public void onCompleted() { + done.countDown(); + } + }); + + for (int i = 0; i < 3; i++) { + requests.onNext(Order.newBuilder().setId("o" + i).setAmountCents(1000).build()); + } + // Half-close: "I have no more requests". Forgetting this is the single most common + // client-streaming bug -- the server's onCompleted never fires and the call hangs until + // the deadline, or forever if there is no deadline. + requests.onCompleted(); + + assertThat(done.await(10, TimeUnit.SECONDS)).isTrue(); + Report.bullet("accepted=%d totalCents=%d", + summary.get().getAccepted(), summary.get().getTotalCents()); + assertThat(summary.get().getAccepted()).isEqualTo(3); + + Report.takeaway("requests.onCompleted() is the half-close. Without it the call hangs until"); + Report.takeaway("the deadline -- and if you did not set a deadline, it hangs forever."); + } + + @Test + void bidirectionalStreaming() throws Exception { + Report.title("Bidirectional streaming: many requests, many responses"); + + List received = new ArrayList<>(); + CountDownLatch done = new CountDownLatch(1); + + StreamObserver requests = asyncStub().sync(new StreamObserver<>() { + @Override + public void onNext(Order value) { + received.add(value); + } + + @Override + public void onError(Throwable t) { + done.countDown(); + } + + @Override + public void onCompleted() { + done.countDown(); + } + }); + + for (int i = 0; i < 4; i++) { + requests.onNext(Order.newBuilder().setId("b" + i).setStatus(Order.Status.PENDING).build()); + } + requests.onCompleted(); + + assertThat(done.await(10, TimeUnit.SECONDS)).isTrue(); + Report.bullet("sent 4, received %d, all status=%s", + received.size(), received.isEmpty() ? "-" : received.getFirst().getStatus()); + assertThat(received).hasSize(4); + assertThat(received).allMatch(o -> o.getStatus() == Order.Status.PAID); + + Report.takeaway("Request and response streams are INDEPENDENT. The server may respond"); + Report.takeaway("before you finish sending, or not at all until you half-close. Do not"); + Report.takeaway("assume a request/response pairing -- that is your protocol's job, not gRPC's."); + Report.takeaway(""); + Report.takeaway("StreamObserver is NOT thread-safe. Calling onNext from two threads without"); + Report.takeaway("synchronisation corrupts the stream, and the symptom is usually a"); + Report.takeaway("deserialization error on the far side rather than anything pointing here."); + } +} diff --git a/src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java b/src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java new file mode 100644 index 0000000..2234431 --- /dev/null +++ b/src/test/java/com/ankurm/grpc/_02_troubleshooting/HardToDiagnoseTest.java @@ -0,0 +1,265 @@ +package com.ankurm.grpc._02_troubleshooting; + +import com.ankurm.grpc.orders.GetOrderRequest; +import com.ankurm.grpc.orders.ListOrdersRequest; +import com.ankurm.grpc.orders.Order; +import com.ankurm.grpc.orders.OrderServiceGrpc; +import com.ankurm.grpc.orders.OrderServiceImpl; +import com.ankurm.grpc.orders.Payload; +import com.ankurm.grpc.orders.SizeRequest; +import com.ankurm.grpc.orders.SlowRequest; +import com.ankurm.grpc.support.GrpcTestBase; +import com.ankurm.grpc.support.Report; +import io.grpc.Metadata; +import io.grpc.Status; +import io.grpc.StatusRuntimeException; +import io.grpc.stub.StreamObserver; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.awaitility.Awaitility.await; + +/** + *

The failures that cost hours

+ * + *

Each test here reproduces a problem that is hard to diagnose from its symptom, then + * shows the fix. They are grouped in one class because they share a server and the suite is faster + * that way; each is independent. + * + *

Marked in the transcript as {@code [PROBLEM]} and {@code [FIX]}. + */ +class HardToDiagnoseTest extends GrpcTestBase { + + @Autowired + OrderServiceImpl service; + + @BeforeEach + void reset() { + service.reset(); + } + + // ======================================================================================= + // 1. The 4 MB message limit + // ======================================================================================= + + /** + * Symptom: {@code RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 4194304}. + * Works in dev, fails in production, or works for 99% of requests and fails for the big ones. + * + *

Why it is confusing: the limit is per message, it applies to the + * receiving side, and client and server are configured separately. So raising + * it on the server does nothing for a large response, which is the client's inbound + * limit. Half the reports of "I raised the limit and it still fails" are this asymmetry. + */ + @Test + void inProcessTransportDoesNotEnforceMessageSizeLimits() { + Report.title("1. The 4 MB limit -- and why your integration tests will not catch it"); + + int fourMb = 4 * 1024 * 1024; + + Report.section("Sending a 4 MB + 1 KB response over the IN-PROCESS transport"); + Payload big = blockingStub().getLargePayload( + SizeRequest.newBuilder().setSizeBytes(fourMb + 1024).build()); + Report.bullet("received %,d bytes -- no error", big.getData().size()); + + assertThat(big.getData().size()) + .as("in-process transport happily delivers a message over the default limit") + .isEqualTo(fourMb + 1024); + + Report.problem("This message is OVER the 4 MB default limit and nothing complained."); + Report.problem("The in-process transport passes message objects BY REFERENCE -- it never"); + Report.problem("serialises them -- so there are no bytes to measure and the limit cannot"); + Report.problem("be applied. Over a real Netty transport the same call fails with:"); + Report.problem(" RESOURCE_EXHAUSTED: gRPC message exceeds maximum size 4194304"); + + Report.takeaway("This is a genuinely nasty testing trap. In-process transport is otherwise"); + Report.takeaway("an excellent test transport -- interceptors, statuses, deadlines and"); + Report.takeaway("cancellation all behave correctly -- but message size limits, compression"); + Report.takeaway("and anything else that depends on the wire format do NOT apply."); + Report.takeaway("A payload-size regression will pass every in-process test you have."); + Report.takeaway(""); + Report.takeaway("Test size limits against a real port, or not at all -- but do not believe"); + Report.takeaway("a green in-process suite on this point."); + + Report.section("The fix in production, for when you do hit it"); + Report.fix("The limit applies to the RECEIVER, and the two sides are configured separately:"); + Report.fix(" large RESPONSE -> client inbound limit"); + Report.fix(" spring.grpc.client.channel..inbound.message.max-size=16MB"); + Report.fix(" large REQUEST -> server inbound limit"); + Report.fix(" spring.grpc.server.inbound.message.max-size=16MB"); + Report.fix("Per call, without touching configuration:"); + Report.fix(" stub.withMaxInboundMessageSize(16 * 1024 * 1024)"); + Report.fix(""); + Report.fix("Half of all 'I raised the limit and it still fails' reports are this"); + Report.fix("asymmetry: raising it on the server does nothing for a large RESPONSE."); + Report.fix(""); + Report.fix("Do not raise it globally. A large limit turns a malformed request into an OOM."); + Report.fix("Prefer streaming: the limit is per MESSAGE, so 1000 small messages are fine."); + } + + // ======================================================================================= + // 2. No deadline + // ======================================================================================= + + /** + * Symptom: a thread pool fills up and the service stops responding, with no errors in + * the logs. Or: a retry storm during a downstream slowdown. + * + *

Why it is confusing: gRPC calls have no default deadline. A call with no + * deadline waits forever. Nothing warns you, and in testing the server always responds quickly, + * so it never surfaces until the day something downstream is slow. + */ + @Test + void callsHaveNoDefaultDeadline() { + Report.title("2. gRPC has NO default deadline -- a hung server hangs you forever"); + + Report.section("A call with an explicit deadline shorter than the server takes"); + long start = System.nanoTime(); + assertThatThrownBy(() -> blockingStub() + .withDeadlineAfter(300, TimeUnit.MILLISECONDS) + .slowCall(SlowRequest.newBuilder().setSleepMillis(3_000).build())) + .isInstanceOf(StatusRuntimeException.class) + .satisfies(e -> { + Status s = ((StatusRuntimeException) e).getStatus(); + Report.bullet("%s after %d ms", s.getCode(), + (System.nanoTime() - start) / 1_000_000); + assertThat(s.getCode()).isEqualTo(Status.Code.DEADLINE_EXCEEDED); + }); + + Report.problem("Without .withDeadlineAfter(...) that call would have blocked for the full"); + Report.problem("3 seconds -- and if the server never responded, forever."); + + Report.section("The fix"); + Report.fix("Per call: stub.withDeadlineAfter(2, TimeUnit.SECONDS)"); + Report.fix("Per channel: register a DefaultDeadlineSetupClientInterceptor, or set"); + Report.fix(" spring.grpc.client.channel..default.deadline=2s"); + Report.fix("Treat a missing deadline as a code-review failure, like a missing timeout"); + Report.fix("on an HTTP client."); + + Report.takeaway("A deadline is ABSOLUTE and propagates: if service A calls B with 2s"); + Report.takeaway("remaining, B sees 2s, not a fresh 2s. That is the property that stops a"); + Report.takeaway("deep call chain from multiplying its timeouts -- and the reason you should"); + Report.takeaway("set the deadline at the EDGE, not re-set it at every hop."); + } + + // ======================================================================================= + // 3. Cancellation the server never notices + // ======================================================================================= + + /** + * Symptom: the client timed out ten minutes ago, and the server is still burning CPU and + * database connections on its request. + * + *

Why it is confusing: gRPC does not interrupt your thread when a call is cancelled. + * It flips a flag on the {@link io.grpc.Context}. A server that never checks it keeps working. + */ + @Test + void serverMustCheckForCancellationItself() throws Exception { + Report.title("3. Cancellation does not interrupt your thread -- you must check for it"); + + Report.section("Client starts a 200-message stream, then cancels after a few messages"); + CountDownLatch got3 = new CountDownLatch(3); + AtomicReference clientError = new AtomicReference<>(); + + io.grpc.stub.ClientCallStreamObserver requestStream = + new io.grpc.stub.ClientCallStreamObserver<>() { + @Override public boolean isReady() { return true; } + @Override public void setOnReadyHandler(Runnable r) { } + @Override public void disableAutoInboundFlowControl() { } + @Override public void request(int count) { } + @Override public void setMessageCompression(boolean enable) { } + @Override public void cancel(String message, Throwable cause) { } + @Override public void onNext(ListOrdersRequest value) { } + @Override public void onError(Throwable t) { } + @Override public void onCompleted() { } + }; + + io.grpc.Context.CancellableContext ctx = io.grpc.Context.current().withCancellation(); + ctx.run(() -> asyncStub().listOrders( + ListOrdersRequest.newBuilder().setCount(200).setDelayMillis(20).build(), + new StreamObserver<>() { + @Override public void onNext(Order value) { got3.countDown(); } + @Override public void onError(Throwable t) { clientError.set(t); } + @Override public void onCompleted() { } + })); + + assertThat(got3.await(10, TimeUnit.SECONDS)).as("stream started").isTrue(); + Report.bullet("received 3 messages, now cancelling"); + ctx.cancel(new RuntimeException("client gave up")); + + await().atMost(15, TimeUnit.SECONDS) + .until(() -> service.cancellationsObserved() > 0); + + Report.bullet("server noticed cancellation after %d of 200 messages", service.lastListProgress()); + assertThat(service.lastListProgress()).isLessThan(200); + + Report.fix("The server loop checks Context.current().isCancelled() between messages and"); + Report.fix("returns. Without that check it would have produced all 200 messages into a"); + Report.fix("dead stream, doing every database read and every serialization for nothing."); + + Report.takeaway("Any server handler that loops, or that does work in stages, should check"); + Report.takeaway("Context.current().isCancelled(). For blocking work, propagate the Context"); + Report.takeaway("to worker threads with Context.current().wrap(runnable) -- otherwise the"); + Report.takeaway("cancellation flag is invisible to them."); + Report.takeaway(""); + Report.takeaway("Also: after cancellation the stream is CLOSED. Calling onNext/onCompleted"); + Report.takeaway("on it throws IllegalStateException. Just return."); + } + + // ======================================================================================= + // 4. Errors that tell the caller nothing + // ======================================================================================= + + /** + * Symptom: {@code UNKNOWN} statuses everywhere, and callers that cannot tell a retryable + * failure from a permanent one. + * + *

Why it is confusing: any exception that escapes a gRPC handler becomes + * {@code UNKNOWN} with no message, because leaking exception text across a service + * boundary would be an information disclosure risk. So the useful part of your error is + * discarded by default. + */ + @Test + void statusCodesAndTrailersCarryTheActualReason() { + Report.title("4. Errors: UNKNOWN by default, useful only if you make them so"); + + Report.section("A handler that fails with a proper status and trailing metadata"); + assertThatThrownBy(() -> blockingStub() + .alwaysFails(GetOrderRequest.newBuilder().setId("o-42").build())) + .isInstanceOf(StatusRuntimeException.class) + .satisfies(e -> { + StatusRuntimeException sre = (StatusRuntimeException) e; + Report.bullet("code : %s", sre.getStatus().getCode()); + Report.bullet("description : %s", sre.getStatus().getDescription()); + + Metadata trailers = sre.getTrailers(); + String reason = trailers == null ? null : trailers.get(OrderServiceImpl.REASON_KEY); + String retryAfter = trailers == null ? null : trailers.get(OrderServiceImpl.RETRY_AFTER_KEY); + Report.bullet("x-failure-reason : %s", reason); + Report.bullet("x-retry-after-seconds : %s", retryAfter); + + assertThat(sre.getStatus().getCode()).isEqualTo(Status.Code.FAILED_PRECONDITION); + assertThat(reason).isEqualTo("ORDER_LOCKED"); + }); + + Report.problem("If the handler had simply thrown IllegalStateException, the caller would"); + Report.problem("have received UNKNOWN with a null description. Nothing actionable at all."); + + Report.fix("Throw StatusRuntimeException with a code that means something:"); + Report.fix(" NOT_FOUND / INVALID_ARGUMENT / FAILED_PRECONDITION -- do NOT retry"); + Report.fix(" UNAVAILABLE / DEADLINE_EXCEEDED / RESOURCE_EXHAUSTED -- retry may help"); + Report.fix("Attach machine-readable context as trailers, not in the description string."); + + Report.takeaway("The status code is the contract. Clients, retry policies, circuit breakers"); + Report.takeaway("and dashboards all key off it, so choosing it carelessly makes every"); + Report.takeaway("downstream behaviour wrong -- most damagingly, it makes non-retryable"); + Report.takeaway("failures look retryable and turns one bad request into a storm."); + } +} diff --git a/src/test/java/com/ankurm/grpc/support/GrpcTestBase.java b/src/test/java/com/ankurm/grpc/support/GrpcTestBase.java new file mode 100644 index 0000000..cdc6a28 --- /dev/null +++ b/src/test/java/com/ankurm/grpc/support/GrpcTestBase.java @@ -0,0 +1,95 @@ +package com.ankurm.grpc.support; + +import com.ankurm.grpc.Application; +import com.ankurm.grpc.orders.OrderServiceGrpc; +import io.grpc.ManagedChannel; +import org.junit.jupiter.api.AfterEach; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.grpc.client.GrpcChannelFactory; +import org.springframework.test.context.TestPropertySource; + +import java.util.ArrayList; +import java.util.List; + +/** + * Base class for every demonstration: boots the application with gRPC on the in-process + * transport and hands out stubs built on channels the test controls. + * + *

Why in-process rather than a real port

+ * The in-process transport is a real gRPC stack -- interceptors, metadata, statuses, deadlines, + * flow control and cancellation all behave normally -- but it skips sockets and TLS. That makes + * tests fast and free of port conflicts, and it is the right default for demonstrating protocol + * behaviour. + * + *

What it does NOT reproduce, and therefore what these tests deliberately cannot show: + * anything about the network. Keepalives, GOAWAY frames, idle-timeout disconnects, load-balancer + * behaviour, TLS negotiation and DNS re-resolution all need a socket. Those are documented in + * {@code docs/06-network-and-keepalive.md} rather than asserted here, because a test that claims to + * prove them over in-process transport would be lying. + * + *

Note also that message size limits are enforced by the in-process transport even though + * no bytes cross a socket, because the limit is applied at the message layer. So the size-limit + * demonstrations are genuine. + */ +@SpringBootTest(classes = Application.class) +@TestPropertySource(properties = { + // Bind gRPC to the in-process transport only. Without this Boot starts a Netty server on + // a real port, which works but makes the suite slower and flakier on CI. + // + // NOTE the exact property names -- they are easy to get wrong and the failure is opaque: + // spring.grpc.server.inprocess.name (NOT "in-process") + // spring.grpc.client.channel..target (SINGULAR "channel"; and "target", NOT "address") + // + // Getting the client one wrong produces no configuration error at all. An unmatched channel + // name is simply passed to the default name resolver as a DNS target, so you get + // UNAVAILABLE: Unable to resolve host orders + // Caused by: java.net.UnknownHostException: orders + // which sends you to look at networking, DNS and service discovery rather than at a typo in + // a property name. Both mistakes in this file cost real time while writing it. + "spring.grpc.server.inprocess.name=orders-test", + "spring.grpc.server.port=-1", + "spring.grpc.client.inprocess.enabled=true", + "spring.grpc.client.channel.orders.target=in-process:orders-test" +}) +public abstract class GrpcTestBase { + + @Autowired + protected GrpcChannelFactory channels; + + private final List opened = new ArrayList<>(); + + /** The in-process target this suite's server listens on. */ + protected static final String TARGET = "in-process:orders-test"; + + /** + * A channel to the test server, tracked so it is shut down after the test. + * + *

Note that this passes the full target rather than the logical channel name. + * {@code GrpcChannelFactory} is a composite: it asks each registered factory whether it + * {@code supports(target)} before any named-channel indirection is applied, so passing a bare + * name to {@code createChannel} in this configuration fails with + * {@code IllegalStateException: No grpc channel factory found that supports target : orders}. + * Named channels are still the right thing in application code, where the property-configured + * target is resolved for you -- see {@code NamedChannelTest}. + */ + protected ManagedChannel channel() { + ManagedChannel channel = channels.createChannel(TARGET); + opened.add(channel); + return channel; + } + + protected OrderServiceGrpc.OrderServiceBlockingStub blockingStub() { + return OrderServiceGrpc.newBlockingStub(channel()); + } + + protected OrderServiceGrpc.OrderServiceStub asyncStub() { + return OrderServiceGrpc.newStub(channel()); + } + + @AfterEach + void closeChannels() { + opened.forEach(ManagedChannel::shutdownNow); + opened.clear(); + } +} diff --git a/src/test/java/com/ankurm/grpc/support/Report.java b/src/test/java/com/ankurm/grpc/support/Report.java new file mode 100644 index 0000000..2649dc2 --- /dev/null +++ b/src/test/java/com/ankurm/grpc/support/Report.java @@ -0,0 +1,38 @@ +package com.ankurm.grpc.support; + +/** Formatting helper so every demonstration prints a readable, diffable transcript. */ +public final class Report { + + private Report() { + } + + public static void title(String text) { + System.out.println(); + System.out.println("=".repeat(78)); + System.out.println(text); + System.out.println("=".repeat(78)); + } + + public static void section(String text) { + System.out.println(); + System.out.println("-- " + text + " " + "-".repeat(Math.max(0, 74 - text.length()))); + } + + public static void bullet(String format, Object... args) { + System.out.printf(" " + format + "%n", args); + } + + /** A problem being reproduced. */ + public static void problem(String format, Object... args) { + System.out.printf(" [PROBLEM] " + format + "%n", args); + } + + /** The fix for the problem just reproduced. */ + public static void fix(String format, Object... args) { + System.out.printf(" [FIX] " + format + "%n", args); + } + + public static void takeaway(String format, Object... args) { + System.out.printf(">> " + format + "%n", args); + } +}