mirror of
https://github.com/appwrite/appwrite.git
synced 2026-05-26 13:51:13 +00:00
Implement task queue and wait() for graphql-php compatibility
- Add task queue mechanism to SwoolePromise matching SyncPromise behavior - Callbacks are enqueued instead of executed immediately - Add wait() method to Swoole adapter to process queue until promise settles - This fixes hanging tests by properly executing deferred promise callbacks Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.5
parent
d7174d0de3
commit
233c4761d6
@@ -8,6 +8,37 @@ use GraphQL\Executor\Promise\Promise as GQLPromise;
|
||||
|
||||
class Swoole extends Adapter
|
||||
{
|
||||
/**
|
||||
* Synchronously wait for promise completion by running the task queue.
|
||||
*
|
||||
* @param GQLPromise $promise
|
||||
* @return mixed
|
||||
* @throws \Throwable
|
||||
*/
|
||||
public function wait(GQLPromise $promise): mixed
|
||||
{
|
||||
/** @var SwoolePromise $swoolePromise */
|
||||
$swoolePromise = $promise->adoptedPromise;
|
||||
$taskQueue = SwoolePromise::getQueue();
|
||||
|
||||
while (
|
||||
$swoolePromise->state === SwoolePromise::PENDING
|
||||
&& !$taskQueue->isEmpty()
|
||||
) {
|
||||
SwoolePromise::runQueue();
|
||||
}
|
||||
|
||||
if ($swoolePromise->state === SwoolePromise::FULFILLED) {
|
||||
return $swoolePromise->result;
|
||||
}
|
||||
|
||||
if ($swoolePromise->state === SwoolePromise::REJECTED) {
|
||||
throw $swoolePromise->result;
|
||||
}
|
||||
|
||||
throw new \Exception('Could not resolve promise');
|
||||
}
|
||||
|
||||
public function create(callable $resolver): GQLPromise
|
||||
{
|
||||
$promise = new SwoolePromise(function ($resolve, $reject) use ($resolver) {
|
||||
|
||||
@@ -3,8 +3,8 @@
|
||||
namespace Appwrite\Promises;
|
||||
|
||||
/**
|
||||
* Swoole-compatible promise implementation that runs synchronously
|
||||
* but can wait for external async operations via Swoole channels.
|
||||
* Swoole-compatible promise implementation that uses a task queue
|
||||
* for deferred callback execution, similar to graphql-php's SyncPromise.
|
||||
*/
|
||||
class Swoole extends Promise
|
||||
{
|
||||
@@ -12,30 +12,56 @@ class Swoole extends Promise
|
||||
public const FULFILLED = 'fulfilled';
|
||||
public const REJECTED = 'rejected';
|
||||
|
||||
public string $promiseState = self::PENDING;
|
||||
public mixed $promiseResult = null;
|
||||
public string $state = self::PENDING;
|
||||
public mixed $result = null;
|
||||
|
||||
/**
|
||||
* Callbacks waiting for this promise to settle
|
||||
* Promises created in `then` method of this promise and awaiting resolution
|
||||
*
|
||||
* @var array<array{self, callable|null, callable|null}>
|
||||
*/
|
||||
protected array $waiting = [];
|
||||
|
||||
/**
|
||||
* Run all tasks in the queue
|
||||
*/
|
||||
public static function runQueue(): void
|
||||
{
|
||||
$q = self::getQueue();
|
||||
while (!$q->isEmpty()) {
|
||||
$task = $q->dequeue();
|
||||
$task();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the shared task queue
|
||||
*
|
||||
* @return \SplQueue<callable(): void>
|
||||
*/
|
||||
public static function getQueue(): \SplQueue
|
||||
{
|
||||
static $queue;
|
||||
|
||||
return $queue ??= new \SplQueue();
|
||||
}
|
||||
|
||||
public function __construct(?callable $executor = null)
|
||||
{
|
||||
if ($executor === null) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
$executor(
|
||||
fn ($value) => $this->doResolve($value),
|
||||
fn ($reason) => $this->doReject($reason)
|
||||
);
|
||||
} catch (\Throwable $e) {
|
||||
$this->doReject($e);
|
||||
}
|
||||
self::getQueue()->enqueue(function () use ($executor): void {
|
||||
try {
|
||||
$executor(
|
||||
fn ($value) => $this->resolve($value),
|
||||
fn ($reason) => $this->reject($reason)
|
||||
);
|
||||
} catch (\Throwable $e) {
|
||||
$this->reject($e);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
protected function execute(
|
||||
@@ -43,60 +69,77 @@ class Swoole extends Promise
|
||||
callable $resolve,
|
||||
callable $reject
|
||||
): void {
|
||||
// Not used - we execute synchronously in constructor
|
||||
// Not used - we use the task queue mechanism instead
|
||||
}
|
||||
|
||||
protected function doResolve(mixed $value): void
|
||||
/**
|
||||
* Resolve the promise with a value
|
||||
*/
|
||||
public function resolve(mixed $value): self
|
||||
{
|
||||
if ($this->promiseState !== self::PENDING) {
|
||||
return;
|
||||
if ($this->state !== self::PENDING) {
|
||||
return $this;
|
||||
}
|
||||
|
||||
// Handle thenable values
|
||||
if (\is_object($value) && \method_exists($value, 'then')) {
|
||||
$value->then(
|
||||
fn ($v) => $this->doResolve($v),
|
||||
fn ($r) => $this->doReject($r)
|
||||
fn ($v) => $this->resolve($v),
|
||||
fn ($r) => $this->reject($r)
|
||||
);
|
||||
return;
|
||||
return $this;
|
||||
}
|
||||
|
||||
$this->promiseState = self::FULFILLED;
|
||||
$this->promiseResult = $value;
|
||||
$this->processWaiting();
|
||||
$this->state = self::FULFILLED;
|
||||
$this->result = $value;
|
||||
$this->enqueueWaitingPromises();
|
||||
|
||||
return $this;
|
||||
}
|
||||
|
||||
protected function doReject(mixed $reason): void
|
||||
/**
|
||||
* Reject the promise with a reason
|
||||
*/
|
||||
public function reject(mixed $reason): self
|
||||
{
|
||||
if ($this->promiseState !== self::PENDING) {
|
||||
return;
|
||||
if ($this->state !== self::PENDING) {
|
||||
return $this;
|
||||
}
|
||||
|
||||
$this->promiseState = self::REJECTED;
|
||||
$this->promiseResult = $reason;
|
||||
$this->processWaiting();
|
||||
$this->state = self::REJECTED;
|
||||
$this->result = $reason;
|
||||
$this->enqueueWaitingPromises();
|
||||
|
||||
return $this;
|
||||
}
|
||||
|
||||
protected function processWaiting(): void
|
||||
/**
|
||||
* Enqueue callbacks for waiting promises
|
||||
*/
|
||||
protected function enqueueWaitingPromises(): void
|
||||
{
|
||||
foreach ($this->waiting as [$promise, $onFulfilled, $onRejected]) {
|
||||
$callback = $this->promiseState === self::FULFILLED ? $onFulfilled : $onRejected;
|
||||
|
||||
if ($callback === null) {
|
||||
if ($this->promiseState === self::FULFILLED) {
|
||||
$promise->doResolve($this->promiseResult);
|
||||
} else {
|
||||
$promise->doReject($this->promiseResult);
|
||||
self::getQueue()->enqueue(function () use ($promise, $onFulfilled, $onRejected): void {
|
||||
if ($this->state === self::FULFILLED) {
|
||||
try {
|
||||
$promise->resolve($onFulfilled === null ? $this->result : $onFulfilled($this->result));
|
||||
} catch (\Throwable $e) {
|
||||
$promise->reject($e);
|
||||
}
|
||||
} elseif ($this->state === self::REJECTED) {
|
||||
try {
|
||||
if ($onRejected === null) {
|
||||
$promise->reject($this->result);
|
||||
} else {
|
||||
$promise->resolve($onRejected($this->result));
|
||||
}
|
||||
} catch (\Throwable $e) {
|
||||
$promise->reject($e);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
try {
|
||||
$result = $callback($this->promiseResult);
|
||||
$promise->doResolve($result);
|
||||
} catch (\Throwable $e) {
|
||||
$promise->doReject($e);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
$this->waiting = [];
|
||||
}
|
||||
|
||||
@@ -104,62 +147,42 @@ class Swoole extends Promise
|
||||
?callable $onFulfilled = null,
|
||||
?callable $onRejected = null
|
||||
): self {
|
||||
if ($this->state === self::REJECTED && $onRejected === null) {
|
||||
return $this;
|
||||
}
|
||||
|
||||
if ($this->state === self::FULFILLED && $onFulfilled === null) {
|
||||
return $this;
|
||||
}
|
||||
|
||||
$promise = new self();
|
||||
$this->waiting[] = [$promise, $onFulfilled, $onRejected];
|
||||
|
||||
if ($this->promiseState === self::PENDING) {
|
||||
$this->waiting[] = [$promise, $onFulfilled, $onRejected];
|
||||
} else {
|
||||
$callback = $this->promiseState === self::FULFILLED ? $onFulfilled : $onRejected;
|
||||
|
||||
if ($callback === null) {
|
||||
if ($this->promiseState === self::FULFILLED) {
|
||||
$promise->doResolve($this->promiseResult);
|
||||
} else {
|
||||
$promise->doReject($this->promiseResult);
|
||||
}
|
||||
} else {
|
||||
try {
|
||||
$result = $callback($this->promiseResult);
|
||||
$promise->doResolve($result);
|
||||
} catch (\Throwable $e) {
|
||||
$promise->doReject($e);
|
||||
}
|
||||
}
|
||||
if ($this->state !== self::PENDING) {
|
||||
$this->enqueueWaitingPromises();
|
||||
}
|
||||
|
||||
return $promise;
|
||||
}
|
||||
|
||||
public function resolve(mixed $value): self
|
||||
{
|
||||
$this->doResolve($value);
|
||||
return $this;
|
||||
}
|
||||
|
||||
public function reject(mixed $reason): self
|
||||
{
|
||||
$this->doReject($reason);
|
||||
return $this;
|
||||
}
|
||||
|
||||
public function isPending(): bool
|
||||
{
|
||||
return $this->promiseState === self::PENDING;
|
||||
return $this->state === self::PENDING;
|
||||
}
|
||||
|
||||
public function isFulfilled(): bool
|
||||
{
|
||||
return $this->promiseState === self::FULFILLED;
|
||||
return $this->state === self::FULFILLED;
|
||||
}
|
||||
|
||||
public function isRejected(): bool
|
||||
{
|
||||
return $this->promiseState === self::REJECTED;
|
||||
return $this->state === self::REJECTED;
|
||||
}
|
||||
|
||||
public function getResult(): mixed
|
||||
{
|
||||
return $this->promiseResult;
|
||||
return $this->result;
|
||||
}
|
||||
|
||||
public static function all(iterable $promisesOrValues): self
|
||||
@@ -169,21 +192,19 @@ class Swoole extends Promise
|
||||
: \iterator_to_array($promisesOrValues);
|
||||
|
||||
$total = \count($promisesOrValues);
|
||||
$promise = new self();
|
||||
$all = new self();
|
||||
|
||||
if ($total === 0) {
|
||||
$promise->doResolve([]);
|
||||
return $promise;
|
||||
$all->resolve([]);
|
||||
return $all;
|
||||
}
|
||||
|
||||
$count = 0;
|
||||
$result = [];
|
||||
$rejected = false;
|
||||
|
||||
$checkComplete = static function () use (&$count, $total, &$result, &$rejected, $promise): void {
|
||||
if (!$rejected && $count === $total) {
|
||||
\ksort($result);
|
||||
$promise->doResolve($result);
|
||||
$resolveAllWhenFinished = static function () use (&$count, $total, $all, &$result): void {
|
||||
if ($count === $total) {
|
||||
$all->resolve($result);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -191,18 +212,12 @@ class Swoole extends Promise
|
||||
if ($promiseOrValue instanceof self) {
|
||||
$result[$index] = null;
|
||||
$promiseOrValue->then(
|
||||
static function ($value) use (&$result, $index, &$count, $checkComplete) {
|
||||
static function ($value) use (&$result, $index, &$count, $resolveAllWhenFinished) {
|
||||
$result[$index] = $value;
|
||||
++$count;
|
||||
$checkComplete();
|
||||
return $value;
|
||||
$resolveAllWhenFinished();
|
||||
},
|
||||
static function ($error) use (&$rejected, $promise) {
|
||||
if (!$rejected) {
|
||||
$rejected = true;
|
||||
$promise->doReject($error);
|
||||
}
|
||||
}
|
||||
[$all, 'reject']
|
||||
);
|
||||
} else {
|
||||
$result[$index] = $promiseOrValue;
|
||||
@@ -210,8 +225,8 @@ class Swoole extends Promise
|
||||
}
|
||||
}
|
||||
|
||||
$checkComplete();
|
||||
$resolveAllWhenFinished();
|
||||
|
||||
return $promise;
|
||||
return $all;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user