From 4f724d57c9db7e3aa00796a4c422a458a2c168df Mon Sep 17 00:00:00 2001 From: nexxo Date: Sat, 25 Jul 2026 09:54:52 +0200 Subject: [PATCH] feat(engine): orchestrator core (state machine + tick + lock) StepResult, ProvisioningStep contract, PipelineRegistry, RunRunner (per-run lock, advance/retry/fail, backoff, timeout, append-only events + Reverb StepAdvanced), AdvanceRunJob (provisioning queue), minutely Tick, admin.runs channel. 21 tests. Co-Authored-By: Claude Opus 4.8 --- app/Console/TickProvisioning.php | 25 +++ app/Providers/AppServiceProvider.php | 16 +- .../Contracts/ProvisioningStep.php | 24 +++ app/Provisioning/Events/StepAdvanced.php | 36 +++++ app/Provisioning/Jobs/AdvanceRunJob.php | 30 ++++ app/Provisioning/PipelineRegistry.php | 41 +++++ app/Provisioning/RunRunner.php | 146 ++++++++++++++++++ app/Provisioning/StepResult.php | 35 +++++ config/provisioning.php | 47 ++++++ routes/channels.php | 3 + routes/console.php | 8 + .../Provisioning/PipelineRegistryTest.php | 21 +++ tests/Feature/Provisioning/RunRunnerTest.php | 132 ++++++++++++++++ .../Provisioning/TickProvisioningTest.php | 20 +++ tests/Support/Steps/FakeAdvanceStep.php | 30 ++++ tests/Support/Steps/FakeFailStep.php | 30 ++++ tests/Support/Steps/FakeRetryStep.php | 30 ++++ tests/Support/Steps/FakeThrowStep.php | 31 ++++ tests/Unit/StepResultTest.php | 14 ++ 19 files changed, 718 insertions(+), 1 deletion(-) create mode 100644 app/Console/TickProvisioning.php create mode 100644 app/Provisioning/Contracts/ProvisioningStep.php create mode 100644 app/Provisioning/Events/StepAdvanced.php create mode 100644 app/Provisioning/Jobs/AdvanceRunJob.php create mode 100644 app/Provisioning/PipelineRegistry.php create mode 100644 app/Provisioning/RunRunner.php create mode 100644 app/Provisioning/StepResult.php create mode 100644 config/provisioning.php create mode 100644 tests/Feature/Provisioning/PipelineRegistryTest.php create mode 100644 tests/Feature/Provisioning/RunRunnerTest.php create mode 100644 tests/Feature/Provisioning/TickProvisioningTest.php create mode 100644 tests/Support/Steps/FakeAdvanceStep.php create mode 100644 tests/Support/Steps/FakeFailStep.php create mode 100644 tests/Support/Steps/FakeRetryStep.php create mode 100644 tests/Support/Steps/FakeThrowStep.php create mode 100644 tests/Unit/StepResultTest.php diff --git a/app/Console/TickProvisioning.php b/app/Console/TickProvisioning.php new file mode 100644 index 0000000..614a005 --- /dev/null +++ b/app/Console/TickProvisioning.php @@ -0,0 +1,25 @@ +whereIn('status', [ProvisioningRun::STATUS_RUNNING, ProvisioningRun::STATUS_WAITING]) + ->where(function ($q) { + $q->whereNull('next_attempt_at')->orWhere('next_attempt_at', '<=', now()); + }) + ->get() + ->each(fn (ProvisioningRun $run) => AdvanceRunJob::dispatch($run->uuid)); + } +} diff --git a/app/Providers/AppServiceProvider.php b/app/Providers/AppServiceProvider.php index 452e6b6..2d6bdf7 100644 --- a/app/Providers/AppServiceProvider.php +++ b/app/Providers/AppServiceProvider.php @@ -2,6 +2,13 @@ namespace App\Providers; +use App\Provisioning\PipelineRegistry; +use App\Services\Proxmox\HttpProxmoxClient; +use App\Services\Proxmox\ProxmoxClient; +use App\Services\Ssh\PhpseclibRemoteShell; +use App\Services\Ssh\RemoteShell; +use App\Services\Wireguard\LocalWireguardHub; +use App\Services\Wireguard\WireguardHub; use Illuminate\Support\ServiceProvider; class AppServiceProvider extends ServiceProvider @@ -11,7 +18,14 @@ class AppServiceProvider extends ServiceProvider */ public function register(): void { - // + $this->app->singleton(PipelineRegistry::class, fn () => new PipelineRegistry( + config('provisioning.pipelines', []), + )); + + // Real I/O implementations; tests swap in fakes via app()->instance(). + $this->app->bind(RemoteShell::class, PhpseclibRemoteShell::class); + $this->app->bind(WireguardHub::class, LocalWireguardHub::class); + $this->app->bind(ProxmoxClient::class, HttpProxmoxClient::class); } /** diff --git a/app/Provisioning/Contracts/ProvisioningStep.php b/app/Provisioning/Contracts/ProvisioningStep.php new file mode 100644 index 0000000..4e55334 --- /dev/null +++ b/app/Provisioning/Contracts/ProvisioningStep.php @@ -0,0 +1,24 @@ +. */ + public function label(): string; + + /** Seconds this step may run before the runner treats it as timed out. */ + public function maxDuration(): int; + + public function execute(ProvisioningRun $run): StepResult; +} diff --git a/app/Provisioning/Events/StepAdvanced.php b/app/Provisioning/Events/StepAdvanced.php new file mode 100644 index 0000000..6416b7d --- /dev/null +++ b/app/Provisioning/Events/StepAdvanced.php @@ -0,0 +1,36 @@ + */ + public function broadcastOn(): array + { + return [new PrivateChannel('admin.runs')]; + } + + public function broadcastAs(): string + { + return 'StepAdvanced'; + } +} diff --git a/app/Provisioning/Jobs/AdvanceRunJob.php b/app/Provisioning/Jobs/AdvanceRunJob.php new file mode 100644 index 0000000..1ac5f63 --- /dev/null +++ b/app/Provisioning/Jobs/AdvanceRunJob.php @@ -0,0 +1,30 @@ +onQueue('provisioning'); + } + + public function handle(RunRunner $runner): void + { + $run = ProvisioningRun::query()->where('uuid', $this->runUuid)->first(); + + if ($run !== null) { + $runner->advance($run); + } + } +} diff --git a/app/Provisioning/PipelineRegistry.php b/app/Provisioning/PipelineRegistry.php new file mode 100644 index 0000000..a982526 --- /dev/null +++ b/app/Provisioning/PipelineRegistry.php @@ -0,0 +1,41 @@ +>> $pipelines */ + public function __construct(private array $pipelines) {} + + /** @return array> */ + public function stepsFor(string $pipeline): array + { + if (! isset($this->pipelines[$pipeline])) { + throw new InvalidArgumentException("Unknown pipeline [{$pipeline}]."); + } + + return $this->pipelines[$pipeline]; + } + + public function count(string $pipeline): int + { + return count($this->stepsFor($pipeline)); + } + + public function resolve(string $pipeline, int $index): ProvisioningStep + { + $steps = $this->stepsFor($pipeline); + + if (! isset($steps[$index])) { + throw new InvalidArgumentException("No step at index {$index} for pipeline [{$pipeline}]."); + } + + return app($steps[$index]); + } +} diff --git a/app/Provisioning/RunRunner.php b/app/Provisioning/RunRunner.php new file mode 100644 index 0000000..ab7197a --- /dev/null +++ b/app/Provisioning/RunRunner.php @@ -0,0 +1,146 @@ +uuid, 120); + + if (! $lock->get()) { + return; // another worker already holds this run + } + + try { + $this->runLocked($run); + } finally { + $lock->release(); + } + } + + private function runLocked(ProvisioningRun $run): void + { + $run->refresh(); + + if (! in_array($run->status, [ + ProvisioningRun::STATUS_PENDING, + ProvisioningRun::STATUS_RUNNING, + ProvisioningRun::STATUS_WAITING, + ], true)) { + return; + } + + // started_at marks when the current step began (drives the timeout). + if ($run->started_at === null || $run->status === ProvisioningRun::STATUS_PENDING) { + $run->started_at ??= now(); + } + $run->status = ProvisioningRun::STATUS_RUNNING; + $run->save(); + + $step = $this->registry->resolve($run->pipeline, $run->current_step); + + $timedOut = $run->started_at !== null + && $run->started_at->copy()->addSeconds($step->maxDuration())->isPast(); + + try { + $result = $timedOut + ? StepResult::retry(0, 'step timed out after '.$step->maxDuration().'s') + : $step->execute($run); + } catch (Throwable $e) { + $result = StepResult::retry($this->backoff($run->attempt), 'exception: '.$e->getMessage()); + } + + $this->apply($run, $step, $result); + } + + private function apply(ProvisioningRun $run, ProvisioningStep $step, StepResult $result): void + { + match ($result->type) { + StepResult::ADVANCE => $this->onAdvance($run, $step), + StepResult::RETRY => $this->onRetry($run, $step, $result), + StepResult::FAIL => $this->onFail($run, $step, $result), + }; + } + + private function onAdvance(ProvisioningRun $run, ProvisioningStep $step): void + { + $isLast = $run->current_step >= $this->registry->count($run->pipeline) - 1; + + if ($isLast) { + $run->status = ProvisioningRun::STATUS_COMPLETED; + $run->finished_at = now(); + $run->attempt = 0; + $run->save(); + $this->record($run, $step, 'advanced', null); + + return; + } + + $run->current_step += 1; + $run->attempt = 0; + $run->started_at = now(); + $run->status = ProvisioningRun::STATUS_RUNNING; + $run->save(); + $this->record($run, $step, 'advanced', null); + + AdvanceRunJob::dispatch($run->uuid); + } + + private function onRetry(ProvisioningRun $run, ProvisioningStep $step, StepResult $result): void + { + $run->attempt += 1; + + if ($run->attempt >= $run->max_attempts) { + $run->status = ProvisioningRun::STATUS_FAILED; + $run->error = $result->reason; + $run->save(); + $this->record($run, $step, 'failed', $result->reason); + + return; + } + + $run->status = ProvisioningRun::STATUS_WAITING; + $run->next_attempt_at = now()->addSeconds($result->afterSeconds); + $run->save(); + $this->record($run, $step, 'retry', $result->reason); + } + + private function onFail(ProvisioningRun $run, ProvisioningStep $step, StepResult $result): void + { + $run->status = ProvisioningRun::STATUS_FAILED; + $run->error = $result->reason; + $run->save(); + $this->record($run, $step, 'failed', $result->reason); + } + + private function record(ProvisioningRun $run, ProvisioningStep $step, string $outcome, ?string $message): void + { + $run->events()->create([ + 'step' => $step->key(), + 'attempt' => $run->attempt, + 'outcome' => $outcome, + 'message' => $message, + ]); + + StepAdvanced::dispatch($run->uuid, $step->key(), $outcome, $run->current_step, $run->status); + } + + private function backoff(int $attempt): int + { + return (int) min(300, 15 * (2 ** $attempt)); + } +} diff --git a/app/Provisioning/StepResult.php b/app/Provisioning/StepResult.php new file mode 100644 index 0000000..5b34288 --- /dev/null +++ b/app/Provisioning/StepResult.php @@ -0,0 +1,35 @@ + [ + 'host' => [ + Host\ValidateHostInput::class, + Host\EstablishSshTrust::class, + Host\PrepareBaseSystem::class, + Host\ConfigureWireguard::class, + Host\InstallProxmoxVe::class, + Host\RebootIntoPveKernel::class, + Host\ConfigureProxmox::class, + Host\CreateAutomationToken::class, + Host\VerifyProxmoxApi::class, + Host\RegisterCapacity::class, + Host\CompleteHostOnboarding::class, + ], + ], + + // CluPilot VM acts as the WireGuard hub; hosts join it during onboarding. + 'wireguard' => [ + 'subnet' => env('CLUPILOT_WG_SUBNET', '10.66.0.0/24'), + 'endpoint' => env('CLUPILOT_WG_ENDPOINT', ''), // host:port reachable by peers + 'hub_public_key' => env('CLUPILOT_WG_HUB_PUBKEY', ''), + 'config_path' => env('CLUPILOT_WG_CONFIG_PATH', '/etc/wireguard/wg0.conf'), + ], + + // SSH identity CluPilot deploys to each host after first password login. + 'ssh' => [ + 'public_key' => env('CLUPILOT_SSH_PUBLIC_KEY', ''), + 'private_key' => env('CLUPILOT_SSH_PRIVATE_KEY', ''), + ], + + // Proxmox automation role/user created on each host. + 'proxmox' => [ + 'role_id' => 'CluPilotAutomation', + 'role_privs' => 'VM.Allocate,VM.Clone,VM.Config.Disk,VM.Config.CPU,VM.Config.Memory,VM.Config.Network,VM.Config.Options,VM.Config.Cloudinit,VM.PowerMgmt,VM.Monitor,VM.Audit,VM.GuestAgent.Audit,VM.GuestAgent.Unrestricted,Datastore.AllocateSpace,Datastore.Audit,Sys.Audit', + 'user' => 'automation@pve', + 'token_name' => 'clupilot', + ], +]; diff --git a/routes/channels.php b/routes/channels.php index df2ad28..18b4ea1 100644 --- a/routes/channels.php +++ b/routes/channels.php @@ -5,3 +5,6 @@ use Illuminate\Support\Facades\Broadcast; Broadcast::channel('App.Models.User.{id}', function ($user, $id) { return (int) $user->id === (int) $id; }); + +// Operator console live provisioning feed — admins only. +Broadcast::channel('admin.runs', fn ($user) => (bool) $user->is_admin); diff --git a/routes/console.php b/routes/console.php index 3c9adf1..368dd2e 100644 --- a/routes/console.php +++ b/routes/console.php @@ -1,8 +1,16 @@ comment(Inspiring::quote()); })->purpose('Display an inspiring quote'); + +// Advance every due provisioning run once a minute (runs in the scheduler service). +Schedule::call(fn () => app(TickProvisioning::class)()) + ->everyMinute() + ->name('provisioning-tick') + ->withoutOverlapping(); diff --git a/tests/Feature/Provisioning/PipelineRegistryTest.php b/tests/Feature/Provisioning/PipelineRegistryTest.php new file mode 100644 index 0000000..ebee57d --- /dev/null +++ b/tests/Feature/Provisioning/PipelineRegistryTest.php @@ -0,0 +1,21 @@ + [FakeAdvanceStep::class, FakeRetryStep::class]]); + + expect($registry->count('host'))->toBe(2) + ->and($registry->resolve('host', 0))->toBeInstanceOf(FakeAdvanceStep::class) + ->and($registry->resolve('host', 1))->toBeInstanceOf(FakeRetryStep::class); +}); + +it('throws for an unknown pipeline', function () { + (new PipelineRegistry([]))->stepsFor('nope'); +})->throws(InvalidArgumentException::class); + +it('throws for an out-of-range step index', function () { + (new PipelineRegistry(['host' => [FakeAdvanceStep::class]]))->resolve('host', 5); +})->throws(InvalidArgumentException::class); diff --git a/tests/Feature/Provisioning/RunRunnerTest.php b/tests/Feature/Provisioning/RunRunnerTest.php new file mode 100644 index 0000000..3abadb0 --- /dev/null +++ b/tests/Feature/Provisioning/RunRunnerTest.php @@ -0,0 +1,132 @@ +instance(PipelineRegistry::class, new PipelineRegistry($pipelines)); +} + +it('advances to the next step and queues the follow-up job', function () { + Queue::fake(); + bindPipeline(['test' => [FakeAdvanceStep::class, FakeAdvanceStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'current_step' => 0, 'status' => 'pending']); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->current_step)->toBe(1) + ->and($run->status)->toBe('running') + ->and($run->events()->where('outcome', 'advanced')->count())->toBe(1); + Queue::assertPushed(AdvanceRunJob::class); +}); + +it('completes the run on the last step', function () { + Queue::fake(); + bindPipeline(['test' => [FakeAdvanceStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'status' => 'pending']); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->status)->toBe('completed')->and($run->finished_at)->not->toBeNull(); + Queue::assertNotPushed(AdvanceRunJob::class); +}); + +it('waits and increments attempt on retry', function () { + Queue::fake(); + bindPipeline(['test' => [FakeRetryStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'max_attempts' => 5]); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->status)->toBe('waiting') + ->and($run->attempt)->toBe(1) + ->and($run->next_attempt_at)->not->toBeNull(); + Queue::assertNotPushed(AdvanceRunJob::class); +}); + +it('fails once retries are exhausted', function () { + bindPipeline(['test' => [FakeRetryStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'max_attempts' => 1]); + + app(RunRunner::class)->advance($run); + + expect($run->refresh()->status)->toBe('failed'); +}); + +it('fails immediately on a fail result', function () { + bindPipeline(['test' => [FakeFailStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test']); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->status)->toBe('failed')->and($run->error)->toBe('boom'); +}); + +it('treats a thrown exception as a retry', function () { + bindPipeline(['test' => [FakeThrowStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'max_attempts' => 5]); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->status)->toBe('waiting') + ->and($run->attempt)->toBe(1) + ->and($run->events()->where('outcome', 'retry')->exists())->toBeTrue(); +}); + +it('retries when a step exceeds its max duration', function () { + bindPipeline(['test' => [FakeAdvanceStep::class]]); + $run = ProvisioningRun::factory()->create([ + 'pipeline' => 'test', + 'status' => 'running', + 'started_at' => now()->subMinutes(10), + 'max_attempts' => 5, + ]); + + app(RunRunner::class)->advance($run); + $run->refresh(); + + expect($run->status)->toBe('waiting') + ->and($run->events()->where('message', 'like', '%timed out%')->exists())->toBeTrue(); +}); + +it('is a no-op while the run lock is held', function () { + bindPipeline(['test' => [FakeAdvanceStep::class, FakeAdvanceStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test', 'current_step' => 0]); + + $lock = Cache::lock('run:'.$run->uuid, 120); + expect($lock->get())->toBeTrue(); + + app(RunRunner::class)->advance($run); + expect($run->refresh()->current_step)->toBe(0); + + $lock->release(); +}); + +it('broadcasts a StepAdvanced event', function () { + Event::fake([StepAdvanced::class]); + bindPipeline(['test' => [FakeAdvanceStep::class]]); + $run = ProvisioningRun::factory()->create(['pipeline' => 'test']); + + app(RunRunner::class)->advance($run); + + Event::assertDispatched( + StepAdvanced::class, + fn ($e) => $e->runUuid === $run->uuid && $e->status === 'completed', + ); +}); diff --git a/tests/Feature/Provisioning/TickProvisioningTest.php b/tests/Feature/Provisioning/TickProvisioningTest.php new file mode 100644 index 0000000..97e6f9c --- /dev/null +++ b/tests/Feature/Provisioning/TickProvisioningTest.php @@ -0,0 +1,20 @@ +create(['status' => 'waiting', 'next_attempt_at' => now()->subMinute()]); + ProvisioningRun::factory()->create(['status' => 'running', 'next_attempt_at' => null]); + ProvisioningRun::factory()->create(['status' => 'waiting', 'next_attempt_at' => now()->addHour()]); + ProvisioningRun::factory()->create(['status' => 'completed']); + ProvisioningRun::factory()->create(['status' => 'paused']); + + app(TickProvisioning::class)(); + + Queue::assertPushed(AdvanceRunJob::class, 2); +}); diff --git a/tests/Support/Steps/FakeAdvanceStep.php b/tests/Support/Steps/FakeAdvanceStep.php new file mode 100644 index 0000000..ec65836 --- /dev/null +++ b/tests/Support/Steps/FakeAdvanceStep.php @@ -0,0 +1,30 @@ +type)->toBe(StepResult::ADVANCE); + + $retry = StepResult::retry(30, 'wait'); + expect($retry->type)->toBe(StepResult::RETRY) + ->and($retry->afterSeconds)->toBe(30) + ->and($retry->reason)->toBe('wait'); + + expect(StepResult::fail('boom')->reason)->toBe('boom'); +});