GRPCStreamStateMachine.swift 60 KB

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