Pusher.php 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203
  1. <?php
  2. namespace think\swoole\websocket;
  3. use Swoole\Server;
  4. /**
  5. * Class Pusher
  6. */
  7. class Pusher
  8. {
  9. /**
  10. * @var Server|\Swoole\WebSocket\Server
  11. */
  12. protected $server;
  13. /**
  14. * @var int
  15. */
  16. protected $sender;
  17. /**
  18. * @var array
  19. */
  20. protected $descriptors;
  21. /**
  22. * @var bool
  23. */
  24. protected $broadcast;
  25. /**
  26. * @var bool
  27. */
  28. protected $assigned;
  29. /**
  30. * @var string
  31. */
  32. protected $payload;
  33. /**
  34. * Push constructor.
  35. *
  36. * @param Server $server
  37. * @param int $sender
  38. * @param array $descriptors
  39. * @param bool $broadcast
  40. * @param bool $assigned
  41. * @param string $payload
  42. */
  43. public function __construct(
  44. Server $server,
  45. string $payload,
  46. int $sender = 0,
  47. array $descriptors = [],
  48. bool $broadcast = false,
  49. bool $assigned = false
  50. )
  51. {
  52. $this->sender = $sender;
  53. $this->descriptors = $descriptors;
  54. $this->broadcast = $broadcast;
  55. $this->assigned = $assigned;
  56. $this->payload = $payload;
  57. $this->server = $server;
  58. }
  59. /**
  60. * @return int
  61. */
  62. public function getSender(): int
  63. {
  64. return $this->sender;
  65. }
  66. /**
  67. * @return array
  68. */
  69. public function getDescriptors(): array
  70. {
  71. return $this->descriptors;
  72. }
  73. /**
  74. * @param int $descriptor
  75. *
  76. * @return self
  77. */
  78. public function addDescriptor($descriptor): self
  79. {
  80. return $this->addDescriptors([$descriptor]);
  81. }
  82. /**
  83. * @param array $descriptors
  84. *
  85. * @return self
  86. */
  87. public function addDescriptors(array $descriptors): self
  88. {
  89. $this->descriptors = array_values(
  90. array_unique(
  91. array_merge($this->descriptors, $descriptors)
  92. )
  93. );
  94. return $this;
  95. }
  96. /**
  97. * @param int $descriptor
  98. *
  99. * @return bool
  100. */
  101. public function hasDescriptor(int $descriptor): bool
  102. {
  103. return in_array($descriptor, $this->descriptors);
  104. }
  105. /**
  106. * @return bool
  107. */
  108. public function isBroadcast(): bool
  109. {
  110. return $this->broadcast;
  111. }
  112. /**
  113. * @return bool
  114. */
  115. public function isAssigned(): bool
  116. {
  117. return $this->assigned;
  118. }
  119. /**
  120. * @return string
  121. */
  122. public function getPayload(): string
  123. {
  124. return $this->payload;
  125. }
  126. /**
  127. * @return bool
  128. */
  129. public function shouldBroadcast(): bool
  130. {
  131. return $this->broadcast && empty($this->descriptors) && !$this->assigned;
  132. }
  133. /**
  134. * Returns all descriptors that are websocket
  135. *
  136. * @return array
  137. */
  138. protected function getWebsocketConnections(): array
  139. {
  140. return array_filter(iterator_to_array($this->server->connections), function ($fd) {
  141. return (bool) ($this->server->getClientInfo($fd)['websocket_status'] ?? false);
  142. });
  143. }
  144. /**
  145. * @param int $fd
  146. *
  147. * @return bool
  148. */
  149. protected function shouldPushToDescriptor(int $fd): bool
  150. {
  151. if (!$this->server->isEstablished($fd)) {
  152. return false;
  153. }
  154. return !$this->broadcast || $this->sender !== (int) $fd;
  155. }
  156. /**
  157. * Push message to related descriptors
  158. * @return void
  159. */
  160. public function push(): void
  161. {
  162. // attach sender if not broadcast
  163. if (!$this->broadcast && $this->sender && !$this->hasDescriptor($this->sender)) {
  164. $this->addDescriptor($this->sender);
  165. }
  166. // check if to broadcast to other clients
  167. if ($this->shouldBroadcast()) {
  168. $this->addDescriptors($this->getWebsocketConnections());
  169. }
  170. // push message to designated fds
  171. foreach ($this->descriptors as $descriptor) {
  172. if ($this->shouldPushToDescriptor($descriptor)) {
  173. $this->server->push($descriptor, $this->payload);
  174. }
  175. }
  176. }
  177. }