GRPCStreamStateMachine.swift 49 KB

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