1
0
mirror of https://github.com/danog/parallel.git synced 2024-12-12 09:09:35 +01:00
parallel/src/Threading/Thread.php

189 lines
4.9 KiB
PHP
Raw Normal View History

2015-07-14 00:30:59 +02:00
<?php
namespace Icicle\Concurrent\Threading;
use Icicle\Concurrent\ChannelInterface;
use Icicle\Concurrent\Exception\InvalidArgumentError;
use Icicle\Concurrent\Exception\SynchronizationError;
use Icicle\Concurrent\Sync\Channel;
use Icicle\Concurrent\Sync\ExitStatusInterface;
use Icicle\Coroutine;
2015-08-25 16:37:22 +02:00
use Icicle\Socket\Stream\DuplexStream;
2015-07-27 00:53:00 +02:00
/**
* Implements an execution context using native multi-threading.
2015-08-11 00:38:58 +02:00
*
* The thread context is not itself threaded. A local instance of the context is
* maintained both in the context that creates the thread and in the thread
* itself.
*/
class Thread implements ChannelInterface
2015-07-14 00:30:59 +02:00
{
2015-08-25 16:39:58 +02:00
const LATENCY_TIMEOUT = 0.01; // 10 ms
2015-07-27 00:53:00 +02:00
/**
* @var \Icicle\Concurrent\Threading\InternalThread An internal thread instance.
2015-07-27 00:53:00 +02:00
*/
private $thread;
2015-07-27 00:53:00 +02:00
/**
* @var \Icicle\Concurrent\Sync\Channel A channel for communicating with the thread.
*/
private $channel;
/**
* @var resource
*/
private $socket;
2015-08-07 01:59:25 +02:00
/**
* Spawns a new thread and runs it.
*
* @param callable $function A callable to invoke in the thread.
*
* @return \Icicle\Concurrent\Threading\Thread The thread object that was spawned.
2015-08-07 01:59:25 +02:00
*/
public static function spawn(callable $function /* , ...$args */)
{
2015-08-24 17:47:36 +02:00
$class = new \ReflectionClass(__CLASS__);
$thread = $class->newInstanceArgs(func_get_args());
$thread->start();
return $thread;
}
2015-08-18 17:12:06 +02:00
/**
* Creates a new thread context from a thread.
*
* @param callable $function
2015-08-18 17:12:06 +02:00
*/
public function __construct(callable $function /* , ...$args */)
{
$args = array_slice(func_get_args(), 1);
list($channel, $this->socket) = Channel::createSocketPair();
$this->thread = new InternalThread($this->socket, $function, $args);
2015-08-25 16:37:22 +02:00
$this->channel = new Channel(new DuplexStream($channel));
}
2015-08-18 17:12:06 +02:00
2015-07-27 00:53:00 +02:00
/**
* Checks if the context is running.
2015-07-27 00:53:00 +02:00
*
* @return bool True if the context is running, otherwise false.
2015-07-27 00:53:00 +02:00
*/
public function isRunning()
{
return $this->thread->isRunning();
2015-08-07 01:59:25 +02:00
}
/**
* Starts the context execution.
*/
public function start()
2015-07-14 00:30:59 +02:00
{
if ($this->isRunning()) {
throw new SynchronizationError('The thread has already been started.');
}
$this->thread->start(PTHREADS_INHERIT_ALL);
}
2015-08-11 00:38:58 +02:00
/**
* Immediately kills the context.
*/
public function kill()
{
$this->thread->kill();
$this->channel->close();
2015-08-25 16:37:22 +02:00
fclose($this->socket);
}
2015-08-18 17:12:06 +02:00
/**
* @coroutine
*
* Gets a promise that resolves when the context ends and joins with the
* parent context.
2015-08-18 17:12:06 +02:00
*
* @return \Generator
*
* @resolve mixed Resolved with the return or resolution value of the context once it has completed execution.
*
* @throws \Icicle\Concurrent\Exception\SynchronizationError Thrown if an exit status object is not received.
2015-08-18 17:12:06 +02:00
*/
public function join()
2015-08-18 17:12:06 +02:00
{
if (!$this->isRunning()) {
throw new SynchronizationError('The thread has not been started or has already finished.');
2015-08-18 17:12:06 +02:00
}
try {
$response = (yield $this->channel->receive());
if (!$response instanceof ExitStatusInterface) {
throw new SynchronizationError('Did not receive an exit status from thread.');
2015-08-18 17:12:06 +02:00
}
yield $response->getResult();
2015-08-18 17:12:06 +02:00
} finally {
$this->thread->join();
$this->channel->close();
fclose($this->socket);
2015-08-18 17:12:06 +02:00
}
}
/**
* {@inheritdoc}
2015-08-18 17:12:06 +02:00
*/
public function receive()
2015-08-18 17:12:06 +02:00
{
if (!$this->isRunning()) {
throw new SynchronizationError('The thread has not been started or has already finished.');
}
$data = (yield $this->channel->receive());
if ($data instanceof ExitStatusInterface) {
$data = $data->getResult();
throw new SynchronizationError(sprintf(
'Thread unexpectedly exited with result of type: %s',
is_object($data) ? get_class($data) : gettype($data)
));
}
yield $data;
2015-08-18 17:12:06 +02:00
}
2015-08-07 01:59:25 +02:00
/**
* {@inheritdoc}
2015-08-07 01:59:25 +02:00
*/
public function send($data)
{
if (!$this->isRunning()) {
throw new SynchronizationError('The thread has not been started or has already finished.');
}
2015-08-07 07:07:53 +02:00
if ($data instanceof ExitStatusInterface) {
throw new InvalidArgumentError('Cannot send exit status objects.');
2015-07-27 00:53:00 +02:00
}
2015-07-14 00:30:59 +02:00
return $this->channel->send($data);
}
/**
* @param callable $callback
*
* @return \Generator
*/
public function synchronized(callable $callback)
{
while (!$this->thread->tsl()) {
2015-08-25 16:39:58 +02:00
yield Coroutine\sleep(self::LATENCY_TIMEOUT);
2015-08-19 02:12:58 +02:00
}
try {
yield $callback($this);
} finally {
$this->thread->release();
}
}
2015-07-14 00:30:59 +02:00
}