| @@ -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 | /** |