GRPCStreamStateMachine.swift 65 KB

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