InteractsWithWebsocket.php 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192
  1. <?php
  2. namespace think\swoole\concerns;
  3. use Swoole\Http\Request;
  4. use Swoole\Websocket\Frame;
  5. use Swoole\Websocket\Server;
  6. use think\App;
  7. use think\Container;
  8. use think\helper\Str;
  9. use think\swoole\contract\websocket\RoomInterface;
  10. use think\swoole\Middleware;
  11. use think\swoole\Websocket;
  12. use think\swoole\websocket\Room;
  13. /**
  14. * Trait InteractsWithWebsocket
  15. * @package think\swoole\concerns
  16. *
  17. * @property App $app
  18. * @property Container $container
  19. * @method \Swoole\Server getServer()
  20. */
  21. trait InteractsWithWebsocket
  22. {
  23. /**
  24. * @var boolean
  25. */
  26. protected $isWebsocketServer = false;
  27. /**
  28. * @var RoomInterface
  29. */
  30. protected $websocketRoom;
  31. /**
  32. * Websocket server events.
  33. *
  34. * @var array
  35. */
  36. protected $wsEvents = ['open', 'message', 'close'];
  37. /**
  38. * "onOpen" listener.
  39. *
  40. * @param Server $server
  41. * @param Request $req
  42. */
  43. public function onOpen($server, $req)
  44. {
  45. $this->waitCoordinator('workerStart');
  46. $this->runInSandbox(function (App $app, Websocket $websocket) use ($req) {
  47. $request = $this->prepareRequest($req);
  48. $app->instance('request', $request);
  49. $request = $this->setRequestThroughMiddleware($app, $request);
  50. $websocket->setSender($req->fd);
  51. $websocket->onOpen($req->fd, $request);
  52. }, $req->fd, true);
  53. }
  54. /**
  55. * "onMessage" listener.
  56. *
  57. * @param Server $server
  58. * @param Frame $frame
  59. */
  60. public function onMessage($server, $frame)
  61. {
  62. $this->runInSandbox(function (Websocket $websocket) use ($frame) {
  63. $websocket->setSender($frame->fd);
  64. $websocket->onMessage($frame);
  65. }, $frame->fd, true);
  66. }
  67. /**
  68. * "onClose" listener.
  69. *
  70. * @param Server $server
  71. * @param int $fd
  72. * @param int $reactorId
  73. */
  74. public function onClose($server, $fd, $reactorId)
  75. {
  76. if (!$server instanceof Server || !$this->isWebsocketServer($fd)) {
  77. return;
  78. }
  79. $this->runInSandbox(function (Websocket $websocket) use ($fd, $reactorId) {
  80. $websocket->setSender($fd);
  81. try {
  82. $websocket->onClose($fd, $reactorId);
  83. } finally {
  84. // leave all rooms
  85. $websocket->leave();
  86. }
  87. }, $fd);
  88. }
  89. /**
  90. * @param App $app
  91. * @param \think\Request $request
  92. * @return \think\Request
  93. */
  94. protected function setRequestThroughMiddleware(App $app, \think\Request $request)
  95. {
  96. return Middleware::make($app, $this->getConfig('websocket.middleware', []))
  97. ->pipeline()
  98. ->send($request)
  99. ->then(function ($request) {
  100. return $request;
  101. });
  102. }
  103. /**
  104. * Prepare settings if websocket is enabled.
  105. */
  106. protected function prepareWebsocket()
  107. {
  108. if (!$this->isWebsocketServer = $this->getConfig('websocket.enable', false)) {
  109. return;
  110. }
  111. $this->events = array_merge($this->events ?? [], $this->wsEvents);
  112. $this->prepareWebsocketRoom();
  113. $this->onEvent('workerStart', function () {
  114. $this->bindWebsocketRoom();
  115. $this->bindWebsocketHandler();
  116. $this->prepareWebsocketListener();
  117. });
  118. }
  119. /**
  120. * Check if it's a websocket fd.
  121. *
  122. * @param int $fd
  123. *
  124. * @return bool
  125. */
  126. protected function isWebsocketServer(int $fd): bool
  127. {
  128. return $this->getServer()->getClientInfo($fd)['websocket_status'] ?? false;
  129. }
  130. /**
  131. * Prepare websocket room.
  132. */
  133. protected function prepareWebsocketRoom()
  134. {
  135. // create room instance and initialize
  136. $this->websocketRoom = $this->container->make(Room::class);
  137. $this->websocketRoom->prepare();
  138. }
  139. protected function prepareWebsocketListener()
  140. {
  141. $listeners = $this->getConfig('websocket.listen', []);
  142. foreach ($listeners as $event => $listener) {
  143. $this->app->event->listen('swoole.websocket.' . Str::studly($event), $listener);
  144. }
  145. $subscribers = $this->getConfig('websocket.subscribe', []);
  146. foreach ($subscribers as $subscriber) {
  147. $this->app->event->observe($subscriber, 'swoole.websocket.');
  148. }
  149. }
  150. /**
  151. * Prepare websocket handler for onOpen and onClose callback
  152. */
  153. protected function bindWebsocketHandler()
  154. {
  155. $handlerClass = $this->getConfig('websocket.handler');
  156. if ($handlerClass && is_subclass_of($handlerClass, Websocket::class)) {
  157. $this->app->bind(Websocket::class, $handlerClass);
  158. }
  159. }
  160. /**
  161. * Bind room instance to app container.
  162. */
  163. protected function bindWebsocketRoom(): void
  164. {
  165. $this->app->instance(Room::class, $this->websocketRoom);
  166. }
  167. }