GRPCStreamStateMachine.swift 66 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946
  1. /*
  2. * Copyright 2024, gRPC Authors All rights reserved.
  3. *
  4. * Licensed under the Apache License, Version 2.0 (the "License");
  5. * you may not use this file except in compliance with the License.
  6. * You may obtain a copy of the License at
  7. *
  8. * http://www.apache.org/licenses/LICENSE-2.0
  9. *
  10. * Unless required by applicable law or agreed to in writing, software
  11. * distributed under the License is distributed on an "AS IS" BASIS,
  12. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  13. * See the License for the specific language governing permissions and
  14. * limitations under the License.
  15. */
  16. import GRPCCore
  17. import NIOCore
  18. import NIOHPACK
  19. import NIOHTTP1
  20. enum GRPCStreamStateMachineConfiguration {
  21. case client(ClientConfiguration)
  22. case server(ServerConfiguration)
  23. enum Scheme: String {
  24. case http
  25. case https
  26. }
  27. struct ClientConfiguration {
  28. var methodDescriptor: MethodDescriptor
  29. var scheme: Scheme
  30. var outboundEncoding: CompressionAlgorithm
  31. var acceptedEncodings: CompressionAlgorithmSet
  32. init(
  33. methodDescriptor: MethodDescriptor,
  34. scheme: Scheme,
  35. outboundEncoding: CompressionAlgorithm,
  36. acceptedEncodings: CompressionAlgorithmSet
  37. ) {
  38. self.methodDescriptor = methodDescriptor
  39. self.scheme = scheme
  40. self.outboundEncoding = outboundEncoding
  41. self.acceptedEncodings = acceptedEncodings.union(.none)
  42. }
  43. }
  44. struct ServerConfiguration {
  45. var scheme: Scheme
  46. var acceptedEncodings: CompressionAlgorithmSet
  47. init(scheme: Scheme, acceptedEncodings: CompressionAlgorithmSet) {
  48. self.scheme = scheme
  49. self.acceptedEncodings = acceptedEncodings.union(.none)
  50. }
  51. }
  52. }
  53. private enum GRPCStreamStateMachineState {
  54. case clientIdleServerIdle(ClientIdleServerIdleState)
  55. case clientOpenServerIdle(ClientOpenServerIdleState)
  56. case clientOpenServerOpen(ClientOpenServerOpenState)
  57. case clientOpenServerClosed(ClientOpenServerClosedState)
  58. case clientClosedServerIdle(ClientClosedServerIdleState)
  59. case clientClosedServerOpen(ClientClosedServerOpenState)
  60. case clientClosedServerClosed(ClientClosedServerClosedState)
  61. case _modifying
  62. struct ClientIdleServerIdleState {
  63. let maximumPayloadSize: Int
  64. }
  65. struct ClientOpenServerIdleState {
  66. let maximumPayloadSize: Int
  67. var framer: GRPCMessageFramer
  68. var compressor: Zlib.Compressor?
  69. var outboundCompression: CompressionAlgorithm
  70. // The deframer must be optional because the client will not have one configured
  71. // until the server opens and sends a grpc-encoding header.
  72. // It will be present for the server though, because even though it's idle,
  73. // it can still receive compressed messages from the client.
  74. let deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>?
  75. var decompressor: Zlib.Decompressor?
  76. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  77. // Store the headers received from the remote peer, its storage can be reused when sending
  78. // headers back to the remote peer.
  79. var headers: HPACKHeaders
  80. init(
  81. previousState: ClientIdleServerIdleState,
  82. compressor: Zlib.Compressor?,
  83. outboundCompression: CompressionAlgorithm,
  84. framer: GRPCMessageFramer,
  85. decompressor: Zlib.Decompressor?,
  86. deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>?,
  87. headers: HPACKHeaders
  88. ) {
  89. self.maximumPayloadSize = previousState.maximumPayloadSize
  90. self.compressor = compressor
  91. self.outboundCompression = outboundCompression
  92. self.framer = framer
  93. self.decompressor = decompressor
  94. self.deframer = deframer
  95. self.inboundMessageBuffer = .init()
  96. self.headers = headers
  97. }
  98. }
  99. struct ClientOpenServerOpenState {
  100. var framer: GRPCMessageFramer
  101. var compressor: Zlib.Compressor?
  102. var outboundCompression: CompressionAlgorithm
  103. let deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>
  104. var decompressor: Zlib.Decompressor?
  105. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  106. // Store the headers received from the remote peer, its storage can be reused when sending
  107. // headers back to the remote peer.
  108. var headers: HPACKHeaders
  109. init(
  110. previousState: ClientOpenServerIdleState,
  111. deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>,
  112. decompressor: Zlib.Decompressor?
  113. ) {
  114. self.framer = previousState.framer
  115. self.compressor = previousState.compressor
  116. self.outboundCompression = previousState.outboundCompression
  117. self.deframer = deframer
  118. self.decompressor = decompressor
  119. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  120. self.headers = previousState.headers
  121. }
  122. }
  123. struct ClientOpenServerClosedState {
  124. var framer: GRPCMessageFramer?
  125. var compressor: Zlib.Compressor?
  126. var outboundCompression: CompressionAlgorithm
  127. let deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>?
  128. var decompressor: Zlib.Decompressor?
  129. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  130. // This transition should only happen on the server-side when, upon receiving
  131. // initial client metadata, some of the headers are invalid and we must reject
  132. // the RPC.
  133. // We will mark the client as open (because it sent initial metadata albeit
  134. // invalid) but we'll close the server, meaning all future messages sent from
  135. // the client will be ignored. Because of this, we won't need to frame or
  136. // deframe any messages, as we won't be reading or writing any messages.
  137. init(previousState: ClientIdleServerIdleState) {
  138. self.framer = nil
  139. self.compressor = nil
  140. self.outboundCompression = .none
  141. self.deframer = nil
  142. self.decompressor = nil
  143. self.inboundMessageBuffer = .init()
  144. }
  145. init(previousState: ClientOpenServerOpenState) {
  146. self.framer = previousState.framer
  147. self.compressor = previousState.compressor
  148. self.outboundCompression = previousState.outboundCompression
  149. self.deframer = previousState.deframer
  150. self.decompressor = previousState.decompressor
  151. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  152. }
  153. init(previousState: ClientOpenServerIdleState) {
  154. self.framer = previousState.framer
  155. self.compressor = previousState.compressor
  156. self.outboundCompression = previousState.outboundCompression
  157. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  158. // The server went directly from idle to closed - this means it sent a
  159. // trailers-only response:
  160. // - if we're the client, the previous state was a nil deframer, but that
  161. // is okay because we don't need a deframer as the server won't be sending
  162. // any messages;
  163. // - if we're the server, we'll keep whatever deframer we had.
  164. self.deframer = previousState.deframer
  165. self.decompressor = previousState.decompressor
  166. }
  167. }
  168. struct ClientClosedServerIdleState {
  169. let maximumPayloadSize: Int
  170. var framer: GRPCMessageFramer
  171. var compressor: Zlib.Compressor?
  172. var outboundCompression: CompressionAlgorithm
  173. let deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>?
  174. var decompressor: Zlib.Decompressor?
  175. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  176. // Store the headers received from the remote peer, its storage can be reused when sending
  177. // headers back to the remote peer.
  178. var headers: HPACKHeaders
  179. /// This transition should only happen on the client-side.
  180. /// It can happen if the request times out before the client outbound can be opened, or if the stream is
  181. /// unexpectedly closed for some other reason on the client before it can transition to open.
  182. init(previousState: ClientIdleServerIdleState) {
  183. self.maximumPayloadSize = previousState.maximumPayloadSize
  184. // We don't need a compressor since we won't be sending any messages.
  185. self.framer = GRPCMessageFramer()
  186. self.compressor = nil
  187. self.outboundCompression = .none
  188. // We haven't received anything from the server.
  189. self.deframer = nil
  190. self.decompressor = nil
  191. self.inboundMessageBuffer = .init()
  192. self.headers = [:]
  193. }
  194. /// This transition should only happen on the server-side.
  195. /// We are closing the client as soon as it opens (i.e., endStream was set when receiving the client's
  196. /// initial metadata). We don't need to know a decompression algorithm, since we won't receive
  197. /// any more messages from the client anyways, as it's closed.
  198. init(
  199. previousState: ClientIdleServerIdleState,
  200. compressionAlgorithm: CompressionAlgorithm,
  201. headers: HPACKHeaders
  202. ) {
  203. self.maximumPayloadSize = previousState.maximumPayloadSize
  204. if let zlibMethod = Zlib.Method(encoding: compressionAlgorithm) {
  205. self.compressor = Zlib.Compressor(method: zlibMethod)
  206. self.outboundCompression = compressionAlgorithm
  207. } else {
  208. self.compressor = nil
  209. self.outboundCompression = .none
  210. }
  211. self.framer = GRPCMessageFramer()
  212. // We don't need a deframer since we won't receive any messages from the
  213. // client: it's closed.
  214. self.deframer = nil
  215. self.inboundMessageBuffer = .init()
  216. self.headers = headers
  217. }
  218. init(previousState: ClientOpenServerIdleState) {
  219. self.maximumPayloadSize = previousState.maximumPayloadSize
  220. self.framer = previousState.framer
  221. self.compressor = previousState.compressor
  222. self.outboundCompression = previousState.outboundCompression
  223. self.deframer = previousState.deframer
  224. self.decompressor = previousState.decompressor
  225. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  226. self.headers = previousState.headers
  227. }
  228. }
  229. struct ClientClosedServerOpenState {
  230. var framer: GRPCMessageFramer
  231. var compressor: Zlib.Compressor?
  232. var outboundCompression: CompressionAlgorithm
  233. let deframer: NIOSingleStepByteToMessageProcessor<GRPCMessageDeframer>?
  234. var decompressor: Zlib.Decompressor?
  235. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  236. // Store the headers received from the remote peer, its storage can be reused when sending
  237. // headers back to the remote peer.
  238. var headers: HPACKHeaders
  239. init(previousState: ClientOpenServerOpenState) {
  240. self.framer = previousState.framer
  241. self.compressor = previousState.compressor
  242. self.outboundCompression = previousState.outboundCompression
  243. self.deframer = previousState.deframer
  244. self.decompressor = previousState.decompressor
  245. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  246. self.headers = previousState.headers
  247. }
  248. /// This should be called from the server path, as the deframer will already be configured in this scenario.
  249. init(previousState: ClientClosedServerIdleState) {
  250. self.framer = previousState.framer
  251. self.compressor = previousState.compressor
  252. self.outboundCompression = previousState.outboundCompression
  253. // In the case of the server, we don't need to deframe/decompress any more
  254. // messages, since the client's closed.
  255. self.deframer = nil
  256. self.decompressor = nil
  257. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  258. self.headers = previousState.headers
  259. }
  260. /// This should only be called from the client path, as the deframer has not yet been set up.
  261. init(
  262. previousState: ClientClosedServerIdleState,
  263. decompressionAlgorithm: CompressionAlgorithm
  264. ) {
  265. self.framer = previousState.framer
  266. self.compressor = previousState.compressor
  267. self.outboundCompression = previousState.outboundCompression
  268. // In the case of the client, it will only be able to set up the deframer
  269. // after it receives the chosen encoding from the server.
  270. if let zlibMethod = Zlib.Method(encoding: decompressionAlgorithm) {
  271. self.decompressor = Zlib.Decompressor(method: zlibMethod)
  272. }
  273. let decoder = GRPCMessageDeframer(
  274. maximumPayloadSize: previousState.maximumPayloadSize,
  275. decompressor: self.decompressor
  276. )
  277. self.deframer = NIOSingleStepByteToMessageProcessor(decoder)
  278. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  279. self.headers = previousState.headers
  280. }
  281. }
  282. struct ClientClosedServerClosedState {
  283. // We still need the framer and compressor in case the server has closed
  284. // but its buffer is not yet empty and still needs to send messages out to
  285. // the client.
  286. var framer: GRPCMessageFramer?
  287. var compressor: Zlib.Compressor?
  288. var outboundCompression: CompressionAlgorithm
  289. // These are already deframed, so we don't need the deframer anymore.
  290. var inboundMessageBuffer: OneOrManyQueue<[UInt8]>
  291. // This transition should only happen on the server-side when, upon receiving
  292. // initial client metadata, some of the headers are invalid and we must reject
  293. // the RPC.
  294. // We will mark the client as closed (because it set the EOS flag, even if
  295. // the initial metadata was invalid) and we'll close the server too.
  296. // Because of this, we won't need to frame any messages, as we
  297. // won't be writing any messages.
  298. init(previousState: ClientIdleServerIdleState) {
  299. self.framer = nil
  300. self.compressor = nil
  301. self.outboundCompression = .none
  302. self.inboundMessageBuffer = .init()
  303. }
  304. init(previousState: ClientClosedServerOpenState) {
  305. self.framer = previousState.framer
  306. self.compressor = previousState.compressor
  307. self.outboundCompression = previousState.outboundCompression
  308. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  309. }
  310. init(previousState: ClientClosedServerIdleState) {
  311. self.framer = previousState.framer
  312. self.compressor = previousState.compressor
  313. self.outboundCompression = previousState.outboundCompression
  314. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  315. }
  316. init(previousState: ClientOpenServerIdleState) {
  317. self.framer = previousState.framer
  318. self.compressor = previousState.compressor
  319. self.outboundCompression = previousState.outboundCompression
  320. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  321. }
  322. init(previousState: ClientOpenServerOpenState) {
  323. self.framer = previousState.framer
  324. self.compressor = previousState.compressor
  325. self.outboundCompression = previousState.outboundCompression
  326. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  327. }
  328. init(previousState: ClientOpenServerClosedState) {
  329. self.framer = previousState.framer
  330. self.compressor = previousState.compressor
  331. self.outboundCompression = previousState.outboundCompression
  332. self.inboundMessageBuffer = previousState.inboundMessageBuffer
  333. }
  334. }
  335. }
  336. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  337. struct GRPCStreamStateMachine {
  338. private var state: GRPCStreamStateMachineState
  339. private var configuration: GRPCStreamStateMachineConfiguration
  340. private var skipAssertions: Bool
  341. struct InvalidState: Error {
  342. var message: String
  343. init(_ message: String) {
  344. self.message = message
  345. }
  346. }
  347. init(
  348. configuration: GRPCStreamStateMachineConfiguration,
  349. maximumPayloadSize: Int,
  350. skipAssertions: Bool = false
  351. ) {
  352. self.state = .clientIdleServerIdle(.init(maximumPayloadSize: maximumPayloadSize))
  353. self.configuration = configuration
  354. self.skipAssertions = skipAssertions
  355. }
  356. mutating func send(metadata: Metadata) throws(InvalidState) -> HPACKHeaders {
  357. switch self.configuration {
  358. case .client(let clientConfiguration):
  359. return try self.clientSend(metadata: metadata, configuration: clientConfiguration)
  360. case .server(let serverConfiguration):
  361. return try self.serverSend(metadata: metadata, configuration: serverConfiguration)
  362. }
  363. }
  364. mutating func send(message: [UInt8], promise: EventLoopPromise<Void>?) throws(InvalidState) {
  365. switch self.configuration {
  366. case .client:
  367. try self.clientSend(message: message, promise: promise)
  368. case .server:
  369. try self.serverSend(message: message, promise: promise)
  370. }
  371. }
  372. mutating func closeOutbound() throws(InvalidState) {
  373. switch self.configuration {
  374. case .client:
  375. try self.clientCloseOutbound()
  376. case .server:
  377. try self.invalidState("Server cannot call close: it must send status and trailers.")
  378. }
  379. }
  380. mutating func send(
  381. status: Status,
  382. metadata: Metadata
  383. ) throws(InvalidState) -> HPACKHeaders {
  384. switch self.configuration {
  385. case .client:
  386. try self.invalidState(
  387. "Client cannot send status and trailer."
  388. )
  389. case .server:
  390. return try self.serverSend(
  391. status: status,
  392. customMetadata: metadata
  393. )
  394. }
  395. }
  396. enum OnMetadataReceived: Equatable {
  397. case receivedMetadata(Metadata, MethodDescriptor?)
  398. case doNothing
  399. // Client-specific actions
  400. case receivedStatusAndMetadata_clientOnly(status: Status, metadata: Metadata)
  401. // Server-specific actions
  402. case rejectRPC_serverOnly(trailers: HPACKHeaders)
  403. case protocolViolation_serverOnly
  404. }
  405. mutating func receive(
  406. headers: HPACKHeaders,
  407. endStream: Bool
  408. ) throws(InvalidState) -> OnMetadataReceived {
  409. switch self.configuration {
  410. case .client(let clientConfiguration):
  411. return try self.clientReceive(
  412. headers: headers,
  413. endStream: endStream,
  414. configuration: clientConfiguration
  415. )
  416. case .server(let serverConfiguration):
  417. return try self.serverReceive(
  418. headers: headers,
  419. endStream: endStream,
  420. configuration: serverConfiguration
  421. )
  422. }
  423. }
  424. enum OnBufferReceivedAction: Equatable {
  425. case readInbound
  426. case doNothing
  427. // This will be returned when the server sends a data frame with EOS set.
  428. // This is invalid as per the protocol specification, because the server
  429. // can only close by sending trailers, not by setting EOS when sending
  430. // a message.
  431. case endRPCAndForwardErrorStatus_clientOnly(Status)
  432. case forwardErrorAndClose_serverOnly(RPCError)
  433. }
  434. mutating func receive(
  435. buffer: ByteBuffer,
  436. endStream: Bool
  437. ) throws(InvalidState) -> OnBufferReceivedAction {
  438. switch self.configuration {
  439. case .client:
  440. return try self.clientReceive(buffer: buffer, endStream: endStream)
  441. case .server:
  442. return try self.serverReceive(buffer: buffer, endStream: endStream)
  443. }
  444. }
  445. /// The result of requesting the next outbound frame, which may contain multiple messages.
  446. enum OnNextOutboundFrame {
  447. /// Either the receiving party is closed, so we shouldn't send any more frames; or the sender is done
  448. /// writing messages (i.e. we are now closed).
  449. case noMoreMessages
  450. /// There isn't a frame ready to be sent, but we could still receive more messages, so keep trying.
  451. case awaitMoreMessages
  452. /// A frame is ready to be sent.
  453. case sendFrame(
  454. frame: ByteBuffer,
  455. promise: EventLoopPromise<Void>?
  456. )
  457. case closeAndFailPromise(EventLoopPromise<Void>?, RPCError)
  458. init(result: Result<ByteBuffer, RPCError>, promise: EventLoopPromise<Void>?) {
  459. switch result {
  460. case .success(let buffer):
  461. self = .sendFrame(frame: buffer, promise: promise)
  462. case .failure(let error):
  463. self = .closeAndFailPromise(promise, error)
  464. }
  465. }
  466. }
  467. mutating func nextOutboundFrame() throws(InvalidState) -> OnNextOutboundFrame {
  468. switch self.configuration {
  469. case .client:
  470. return try self.clientNextOutboundFrame()
  471. case .server:
  472. return try self.serverNextOutboundFrame()
  473. }
  474. }
  475. /// The result of requesting the next inbound message.
  476. enum OnNextInboundMessage: Equatable {
  477. /// The sender is done writing messages and there are no more messages to be received.
  478. case noMoreMessages
  479. /// There isn't a message ready to be sent, but we could still receive more, so keep trying.
  480. case awaitMoreMessages
  481. /// A message has been received.
  482. case receiveMessage([UInt8])
  483. }
  484. mutating func nextInboundMessage() -> OnNextInboundMessage {
  485. switch self.configuration {
  486. case .client:
  487. return self.clientNextInboundMessage()
  488. case .server:
  489. return self.serverNextInboundMessage()
  490. }
  491. }
  492. mutating func tearDown() {
  493. switch self.state {
  494. case .clientIdleServerIdle:
  495. ()
  496. case .clientOpenServerIdle(let state):
  497. state.compressor?.end()
  498. state.decompressor?.end()
  499. case .clientOpenServerOpen(let state):
  500. state.compressor?.end()
  501. state.decompressor?.end()
  502. case .clientOpenServerClosed(let state):
  503. state.compressor?.end()
  504. state.decompressor?.end()
  505. case .clientClosedServerIdle(let state):
  506. state.compressor?.end()
  507. state.decompressor?.end()
  508. case .clientClosedServerOpen(let state):
  509. state.compressor?.end()
  510. state.decompressor?.end()
  511. case .clientClosedServerClosed(let state):
  512. state.compressor?.end()
  513. case ._modifying:
  514. preconditionFailure()
  515. }
  516. }
  517. enum OnUnexpectedInboundClose {
  518. case forwardStatus_clientOnly(Status)
  519. case fireError_serverOnly(any Error)
  520. case doNothing
  521. init(serverCloseReason: UnexpectedInboundCloseReason) {
  522. switch serverCloseReason {
  523. case .streamReset, .channelInactive:
  524. self = .fireError_serverOnly(RPCError(serverCloseReason))
  525. case .errorThrown(let error):
  526. self = .fireError_serverOnly(error)
  527. }
  528. }
  529. }
  530. enum UnexpectedInboundCloseReason {
  531. case streamReset
  532. case channelInactive
  533. case errorThrown(any Error)
  534. }
  535. mutating func unexpectedInboundClose(
  536. reason: UnexpectedInboundCloseReason
  537. ) -> OnUnexpectedInboundClose {
  538. switch self.configuration {
  539. case .client:
  540. return self.clientUnexpectedInboundClose(reason: reason)
  541. case .server:
  542. return self.serverUnexpectedInboundClose(reason: reason)
  543. }
  544. }
  545. }
  546. // - MARK: Client
  547. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  548. extension GRPCStreamStateMachine {
  549. private func makeClientHeaders(
  550. methodDescriptor: MethodDescriptor,
  551. scheme: GRPCStreamStateMachineConfiguration.Scheme,
  552. outboundEncoding: CompressionAlgorithm?,
  553. acceptedEncodings: CompressionAlgorithmSet,
  554. customMetadata: Metadata
  555. ) -> HPACKHeaders {
  556. var headers = HPACKHeaders()
  557. headers.reserveCapacity(7 + customMetadata.count)
  558. // Add required headers.
  559. // See https://github.com/grpc/grpc/blob/master/doc/PROTOCOL-HTTP2.md#requests
  560. // The order is important here: reserved HTTP2 headers (those starting with `:`)
  561. // must come before all other headers.
  562. headers.add("POST", forKey: .method)
  563. headers.add(scheme.rawValue, forKey: .scheme)
  564. headers.add(methodDescriptor.path, forKey: .path)
  565. // Add required gRPC headers.
  566. headers.add(ContentType.grpc.canonicalValue, forKey: .contentType)
  567. headers.add("trailers", forKey: .te) // Used to detect incompatible proxies
  568. if let encoding = outboundEncoding, encoding != .none {
  569. headers.add(encoding.name, forKey: .encoding)
  570. }
  571. for encoding in acceptedEncodings.elements.filter({ $0 != .none }) {
  572. headers.add(encoding.name, forKey: .acceptEncoding)
  573. }
  574. for metadataPair in customMetadata {
  575. headers.add(name: metadataPair.key, value: metadataPair.value.encoded())
  576. }
  577. return headers
  578. }
  579. private mutating func clientSend(
  580. metadata: Metadata,
  581. configuration: GRPCStreamStateMachineConfiguration.ClientConfiguration
  582. ) throws(InvalidState) -> HPACKHeaders {
  583. // Client sends metadata only when opening the stream.
  584. switch self.state {
  585. case .clientIdleServerIdle(let state):
  586. let outboundEncoding = configuration.outboundEncoding
  587. let compressor = Zlib.Method(encoding: outboundEncoding)
  588. .flatMap { Zlib.Compressor(method: $0) }
  589. self.state = .clientOpenServerIdle(
  590. .init(
  591. previousState: state,
  592. compressor: compressor,
  593. outboundCompression: outboundEncoding,
  594. framer: GRPCMessageFramer(),
  595. decompressor: nil,
  596. deframer: nil,
  597. headers: [:]
  598. )
  599. )
  600. return self.makeClientHeaders(
  601. methodDescriptor: configuration.methodDescriptor,
  602. scheme: configuration.scheme,
  603. outboundEncoding: configuration.outboundEncoding,
  604. acceptedEncodings: configuration.acceptedEncodings,
  605. customMetadata: metadata
  606. )
  607. case .clientOpenServerIdle, .clientOpenServerOpen, .clientOpenServerClosed:
  608. try self.invalidState(
  609. "Client is already open: shouldn't be sending metadata."
  610. )
  611. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  612. try self.invalidState(
  613. "Client is closed: can't send metadata."
  614. )
  615. case ._modifying:
  616. preconditionFailure()
  617. }
  618. }
  619. private mutating func clientSend(
  620. message: [UInt8],
  621. promise: EventLoopPromise<Void>?
  622. ) throws(InvalidState) {
  623. switch self.state {
  624. case .clientIdleServerIdle:
  625. try self.invalidState("Client not yet open.")
  626. case .clientOpenServerIdle(var state):
  627. self.state = ._modifying
  628. state.framer.append(message, promise: promise)
  629. self.state = .clientOpenServerIdle(state)
  630. case .clientOpenServerOpen(var state):
  631. self.state = ._modifying
  632. state.framer.append(message, promise: promise)
  633. self.state = .clientOpenServerOpen(state)
  634. case .clientOpenServerClosed:
  635. // The server has closed, so it makes no sense to send the rest of the request.
  636. ()
  637. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  638. try self.invalidState(
  639. "Client is closed, cannot send a message."
  640. )
  641. case ._modifying:
  642. preconditionFailure()
  643. }
  644. }
  645. private mutating func clientCloseOutbound() throws(InvalidState) {
  646. switch self.state {
  647. case .clientIdleServerIdle(let state):
  648. self.state = .clientClosedServerIdle(.init(previousState: state))
  649. case .clientOpenServerIdle(let state):
  650. self.state = .clientClosedServerIdle(.init(previousState: state))
  651. case .clientOpenServerOpen(let state):
  652. self.state = .clientClosedServerOpen(.init(previousState: state))
  653. case .clientOpenServerClosed(let state):
  654. self.state = .clientClosedServerClosed(.init(previousState: state))
  655. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  656. // Client is already closed - nothing to do.
  657. ()
  658. case ._modifying:
  659. preconditionFailure()
  660. }
  661. }
  662. /// Returns the client's next request to the server.
  663. /// - Returns: The request to be made to the server.
  664. private mutating func clientNextOutboundFrame() throws(InvalidState) -> OnNextOutboundFrame {
  665. switch self.state {
  666. case .clientIdleServerIdle:
  667. try self.invalidState("Client is not open yet.")
  668. case .clientOpenServerIdle(var state):
  669. self.state = ._modifying
  670. let next = state.framer.nextResult(compressor: state.compressor)
  671. self.state = .clientOpenServerIdle(state)
  672. if let next = next {
  673. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  674. } else {
  675. return .awaitMoreMessages
  676. }
  677. case .clientOpenServerOpen(var state):
  678. self.state = ._modifying
  679. let next = state.framer.nextResult(compressor: state.compressor)
  680. self.state = .clientOpenServerOpen(state)
  681. if let next = next {
  682. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  683. } else {
  684. return .awaitMoreMessages
  685. }
  686. case .clientClosedServerIdle(var state):
  687. self.state = ._modifying
  688. let next = state.framer.nextResult(compressor: state.compressor)
  689. self.state = .clientClosedServerIdle(state)
  690. if let next = next {
  691. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  692. } else {
  693. return .noMoreMessages
  694. }
  695. case .clientClosedServerOpen(var state):
  696. self.state = ._modifying
  697. let next = state.framer.nextResult(compressor: state.compressor)
  698. self.state = .clientClosedServerOpen(state)
  699. if let next = next {
  700. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  701. } else {
  702. return .noMoreMessages
  703. }
  704. case .clientOpenServerClosed, .clientClosedServerClosed:
  705. // No point in sending any more requests if the server is closed.
  706. return .noMoreMessages
  707. case ._modifying:
  708. preconditionFailure()
  709. }
  710. }
  711. private enum ServerHeadersValidationResult {
  712. case valid
  713. case invalid(OnMetadataReceived)
  714. }
  715. private mutating func clientValidateHeadersReceivedFromServer(
  716. _ metadata: HPACKHeaders
  717. ) -> ServerHeadersValidationResult {
  718. var httpStatus: String? {
  719. metadata.firstString(forKey: .status)
  720. }
  721. var grpcStatus: Status.Code? {
  722. metadata.firstString(forKey: .grpcStatus)
  723. .flatMap { Int($0) }
  724. .flatMap { Status.Code(rawValue: $0) }
  725. }
  726. guard httpStatus == "200" || grpcStatus != nil else {
  727. let httpStatusCode =
  728. httpStatus
  729. .flatMap { Int($0) }
  730. .map { HTTPResponseStatus(statusCode: $0) }
  731. guard let httpStatusCode else {
  732. return .invalid(
  733. .receivedStatusAndMetadata_clientOnly(
  734. status: .init(code: .unknown, message: "HTTP Status Code is missing."),
  735. metadata: Metadata(headers: metadata)
  736. )
  737. )
  738. }
  739. if (100 ... 199).contains(httpStatusCode.code) {
  740. // For 1xx status codes, the entire header should be skipped and a
  741. // subsequent header should be read.
  742. // See https://github.com/grpc/grpc/blob/master/doc/http-grpc-status-mapping.md
  743. return .invalid(.doNothing)
  744. }
  745. // Forward the mapped status code.
  746. return .invalid(
  747. .receivedStatusAndMetadata_clientOnly(
  748. status: .init(
  749. code: Status.Code(httpStatusCode: httpStatusCode),
  750. message: "Unexpected non-200 HTTP Status Code."
  751. ),
  752. metadata: Metadata(headers: metadata)
  753. )
  754. )
  755. }
  756. let contentTypeHeader = metadata.first(name: GRPCHTTP2Keys.contentType.rawValue)
  757. guard contentTypeHeader.flatMap(ContentType.init) != nil else {
  758. return .invalid(
  759. .receivedStatusAndMetadata_clientOnly(
  760. status: .init(
  761. code: .internalError,
  762. message: "Missing \(GRPCHTTP2Keys.contentType.rawValue) header"
  763. ),
  764. metadata: Metadata(headers: metadata)
  765. )
  766. )
  767. }
  768. return .valid
  769. }
  770. private enum ProcessInboundEncodingResult {
  771. case error(OnMetadataReceived)
  772. case success(CompressionAlgorithm)
  773. }
  774. private func processInboundEncoding(
  775. headers: HPACKHeaders,
  776. configuration: GRPCStreamStateMachineConfiguration.ClientConfiguration
  777. ) -> ProcessInboundEncodingResult {
  778. let inboundEncoding: CompressionAlgorithm
  779. if let serverEncoding = headers.first(name: GRPCHTTP2Keys.encoding.rawValue) {
  780. guard let parsedEncoding = CompressionAlgorithm(name: serverEncoding),
  781. configuration.acceptedEncodings.contains(parsedEncoding)
  782. else {
  783. return .error(
  784. .receivedStatusAndMetadata_clientOnly(
  785. status: .init(
  786. code: .internalError,
  787. message:
  788. "The server picked a compression algorithm ('\(serverEncoding)') the client does not know about."
  789. ),
  790. metadata: Metadata(headers: headers)
  791. )
  792. )
  793. }
  794. inboundEncoding = parsedEncoding
  795. } else {
  796. inboundEncoding = .none
  797. }
  798. return .success(inboundEncoding)
  799. }
  800. private func validateTrailers(
  801. _ trailers: HPACKHeaders
  802. ) throws(InvalidState) -> OnMetadataReceived {
  803. let statusValue = trailers.firstString(forKey: .grpcStatus)
  804. let statusCode = statusValue.flatMap {
  805. Int($0)
  806. }.flatMap {
  807. Status.Code(rawValue: $0)
  808. }
  809. let status: Status
  810. if let code = statusCode {
  811. let messageFieldValue = trailers.firstString(forKey: .grpcStatusMessage, canonicalForm: false)
  812. let message = messageFieldValue.map { GRPCStatusMessageMarshaller.unmarshall($0) } ?? ""
  813. status = Status(code: code, message: message)
  814. } else {
  815. let message: String
  816. if let statusValue = statusValue {
  817. message = "Invalid 'grpc-status' in trailers (\(statusValue))"
  818. } else {
  819. message = "No 'grpc-status' value in trailers"
  820. }
  821. status = Status(code: .unknown, message: message)
  822. }
  823. var convertedMetadata = Metadata(headers: trailers)
  824. convertedMetadata.removeAllValues(forKey: GRPCHTTP2Keys.grpcStatus.rawValue)
  825. convertedMetadata.removeAllValues(forKey: GRPCHTTP2Keys.grpcStatusMessage.rawValue)
  826. return .receivedStatusAndMetadata_clientOnly(status: status, metadata: convertedMetadata)
  827. }
  828. private mutating func clientReceive(
  829. headers: HPACKHeaders,
  830. endStream: Bool,
  831. configuration: GRPCStreamStateMachineConfiguration.ClientConfiguration
  832. ) throws(InvalidState) -> OnMetadataReceived {
  833. switch self.state {
  834. case .clientOpenServerIdle(let state):
  835. switch (self.clientValidateHeadersReceivedFromServer(headers), endStream) {
  836. case (.invalid(let action), true):
  837. // The headers are invalid, but the server signalled that it was
  838. // closing the stream, so close both client and server.
  839. self.state = .clientClosedServerClosed(.init(previousState: state))
  840. return action
  841. case (.invalid(let action), false):
  842. self.state = .clientClosedServerIdle(.init(previousState: state))
  843. return action
  844. case (.valid, true):
  845. // This is a trailers-only response: close server.
  846. self.state = .clientOpenServerClosed(.init(previousState: state))
  847. return try self.validateTrailers(headers)
  848. case (.valid, false):
  849. switch self.processInboundEncoding(headers: headers, configuration: configuration) {
  850. case .error(let failure):
  851. return failure
  852. case .success(let inboundEncoding):
  853. let decompressor = Zlib.Method(encoding: inboundEncoding)
  854. .flatMap { Zlib.Decompressor(method: $0) }
  855. let deframer = GRPCMessageDeframer(
  856. maximumPayloadSize: state.maximumPayloadSize,
  857. decompressor: decompressor
  858. )
  859. self.state = .clientOpenServerOpen(
  860. .init(
  861. previousState: state,
  862. deframer: NIOSingleStepByteToMessageProcessor(deframer),
  863. decompressor: decompressor
  864. )
  865. )
  866. return .receivedMetadata(Metadata(headers: headers), nil)
  867. }
  868. }
  869. case .clientOpenServerOpen(let state):
  870. // This state is valid even if endStream is not set: server can send
  871. // trailing metadata without END_STREAM set, and follow it with an
  872. // empty message frame where it is set.
  873. // However, we must make sure that grpc-status is set, otherwise this
  874. // is an invalid state.
  875. if endStream {
  876. self.state = .clientOpenServerClosed(.init(previousState: state))
  877. }
  878. return try self.validateTrailers(headers)
  879. case .clientClosedServerIdle(let state):
  880. switch (self.clientValidateHeadersReceivedFromServer(headers), endStream) {
  881. case (.invalid(let action), true):
  882. // The headers are invalid, but the server signalled that it was
  883. // closing the stream, so close the server side too.
  884. self.state = .clientClosedServerClosed(.init(previousState: state))
  885. return action
  886. case (.invalid(let action), false):
  887. // Client is already closed, so we don't need to update our state.
  888. return action
  889. case (.valid, true):
  890. // This is a trailers-only response: close server.
  891. self.state = .clientClosedServerClosed(.init(previousState: state))
  892. return try self.validateTrailers(headers)
  893. case (.valid, false):
  894. switch self.processInboundEncoding(headers: headers, configuration: configuration) {
  895. case .error(let failure):
  896. return failure
  897. case .success(let inboundEncoding):
  898. self.state = .clientClosedServerOpen(
  899. .init(
  900. previousState: state,
  901. decompressionAlgorithm: inboundEncoding
  902. )
  903. )
  904. return .receivedMetadata(Metadata(headers: headers), nil)
  905. }
  906. }
  907. case .clientClosedServerOpen(let state):
  908. // This state is valid even if endStream is not set: server can send
  909. // trailing metadata without END_STREAM set, and follow it with an
  910. // empty message frame where it is set.
  911. // However, we must make sure that grpc-status is set, otherwise this
  912. // is an invalid state.
  913. if endStream {
  914. self.state = .clientClosedServerClosed(.init(previousState: state))
  915. }
  916. return try self.validateTrailers(headers)
  917. case .clientClosedServerClosed:
  918. // We could end up here if we received a grpc-status header in a previous
  919. // frame (which would have already close the server) and then we receive
  920. // an empty frame with EOS set.
  921. // We wouldn't want to throw in that scenario, so we just ignore it.
  922. // Note that we don't want to ignore it if EOS is not set here though, as
  923. // then it would be an invalid payload.
  924. if !endStream || headers.count > 0 {
  925. try self.invalidState(
  926. "Server is closed, nothing could have been sent."
  927. )
  928. }
  929. return .doNothing
  930. case .clientIdleServerIdle:
  931. try self.invalidState(
  932. "Server cannot have sent metadata if the client is idle."
  933. )
  934. case .clientOpenServerClosed:
  935. try self.invalidState(
  936. "Server is closed, nothing could have been sent."
  937. )
  938. case ._modifying:
  939. preconditionFailure()
  940. }
  941. }
  942. private mutating func clientReceive(
  943. buffer: ByteBuffer,
  944. endStream: Bool
  945. ) throws(InvalidState) -> OnBufferReceivedAction {
  946. // This is a message received by the client, from the server.
  947. switch self.state {
  948. case .clientIdleServerIdle:
  949. try self.invalidState(
  950. "Cannot have received anything from server if client is not yet open."
  951. )
  952. case .clientOpenServerIdle, .clientClosedServerIdle:
  953. try self.invalidState(
  954. "Server cannot have sent a message before sending the initial metadata."
  955. )
  956. case .clientOpenServerOpen(var state):
  957. self.state = ._modifying
  958. if endStream {
  959. // This is invalid as per the protocol specification, because the server
  960. // can only close by sending trailers, not by setting EOS when sending
  961. // a message.
  962. self.state = .clientClosedServerClosed(.init(previousState: state))
  963. return .endRPCAndForwardErrorStatus_clientOnly(
  964. Status(
  965. code: .internalError,
  966. message: """
  967. Server sent EOS alongside a data frame, but server is only allowed \
  968. to close by sending status and trailers.
  969. """
  970. )
  971. )
  972. }
  973. do {
  974. try state.deframer.process(buffer: buffer) { deframedMessage in
  975. state.inboundMessageBuffer.append(deframedMessage)
  976. }
  977. self.state = .clientOpenServerOpen(state)
  978. return .readInbound
  979. } catch {
  980. self.state = .clientOpenServerOpen(state)
  981. let status = Status(code: .internalError, message: "Failed to decode message")
  982. return .endRPCAndForwardErrorStatus_clientOnly(status)
  983. }
  984. case .clientClosedServerOpen(var state):
  985. self.state = ._modifying
  986. if endStream {
  987. self.state = .clientClosedServerClosed(.init(previousState: state))
  988. return .endRPCAndForwardErrorStatus_clientOnly(
  989. Status(
  990. code: .internalError,
  991. message: """
  992. Server sent EOS alongside a data frame, but server is only allowed \
  993. to close by sending status and trailers.
  994. """
  995. )
  996. )
  997. }
  998. // The client may have sent the end stream and thus it's closed,
  999. // but the server may still be responding.
  1000. // The client must have a deframer set up, so force-unwrap is okay.
  1001. do {
  1002. try state.deframer!.process(buffer: buffer) { deframedMessage in
  1003. state.inboundMessageBuffer.append(deframedMessage)
  1004. }
  1005. self.state = .clientClosedServerOpen(state)
  1006. return .readInbound
  1007. } catch {
  1008. self.state = .clientClosedServerOpen(state)
  1009. let status = Status(code: .internalError, message: "Failed to decode message")
  1010. return .endRPCAndForwardErrorStatus_clientOnly(status)
  1011. }
  1012. case .clientOpenServerClosed, .clientClosedServerClosed:
  1013. try self.invalidState(
  1014. "Cannot have received anything from a closed server."
  1015. )
  1016. case ._modifying:
  1017. preconditionFailure()
  1018. }
  1019. }
  1020. private mutating func clientNextInboundMessage() -> OnNextInboundMessage {
  1021. switch self.state {
  1022. case .clientOpenServerOpen(var state):
  1023. self.state = ._modifying
  1024. let message = state.inboundMessageBuffer.pop()
  1025. self.state = .clientOpenServerOpen(state)
  1026. return message.map { .receiveMessage($0) } ?? .awaitMoreMessages
  1027. case .clientOpenServerClosed(var state):
  1028. self.state = ._modifying
  1029. let message = state.inboundMessageBuffer.pop()
  1030. self.state = .clientOpenServerClosed(state)
  1031. return message.map { .receiveMessage($0) } ?? .noMoreMessages
  1032. case .clientClosedServerOpen(var state):
  1033. self.state = ._modifying
  1034. let message = state.inboundMessageBuffer.pop()
  1035. self.state = .clientClosedServerOpen(state)
  1036. return message.map { .receiveMessage($0) } ?? .awaitMoreMessages
  1037. case .clientClosedServerClosed(var state):
  1038. self.state = ._modifying
  1039. let message = state.inboundMessageBuffer.pop()
  1040. self.state = .clientClosedServerClosed(state)
  1041. return message.map { .receiveMessage($0) } ?? .noMoreMessages
  1042. case .clientIdleServerIdle,
  1043. .clientOpenServerIdle,
  1044. .clientClosedServerIdle:
  1045. return .awaitMoreMessages
  1046. case ._modifying:
  1047. preconditionFailure()
  1048. }
  1049. }
  1050. private func invalidState(_ message: String, line: UInt = #line) throws(InvalidState) -> Never {
  1051. if !self.skipAssertions {
  1052. assertionFailure(message, line: line)
  1053. }
  1054. throw InvalidState(message)
  1055. }
  1056. private mutating func clientUnexpectedInboundClose(
  1057. reason: UnexpectedInboundCloseReason
  1058. ) -> OnUnexpectedInboundClose {
  1059. switch self.state {
  1060. case .clientIdleServerIdle(let state):
  1061. self.state = .clientClosedServerClosed(.init(previousState: state))
  1062. return .forwardStatus_clientOnly(Status(RPCError(reason)))
  1063. case .clientOpenServerIdle(let state):
  1064. self.state = .clientClosedServerClosed(.init(previousState: state))
  1065. return .forwardStatus_clientOnly(Status(RPCError(reason)))
  1066. case .clientClosedServerIdle(let state):
  1067. self.state = .clientClosedServerClosed(.init(previousState: state))
  1068. return .forwardStatus_clientOnly(Status(RPCError(reason)))
  1069. case .clientOpenServerOpen(let state):
  1070. self.state = .clientClosedServerClosed(.init(previousState: state))
  1071. return .forwardStatus_clientOnly(Status(RPCError(reason)))
  1072. case .clientClosedServerOpen(let state):
  1073. self.state = .clientClosedServerClosed(.init(previousState: state))
  1074. return .forwardStatus_clientOnly(Status(RPCError(reason)))
  1075. case .clientOpenServerClosed, .clientClosedServerClosed:
  1076. return .doNothing
  1077. case ._modifying:
  1078. preconditionFailure()
  1079. }
  1080. }
  1081. }
  1082. // - MARK: Server
  1083. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  1084. extension GRPCStreamStateMachine {
  1085. private func formResponseHeaders(
  1086. in headers: inout HPACKHeaders,
  1087. outboundEncoding: CompressionAlgorithm?,
  1088. configuration: GRPCStreamStateMachineConfiguration.ServerConfiguration,
  1089. customMetadata: Metadata
  1090. ) {
  1091. headers.removeAll(keepingCapacity: true)
  1092. // Response headers always contain :status (HTTP Status 200) and content-type.
  1093. // They may also contain grpc-encoding, grpc-accept-encoding, and custom metadata.
  1094. headers.reserveCapacity(4 + customMetadata.count)
  1095. headers.add("200", forKey: .status)
  1096. headers.add(ContentType.grpc.canonicalValue, forKey: .contentType)
  1097. if let outboundEncoding, outboundEncoding != .none {
  1098. headers.add(outboundEncoding.name, forKey: .encoding)
  1099. }
  1100. for metadataPair in customMetadata {
  1101. headers.add(name: metadataPair.key, value: metadataPair.value.encoded())
  1102. }
  1103. }
  1104. private mutating func serverSend(
  1105. metadata: Metadata,
  1106. configuration: GRPCStreamStateMachineConfiguration.ServerConfiguration
  1107. ) throws(InvalidState) -> HPACKHeaders {
  1108. // Server sends initial metadata
  1109. switch self.state {
  1110. case .clientOpenServerIdle(var state):
  1111. self.state = ._modifying
  1112. let outboundEncoding = state.outboundCompression
  1113. self.formResponseHeaders(
  1114. in: &state.headers,
  1115. outboundEncoding: outboundEncoding,
  1116. configuration: configuration,
  1117. customMetadata: metadata
  1118. )
  1119. self.state = .clientOpenServerOpen(
  1120. .init(
  1121. previousState: state,
  1122. // In the case of the server, it will already have a deframer set up,
  1123. // because it already knows what encoding the client is using:
  1124. // it's okay to force-unwrap.
  1125. deframer: state.deframer!,
  1126. decompressor: state.decompressor
  1127. )
  1128. )
  1129. return state.headers
  1130. case .clientClosedServerIdle(var state):
  1131. self.state = ._modifying
  1132. let outboundEncoding = state.outboundCompression
  1133. self.formResponseHeaders(
  1134. in: &state.headers,
  1135. outboundEncoding: outboundEncoding,
  1136. configuration: configuration,
  1137. customMetadata: metadata
  1138. )
  1139. self.state = .clientClosedServerOpen(.init(previousState: state))
  1140. return state.headers
  1141. case .clientIdleServerIdle:
  1142. try self.invalidState(
  1143. "Client cannot be idle if server is sending initial metadata: it must have opened."
  1144. )
  1145. case .clientOpenServerClosed, .clientClosedServerClosed:
  1146. try self.invalidState(
  1147. "Server cannot send metadata if closed."
  1148. )
  1149. case .clientOpenServerOpen, .clientClosedServerOpen:
  1150. try self.invalidState(
  1151. "Server has already sent initial metadata."
  1152. )
  1153. case ._modifying:
  1154. preconditionFailure()
  1155. }
  1156. }
  1157. private mutating func serverSend(
  1158. message: [UInt8],
  1159. promise: EventLoopPromise<Void>?
  1160. ) throws(InvalidState) {
  1161. switch self.state {
  1162. case .clientIdleServerIdle, .clientOpenServerIdle, .clientClosedServerIdle:
  1163. try self.invalidState(
  1164. "Server must have sent initial metadata before sending a message."
  1165. )
  1166. case .clientOpenServerOpen(var state):
  1167. self.state = ._modifying
  1168. state.framer.append(message, promise: promise)
  1169. self.state = .clientOpenServerOpen(state)
  1170. case .clientClosedServerOpen(var state):
  1171. self.state = ._modifying
  1172. state.framer.append(message, promise: promise)
  1173. self.state = .clientClosedServerOpen(state)
  1174. case .clientOpenServerClosed, .clientClosedServerClosed:
  1175. try self.invalidState(
  1176. "Server can't send a message if it's closed."
  1177. )
  1178. case ._modifying:
  1179. preconditionFailure()
  1180. }
  1181. }
  1182. private mutating func serverSend(
  1183. status: Status,
  1184. customMetadata: Metadata
  1185. ) throws(InvalidState) -> HPACKHeaders {
  1186. // Close the server.
  1187. switch self.state {
  1188. case .clientOpenServerOpen(var state):
  1189. self.state = ._modifying
  1190. state.headers.formTrailers(status: status, metadata: customMetadata)
  1191. self.state = .clientOpenServerClosed(.init(previousState: state))
  1192. return state.headers
  1193. case .clientClosedServerOpen(var state):
  1194. self.state = ._modifying
  1195. state.headers.formTrailers(status: status, metadata: customMetadata)
  1196. self.state = .clientClosedServerClosed(.init(previousState: state))
  1197. return state.headers
  1198. case .clientOpenServerIdle(var state):
  1199. self.state = ._modifying
  1200. state.headers.formTrailersOnly(status: status, metadata: customMetadata)
  1201. self.state = .clientOpenServerClosed(.init(previousState: state))
  1202. return state.headers
  1203. case .clientClosedServerIdle(var state):
  1204. self.state = ._modifying
  1205. state.headers.formTrailersOnly(status: status, metadata: customMetadata)
  1206. self.state = .clientClosedServerClosed(.init(previousState: state))
  1207. return state.headers
  1208. case .clientIdleServerIdle:
  1209. try self.invalidState(
  1210. "Server can't send status if client is idle."
  1211. )
  1212. case .clientOpenServerClosed, .clientClosedServerClosed:
  1213. try self.invalidState(
  1214. "Server can't send anything if closed."
  1215. )
  1216. case ._modifying:
  1217. preconditionFailure()
  1218. }
  1219. }
  1220. private mutating func serverReceive(
  1221. headers: HPACKHeaders,
  1222. endStream: Bool,
  1223. configuration: GRPCStreamStateMachineConfiguration.ServerConfiguration
  1224. ) throws(InvalidState) -> OnMetadataReceived {
  1225. func closeServer(
  1226. from state: GRPCStreamStateMachineState.ClientIdleServerIdleState,
  1227. endStream: Bool
  1228. ) -> GRPCStreamStateMachineState {
  1229. if endStream {
  1230. return .clientClosedServerClosed(.init(previousState: state))
  1231. } else {
  1232. return .clientOpenServerClosed(.init(previousState: state))
  1233. }
  1234. }
  1235. switch self.state {
  1236. case .clientIdleServerIdle(let state):
  1237. let contentType = headers.firstString(forKey: .contentType)
  1238. .flatMap { ContentType(value: $0) }
  1239. if contentType == nil {
  1240. self.state = .clientOpenServerClosed(.init(previousState: state))
  1241. // Respond with HTTP-level Unsupported Media Type status code.
  1242. var trailers = HPACKHeaders()
  1243. trailers.add("415", forKey: .status)
  1244. return .rejectRPC_serverOnly(trailers: trailers)
  1245. }
  1246. guard let pathHeader = headers.firstString(forKey: .path) else {
  1247. self.state = closeServer(from: state, endStream: endStream)
  1248. return .rejectRPC_serverOnly(
  1249. trailers: .trailersOnly(
  1250. code: .invalidArgument,
  1251. message: "No \(GRPCHTTP2Keys.path.rawValue) header has been set."
  1252. )
  1253. )
  1254. }
  1255. guard let path = MethodDescriptor(path: pathHeader) else {
  1256. self.state = closeServer(from: state, endStream: endStream)
  1257. return .rejectRPC_serverOnly(
  1258. trailers: .trailersOnly(
  1259. code: .unimplemented,
  1260. message:
  1261. "The given \(GRPCHTTP2Keys.path.rawValue) (\(pathHeader)) does not correspond to a valid method."
  1262. )
  1263. )
  1264. }
  1265. let scheme = headers.firstString(forKey: .scheme)
  1266. .flatMap { GRPCStreamStateMachineConfiguration.Scheme(rawValue: $0) }
  1267. if scheme == nil {
  1268. self.state = closeServer(from: state, endStream: endStream)
  1269. return .rejectRPC_serverOnly(
  1270. trailers: .trailersOnly(
  1271. code: .invalidArgument,
  1272. message: ":scheme header must be present and one of \"http\" or \"https\"."
  1273. )
  1274. )
  1275. }
  1276. guard let method = headers.firstString(forKey: .method), method == "POST" else {
  1277. self.state = closeServer(from: state, endStream: endStream)
  1278. return .rejectRPC_serverOnly(
  1279. trailers: .trailersOnly(
  1280. code: .invalidArgument,
  1281. message: ":method header is expected to be present and have a value of \"POST\"."
  1282. )
  1283. )
  1284. }
  1285. guard let te = headers.firstString(forKey: .te), te == "trailers" else {
  1286. self.state = closeServer(from: state, endStream: endStream)
  1287. return .rejectRPC_serverOnly(
  1288. trailers: .trailersOnly(
  1289. code: .invalidArgument,
  1290. message: "\"te\" header is expected to be present and have a value of \"trailers\"."
  1291. )
  1292. )
  1293. }
  1294. // Firstly, find out if we support the client's chosen encoding, and reject
  1295. // the RPC if we don't.
  1296. let inboundEncoding: CompressionAlgorithm
  1297. let encodingValues = headers.values(
  1298. forHeader: GRPCHTTP2Keys.encoding.rawValue,
  1299. canonicalForm: true
  1300. )
  1301. var encodingValuesIterator = encodingValues.makeIterator()
  1302. if let rawEncoding = encodingValuesIterator.next() {
  1303. guard encodingValuesIterator.next() == nil else {
  1304. self.state = closeServer(from: state, endStream: endStream)
  1305. return .rejectRPC_serverOnly(
  1306. trailers: .trailersOnly(
  1307. code: .internalError,
  1308. message: "\(GRPCHTTP2Keys.encoding) must contain no more than one value."
  1309. )
  1310. )
  1311. }
  1312. guard let clientEncoding = CompressionAlgorithm(name: rawEncoding),
  1313. configuration.acceptedEncodings.contains(clientEncoding)
  1314. else {
  1315. self.state = closeServer(from: state, endStream: endStream)
  1316. var trailers = HPACKHeaders.trailersOnly(
  1317. code: .unimplemented,
  1318. message: """
  1319. \(rawEncoding) compression is not supported; \
  1320. supported algorithms are listed in grpc-accept-encoding
  1321. """
  1322. )
  1323. for acceptedEncoding in configuration.acceptedEncodings.elements {
  1324. trailers.add(name: GRPCHTTP2Keys.acceptEncoding.rawValue, value: acceptedEncoding.name)
  1325. }
  1326. return .rejectRPC_serverOnly(trailers: trailers)
  1327. }
  1328. // Server supports client's encoding.
  1329. inboundEncoding = clientEncoding
  1330. } else {
  1331. inboundEncoding = .none
  1332. }
  1333. // Secondly, find a compatible encoding the server can use to compress outbound messages,
  1334. // based on the encodings the client has advertised.
  1335. var outboundEncoding: CompressionAlgorithm = .none
  1336. let clientAdvertisedEncodings = headers.values(
  1337. forHeader: GRPCHTTP2Keys.acceptEncoding.rawValue,
  1338. canonicalForm: true
  1339. )
  1340. // Find the preferred encoding and use it to compress responses.
  1341. for clientAdvertisedEncoding in clientAdvertisedEncodings {
  1342. if let algorithm = CompressionAlgorithm(name: clientAdvertisedEncoding),
  1343. configuration.acceptedEncodings.contains(algorithm)
  1344. {
  1345. outboundEncoding = algorithm
  1346. break
  1347. }
  1348. }
  1349. if endStream {
  1350. self.state = .clientClosedServerIdle(
  1351. .init(
  1352. previousState: state,
  1353. compressionAlgorithm: outboundEncoding,
  1354. headers: headers
  1355. )
  1356. )
  1357. } else {
  1358. let compressor = Zlib.Method(encoding: outboundEncoding)
  1359. .flatMap { Zlib.Compressor(method: $0) }
  1360. let decompressor = Zlib.Method(encoding: inboundEncoding)
  1361. .flatMap { Zlib.Decompressor(method: $0) }
  1362. let deframer = GRPCMessageDeframer(
  1363. maximumPayloadSize: state.maximumPayloadSize,
  1364. decompressor: decompressor
  1365. )
  1366. self.state = .clientOpenServerIdle(
  1367. .init(
  1368. previousState: state,
  1369. compressor: compressor,
  1370. outboundCompression: outboundEncoding,
  1371. framer: GRPCMessageFramer(),
  1372. decompressor: decompressor,
  1373. deframer: NIOSingleStepByteToMessageProcessor(deframer),
  1374. headers: headers
  1375. )
  1376. )
  1377. }
  1378. return .receivedMetadata(Metadata(headers: headers), path)
  1379. case .clientOpenServerIdle, .clientOpenServerOpen, .clientOpenServerClosed:
  1380. // Metadata has already been received, should only be sent once by clients.
  1381. return .protocolViolation_serverOnly
  1382. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  1383. try self.invalidState("Client can't have sent metadata if closed.")
  1384. case ._modifying:
  1385. preconditionFailure()
  1386. }
  1387. }
  1388. private mutating func serverReceive(
  1389. buffer: ByteBuffer,
  1390. endStream: Bool
  1391. ) throws(InvalidState) -> OnBufferReceivedAction {
  1392. let action: OnBufferReceivedAction
  1393. switch self.state {
  1394. case .clientIdleServerIdle:
  1395. try self.invalidState("Can't have received a message if client is idle.")
  1396. case .clientOpenServerIdle(var state):
  1397. self.state = ._modifying
  1398. // Deframer must be present on the server side, as we know the decompression
  1399. // algorithm from the moment the client opens.
  1400. do {
  1401. try state.deframer!.process(buffer: buffer) { deframedMessage in
  1402. state.inboundMessageBuffer.append(deframedMessage)
  1403. }
  1404. action = .readInbound
  1405. } catch {
  1406. let error = RPCError(code: .internalError, message: "Failed to decode message")
  1407. action = .forwardErrorAndClose_serverOnly(error)
  1408. }
  1409. if endStream {
  1410. self.state = .clientClosedServerIdle(.init(previousState: state))
  1411. } else {
  1412. self.state = .clientOpenServerIdle(state)
  1413. }
  1414. case .clientOpenServerOpen(var state):
  1415. self.state = ._modifying
  1416. do {
  1417. try state.deframer.process(buffer: buffer) { deframedMessage in
  1418. state.inboundMessageBuffer.append(deframedMessage)
  1419. }
  1420. action = .readInbound
  1421. } catch {
  1422. let error = RPCError(code: .internalError, message: "Failed to decode message")
  1423. action = .forwardErrorAndClose_serverOnly(error)
  1424. }
  1425. if endStream {
  1426. self.state = .clientClosedServerOpen(.init(previousState: state))
  1427. } else {
  1428. self.state = .clientOpenServerOpen(state)
  1429. }
  1430. case .clientOpenServerClosed(let state):
  1431. // Client is not done sending request, but server has already closed.
  1432. // Ignore the rest of the request: do nothing, unless endStream is set,
  1433. // in which case close the client.
  1434. if endStream {
  1435. self.state = .clientClosedServerClosed(.init(previousState: state))
  1436. }
  1437. action = .doNothing
  1438. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  1439. try self.invalidState("Client can't send a message if closed.")
  1440. case ._modifying:
  1441. preconditionFailure()
  1442. }
  1443. return action
  1444. }
  1445. private mutating func serverNextOutboundFrame() throws(InvalidState) -> OnNextOutboundFrame {
  1446. switch self.state {
  1447. case .clientIdleServerIdle, .clientOpenServerIdle, .clientClosedServerIdle:
  1448. try self.invalidState("Server is not open yet.")
  1449. case .clientOpenServerOpen(var state):
  1450. self.state = ._modifying
  1451. let next = state.framer.nextResult(compressor: state.compressor)
  1452. self.state = .clientOpenServerOpen(state)
  1453. if let next = next {
  1454. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  1455. } else {
  1456. return .awaitMoreMessages
  1457. }
  1458. case .clientClosedServerOpen(var state):
  1459. self.state = ._modifying
  1460. let next = state.framer.nextResult(compressor: state.compressor)
  1461. self.state = .clientClosedServerOpen(state)
  1462. if let next = next {
  1463. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  1464. } else {
  1465. return .awaitMoreMessages
  1466. }
  1467. case .clientOpenServerClosed(var state):
  1468. self.state = ._modifying
  1469. let next = state.framer?.nextResult(compressor: state.compressor)
  1470. self.state = .clientOpenServerClosed(state)
  1471. if let next = next {
  1472. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  1473. } else {
  1474. return .noMoreMessages
  1475. }
  1476. case .clientClosedServerClosed(var state):
  1477. self.state = ._modifying
  1478. let next = state.framer?.nextResult(compressor: state.compressor)
  1479. self.state = .clientClosedServerClosed(state)
  1480. if let next = next {
  1481. return OnNextOutboundFrame(result: next.result, promise: next.promise)
  1482. } else {
  1483. return .noMoreMessages
  1484. }
  1485. case ._modifying:
  1486. preconditionFailure()
  1487. }
  1488. }
  1489. private mutating func serverNextInboundMessage() -> OnNextInboundMessage {
  1490. switch self.state {
  1491. case .clientOpenServerIdle(var state):
  1492. self.state = ._modifying
  1493. let request = state.inboundMessageBuffer.pop()
  1494. self.state = .clientOpenServerIdle(state)
  1495. return request.map { .receiveMessage($0) } ?? .awaitMoreMessages
  1496. case .clientOpenServerOpen(var state):
  1497. self.state = ._modifying
  1498. let request = state.inboundMessageBuffer.pop()
  1499. self.state = .clientOpenServerOpen(state)
  1500. return request.map { .receiveMessage($0) } ?? .awaitMoreMessages
  1501. case .clientClosedServerIdle(var state):
  1502. self.state = ._modifying
  1503. let request = state.inboundMessageBuffer.pop()
  1504. self.state = .clientClosedServerIdle(state)
  1505. return request.map { .receiveMessage($0) } ?? .noMoreMessages
  1506. case .clientClosedServerOpen(var state):
  1507. self.state = ._modifying
  1508. let request = state.inboundMessageBuffer.pop()
  1509. self.state = .clientClosedServerOpen(state)
  1510. return request.map { .receiveMessage($0) } ?? .noMoreMessages
  1511. case .clientOpenServerClosed, .clientClosedServerClosed:
  1512. // Server has closed, no need to read.
  1513. return .noMoreMessages
  1514. case .clientIdleServerIdle:
  1515. return .awaitMoreMessages
  1516. case ._modifying:
  1517. preconditionFailure()
  1518. }
  1519. }
  1520. private mutating func serverUnexpectedInboundClose(
  1521. reason: UnexpectedInboundCloseReason
  1522. ) -> OnUnexpectedInboundClose {
  1523. switch self.state {
  1524. case .clientIdleServerIdle(let state):
  1525. self.state = .clientClosedServerClosed(.init(previousState: state))
  1526. return OnUnexpectedInboundClose(serverCloseReason: reason)
  1527. case .clientOpenServerIdle(let state):
  1528. self.state = .clientClosedServerClosed(.init(previousState: state))
  1529. return OnUnexpectedInboundClose(serverCloseReason: reason)
  1530. case .clientOpenServerOpen(let state):
  1531. self.state = .clientClosedServerClosed(.init(previousState: state))
  1532. return OnUnexpectedInboundClose(serverCloseReason: reason)
  1533. case .clientOpenServerClosed(let state):
  1534. self.state = .clientClosedServerClosed(.init(previousState: state))
  1535. return OnUnexpectedInboundClose(serverCloseReason: reason)
  1536. case .clientClosedServerIdle, .clientClosedServerOpen, .clientClosedServerClosed:
  1537. return .doNothing
  1538. case ._modifying:
  1539. preconditionFailure()
  1540. }
  1541. }
  1542. }
  1543. extension MethodDescriptor {
  1544. init?(path: String) {
  1545. var view = path[...]
  1546. guard view.popFirst() == "/" else { return nil }
  1547. // Find the index of the "/" separating the service and method names.
  1548. guard var index = view.firstIndex(of: "/") else { return nil }
  1549. let service = String(view[..<index])
  1550. view.formIndex(after: &index)
  1551. let method = String(view[index...])
  1552. self.init(service: service, method: method)
  1553. }
  1554. }
  1555. internal enum GRPCHTTP2Keys: String {
  1556. case path = ":path"
  1557. case contentType = "content-type"
  1558. case encoding = "grpc-encoding"
  1559. case acceptEncoding = "grpc-accept-encoding"
  1560. case scheme = ":scheme"
  1561. case method = ":method"
  1562. case te = "te"
  1563. case status = ":status"
  1564. case grpcStatus = "grpc-status"
  1565. case grpcStatusMessage = "grpc-message"
  1566. }
  1567. extension HPACKHeaders {
  1568. func firstString(forKey key: GRPCHTTP2Keys, canonicalForm: Bool = true) -> String? {
  1569. self.values(forHeader: key.rawValue, canonicalForm: canonicalForm).first(where: { _ in true })
  1570. .map {
  1571. String($0)
  1572. }
  1573. }
  1574. fileprivate mutating func add(_ value: String, forKey key: GRPCHTTP2Keys) {
  1575. self.add(name: key.rawValue, value: value)
  1576. }
  1577. fileprivate static func trailersOnly(code: Status.Code, message: String) -> Self {
  1578. var trailers = HPACKHeaders()
  1579. HPACKHeaders.formTrailers(
  1580. &trailers,
  1581. isTrailersOnly: true,
  1582. status: Status(code: code, message: message),
  1583. metadata: [:]
  1584. )
  1585. return trailers
  1586. }
  1587. fileprivate mutating func formTrailersOnly(status: Status, metadata: Metadata = [:]) {
  1588. Self.formTrailers(&self, isTrailersOnly: true, status: status, metadata: metadata)
  1589. }
  1590. fileprivate mutating func formTrailers(status: Status, metadata: Metadata = [:]) {
  1591. Self.formTrailers(&self, isTrailersOnly: false, status: status, metadata: metadata)
  1592. }
  1593. private static func formTrailers(
  1594. _ trailers: inout HPACKHeaders,
  1595. isTrailersOnly: Bool,
  1596. status: Status,
  1597. metadata: Metadata
  1598. ) {
  1599. trailers.removeAll(keepingCapacity: true)
  1600. if isTrailersOnly {
  1601. trailers.reserveCapacity(4 + metadata.count)
  1602. trailers.add("200", forKey: .status)
  1603. trailers.add(ContentType.grpc.canonicalValue, forKey: .contentType)
  1604. } else {
  1605. trailers.reserveCapacity(2 + metadata.count)
  1606. }
  1607. trailers.add(String(status.code.rawValue), forKey: .grpcStatus)
  1608. if !status.message.isEmpty, let encoded = GRPCStatusMessageMarshaller.marshall(status.message) {
  1609. trailers.add(encoded, forKey: .grpcStatusMessage)
  1610. }
  1611. for (key, value) in metadata {
  1612. trailers.add(name: key, value: value.encoded())
  1613. }
  1614. }
  1615. }
  1616. extension Zlib.Method {
  1617. init?(encoding: CompressionAlgorithm) {
  1618. switch encoding {
  1619. case .none:
  1620. return nil
  1621. case .deflate:
  1622. self = .deflate
  1623. case .gzip:
  1624. self = .gzip
  1625. default:
  1626. return nil
  1627. }
  1628. }
  1629. }
  1630. extension Metadata {
  1631. init(headers: HPACKHeaders) {
  1632. var metadata = Metadata()
  1633. metadata.reserveCapacity(headers.count)
  1634. for header in headers {
  1635. if header.name.hasSuffix("-bin") {
  1636. do {
  1637. let decodedBinary = try header.value.base64Decoded()
  1638. metadata.addBinary(decodedBinary, forKey: header.name)
  1639. } catch {
  1640. metadata.addString(header.value, forKey: header.name)
  1641. }
  1642. } else {
  1643. metadata.addString(header.value, forKey: header.name)
  1644. }
  1645. }
  1646. self = metadata
  1647. }
  1648. }
  1649. extension Status.Code {
  1650. // See https://github.com/grpc/grpc/blob/master/doc/http-grpc-status-mapping.md
  1651. init(httpStatusCode: HTTPResponseStatus) {
  1652. switch httpStatusCode {
  1653. case .badRequest:
  1654. self = .internalError
  1655. case .unauthorized:
  1656. self = .unauthenticated
  1657. case .forbidden:
  1658. self = .permissionDenied
  1659. case .notFound:
  1660. self = .unimplemented
  1661. case .tooManyRequests, .badGateway, .serviceUnavailable, .gatewayTimeout:
  1662. self = .unavailable
  1663. default:
  1664. self = .unknown
  1665. }
  1666. }
  1667. }
  1668. extension MethodDescriptor {
  1669. var path: String {
  1670. return "/\(self.service)/\(self.method)"
  1671. }
  1672. }
  1673. extension RPCError {
  1674. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  1675. fileprivate init(_ reason: GRPCStreamStateMachine.UnexpectedInboundCloseReason) {
  1676. switch reason {
  1677. case .streamReset:
  1678. self = RPCError(
  1679. code: .unavailable,
  1680. message: "Stream unexpectedly closed: a RST_STREAM frame was received."
  1681. )
  1682. case .channelInactive:
  1683. self = RPCError(code: .unavailable, message: "Stream unexpectedly closed.")
  1684. case .errorThrown:
  1685. self = RPCError(code: .unavailable, message: "Stream unexpectedly closed with error.")
  1686. }
  1687. }
  1688. }
  1689. extension Status {
  1690. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  1691. fileprivate init(_ error: RPCError) {
  1692. self = Status(code: Status.Code(error.code), message: error.message)
  1693. }
  1694. }
  1695. extension RPCError {
  1696. @available(macOS 13.0, iOS 16.0, watchOS 9.0, tvOS 16.0, *)
  1697. init(_ invalidState: GRPCStreamStateMachine.InvalidState) {
  1698. self = RPCError(code: .internalError, message: "Invalid state", cause: invalidState)
  1699. }
  1700. }