GRPCStreamStateMachine.swift 57 KB

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