Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- **`ThreadChannel::send()` and `ThreadChannel::recv()` accepted a cancellation token and ignored it.** Both methods parsed the `?Async\Completable` argument and neither passed it on: `recv()` handed `NULL` to a receive that has supported cancellation all along, and `send()` had no parameter to take one. A parked call could then be broken only by cancelling the coroutine, so every bounded wait had to be hand-built from a cancel race — while `Async\Channel` honoured the same argument on the same signature. Measured before the fix: `recv(Async\timeout(300))` was still parked at 1000 ms. The token now ends either wait with `OperationCanceledException`, the exception `Channel::recv()` raises, carrying the token's own error as its previous, so one `catch (AsyncCancellation)` covers both classes. A token that fired before the call ends it before it waits. A wake is attributed to the token by asking the token: one freed slot wakes every parked sender and one sent value wakes every parked receiver, so the losers of that race park again instead of reporting a cancellation nobody requested — which would leave the method throwing an exception that was never raised (`ZEND_ASSERT(EG(exception))`, SIGABRT on a debug build, a silent `NULL` on a release one).
- **Sending a value that cannot cross a thread boundary aborted the process instead of only throwing.** `ThreadChannel::send()` transfers its argument into persistent memory before it takes the lock, and the transfer refuses what it cannot copy — a resource, an object with dynamic properties — by releasing the partial graph, leaving the destination `IS_UNDEF` and throwing. The send did not check for that: it pushed the undefined slot into the buffer and reported success, so the caller got the right exception while the buffer held a value no receiver can interpret. `ThreadPool` reads the task as an array, so a debug build died on the assertion at `thread_pool.c:347` and a release build, where that assertion is compiled out, reads array fields from a value that is not one. The send now leaves the buffer untouched and returns false, which every caller already handles: `ThreadChannel::send()` rethrows, and `ThreadPool::submit()` and `map()` release the snapshot and the future first. Reproduced with `$pool->submit(fn () => 1, fopen('php://memory', 'r'))`: SIGABRT before, a caught `Error` and exit 0 after.

## [0.9.3] - 2026-08-13
Expand Down
66 changes: 66 additions & 0 deletions tests/thread_channel/046-cancellation_token.phpt
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
--TEST--
ThreadChannel: the cancellation token ends a wait
--SKIPIF--
<?php
if (!PHP_ZTS) die('skip ZTS required');
?>
--FILE--
<?php

use Async\ThreadChannel;

use function Async\spawn;

// A fired token ends a parked recv() or send() with the exception Channel::recv()
// raises, so one catch (AsyncCancellation) covers both channel classes.
function timed(string $label, callable $body): void
{
$c = spawn(static function () use ($body, $label): void {
try {
$body();
echo "$label: returned\n";
} catch (Async\AsyncCancellation $e) {
echo "$label: " . $e::class . "\n";
}
});

// Watchdog: a wait that did not end on the token shows up as STILL PARKED.
spawn(static function () use ($c, $label): void {
Async\delay(1000);
if (!$c->isCompleted()) {
echo "$label: STILL PARKED\n";
$c->cancel();
}
});
}

$empty = new ThreadChannel(4);
timed('recv', static fn() => $empty->recv(Async\timeout(100)));

$full = new ThreadChannel(1);
$full->send('taken');
timed('send', static fn() => $full->send('overflow', Async\timeout(100)));

// A token that never fires ends nothing: both calls complete on their own.
$open = new ThreadChannel(4);
$open->send('value', Async\timeout(5000));
var_dump($open->recv(Async\timeout(5000)));

// The token's own exception is chained as previous, as Async\Channel chains it.
$chained = new ThreadChannel(4);

try {
$chained->recv(Async\timeout(100));
} catch (Async\AsyncCancellation $e) {
echo "previous: ", $e->getPrevious()?->getMessage() ?? 'NULL', "\n";
}

Async\delay(1200);
echo "Done\n";
?>
--EXPECT--
string(5) "value"
previous: Timeout occurred after 100 milliseconds
recv: Async\OperationCanceledException
send: Async\OperationCanceledException
Done
64 changes: 64 additions & 0 deletions tests/thread_channel/047-cancellation_token_spurious_wake.phpt
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
--TEST--
ThreadChannel: a wake that is not the token parks the call again
--SKIPIF--
<?php
if (!PHP_ZTS) die('skip ZTS required');
?>
--FILE--
<?php

use Async\ThreadChannel;

use function Async\spawn;

// One freed slot wakes every parked sender and one value wakes every parked
// receiver, but only one can proceed. The losers were not cancelled: treating the
// wake as a cancellation would end them with no value and no exception.
$full = new ThreadChannel(1);
$full->send('taken');

$sent = [];

foreach (['b', 'c'] as $value) {
spawn(static function () use ($full, $value, &$sent): void {
$full->send($value, Async\timeout(5000));
$sent[] = $value;
});
}

