2
0

GRPCStreamStateMachine.swift 55 KB

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