PluginProbe
Media Cloud Sync / 1.4.2
Media Cloud Sync v1.4.2
1.4.2 1.4.1 1.4.0 1.3.12 1.3.11 1.3.10 trunk 1.0.0 1.0.1 1.0.2 1.0.3 1.1.0 1.1.1 1.2.0 1.2.10 1.2.11 1.2.12 1.2.13 1.2.2 1.2.3 1.2.4 1.2.5 1.2.6 1.2.7 1.2.8 All 36 releases
← All changes | includes/sdk/google/google/gax/src/BidiStream.php +44 -17 1.2.7 → 1.4.2 View file →
@@ -29,34 +29,42 @@
29 29 * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
30 30 * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
31 31 * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
32 32 */
33 -namespace Dudlewebs\WPMCS\Google\ApiCore;
33 +namespace Dudlewebs\WPMCS\GCP\Google\ApiCore;
34 34
35 -use Dudlewebs\WPMCS\Google\Rpc\Code;
36 -use Dudlewebs\WPMCS\Grpc\BidiStreamingCall;
35 +use Dudlewebs\WPMCS\GCP\Google\Auth\Logging\LoggingTrait;
36 +use Dudlewebs\WPMCS\GCP\Google\Auth\Logging\RpcLogEvent;
37 +use Dudlewebs\WPMCS\GCP\Google\Protobuf\Internal\Message;
38 +use Dudlewebs\WPMCS\GCP\Google\Rpc\Code;
39 +use Dudlewebs\WPMCS\GCP\Grpc\BidiStreamingCall;
40 +use Dudlewebs\WPMCS\GCP\Psr\Log\LoggerInterface;
37 41 /**
38 42 * BidiStream is the response object from a gRPC bidirectional streaming API call.
39 43 */
40 44 class BidiStream
41 45 {
46 + use LoggingTrait;
42 47 private $call;
43 48 private $isComplete = \false;
44 49 private $writesClosed = \false;
45 50 private $resourcesGetMethod = null;
46 51 private $pendingResources = [];
52 + private null|LoggerInterface $logger = null;
47 53 /**
48 54 * BidiStream constructor.
49 55 *
50 56 * @param BidiStreamingCall $bidiStreamingCall The gRPC bidirectional streaming call object
51 57 * @param array $streamingDescriptor
58 + * @param null|LoggerInterface $logger
52 59 */
53 - public function __construct(BidiStreamingCall $bidiStreamingCall, array $streamingDescriptor = [])
60 + public function __construct(BidiStreamingCall $bidiStreamingCall, array $streamingDescriptor = [], null|LoggerInterface $logger = null)
54 61 {
55 62 $this->call = $bidiStreamingCall;
56 - if (array_key_exists('resourcesGetMethod', $streamingDescriptor)) {
63 + if (\array_key_exists('resourcesGetMethod', $streamingDescriptor)) {
57 64 $this->resourcesGetMethod = $streamingDescriptor['resourcesGetMethod'];
58 65 }
66 + $this->logger = $logger;
59 67 }
60 68 /**
61 69 * Write request to the server.
62 70 *
@@ -65,13 +73,21 @@
65 73 */
66 74 public function write($request)
67 75 {
68 76 if ($this->isComplete) {
69 - throw new ValidationException("Cannot call write() after streaming call is complete.");
77 + throw new ValidationException('Cannot call write() after streaming call is complete.');
70 78 }
71 79 if ($this->writesClosed) {
72 - throw new ValidationException("Cannot call write() after calling closeWrite().");
80 + throw new ValidationException('Cannot call write() after calling closeWrite().');
73 81 }
82 + if ($this->logger && $request instanceof Message) {
83 + $logEvent = new RpcLogEvent();
84 + $logEvent->headers = null;
85 + $logEvent->payload = $request->serializeToJsonString();
86 + $logEvent->processId = (int) \getmypid();
87 + $logEvent->requestId = \crc32((string) \spl_object_id($this) . \getmypid());
88 + $this->logRequest($logEvent);
89 + }
74 90 $this->call->write($request);
75 91 }
76 92 /**
77 93 * Write all requests in $requests.
@@ -93,9 +109,9 @@
93 109 */
94 110 public function closeWrite()
95 111 {
96 112 if ($this->isComplete) {
97 - throw new ValidationException("Cannot call closeWrite() after streaming call is complete.");
113 + throw new ValidationException('Cannot call closeWrite() after streaming call is complete.');
98 114 }
99 115 if (!$this->writesClosed) {
100 116 $this->call->writesDone();
101 117 $this->writesClosed = \true;
@@ -111,27 +127,27 @@
111 127 */
112 128 public function read()
113 129 {
114 130 if ($this->isComplete) {
115 - throw new ValidationException("Cannot call read() after streaming call is complete.");
131 + throw new ValidationException('Cannot call read() after streaming call is complete.');
116 132 }
117 133 $resourcesGetMethod = $this->resourcesGetMethod;
118 - if (!is_null($resourcesGetMethod)) {
119 - if (count($this->pendingResources) === 0) {
134 + if (!\is_null($resourcesGetMethod)) {
135 + if (\count($this->pendingResources) === 0) {
120 136 $response = $this->call->read();
121 - if (!is_null($response)) {
137 + if (!\is_null($response)) {
122 138 $pendingResources = [];
123 139 foreach ($response->{$resourcesGetMethod}() as $resource) {
124 140 $pendingResources[] = $resource;
125 141 }
126 - $this->pendingResources = array_reverse($pendingResources);
142 + $this->pendingResources = \array_reverse($pendingResources);
127 143 }
128 144 }
129 - $result = array_pop($this->pendingResources);
145 + $result = \array_pop($this->pendingResources);
130 146 } else {
131 147 $result = $this->call->read();
132 148 }
133 - if (is_null($result)) {
149 + if (\is_null($result)) {
134 150 $status = $this->call->getStatus();
135 151 $this->isComplete = \true;
136 152 if (!($status->code == Code::OK)) {
137 153 throw ApiException::createFromStdClass($status);
@@ -136,8 +152,19 @@
136 152 if (!($status->code == Code::OK)) {
137 153 throw ApiException::createFromStdClass($status);
138 154 }
139 155 }
156 + if ($this->logger) {
157 + $responseEvent = new RpcLogEvent();
158 + $responseEvent->headers = $this->call->getMetadata();
159 + $responseEvent->status = $status->code ?? null;
160 + $responseEvent->processId = (int) \getmypid();
161 + $responseEvent->requestId = \crc32((string) \spl_object_id($this) . \getmypid());
162 + if ($result instanceof Message) {
163 + $responseEvent->payload = $result->serializeToJsonString();
164 + }
165 + $this->logResponse($responseEvent);
166 + }
140 167 return $result;
141 168 }
142 169 /**
143 170 * Call closeWrite(), and read all responses from the server, until the streaming call is
@@ -150,10 +177,10 @@
150 177 public function closeWriteAndReadAll()
151 178 {
152 179 $this->closeWrite();
153 180 $response = $this->read();
154 - while (!is_null($response)) {
155 - yield $response;
181 + while (!\is_null($response)) {
182 + (yield $response);
156 183 $response = $this->read();
157 184 }
158 185 }
159 186 /**