Async\delay(50);
$taken = [$full->recv()]; // one free slot, both senders woken
Async\delay(50);
$taken[] = $full->recv();
Async\delay(50);
$taken[] = $full->recv();

\sort($sent);
\sort($taken);
echo "sent: ", \implode(',', $sent), "\n";
echo "taken: ", \implode(',', $taken), "\n";

$empty = new ThreadChannel(4);
$received = [];

foreach (['x', 'y'] as $label) {
spawn(static function () use ($empty, &$received): void {
$received[] = $empty->recv(Async\timeout(5000));
});
}

Async\delay(50);
$empty->send('one'); // one value, both receivers woken
Async\delay(50);
$empty->send('two');
Async\delay(50);

\sort($received);
echo "received: ", \implode(',', $received), "\n";
echo "Done\n";
?>
--EXPECT--
sent: b,c
taken: b,c,taken
received: one,two
Done
115 changes: 103 additions & 12 deletions thread_channel.c
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,12 @@
} \
}

/* Mark Future token as used (non-blocking path) and throw if already fired. */
#define CANCELLATION_TOKEN_PREPARE(ct) \
if ((ct) != NULL && UNEXPECTED(async_resolve_cancel_token(ct))) { \
RETURN_THROWS(); \
}

zend_class_entry *async_ce_thread_channel = NULL;
zend_class_entry *async_ce_thread_channel_exception = NULL;
static zend_object_handlers async_thread_channel_handlers;
Expand Down Expand Up @@ -95,7 +101,11 @@ void async_thread_channel_close_owned(void)
static bool thread_channel_send(zend_async_channel_t *channel, zval *value);
static bool thread_channel_receive(zend_async_channel_t *channel, zval *result, zend_async_event_t *cancellation);

static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
/* Send, with an optional event that ends a wait on a full buffer. False means the
* value was refused, the channel is closed, or the event fired; EG(exception) does
* not tell them apart (a timeout token leaves one pending), the event's state does. */
static bool thread_channel_send_ex(
zend_async_channel_t *channel, zval *value, zend_async_event_t *cancellation)
{
async_thread_channel_t *ch = (async_thread_channel_t *) channel;
zend_async_trigger_event_t *trigger = NULL;
Expand All @@ -105,10 +115,8 @@ static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
async_thread_transfer_zval(&persistent_copy, value);

if (UNEXPECTED(Z_TYPE(persistent_copy) == IS_UNDEF)) {
/* The value is not transferable between threads. async_thread_transfer_zval
* has already released the partial graph and thrown; leave the buffer
* untouched, because a receiver has no way to tell an undefined slot from
* a real message and every reader asserts on the type it expects. */
/* Not transferable; async_thread_transfer_zval already threw. Nothing goes
* into the buffer: a receiver cannot tell an undefined slot from a message. */
return false;
}

Expand Down Expand Up @@ -147,6 +155,25 @@ static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
zend_async_resume_when(ZEND_ASYNC_CURRENT_COROUTINE,
&trigger->base, false, zend_async_waker_callback_resolve, NULL);

if (cancellation != NULL) {
if (UNEXPECTED(ZEND_ASYNC_EVENT_IS_CLOSED(cancellation))) {
/* zend_async_resume_when refuses a closed event, so suspending would arm
* the channel trigger alone and never time out. No race with the register
* below: an event closes from a loop callback, and the loop is not running. */
ASYNC_MUTEX_LOCK(ch->mutex);
zend_hash_index_del(&ch->sender_triggers, (zend_ulong)(uintptr_t) trigger);
ASYNC_MUTEX_UNLOCK(ch->mutex);
ZEND_ASYNC_WAKER_DESTROY(ZEND_ASYNC_CURRENT_COROUTINE);
async_thread_release_transferred_zval(&persistent_copy);
trigger->base.dispose(&trigger->base);

return false;
}

zend_async_resume_when(ZEND_ASYNC_CURRENT_COROUTINE,
cancellation, false, zend_async_waker_callback_resolve, NULL);
}

/* A bailout through SUSPEND would skip the dispose paths below and leak the
* trigger (open uv_async blocks uv_loop_close). Catch, dispose, re-raise. */
bool channel_bailed = false;
Expand All @@ -171,6 +198,8 @@ static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
/* Woke up — remove from sender queue */
ASYNC_MUTEX_LOCK(ch->mutex);
zend_hash_index_del(&ch->sender_triggers, (zend_ulong)(uintptr_t) trigger);
const bool closed = ZEND_ASYNC_EVENT_IS_CLOSED(&ch->channel.event);
const bool still_full = circular_buffer_count(&ch->buffer) >= (size_t) ch->capacity;
ASYNC_MUTEX_UNLOCK(ch->mutex);

if (EG(exception)) {
Expand All @@ -179,9 +208,22 @@ static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
return false;
}

if (cancellation != NULL && ZEND_ASYNC_EVENT_IS_CLOSED(cancellation) && closed == false && still_full) {
/* One freed slot wakes every parked sender, so the buffer alone cannot say
* this call was cancelled: the losers of that race retry instead of reporting. */
async_thread_release_transferred_zval(&persistent_copy);
trigger->base.dispose(&trigger->base);
return false;
}

goto retry;
}

