Repository navigation
Expand file tree
/
Copy pathConcurrent.php
More file actions
90 lines (76 loc) · 2.27 KB
/
Copy pathConcurrent.php
File metadata and controls
90 lines (76 loc) · 2.27 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
<?php
declare(strict_types=1);
/**
* This file is part of Hyperf.
*
* @link https://www.hyperf.io
* @document https://hyperf.wiki
* @contact [email protected]
* @license https://github.com/hyperf/hyperf/blob/master/LICENSE
*/
namespace Hyperf\Coroutine;
use Hyperf\Context\ApplicationContext;
use Hyperf\Contract\StdoutLoggerInterface;
use Hyperf\Coroutine\Exception\InvalidArgumentException;
use Hyperf\Engine\Channel;
use Hyperf\ExceptionHandler\Formatter\FormatterInterface;
use Throwable;
/**
* @method bool isFull()
* @method bool isEmpty()
*/
class Concurrent
{
protected Channel $channel;
public function __construct(protected int $limit)
{
$this->channel = new Channel($limit);
}
public function __call($name, $arguments)
{
if (in_array($name, ['isFull', 'isEmpty'])) {
return $this->channel->{$name}(...$arguments);
}
throw new InvalidArgumentException(sprintf('The method %s is not supported.', $name));
}
public function getLimit(): int
{
return $this->limit;
}
public function length(): int
{
return $this->channel->getLength();
}
public function getLength(): int
{
return $this->channel->getLength();
}
public function getRunningCoroutineCount(): int
{
return $this->getLength();
}
public function getChannel(): Channel
{
return $this->channel;
}
public function create(callable $callable): void
{
$this->channel->push(true);
Coroutine::create(function () use ($callable) {
try {
$callable();
} catch (Throwable $exception) {
if (ApplicationContext::hasContainer()) {
$container = ApplicationContext::getContainer();
if ($container->has(StdoutLoggerInterface::class) && $container->has(FormatterInterface::class)) {
$logger = $container->get(StdoutLoggerInterface::class);
$formatter = $container->get(FormatterInterface::class);
$logger->error($formatter->format($exception));
}
}
} finally {
$this->channel->pop();
}
});
}
}