2
0

GRPCStreamStateMachine.swift 56 KB

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