static bool thread_channel_send(zend_async_channel_t *channel, zval *value)
{
return thread_channel_send_ex(channel, value, NULL);
}

static bool thread_channel_receive(
zend_async_channel_t *channel, zval *result, zend_async_event_t *cancellation)
{
Expand Down Expand Up @@ -231,7 +273,20 @@ static bool thread_channel_receive(

zend_async_resume_when(ZEND_ASYNC_CURRENT_COROUTINE,
&trigger->base, false, zend_async_waker_callback_resolve, NULL);

if (cancellation != NULL) {
if (UNEXPECTED(ZEND_ASYNC_EVENT_IS_CLOSED(cancellation))) {
/* Same guard as in thread_channel_send_ex: zend_async_resume_when refuses
* a closed event, and suspending without it never times out. */
ASYNC_MUTEX_LOCK(ch->mutex);
zend_hash_index_del(&ch->receiver_triggers, (zend_ulong)(uintptr_t) trigger);
ASYNC_MUTEX_UNLOCK(ch->mutex);
ZEND_ASYNC_WAKER_DESTROY(ZEND_ASYNC_CURRENT_COROUTINE);
trigger->base.dispose(&trigger->base);

return false;
}

zend_async_resume_when(ZEND_ASYNC_CURRENT_COROUTINE,
cancellation, false, zend_async_waker_callback_resolve, NULL);
}
Expand Down Expand Up @@ -274,14 +329,14 @@ static bool thread_channel_receive(
return !closed;
}

if (cancellation != NULL && closed == false) {
/* Non-wait_only call: data still wasn't ready (else the retry's
* pop branch would have caught it), and channel isn't closed, so
* the cancellation event is what woke us. Return false without
* exception — caller distinguishes from the closed-channel path. */
if (cancellation != NULL && ZEND_ASYNC_EVENT_IS_CLOSED(cancellation) && closed == false) {
/* Non-wait_only call: return false without an exception, which the caller
* distinguishes from the closed-channel path. One send wakes every parked
* receiver, so the losers of that race park again instead of reporting. */
ASYNC_MUTEX_LOCK(ch->mutex);
const bool still_empty = !circular_buffer_is_not_empty(&ch->buffer);
ASYNC_MUTEX_UNLOCK(ch->mutex);

if (still_empty) {
trigger->base.dispose(&trigger->base);
return false;
Expand Down Expand Up @@ -498,6 +553,30 @@ METHOD(__construct)
obj->channel = async_thread_channel_create((int32_t) capacity);
}

/* Throw what `Channel::recv()` throws for a fired token: OperationCanceledException
* with the token's own exception as previous, so one catch covers both channel
* classes. A token that never fired is left alone — that wait ended some other way. */
static void report_cancellation(zend_object *token)
{
if (token == NULL || !ZEND_ASYNC_EVENT_IS_CLOSED(ZEND_ASYNC_OBJECT_TO_EVENT(token))) {
return;
}

/* A timeout token leaves its TimeoutException in EG() rather than on the event,
* where async_resolve_cancel_token would not find it and would replace it. */
zend_object *raised = EG(exception);

if (raised != NULL) {
GC_ADDREF(raised);
zend_clear_exception();
}

async_resolve_cancel_token(token);

ZEND_ASSERT(EG(exception) != NULL);
zend_exception_set_previous(EG(exception), raised);
}

METHOD(send)
{
zval *value;
Expand All @@ -511,7 +590,13 @@ METHOD(send)

ENSURE_COROUTINE_CONTEXT

if (!THIS_CHANNEL()->channel.send(&THIS_CHANNEL()->channel, value)) {
CANCELLATION_TOKEN_PREPARE(cancellation_token)

zend_async_event_t *cancellation =
cancellation_token != NULL ? ZEND_ASYNC_OBJECT_TO_EVENT(cancellation_token) : NULL;

if (!thread_channel_send_ex(&THIS_CHANNEL()->channel, value, cancellation)) {
report_cancellation(cancellation_token);
RETURN_THROWS();
}
}
Expand All @@ -527,7 +612,13 @@ METHOD(recv)

ENSURE_COROUTINE_CONTEXT

if (!THIS_CHANNEL()->channel.receive(&THIS_CHANNEL()->channel, return_value, NULL)) {
CANCELLATION_TOKEN_PREPARE(cancellation_token)

zend_async_event_t *cancellation =
cancellation_token != NULL ? ZEND_ASYNC_OBJECT_TO_EVENT(cancellation_token) : NULL;

if (!THIS_CHANNEL()->channel.receive(&THIS_CHANNEL()->channel, return_value, cancellation)) {
report_cancellation(cancellation_token);
RETURN_THROWS();
}
}
Expand Down