Call.swift 10.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266
  1. /*
  2. *
  3. * Copyright 2016, Google Inc.
  4. * All rights reserved.
  5. *
  6. * Redistribution and use in source and binary forms, with or without
  7. * modification, are permitted provided that the following conditions are
  8. * met:
  9. *
  10. * * Redistributions of source code must retain the above copyright
  11. * notice, this list of conditions and the following disclaimer.
  12. * * Redistributions in binary form must reproduce the above
  13. * copyright notice, this list of conditions and the following disclaimer
  14. * in the documentation and/or other materials provided with the
  15. * distribution.
  16. * * Neither the name of Google Inc. nor the names of its
  17. * contributors may be used to endorse or promote products derived from
  18. * this software without specific prior written permission.
  19. *
  20. * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
  21. * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
  22. * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
  23. * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
  24. * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
  25. * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
  26. * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
  27. * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
  28. * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
  29. * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
  30. * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
  31. *
  32. */
  33. #if SWIFT_PACKAGE
  34. import CgRPC
  35. #endif
  36. import Foundation
  37. /// Singleton class that provides a mutex for synchronizing calls to cgrpc_call_perform()
  38. private class CallLock {
  39. var mutex : Mutex
  40. init() {
  41. mutex = Mutex()
  42. }
  43. static let sharedInstance = CallLock()
  44. }
  45. /// A gRPC API call
  46. public class Call {
  47. /// Pointer to underlying C representation
  48. private var call : UnsafeMutableRawPointer!
  49. /// Completion queue used for call
  50. private var completionQueue: CompletionQueue
  51. /// True if this instance is responsible for deleting the underlying C representation
  52. private var owned : Bool
  53. /// A queue of pending messages to send over the call
  54. private var pendingMessages : Array<Data>
  55. /// True if a message write operation is underway
  56. private var writing : Bool
  57. /// Initializes a Call representation
  58. ///
  59. /// - Parameter call: the underlying C representation
  60. /// - Parameter owned: true if this instance is responsible for deleting the underlying call
  61. init(call: UnsafeMutableRawPointer, owned: Bool, completionQueue: CompletionQueue) {
  62. self.call = call
  63. self.owned = owned
  64. self.completionQueue = completionQueue
  65. self.pendingMessages = []
  66. self.writing = false
  67. }
  68. deinit {
  69. if (owned) {
  70. cgrpc_call_destroy(call)
  71. }
  72. }
  73. /// Initiate performance of a call without waiting for completion
  74. ///
  75. /// - Parameter operations: array of operations to be performed
  76. /// - Parameter tag: integer tag that will be attached to these operations
  77. /// - Returns: the result of initiating the call
  78. func performOperations(operations: OperationGroup,
  79. tag: Int64,
  80. completionQueue: CompletionQueue)
  81. -> grpc_call_error {
  82. let mutex = CallLock.sharedInstance.mutex
  83. mutex.lock()
  84. let error = cgrpc_call_perform(call, operations.operations, tag)
  85. mutex.unlock()
  86. return error
  87. }
  88. /// Performs a nonstreaming gRPC API call
  89. ///
  90. /// - Parameter message: a ByteBuffer containing the message to send
  91. /// - Parameter metadata: metadata to send with the call
  92. /// - Returns: a CallResponse object containing results of the call
  93. public func performNonStreamingCall(messageData: Data,
  94. metadata: Metadata,
  95. completion: @escaping ((CallResponse) -> Void)) -> Void {
  96. let messageBuffer = ByteBuffer(data:messageData)
  97. let operation_sendInitialMetadata = Operation_SendInitialMetadata(metadata:metadata);
  98. let operation_sendMessage = Operation_SendMessage(message:messageBuffer)
  99. let operation_sendCloseFromClient = Operation_SendCloseFromClient()
  100. let operation_receiveInitialMetadata = Operation_ReceiveInitialMetadata()
  101. let operation_receiveStatusOnClient = Operation_ReceiveStatusOnClient()
  102. let operation_receiveMessage = Operation_ReceiveMessage()
  103. let group = OperationGroup(call:self,
  104. operations:[operation_sendInitialMetadata,
  105. operation_sendMessage,
  106. operation_sendCloseFromClient,
  107. operation_receiveInitialMetadata,
  108. operation_receiveStatusOnClient,
  109. operation_receiveMessage])
  110. { (event) in
  111. if (event.type == GRPC_OP_COMPLETE) {
  112. let response = CallResponse(status:operation_receiveStatusOnClient.status(),
  113. statusDetails:operation_receiveStatusOnClient.statusDetails(),
  114. message:operation_receiveMessage.message(),
  115. initialMetadata:operation_receiveInitialMetadata.metadata(),
  116. trailingMetadata:operation_receiveStatusOnClient.metadata())
  117. completion(response)
  118. } else {
  119. completion(CallResponse(completion: event.type))
  120. }
  121. }
  122. let call_error = self.perform(call: self, operations: group)
  123. if call_error != GRPC_CALL_OK {
  124. print ("call error = \(call_error)")
  125. }
  126. }
  127. // perform a group of operations (used internally)
  128. private func perform(call: Call, operations: OperationGroup) -> grpc_call_error {
  129. self.completionQueue.operationGroups[operations.tag] = operations
  130. return call.performOperations(operations:operations,
  131. tag:operations.tag,
  132. completionQueue: self.completionQueue)
  133. }
  134. // start a streaming connection
  135. public func start(metadata: Metadata) {
  136. self.sendInitialMetadata(metadata: metadata)
  137. self.receiveInitialMetadata()
  138. self.receiveStatus()
  139. }
  140. // send a message over a streaming connection
  141. public func sendMessage(data: Data) {
  142. DispatchQueue.main.async {
  143. if self.writing {
  144. self.pendingMessages.append(data)
  145. } else {
  146. self.writing = true
  147. self.sendWithoutBlocking(data: data)
  148. }
  149. }
  150. }
  151. private func sendWithoutBlocking(data: Data) {
  152. let messageBuffer = ByteBuffer(data:data)
  153. let operation_sendMessage = Operation_SendMessage(message:messageBuffer)
  154. let operations = OperationGroup(call:self, operations:[operation_sendMessage])
  155. { (event) in
  156. DispatchQueue.main.async {
  157. if self.pendingMessages.count > 0 {
  158. let nextMessage = self.pendingMessages.first!
  159. self.pendingMessages.removeFirst()
  160. self.sendWithoutBlocking(data: nextMessage)
  161. } else {
  162. self.writing = false
  163. }
  164. }
  165. }
  166. let call_error = self.perform(call:self, operations:operations)
  167. if call_error != GRPC_CALL_OK {
  168. print("call error \(call_error)")
  169. }
  170. }
  171. // receive a message over a streaming connection
  172. public func receiveMessage(callback:@escaping ((Data!) -> Void)) -> Void {
  173. let operation_receiveMessage = Operation_ReceiveMessage()
  174. let operations = OperationGroup(call:self, operations:[operation_receiveMessage])
  175. { (event) in
  176. if let messageBuffer = operation_receiveMessage.message() {
  177. callback(messageBuffer.data())
  178. }
  179. }
  180. let call_error = self.perform(call:self, operations:operations)
  181. if call_error != GRPC_CALL_OK {
  182. print("call error \(call_error)")
  183. }
  184. }
  185. // send initial metadata over a streaming connection
  186. private func sendInitialMetadata(metadata: Metadata) {
  187. let operation_sendInitialMetadata = Operation_SendInitialMetadata(metadata:metadata);
  188. let operations = OperationGroup(call:self, operations:[operation_sendInitialMetadata])
  189. { (event) in
  190. if (event.type == GRPC_OP_COMPLETE) {
  191. print("call status \(event.type) \(event.tag)")
  192. } else {
  193. return
  194. }
  195. }
  196. let call_error = self.perform(call:self, operations:operations)
  197. if call_error != GRPC_CALL_OK {
  198. print("call error: \(call_error)")
  199. }
  200. }
  201. // receive initial metadata from a streaming connection
  202. private func receiveInitialMetadata() {
  203. let operation_receiveInitialMetadata = Operation_ReceiveInitialMetadata()
  204. let operations = OperationGroup(call:self, operations:[operation_receiveInitialMetadata])
  205. { (event) in
  206. let initialMetadata = operation_receiveInitialMetadata.metadata()
  207. for j in 0..<initialMetadata.count() {
  208. print("Received initial metadata -> " + initialMetadata.key(index:j) + " : " + initialMetadata.value(index:j))
  209. }
  210. }
  211. let call_error = self.perform(call:self, operations:operations)
  212. if call_error != GRPC_CALL_OK {
  213. print("call error \(call_error)")
  214. }
  215. }
  216. // receive status from a streaming connection
  217. private func receiveStatus() {
  218. let operation_receiveStatus = Operation_ReceiveStatusOnClient()
  219. let operations = OperationGroup(call:self,
  220. operations:[operation_receiveStatus])
  221. { (event) in
  222. print("status = \(operation_receiveStatus.status()), \(operation_receiveStatus.statusDetails())")
  223. }
  224. let call_error = self.perform(call:self, operations:operations)
  225. if call_error != GRPC_CALL_OK {
  226. print("call error \(call_error)")
  227. }
  228. }
  229. // close a streaming connection
  230. public func close(completion:@escaping (() -> Void)) {
  231. let operation_sendCloseFromClient = Operation_SendCloseFromClient()
  232. let operations = OperationGroup(call:self, operations:[operation_sendCloseFromClient])
  233. { (event) in
  234. completion()
  235. }
  236. let call_error = self.perform(call:self, operations:operations)
  237. if call_error != GRPC_CALL_OK {
  238. print("call error \(call_error)")
  239. }
  240. }
  241. }