33namespace Utopia \Queue \Adapter ;
44
55use Swoole \Coroutine ;
6+ use Swoole \Coroutine \Channel ;
67use Swoole \Process ;
8+ use Utopia \DI \Container ;
79use Utopia \Queue \Adapter ;
810use Utopia \Queue \Consumer ;
11+ use Utopia \Queue \Error \ConsumerFailures ;
12+ use Utopia \Queue \Message ;
913
1014class Swoole extends Adapter
1115{
16+ protected const string CONTEXT_KEY = '__utopia__ ' ;
17+
1218 /** @var Process[] */
1319 protected array $ workers = [];
1420
@@ -23,9 +29,11 @@ public function __construct(
2329 int $ workerNum ,
2430 string $ queue ,
2531 string $ namespace = 'utopia-queue ' ,
32+ protected int $ maxCoroutines = 1 ,
33+ Container $ resources = new Container (),
2634 ) {
27- parent ::__construct ($ workerNum , $ queue , $ namespace );
28- $ this ->consumer = $ consumer ;
35+ parent ::__construct ($ consumer , $ workerNum , $ queue , $ namespace, $ resources );
36+ $ this ->maxCoroutines = \max ( 1 , $ maxCoroutines ) ;
2937 }
3038
3139 public function start (): self
@@ -71,6 +79,65 @@ protected function spawnWorker(int $workerId): void
7179 $ this ->workers [$ pid ] = $ process ;
7280 }
7381
82+ public function consume (callable $ messageCallback , callable $ successCallback , callable $ errorCallback ): void
83+ {
84+ $ messageCallback = function (Message $ message ) use ($ messageCallback ) {
85+ Coroutine::getContext ()[self ::CONTEXT_KEY ] = new Container ($ this ->resources ());
86+
87+ return $ messageCallback ($ message );
88+ };
89+
90+ $ errorCallback = function (?Message $ message , \Throwable $ error ) use ($ errorCallback ) {
91+ if ($ message === null ) {
92+ Coroutine::getContext ()[self ::CONTEXT_KEY ] = new Container ($ this ->resources ());
93+ }
94+
95+ $ errorCallback ($ message , $ error );
96+ };
97+
98+ $ channel = new Channel ($ this ->maxCoroutines );
99+ $ errors = [];
100+
101+ for ($ i = 0 ; $ i < $ this ->maxCoroutines ; $ i ++) {
102+ Coroutine::create (function () use ($ messageCallback , $ successCallback , $ errorCallback , $ channel , &$ errors ) {
103+ try {
104+ $ this ->consumer ->consume (
105+ $ this ->queue ,
106+ $ messageCallback ,
107+ $ successCallback ,
108+ $ errorCallback ,
109+ );
110+ } catch (\Throwable $ error ) {
111+ $ errors [] = $ error ;
112+ $ this ->consumer ->close ();
113+ $ channel ->push (true );
114+ return ;
115+ }
116+
117+ $ channel ->push (true );
118+ });
119+ }
120+
121+ for ($ i = 0 ; $ i < $ this ->maxCoroutines ; $ i ++) {
122+ $ channel ->pop ();
123+ }
124+
125+ $ channel ->close ();
126+
127+ if ($ errors !== []) {
128+ throw new ConsumerFailures ($ errors );
129+ }
130+ }
131+
132+ public function context (): Container
133+ {
134+ if (Coroutine::getCid () !== -1 ) {
135+ return Coroutine::getContext ()[self ::CONTEXT_KEY ] ?? $ this ->resources ();
136+ }
137+
138+ return $ this ->resources ();
139+ }
140+
74141 protected function reap (): void
75142 {
76143 while (($ ret = Process::wait (false )) !== false ) {
0 commit comments