Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

task delay #41

Merged
merged 4 commits into from
Dec 16, 2023
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
40 changes: 21 additions & 19 deletions composer.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

69 changes: 59 additions & 10 deletions src/TaskManager.php
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,11 @@
use kuaukutsu\poc\task\event\PublisherEvent;
use kuaukutsu\poc\task\event\ProcessEvent;
use kuaukutsu\poc\task\event\ProcessTimeoutEvent;
use kuaukutsu\poc\task\event\ProcessContextEvent;
use kuaukutsu\poc\task\event\ProcessExceptionEvent;
use kuaukutsu\poc\task\exception\ProcessingException;
use kuaukutsu\poc\task\processing\TaskProcess;
use kuaukutsu\poc\task\processing\TaskProcessContext;
use kuaukutsu\poc\task\processing\TaskProcessFactory;
use kuaukutsu\poc\task\processing\TaskProcessing;

Expand All @@ -32,6 +34,11 @@ final class TaskManager implements EventPublisherInterface
*/
private array $processesActive = [];

/**
* @var array<non-empty-string, TaskProcessContext>
*/
private array $processesDelay = [];

/**
* A unique identifier that can be used to cancel, enable or disable the callback.
*/
Expand All @@ -58,7 +65,10 @@ public function run(TaskManagerOptions $options = new TaskManagerOptions()): voi
function () use ($options): void {
$this->trigger(
Event::LoopTick,
new LoopTickEvent(new DateTimeImmutable())
new LoopTickEvent(
count($this->processesActive),
count($this->processesDelay),
)
);

try {
Expand Down Expand Up @@ -88,11 +98,15 @@ function () use ($options): void {
$this->keeperId = EventLoop::repeat(
$options->getKeeperInterval(),
function () use ($options): void {
if ($this->processesActive === []) {
if ($this->processesActive === [] && $this->processesDelay === []) {
$this->keeperDisable();
return;
}

foreach ($this->processesDelay as $context) {
$this->processDelayActive($context, $options);
}

foreach ($this->processesActive as $process) {
if ($process->isRunning() === false) {
$this->trigger(
Expand Down Expand Up @@ -230,12 +244,43 @@ private function processRun(TaskManagerOptions $options): void
&& count($this->processesActive) < $options->getQueueSize()
) {
$context = $this->processing->getTaskProcess();
if (array_key_exists($context->stage, $this->processesActive) === false) {
$process = $this->processFactory->create($context, $options);
$process->start();

$this->processPush($process);
if ($context->timestamp > 0 && $context->timestamp > time()) {
$this->processDelayPush($context);
continue;
}

$this->processStart($context, $options);
}
}

private function processStart(TaskProcessContext $context, TaskManagerOptions $options): void
{
if (array_key_exists($context->getHash(), $this->processesActive) === false) {
$process = $this->processFactory->create($context, $options);
$process->start();

$this->processPush($process);
}
}

private function processDelayPush(TaskProcessContext $context): void
{
if (
array_key_exists($context->getHash(), $this->processesDelay)
&& $this->processesDelay[$context->getHash()]->timestamp > $context->timestamp
) {
$context = $this->processesDelay[$context->getHash()];
}

$this->processesDelay[$context->getHash()] = $context;
}

private function processDelayActive(TaskProcessContext $context, TaskManagerOptions $options): void
{
if ($context->timestamp < time()) {
$this->processStart($context, $options);
$this->trigger(Event::ProcessDelay, new ProcessContextEvent($context));
unset($this->processesDelay[$context->getHash()]);
}
}

Expand All @@ -245,15 +290,19 @@ private function processPush(TaskProcess $process): void
$this->keeperEnable();
}

$this->processesActive[$process->stage] = $process;
$this->processesActive[$process->hash] = $process;
$this->trigger(Event::ProcessPush, new ProcessEvent($process));
}

private function processPull(TaskProcess $process): void
{
$process->stop(0);
unset($this->processesActive[$process->stage]);
if ($this->processesActive === []) {
unset(
$this->processesActive[$process->hash],
$this->processesDelay[$process->hash],
);

if ($this->processesActive === [] && $this->processesDelay === []) {
$this->keeperDisable();
}

Expand Down
2 changes: 2 additions & 0 deletions src/event/Event.php
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ enum Event: string

case ProcessStop = 'process-stop-event';

case ProcessDelay = 'process-delay-event';

case ProcessSuccess = 'process-success-event';

case ProcessError = 'process-error-event';
Expand Down
7 changes: 3 additions & 4 deletions src/event/LoopTickEvent.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,18 +4,17 @@

namespace kuaukutsu\poc\task\event;

use DateTimeImmutable;

/**
* @psalm-immutable
*/
final class LoopTickEvent implements EventInterface
{
private readonly string $message;

public function __construct(DateTimeImmutable $tick)
public function __construct(int $countProcessActive, int $countProcessDelay)
{
$this->message = 'tick: ' . $tick->format('c');
$datetime = gmdate('c');
$this->message = "tick: $datetime, active: $countProcessActive, delay: $countProcessDelay.";
}

public function getMessage(): string
Expand Down
32 changes: 32 additions & 0 deletions src/event/ProcessContextEvent.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
<?php

declare(strict_types=1);

namespace kuaukutsu\poc\task\event;

use kuaukutsu\poc\task\processing\TaskProcessContext;

final class ProcessContextEvent implements EventInterface
{
private readonly string $uuid;

private readonly string $message;

public function __construct(TaskProcessContext $context)
{
$this->uuid = $context->task;

$time = date('c', $context->timestamp);
$this->message = "[{$context->getHash()}] timer $time: " . $context->task;
}

public function getUuid(): string
{
return $this->uuid;
}

public function getMessage(): string
{
return $this->message;
}
}
17 changes: 1 addition & 16 deletions src/handler/StageHandler.php
Original file line number Diff line number Diff line change
Expand Up @@ -32,22 +32,7 @@ public function handle(string $taskUuid, string $stageUuid, ?string $previous =
try {
$this->stdout(
$this->stateSerialize(
$this->handler->complete($taskUuid)
)
);

return TaskProcess::SUCCESS;
} catch (ProcessingException $exception) {
$this->stderr($exception->getMessage());
return TaskProcess::ERROR;
}
}

if ((new TaskCommand($stageUuid))->isState()) {
try {
$this->stdout(
$this->stateSerialize(
$this->handler->state($taskUuid)
$this->handler->complete($taskUuid, $stageUuid)
)
);

Expand Down
Loading