-
-
Notifications
You must be signed in to change notification settings - Fork 2
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Support defining code based input sources and actions.
* Input Source functions are defined via extending the SourceFunction class * Action functions @note: There is no protection against blocking. These functions are more useful for testing engines or simple input/action processing. Additional Changes: * Loop is initialised in the Scheduler constructors instead of delaying till run() * Input and Action processes can be defined using an array of parameters rather than an escaped string
- Loading branch information
1 parent
ac01dee
commit 41b7f5b
Showing
4 changed files
with
263 additions
and
32 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,69 @@ | ||
<?php declare(strict_types=1); | ||
|
||
use Bref\Logger\StderrLogger; | ||
use EdgeTelemetrics\EventCorrelation\Event; | ||
use EdgeTelemetrics\EventCorrelation\Rule; | ||
use EdgeTelemetrics\EventCorrelation\Scheduler; | ||
use Psr\Log\LogLevel; | ||
|
||
require __DIR__ . '/../../vendor/autoload.php'; | ||
|
||
class SendCountToStdOut extends Rule\MatchSingle { | ||
const EVENTS = [['SampleValueEvent']]; | ||
|
||
public function onComplete(): void | ||
{ | ||
$this->emit('data', [New \EdgeTelemetrics\EventCorrelation\Action('echo', ['value' => $this->getFirstEvent()->value])]); | ||
} | ||
} | ||
|
||
$rules = [ | ||
SendCountToStdOut::class, | ||
]; | ||
|
||
$scheduler = new \EdgeTelemetrics\EventCorrelation\Scheduler($rules); | ||
$scheduler->setLogger(new StderrLogger(LogLevel::DEBUG)); | ||
|
||
$scheduler->register_action('echo', function($vars) { | ||
echo 'Next Value: ' . $vars['value'] . PHP_EOL; | ||
}); | ||
|
||
$numberGenClass = new class() extends Scheduler\SourceFunction { | ||
protected \React\EventLoop\TimerInterface $timer; | ||
|
||
function functionStart(): void | ||
{ | ||
$this->timer = $this->loop->addPeriodicTimer(1.0, function () { | ||
static $count = 1; | ||
try { | ||
$event = new Event(['event' => 'SampleValueEvent', 'value' => $count++]); | ||
$this->emit('data', [$event]); | ||
|
||
if ($count > 10) { | ||
$this->exit(); | ||
} | ||
} catch (Throwable $exception) { | ||
$this->emit('error', [$exception]); | ||
$this->exit(255); | ||
return; | ||
} | ||
}); | ||
} | ||
|
||
function exit(int $code = 0): void | ||
{ | ||
$this->running = false; | ||
$this->loop->cancelTimer($this->timer); | ||
unset($this->timer); | ||
$this->emit('exit', [$code]); | ||
} | ||
|
||
function functionStop(): void | ||
{ | ||
$this->exit(); | ||
} | ||
}; | ||
|
||
$scheduler->register_input_process('generator', $numberGenClass, null, [], false); | ||
|
||
$scheduler->run(); |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,31 @@ | ||
<?php declare(strict_types=1); | ||
|
||
/* | ||
* This file is part of the PHP Event Correlation package. | ||
* | ||
* (c) James Lucas <james@lucas.net.au> | ||
* | ||
* For the full copyright and license information, please view the LICENSE | ||
* file that was distributed with this source code. | ||
*/ | ||
|
||
namespace EdgeTelemetrics\EventCorrelation\Scheduler; | ||
|
||
use Closure; | ||
use Psr\Log\LoggerAwareInterface; | ||
use Psr\Log\LoggerAwareTrait; | ||
use Psr\Log\LoggerInterface; | ||
|
||
class ClosureActionWrapper implements LoggerAwareInterface { | ||
/** PSR3 logger provides $this->logger */ | ||
use LoggerAwareTrait; | ||
|
||
public function __construct(private Closure $closure, LoggerInterface $logger) { | ||
$this->setLogger($logger); | ||
} | ||
|
||
public function run(array $args): void | ||
{ | ||
$this->closure->call($this, $args); | ||
} | ||
} |
Oops, something went wrong.