GRPCStreamStateMachine.swift 55 KB

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