From a41f6cbf3184cdddd570113a4c77ab669b29da13 Mon Sep 17 00:00:00 2001 From: hakeemRash Date: Sat, 29 Aug 2026 22:39:04 +0300 Subject: [PATCH 1/4] feat(tests): add comprehensive tests for job payload signature and security gateway --- .../Pipelines/Worker/RejectedJobException.php | 28 +++ tests/Fixtures/JobModule/Provider.php | 16 ++ tests/Fixtures/JobModule/module.json | 21 ++ .../Boot/CompileJobManifestStageTest.php | 124 ++++++++++ .../Unit/Kernel/Boot/KernelBootCacheTest.php | 220 ++++++++++++++++++ .../Worker/JobPayloadSignatureTest.php | 155 ++++++++++++ .../Worker/WorkerLoopSignatureTest.php | 155 ++++++++++++ .../Kernel/Security/SecurityGatewayTest.php | 164 +++++++++++++ 8 files changed, 883 insertions(+) create mode 100644 src/Kernel/Pipelines/Worker/RejectedJobException.php create mode 100644 tests/Fixtures/JobModule/Provider.php create mode 100644 tests/Fixtures/JobModule/module.json create mode 100644 tests/Unit/Kernel/Boot/CompileJobManifestStageTest.php create mode 100644 tests/Unit/Kernel/Boot/KernelBootCacheTest.php create mode 100644 tests/Unit/Kernel/Pipelines/Worker/JobPayloadSignatureTest.php create mode 100644 tests/Unit/Kernel/Pipelines/Worker/WorkerLoopSignatureTest.php create mode 100644 tests/Unit/Kernel/Security/SecurityGatewayTest.php diff --git a/src/Kernel/Pipelines/Worker/RejectedJobException.php b/src/Kernel/Pipelines/Worker/RejectedJobException.php new file mode 100644 index 0000000..129976c --- /dev/null +++ b/src/Kernel/Pipelines/Worker/RejectedJobException.php @@ -0,0 +1,28 @@ +root = sys_get_temp_dir() . '/hkm-jobs-' . bin2hex(random_bytes(6)); + mkdir($this->root . '/var/cache/manifests', 0775, true); + + $this->previousBase = Paths::base(); + $this->previousProject = Paths::project(); + Paths::setBase($this->root); + Paths::setProject($this->root); + } + + protected function tearDown(): void + { + Paths::setBase((string) $this->previousBase); + Paths::setProject($this->previousProject); + + foreach (glob($this->root . '/var/cache/manifests/*') ?: [] as $file) { + @unlink($file); + } + @rmdir($this->root . '/var/cache/manifests'); + @rmdir($this->root . '/var/cache'); + @rmdir($this->root . '/var'); + @rmdir($this->root); + } + + /** @return array> */ + private function compile(): array + { + (new CompileJobManifestStage([JobProvider::class], new ManifestReader()))->run(); + + return ManifestReader::readCompiled('job-manifest.php'); + } + + public function test_it_compiles_the_declared_retry_policy(): void + { + $jobs = $this->compile(); + + self::assertSame( + ['max' => 5, 'strategy' => 'linear', 'base' => 30, 'jitter' => false], + $jobs['jobs.send-invoice']['retry'], + ); + } + + public function test_it_compiles_the_declared_timeout(): void + { + self::assertSame(30, $this->compile()['jobs.send-invoice']['timeout']); + } + + public function test_an_undeclared_retry_stays_null(): void + { + // Not a default: "this job said nothing" must stay distinguishable from + // "this job asked for the defaults", so WorkerLoop's own fallback + // remains the single place the default lives. + $jobs = $this->compile(); + + self::assertNull($jobs['jobs.rebuild-index']['retry']); + self::assertNull($jobs['jobs.rebuild-index']['timeout']); + } + + public function test_the_shorthand_retry_count_is_understood(): void + { + $retry = $this->compile()['jobs.prune']['retry']; + + self::assertSame(5, $retry['max']); + self::assertSame('exponential', $retry['strategy'], 'the default strategy'); + } + + public function test_an_unknown_strategy_falls_back_rather_than_failing(): void + { + $retry = $this->compile()['jobs.report']['retry']; + + self::assertSame('exponential', $retry['strategy']); + } + + public function test_a_zero_max_is_raised_to_one(): void + { + // "retry": {"max": 0} would mean the job can never run, which is never + // what it means. + self::assertSame(1, $this->compile()['jobs.report']['retry']['max']); + } + + public function test_a_nonsensical_timeout_is_dropped(): void + { + self::assertNull($this->compile()['jobs.report']['timeout']); + } + + public function test_the_existing_fields_are_unchanged(): void + { + $job = $this->compile()['jobs.send-invoice']; + + self::assertSame('emails', $job['queue']); + self::assertSame('App\\Jobs\\SendInvoice', $job['handler']); + self::assertSame('jobs.demo', $job['solves']); + self::assertSame(JobProvider::class, $job['module']); + } +} diff --git a/tests/Unit/Kernel/Boot/KernelBootCacheTest.php b/tests/Unit/Kernel/Boot/KernelBootCacheTest.php new file mode 100644 index 0000000..dcedbdb --- /dev/null +++ b/tests/Unit/Kernel/Boot/KernelBootCacheTest.php @@ -0,0 +1,220 @@ +root = sys_get_temp_dir() . '/hkm-bootcache-' . bin2hex(random_bytes(6)); + mkdir($this->root . '/var/cache/manifests', 0775, true); + + $this->previousBase = Paths::base(); + $this->previousProject = Paths::project(); + $this->previousEnv = $_ENV['BOOT_CACHE'] ?? false; + } + + protected function tearDown(): void + { + Paths::setBase((string) $this->previousBase); + Paths::setProject($this->previousProject); + + if ($this->previousEnv === false) { + unset($_ENV['BOOT_CACHE']); + } else { + $_ENV['BOOT_CACHE'] = $this->previousEnv; + } + + foreach (glob($this->root . '/var/cache/manifests/*') ?: [] as $file) { + @unlink($file); + } + @rmdir($this->root . '/var/cache/manifests'); + @rmdir($this->root . '/var/cache'); + @rmdir($this->root . '/var'); + @rmdir($this->root); + } + + /** @param list $essentials */ + private function build(array $essentials): Kernel + { + return Kernel::configure() + ->withBasePath($this->root) + ->withProjectPath($this->root) + ->withPorts([ + DatabasePort::class => self::database(), + CachePort::class => self::cache(), + ]) + ->withSecurity([self::securityLayer()]) + ->withModules([Provider::class]) + ->withEssentialModules($essentials) + ->build(); + } + + /** + * Whether the SECOND build recompiled. ManifestWriter renames a fresh temp + * file into place, so a rewrite always changes the inode — which mtime, at + * one-second granularity, would miss entirely inside one test run. + * + * @param list $essentials + */ + private function secondBuildRecompiled(array $essentials): bool + { + $manifest = $this->root . '/var/cache/manifests/route-manifest.php'; + + $this->build($essentials); + clearstatcache(true, $manifest); + $before = fileinode($manifest); + + $this->build($essentials); + clearstatcache(true, $manifest); + + return fileinode($manifest) !== $before; + } + + // ── The regression ────────────────────────────────────────────────────── + + public function test_essentials_declared_as_domains_still_hit_the_cache(): void + { + $_ENV['BOOT_CACHE'] = '1'; + + // 'prefixed.demo' is the fixture's solves domain — the shape proj.json + // "essentials" uses, and the shape the docs recommend. Hashing the + // resolved provider class instead made this case miss forever. + self::assertFalse( + $this->secondBuildRecompiled(['prefixed.demo']), + 'a second build with domain-form essentials must skip the compile', + ); + } + + public function test_it_hits_with_no_essentials_at_all(): void + { + $_ENV['BOOT_CACHE'] = '1'; + + self::assertFalse($this->secondBuildRecompiled([])); + } + + public function test_it_hits_with_essentials_already_given_as_classes(): void + { + $_ENV['BOOT_CACHE'] = '1'; + + self::assertFalse($this->secondBuildRecompiled([Provider::class])); + } + + // ── The flag still has to mean something ──────────────────────────────── + + public function test_it_recompiles_every_build_when_the_flag_is_off(): void + { + unset($_ENV['BOOT_CACHE']); + + self::assertTrue( + $this->secondBuildRecompiled(['prefixed.demo']), + 'without BOOT_CACHE every build must recompile — that is the default', + ); + } + + public function test_changing_a_builder_input_invalidates_the_cache(): void + { + $_ENV['BOOT_CACHE'] = '1'; + + $manifest = $this->root . '/var/cache/manifests/route-manifest.php'; + + $this->build([]); + clearstatcache(true, $manifest); + $before = fileinode($manifest); + + // A different project route is a different application, cache or not. + Kernel::configure() + ->withBasePath($this->root) + ->withProjectPath($this->root) + ->withPorts([DatabasePort::class => self::database(), CachePort::class => self::cache()]) + ->withSecurity([self::securityLayer()]) + ->withModules([Provider::class]) + ->withRoutes([['method' => 'GET', 'path' => '/added', 'handler' => 'A\\C@index']]) + ->build(); + + clearstatcache(true, $manifest); + self::assertNotSame($before, fileinode($manifest)); + } + + // ── Port stubs ────────────────────────────────────────────────────────── + + private static function database(): DatabasePort + { + return new class implements DatabasePort { + public function query(string $sql, array $params = []): array { return []; } + public function queryOne(string $sql, array $params = []): ?array { return null; } + public function execute(string $sql, array $params = []): int { return 0; } + public function upsert(string $t, array $v, array $c, ?array $u = null): int { return 0; } + public function lastInsertId(?string $sequence = null): string { return '1'; } + public function beginTransaction(): void {} + public function commit(): void {} + public function rollback(): void {} + public function inTransaction(): bool { return false; } + public function driver(): string { return 'sqlite'; } + }; + } + + private static function cache(): CachePort + { + return new class implements CachePort { + public function get(string $key): mixed { return null; } + public function set(string $key, mixed $value, ?int $ttl = null): bool { return true; } + public function delete(string $key): bool { return true; } + public function has(string $key): bool { return false; } + public function remember(string $key, int $ttl, callable $callback): mixed { return $callback(); } + public function increment(string $key, int $by = 1): int { return $by; } + public function deletePattern(string $pattern): int { return 0; } + public function flush(): bool { return true; } + public function lock(string $name, int $seconds = 0, ?string $owner = null): Lock + { + throw new \RuntimeException('not needed for boot'); + } + public function restoreLock(string $name, string $owner): Lock + { + throw new \RuntimeException('not needed for boot'); + } + }; + } + + private static function securityLayer(): SecurityLayerContract + { + return new class implements SecurityLayerContract { + public function check(Request $request): SecurityVerdict + { + return SecurityVerdict::allow($request); + } + }; + } +} diff --git a/tests/Unit/Kernel/Pipelines/Worker/JobPayloadSignatureTest.php b/tests/Unit/Kernel/Pipelines/Worker/JobPayloadSignatureTest.php new file mode 100644 index 0000000..b9a9f17 --- /dev/null +++ b/tests/Unit/Kernel/Pipelines/Worker/JobPayloadSignatureTest.php @@ -0,0 +1,155 @@ + $overrides */ + private static function payload(array $overrides = [], string $signature = ''): JobPayload + { + $fields = [ + 'jobId' => 'job-1', + 'jobClass' => 'App\\Jobs\\SendInvoice', + 'data' => ['invoiceId' => 42, 'to' => 'a@b.test'], + 'queue' => 'emails', + 'attempts' => 0, + 'maxAttempts' => 3, + ...$overrides, + ]; + + return new JobPayload( + jobId: $fields['jobId'], + jobClass: $fields['jobClass'], + data: $fields['data'], + queue: $fields['queue'], + attempts: $fields['attempts'], + maxAttempts: $fields['maxAttempts'], + enqueuedAt: new \DateTimeImmutable('@1700000000'), + signature: $signature, + ); + } + + /** A payload signed exactly as a correct QueuePort adapter would sign it. */ + private static function signed(array $overrides = []): JobPayload + { + $unsigned = self::payload($overrides); + + return self::payload($overrides, $unsigned->signatureFor(self::SECRET)); + } + + // ── The happy path ────────────────────────────────────────────────────── + + public function test_a_correctly_signed_payload_verifies(): void + { + self::assertTrue(self::signed()->isSignatureValid(self::SECRET)); + } + + public function test_a_payload_signed_with_another_key_does_not_verify(): void + { + self::assertFalse(self::signed()->isSignatureValid('a-different-key')); + } + + // ── The regression: every producer-authored field is covered ──────────── + + public function test_swapping_the_job_class_invalidates_a_captured_signature(): void + { + $captured = self::signed()->signature(); + + // The whole attack: keep the signature, change what runs. + $tampered = self::payload(['jobClass' => 'App\\Jobs\\GrantAdmin'], $captured); + + self::assertFalse( + $tampered->isSignatureValid(self::SECRET), + 'jobClass decides which code runs and MUST be signed', + ); + } + + public function test_tampering_with_the_data_invalidates_the_signature(): void + { + $captured = self::signed()->signature(); + $tampered = self::payload(['data' => ['invoiceId' => 99, 'to' => 'attacker@evil.test']], $captured); + + self::assertFalse($tampered->isSignatureValid(self::SECRET)); + } + + public function test_widening_the_retry_budget_invalidates_the_signature(): void + { + $captured = self::signed()->signature(); + $tampered = self::payload(['maxAttempts' => 100_000], $captured); + + self::assertFalse($tampered->isSignatureValid(self::SECRET)); + } + + public function test_moving_the_job_to_another_queue_invalidates_the_signature(): void + { + $captured = self::signed()->signature(); + $tampered = self::payload(['queue' => 'high-priority'], $captured); + + self::assertFalse($tampered->isSignatureValid(self::SECRET)); + } + + // ── attempts is queue-managed, so it is deliberately NOT signed ───────── + + public function test_a_retry_still_verifies_after_the_driver_increments_attempts(): void + { + $captured = self::signed()->signature(); + $retried = self::payload(['attempts' => 2], $captured); + + // Signing this would break every job on its first retry, since release() + // increments it. maxAttempts is signed instead, so the budget still holds. + self::assertTrue($retried->isSignatureValid(self::SECRET)); + } + + // ── Unsigned is never a way past the check ────────────────────────────── + + public function test_an_unsigned_payload_never_verifies(): void + { + self::assertFalse(self::payload()->isSignatureValid(self::SECRET)); + } + + public function test_an_empty_secret_never_verifies(): void + { + // The CALLER decides whether verification applies at all. Once it does, + // an empty key must not be a skeleton key. + self::assertFalse(self::signed()->isSignatureValid('')); + } + + // ── Canonicalisation ──────────────────────────────────────────────────── + + public function test_payload_key_order_does_not_change_the_signature(): void + { + // A driver round-tripping the envelope through JSON is under no + // obligation to preserve key order; an unstable input would make the + // HMAC reject its own legitimate messages. + $a = self::payload(['data' => ['b' => 2, 'a' => 1, 'nested' => ['y' => 1, 'x' => 2]]]); + $b = self::payload(['data' => ['a' => 1, 'nested' => ['x' => 2, 'y' => 1], 'b' => 2]]); + + self::assertSame($a->signatureFor(self::SECRET), $b->signatureFor(self::SECRET)); + } + + public function test_list_order_is_meaningful_and_does_change_the_signature(): void + { + $a = self::payload(['data' => ['steps' => ['charge', 'ship']]]); + $b = self::payload(['data' => ['steps' => ['ship', 'charge']]]); + + self::assertNotSame($a->signatureFor(self::SECRET), $b->signatureFor(self::SECRET)); + } +} diff --git a/tests/Unit/Kernel/Pipelines/Worker/WorkerLoopSignatureTest.php b/tests/Unit/Kernel/Pipelines/Worker/WorkerLoopSignatureTest.php new file mode 100644 index 0000000..023168a --- /dev/null +++ b/tests/Unit/Kernel/Pipelines/Worker/WorkerLoopSignatureTest.php @@ -0,0 +1,155 @@ + 1], + queue: 'default', + attempts: 0, + maxAttempts: 3, + enqueuedAt: new \DateTimeImmutable('@1700000000'), + signature: $signature, + ); + } + + private static function signedPayload(string $jobClass = SpyJob::class): JobPayload + { + $unsigned = self::payload('', $jobClass); + + return self::payload($unsigned->signatureFor(self::SECRET), $jobClass); + } + + protected function setUp(): void + { + SpyJob::$handled = 0; + } + + // ── An unverifiable payload is dead-lettered, never acked ─────────────── + + public function test_an_unsigned_payload_is_failed_not_acked(): void + { + $queue = $this->runOnce(new RecordingQueue([self::payload('')]), self::SECRET); + + self::assertSame(['fail'], $queue->calls, 'a rejected payload must go to the dead-letter queue'); + self::assertSame(0, SpyJob::$handled, 'it must never reach the job'); + } + + public function test_a_forged_payload_is_failed_not_acked(): void + { + $queue = $this->runOnce(new RecordingQueue([self::payload('deadbeef')]), self::SECRET); + + self::assertSame(['fail'], $queue->calls); + self::assertSame(0, SpyJob::$handled); + } + + public function test_a_rejected_payload_is_not_retried(): void + { + // release() would put it back for another attempt; a signature that does + // not verify will not verify on the second attempt either. + $queue = $this->runOnce(new RecordingQueue([self::payload('')]), self::SECRET); + + self::assertNotContains('release', $queue->calls); + } + + // ── A correctly signed payload runs normally ──────────────────────────── + + public function test_a_correctly_signed_payload_runs_and_is_acked(): void + { + $queue = $this->runOnce(new RecordingQueue([self::signedPayload()]), self::SECRET); + + self::assertSame(['ack'], $queue->calls); + self::assertSame(1, SpyJob::$handled); + } + + // ── With no secret configured, verification is off ────────────────────── + + public function test_no_secret_leaves_verification_off(): void + { + // The historical behaviour, and still the default: an application whose + // queue nothing else can write to should not have to sign anything. + $queue = $this->runOnce(new RecordingQueue([self::payload('')]), ''); + + self::assertSame(['ack'], $queue->calls); + self::assertSame(1, SpyJob::$handled); + } + + /** Run exactly one job through a loop wired to $queue. */ + private function runOnce(RecordingQueue $queue, string $secret): RecordingQueue + { + $core = new CoreContainer(); + $core->instance(QueuePort::class, $queue); + + $loop = new WorkerLoop($core, ErrorPipeline::notifiers([]), new WorkerPipeline(), $secret); + $loop->run(maxIterations: 1); + + return $queue; + } +} + +/** Records which lifecycle verb the loop resolved a job with. */ +final class RecordingQueue implements QueuePort +{ + /** @var list */ + public array $calls = []; + + /** @param list $pending */ + public function __construct(private array $pending = []) {} + + public function push(string $jobClass, array $payload, string $queue = 'default', int $delay = 0): string { return 'id'; } + public function later(int $seconds, string $jobClass, array $payload, string $queue = 'default'): string { return 'id'; } + public function size(string $queue = 'default'): int { return count($this->pending); } + + public function pop(string $queue = 'default'): ?JobPayload + { + return array_shift($this->pending); + } + + public function ack(JobPayload $payload): void { $this->calls[] = 'ack'; } + public function release(JobPayload $payload, int $delay = 0): void { $this->calls[] = 'release'; } + public function fail(JobPayload $payload, ?\Throwable $reason = null): void { $this->calls[] = 'fail'; } +} + +final class SpyJob implements JobContract +{ + public static int $handled = 0; + + public function handle(JobPayload $payload): JobResult + { + self::$handled++; + + return JobResult::success(); + } + + public function failed(JobPayload $payload, \Throwable $e): void {} +} diff --git a/tests/Unit/Kernel/Security/SecurityGatewayTest.php b/tests/Unit/Kernel/Security/SecurityGatewayTest.php new file mode 100644 index 0000000..fa29e05 --- /dev/null +++ b/tests/Unit/Kernel/Security/SecurityGatewayTest.php @@ -0,0 +1,164 @@ +identity === null + ? SecurityVerdict::allow($request) + : SecurityVerdict::allowWithIdentity($this->identity); + } + }; + } + + /** A layer that records what identity it was handed, then allows. */ + private static function spy(?Identity &$seen): SecurityLayerContract + { + return new class ($seen) implements SecurityLayerContract { + public function __construct(private ?Identity &$seen) {} + + public function check(Request $request): SecurityVerdict + { + $this->seen = $request->identity(); + + return SecurityVerdict::allow($request); + } + }; + } + + private static function denies(int $code, string $reason): SecurityLayerContract + { + return new class ($code, $reason) implements SecurityLayerContract { + public function __construct(private readonly int $code, private readonly string $reason) {} + + public function check(Request $request): SecurityVerdict + { + return SecurityVerdict::deny($this->code, $this->reason); + } + }; + } + + // ── Identity reaches the verdict ──────────────────────────────────────── + + public function test_the_last_layers_identity_reaches_the_verdict(): void + { + $identity = Identity::asUser('user-7', 'tenant-a'); + + $verdict = (new SecurityGateway([self::allows(), self::allows($identity)])) + ->inspect(self::request()); + + self::assertTrue($verdict->isAllowed()); + self::assertSame($identity, $verdict->identity()); + } + + public function test_a_sole_layers_identity_reaches_the_verdict(): void + { + $identity = Identity::asUser('user-1'); + + $verdict = (new SecurityGateway([self::allows($identity)]))->inspect(self::request()); + + self::assertSame($identity, $verdict->identity()); + } + + public function test_no_layer_resolving_an_identity_yields_none(): void + { + $verdict = (new SecurityGateway([self::allows(), self::allows()]))->inspect(self::request()); + + self::assertTrue($verdict->isAllowed()); + self::assertNull($verdict->identity()); + } + + // ── A later layer must still SEE an earlier one's identity ────────────── + + public function test_a_later_layer_sees_an_earlier_layers_identity(): void + { + $identity = Identity::asAdmin('tenant-b'); + $seen = null; + + // This is the whole reason the clone exists on a non-final layer: a role + // check that runs after authentication must be able to read the user. + (new SecurityGateway([self::allows($identity), self::spy($seen)])) + ->inspect(self::request()); + + self::assertSame($identity, $seen); + } + + public function test_the_first_layer_sees_no_identity(): void + { + $seen = null; + + (new SecurityGateway([self::spy($seen), self::allows(Identity::asUser('u'))])) + ->inspect(self::request()); + + self::assertNull($seen); + } + + // ── Denial short-circuits ─────────────────────────────────────────────── + + public function test_a_denial_stops_every_later_layer(): void + { + $seen = null; + $ran = false; + $after = new class ($ran) implements SecurityLayerContract { + public function __construct(private bool &$ran) {} + + public function check(Request $request): SecurityVerdict + { + $this->ran = true; + + return SecurityVerdict::allow($request); + } + }; + + $verdict = (new SecurityGateway([self::denies(403, 'nope'), $after])) + ->inspect(self::request()); + + self::assertTrue($verdict->isDenied()); + self::assertSame(403, $verdict->statusCode()); + self::assertSame('nope', $verdict->reason()); + self::assertFalse($ran, 'a denied request must never reach a later layer'); + self::assertNull($seen); + } + + public function test_an_empty_stack_allows(): void + { + // BindSecurityStage refuses this at boot; the gateway itself must still + // behave rather than index past the end of an empty list. + $verdict = (new SecurityGateway([]))->inspect(self::request()); + + self::assertTrue($verdict->isAllowed()); + } +} From e3763f43e6bd40a7071ef4e996d81831f0372901 Mon Sep 17 00:00:00 2001 From: hakeemRash Date: Sat, 29 Aug 2026 22:39:27 +0300 Subject: [PATCH 2/4] feat(worker): implement HMAC signature verification for job payloads and enhance retry strategies --- modules/http | 2 +- .../Boot/Stages/CompileJobManifestStage.php | 83 ++++++- src/Kernel/Kernel.php | 58 ++++- .../Pipelines/Http/Stages/ResolveStage.php | 14 +- src/Kernel/Pipelines/Worker/JobPayload.php | 131 +++++++++- src/Kernel/Pipelines/Worker/WorkerLoop.php | 229 +++++++++++++++++- src/Kernel/Security/SecurityGateway.php | 25 +- src/Kernel/Security/SecurityVerdict.php | 14 ++ tests/Unit/Kernel/Http/RequestTest.php | 45 ++++ tools/docs/hkm-cli-usage.md | 0 tools/src/commands/run.zig | 37 ++- 11 files changed, 601 insertions(+), 37 deletions(-) mode change 100644 => 100755 tools/docs/hkm-cli-usage.md diff --git a/modules/http b/modules/http index e60823d..d1b1366 160000 --- a/modules/http +++ b/modules/http @@ -1 +1 @@ -Subproject commit e60823d6479882fac66588ffbf73da96cd5210ee +Subproject commit d1b1366d89c7514d262e8403931d5d71723521d0 diff --git a/src/Kernel/Boot/Stages/CompileJobManifestStage.php b/src/Kernel/Boot/Stages/CompileJobManifestStage.php index 1d96e0e..4c590e0 100644 --- a/src/Kernel/Boot/Stages/CompileJobManifestStage.php +++ b/src/Kernel/Boot/Stages/CompileJobManifestStage.php @@ -4,7 +4,22 @@ use AlfacodeTeam\PhpServicePlatform\Kernel\Boot\{ManifestReader, ManifestWriter}; -/** Reads jobs[] from every module.json -> job-manifest.php. */ +/** + * Reads jobs[] from every module.json -> job-manifest.php. + * + * RETRY AND TIMEOUT ARE PART OF THE DECLARATION. + * + * module.json has always documented them — + * + * { "type": "job", "queue": "emails", "timeout": 30, + * "retry": { "max": 3, "strategy": "exponential", "jitter": true } } + * + * — but this stage used to drop both on the floor, so every job in every + * application silently shared one hardcoded exponential strategy and no timeout + * at all. A declaration that compiles to nothing is worse than no declaration: + * it reads as a guarantee. They are compiled through now, and WorkerLoop honours + * them per job. + */ final class CompileJobManifestStage implements BootStageContract { /** @param list $moduleClasses */ @@ -23,15 +38,73 @@ public function run(): void if ($name === null) { continue; } + $spec = is_array($job) ? $job : []; + $jobs[$name] = [ - 'handler' => is_array($job) ? ($job['handler'] ?? $name) : $name, - 'queue' => is_array($job) ? ($job['queue'] ?? 'default') : 'default', - 'module' => $moduleClass, - 'solves' => $manifest['solves'], + 'handler' => $spec['handler'] ?? $name, + 'queue' => $spec['queue'] ?? 'default', + 'module' => $moduleClass, + 'solves' => $manifest['solves'] ?? '', + // Fall back to the module-wide declaration: a module.json + // whose "type" is "job" states retry/timeout at the top + // level, which is the shape the docs show. + 'retry' => self::retry($spec['retry'] ?? $manifest['retry'] ?? null), + 'timeout' => self::timeout($spec['timeout'] ?? $manifest['timeout'] ?? null), ]; } } ManifestWriter::write('job-manifest.php', $jobs); } + + /** + * Normalise a `retry` declaration into a shape WorkerLoop can act on without + * re-parsing JSON per job. + * + * An UNDECLARED retry compiles to null, not to a default: that keeps "this + * job says nothing" distinguishable from "this job asked for the defaults", + * so the loop's own fallback stays the single place the default lives. + * + * @return array{max: int, strategy: string, base: int, jitter: bool}|null + */ + private static function retry(mixed $spec): ?array + { + if ($spec === null) { + return null; + } + + // "retry": 5 — the shorthand for "just give me five attempts". + if (is_int($spec) || (is_string($spec) && ctype_digit($spec))) { + $spec = ['max' => (int) $spec]; + } + + if (!is_array($spec)) { + return null; + } + + $strategy = strtolower((string) ($spec['strategy'] ?? 'exponential')); + + return [ + // A job may not declare fewer than one attempt — "retry": {"max": 0} + // would mean the job can never run, which is never what it means. + 'max' => max(1, (int) ($spec['max'] ?? 3)), + 'strategy' => in_array($strategy, ['exponential', 'linear', 'fixed'], true) + ? $strategy + : 'exponential', + 'base' => max(1, (int) ($spec['base'] ?? $spec['delay'] ?? 1)), + 'jitter' => (bool) ($spec['jitter'] ?? false), + ]; + } + + /** Seconds a single attempt may run, or null when the job declares none. */ + private static function timeout(mixed $spec): ?int + { + if ($spec === null || !is_numeric($spec)) { + return null; + } + + $seconds = (int) $spec; + + return $seconds > 0 ? $seconds : null; + } } diff --git a/src/Kernel/Kernel.php b/src/Kernel/Kernel.php index 2073f97..d94ddfe 100644 --- a/src/Kernel/Kernel.php +++ b/src/Kernel/Kernel.php @@ -52,6 +52,8 @@ final class Kernel private array $projectGroups = []; /** @var list hosts this project serves (proj.json "domains") */ private array $projectDomains = []; + /** HMAC key every dequeued job payload must carry; null = read the env. */ + private ?string $workerSecret = null; private ?ErrorPipeline $errorPipeline = null; private ?\Closure $errorPipelineFun = null; private ?string $basePath = null; @@ -122,6 +124,32 @@ public function withSecurity(array $layers): self return $this; } + /** + * Require every dequeued job payload to be HMAC-signed with this key. + * + * A queue is an input channel: whoever can write to it is calling into the + * application. WorkerLoop has always had the check — but the kernel never + * had a way to give it a key, so in every deployment it was dead code and + * the worker ran whatever it was handed. + * + * Defaults to `JOB_SIGNING_SECRET`, and stays OFF when that is unset, which + * is the historical behaviour. It deliberately does NOT fall back to + * APP_KEY: that would switch verification on for every existing application + * at once and reject every job already in flight, since no QueuePort adapter + * signs by default. + * + * TURNING IT ON IS A TWO-SIDED CHANGE. The producing adapter must stamp + * {@see \AlfacodeTeam\PhpServicePlatform\Kernel\Pipelines\Worker\JobPayload::signatureFor()} + * onto the envelope at push() time. Until it does, every payload is rejected + * as unsigned and dead-lettered — correct, but not a discovery to make in + * production. Roll it out producer-first. + */ + public function withWorkerSecret(string $secret): self + { + $this->workerSecret = $secret; + return $this; + } + public function withErrorPipeline(ErrorPipeline|callable|\Closure $pipeline): self { if (is_callable($pipeline)) { @@ -357,7 +385,16 @@ public function build(): self // When BOOT_CACHE is on and nothing the compile read has changed, skip // the compilation and keep only the validation stages, which touch no // disk and must still catch a missing port or an unusable layer. - $cached = BootStamp::enabled() ? BootStamp::read($this->buildHash()) : null; + // Computed ONCE, from the RAW builder inputs, and reused for the write + // below. It must NOT be recomputed after resolveEssentialModules(): + // that turns proj.json's "essentials" domains into provider classes, so + // a hash taken afterwards would never equal the one the next build looks + // the stamp up with — the cache would miss on every request, and pay for + // a stamp rewrite on top of the recompile it failed to skip. + $cacheEnabled = BootStamp::enabled(); + $stampHash = $cacheEnabled ? $this->buildHash() : ''; + + $cached = $cacheEnabled ? BootStamp::read($stampHash) : null; if ($cached !== null) { $pipeline->runValidationOnly(); @@ -377,8 +414,9 @@ public function build(): self // pipeline's reader, whose module.json cache is already warm. $this->essentialModules = $this->resolveEssentialModules($reader); - if (BootStamp::enabled()) { - BootStamp::write($this->buildHash(), $reader->files(), $this->essentialModules); + if ($cacheEnabled) { + // $stampHash, not a fresh buildHash() — see the note above. + BootStamp::write($stampHash, $reader->files(), $this->essentialModules); } $this->built = true; @@ -392,6 +430,13 @@ public function build(): self * essentials, disable policy), so hashing these covers a proj.json edit * without stat'ing it — and covers an edit to bootstrap/app.php itself, * which no file-mtime check would catch. + * + * CALL THIS BEFORE resolveEssentialModules(), AND ONLY ONCE PER BUILD. + * `essentials` here is a BUILDER INPUT — the raw list the project passed, + * which for proj.json is domains ('tenancy.routing'). resolveEssentialModules() + * replaces it with the DERIVED provider classes, so a hash taken afterwards + * describes a different array and can never match the one the next build + * reads with. The derived list belongs in the stamp's payload, not its key. */ private function buildHash(): string { @@ -489,7 +534,12 @@ private function materialize(RuntimeMode $mode): void essentialModules: $this->essentialModules, ); $this->cli = new CliPipeline($this->core, $errorPipeline); - $this->workerLoop = new WorkerLoop($this->core, $errorPipeline, $this->workerPipe); + $this->workerLoop = new WorkerLoop( + $this->core, + $errorPipeline, + $this->workerPipe, + $this->workerSecret ?? (string) (env('JOB_SIGNING_SECRET') ?: ''), + ); // Configuration compiled by CompileConfigManifestStage during build(). // Bound BEFORE module boot() so a Provider can read config while wiring. diff --git a/src/Kernel/Pipelines/Http/Stages/ResolveStage.php b/src/Kernel/Pipelines/Http/Stages/ResolveStage.php index 006d1a7..9699fed 100644 --- a/src/Kernel/Pipelines/Http/Stages/ResolveStage.php +++ b/src/Kernel/Pipelines/Http/Stages/ResolveStage.php @@ -50,12 +50,14 @@ public function handle(Request $request, callable $next): Response return Response::notFound(); } - $request = $request - ->withAttribute('route_entry', $match['entry']) - ->withAttribute('route_params', $match['params']) - ->withAttribute('target_service', $match['entry']['solves']); - - return $next($request); + // One clone, not three: chaining withAttribute() would build two + // intermediate requests — each a deep clone of all seven parameter bags + // — that nothing ever reads. + return $next($request->withAttributes([ + 'route_entry' => $match['entry'], + 'route_params' => $match['params'], + 'target_service' => $match['entry']['solves'], + ])); } /** diff --git a/src/Kernel/Pipelines/Worker/JobPayload.php b/src/Kernel/Pipelines/Worker/JobPayload.php index 73880e5..f41794b 100644 --- a/src/Kernel/Pipelines/Worker/JobPayload.php +++ b/src/Kernel/Pipelines/Worker/JobPayload.php @@ -5,6 +5,40 @@ // ─── JobPayload ─────────────────────────────────────────────────────────────── +/** + * A dequeued unit of work, exactly as the queue handed it over. + * + * ── THE SIGNATURE COVERS THE WHOLE ENVELOPE ────────────────────────────────── + * + * A queue is an input channel like any other. Whoever can write to it is, in + * effect, calling into the application — so the worker must be able to tell a + * payload the application produced from one somebody else planted. + * + * The signature therefore covers every PRODUCER-AUTHORED field: + * + * HMAC_SHA256( secret, jobId | jobClass | queue | maxAttempts | canonical(data) ) + * + * Signing `data` ALONE — as this class used to — left `jobClass` unauthenticated, + * which is the field that decides WHICH CODE RUNS. An attacker who captured one + * legitimately signed envelope could keep the signature, swap the class for any + * other JobContract in the application, and have the worker execute it. Nothing + * about that requires forging a signature, only reusing one. + * + * `attempts` is deliberately NOT signed: it is a QUEUE-managed counter that the + * driver increments on every release(), so covering it would invalidate the + * signature on a job's first retry. `maxAttempts` IS signed, so the retry budget + * the producer set cannot be widened in transit. + * + * `data` is canonicalised (keys sorted, recursively) before signing, because a + * driver that round-trips the envelope through JSON is under no obligation to + * preserve key order — and an unstable input makes an HMAC reject its own + * legitimate messages. + * + * VERIFICATION IS OFF UNTIL A SECRET IS CONFIGURED. See + * {@see \AlfacodeTeam\PhpServicePlatform\Kernel\Kernel::withWorkerSecret()}; + * with no secret the worker runs every payload it is given, which is the + * historical behaviour and is only safe when the queue itself is trusted. + */ final readonly class JobPayload { public function __construct( @@ -27,15 +61,106 @@ public function maxAttempts(): int { return $this->maxAttempts; } public function enqueuedAt(): \DateTimeImmutable { return $this->enqueuedAt; } public function signature(): string { return $this->signature; } + /** + * The signature this payload SHOULD carry — what a QueuePort adapter stamps + * onto the envelope at push() time. + * + * Producing and verifying through the same method is the point: two + * hand-rolled HMACs that must agree forever is how signing quietly stops + * working. + */ + public function signatureFor(string $secret): string + { + return self::sign( + $secret, + $this->jobId, + $this->jobClass, + $this->data, + $this->queue, + $this->maxAttempts, + ); + } + + /** + * Sign an envelope's producer-authored fields. + * + * @param array $data + */ + public static function sign( + string $secret, + string $jobId, + string $jobClass, + array $data, + string $queue, + int $maxAttempts, + ): string { + return hash_hmac('sha256', implode('|', [ + $jobId, + $jobClass, + $queue, + (string) $maxAttempts, + self::canonical($data), + ]), $secret); + } + + /** + * Whether this payload was signed with $secret. + * + * An empty secret or an empty signature is NEVER valid: the caller decides + * whether verification applies at all, and once it does, "unsigned" must not + * be a way past it. + */ public function isSignatureValid(string $secret): bool { - $expected = hash_hmac('sha256', json_encode($this->data), $secret); - return hash_equals($expected, $this->signature); + if ($secret === '' || $this->signature === '') { + return false; + } + + return hash_equals($this->signatureFor($secret), $this->signature); } public function hasExceededMaxAttempts(): bool { return $this->attempts >= $this->maxAttempts; } -} + /** + * A stable string for an arbitrary payload array — keys sorted at every + * depth so encoding order cannot change the signed material. + * + * Falls back to serialize() if the data is not JSON-encodable (a resource, + * NAN, malformed UTF-8). Signing SOMETHING deterministic is what matters; + * throwing here would turn an odd payload into a worker crash. + * + * @param array $data + */ + private static function canonical(array $data): string + { + $sorted = self::sortRecursive($data); + + $json = json_encode($sorted, JSON_UNESCAPED_SLASHES | JSON_UNESCAPED_UNICODE); + + return $json === false ? serialize($sorted) : $json; + } + + /** + * @param array $data + * @return array + */ + private static function sortRecursive(array $data): array + { + foreach ($data as $key => $value) { + if (is_array($value)) { + $data[$key] = self::sortRecursive($value); + } + } + + // Lists keep their order — it is meaningful. Only associative keys are + // sorted, because there their order is an encoding artefact. + if (!array_is_list($data)) { + ksort($data); + } + + return $data; + } +} diff --git a/src/Kernel/Pipelines/Worker/WorkerLoop.php b/src/Kernel/Pipelines/Worker/WorkerLoop.php index 0b4cb0a..41dd5a0 100644 --- a/src/Kernel/Pipelines/Worker/WorkerLoop.php +++ b/src/Kernel/Pipelines/Worker/WorkerLoop.php @@ -7,7 +7,12 @@ use AlfacodeTeam\PhpServicePlatform\Kernel\Error\{ErrorPipeline, ErrorContext}; use AlfacodeTeam\PhpServicePlatform\Kernel\Loading\{DependencyGraphCalculator, OnDemandLoader}; use AlfacodeTeam\PhpServicePlatform\Kernel\Pipelines\Worker\Contracts\JobContract; -use AlfacodeTeam\PhpServicePlatform\Kernel\Pipelines\Worker\Retry\{ExponentialRetryStrategy, RetryStrategyContract}; +use AlfacodeTeam\PhpServicePlatform\Kernel\Pipelines\Worker\Retry\{ + ExponentialRetryStrategy, + FixedRetryStrategy, + LinearRetryStrategy, + RetryStrategyContract +}; use AlfacodeTeam\PhpServicePlatform\Kernel\Ports\QueuePort; use AlfacodeTeam\PhpServicePlatform\Kernel\Support\Paths; @@ -40,16 +45,44 @@ final class WorkerLoop private ?DependencyGraphCalculator $calculator = null; private ?OnDemandLoader $loader = null; + /** + * Per-job retry strategies, built once from the compiled manifest. + * + * @var array + */ + private array $retryByJob = []; + + /** Whether a shutdown signal has already been trapped for this loop. */ + private bool $signalsInstalled = false; + public function __construct( private readonly CoreContainer $core, private readonly ErrorPipeline $errorPipeline, private readonly WorkerPipeline $pipeline, + /** + * HMAC key every dequeued payload must be signed with. EMPTY = no + * verification, which is the historical behaviour and is only safe when + * nothing but the application can write to the queue. + * + * Set it via Kernel::withWorkerSecret() (or JOB_SIGNING_SECRET). Turning + * it on requires the QueuePort adapter to stamp + * {@see JobPayload::signatureFor()} onto the envelope at push() time — + * otherwise every job is rejected as unsigned, which is the correct + * behaviour but a surprising way to discover the requirement. + */ private readonly string $signingSecret = '', - /** Backoff applied by release() in port mode. */ + /** Backoff applied by release() when a job declares no retry of its own. */ private readonly RetryStrategyContract $retry = new ExponentialRetryStrategy(), ) { } + /** + * Ask the loop to finish the job it is on and then exit. + * + * Safe to call from a signal handler: it only sets a flag, so the current + * job still reaches its ack/release/fail instead of being torn down + * mid-flight. + */ public function stop(): void { $this->shouldStop = true; @@ -73,9 +106,18 @@ public function stop(): void * @param (callable():?JobPayload)|null $puller null = use the bound QueuePort * @param int $maxIterations 0 = run forever (until stop()). * @param string $queue which queue to drain in port mode + * @param int $memoryLimitMb stop the loop once the process exceeds this + * resident size, 0 to disable. A long-running PHP process accumulates + * fragmentation no amount of correctness prevents, so a supervised + * worker is EXPECTED to exit and be restarted; the only question is + * whether it does so between jobs or by being OOM-killed mid-job. */ - public function run(?callable $puller = null, int $maxIterations = 0, string $queue = 'default'): void - { + public function run( + ?callable $puller = null, + int $maxIterations = 0, + string $queue = 'default', + int $memoryLimitMb = 0, + ): void { $port = $puller === null ? $this->queuePort() : null; if ($puller === null && $port === null) { @@ -85,12 +127,30 @@ public function run(?callable $puller = null, int $maxIterations = 0, string $qu ); } + $this->trapShutdownSignals(); + $iterations = 0; - while (!$this->shouldStop) { + while (true) { if ($maxIterations > 0 && $iterations++ >= $maxIterations) { break; } + // Deliver any SIGTERM/SIGINT that arrived while the last job ran, + // BEFORE deciding whether to take another one. Checking shouldStop + // only in the `while` condition would let one more job start after + // the signal — which is the opposite of a graceful shutdown. + if (\function_exists('pcntl_signal_dispatch')) { + pcntl_signal_dispatch(); + } + + if ($this->shouldStop) { + break; + } + + if ($memoryLimitMb > 0 && memory_get_usage(true) >= $memoryLimitMb * 1048576) { + break; + } + $payload = $port !== null ? $port->pop($queue) : $puller(); if ($payload === null) { usleep(100_000); // idle backoff @@ -123,6 +183,14 @@ private function processWithPort(QueuePort $port, JobPayload $payload): void // Completed, skipped, or already dead-lettered by process() — either // way it must not come back. Removing it is the whole point of ack. $port->ack($payload); + } catch (RejectedJobException $e) { + // A payload the worker refuses to run at all. Retrying cannot help — + // a bad signature never becomes good — so it goes straight to the + // dead-letter queue, where an operator can see it. It must NOT be + // acked: silently deleting the one piece of evidence that someone is + // writing to your queue is the worst possible response to it. + $this->report($e, $payload); + $port->fail($payload, $e); } catch (\Throwable $e) { if ($payload->hasExceededMaxAttempts()) { $port->fail($payload, $e); @@ -130,7 +198,138 @@ private function processWithPort(QueuePort $port, JobPayload $payload): void return; } - $port->release($payload, $this->retry->delayFor($payload->attempts() + 1)); + $port->release( + $payload, + $this->retryFor($payload->jobClass())->delayFor($payload->attempts() + 1), + ); + } + } + + /** + * The retry strategy a job DECLARED in its module.json, or the loop-wide + * default when it declared none. + * + * Built once per job class and reused — the manifest is a deploy-time + * artefact, so the strategy cannot change under a running worker. + */ + private function retryFor(string $jobName): RetryStrategyContract + { + if (isset($this->retryByJob[$jobName])) { + return $this->retryByJob[$jobName]; + } + + // Load it here too: a job that threw before resolveContainer() ran (an + // unknown class, a container failure) still needs a backoff. + $this->jobManifest ??= $this->loadManifest('job-manifest.php', []); + + $spec = $this->jobManifest[$jobName]['retry'] ?? null; + + if (!is_array($spec)) { + return $this->retryByJob[$jobName] = $this->retry; + } + + $base = (int) ($spec['base'] ?? 1); + $jitter = (bool) ($spec['jitter'] ?? false); + + return $this->retryByJob[$jobName] = match ($spec['strategy'] ?? 'exponential') { + 'linear' => new LinearRetryStrategy($base), + 'fixed' => new FixedRetryStrategy($base), + default => new ExponentialRetryStrategy($base, jitter: $jitter), + }; + } + + /** + * Run a job under the `timeout` its module.json declares. + * + * BEST EFFORT, AND THE LIMITS MATTER. pcntl_alarm delivers SIGALRM, which + * PHP dispatches between opcodes — so it interrupts a runaway loop, but a + * job blocked inside a single long DB query or socket read is not preempted + * until that call returns. A timeout that must hold regardless belongs in + * the driver (a statement timeout, a socket timeout), not here. + * + * It is still worth having: the failure this catches — a job that spins and + * pins a worker until someone notices — is the one that takes a queue down. + * + * Without ext-pcntl, or with no declared timeout, the job runs unbounded, + * exactly as before. + */ + private function runWithTimeout(JobContract $job, JobPayload $payload): JobResult + { + $seconds = ($this->jobManifest ?? [])[$payload->jobClass()]['timeout'] ?? null; + + if (!is_int($seconds) || $seconds <= 0 || !\function_exists('pcntl_alarm')) { + return $job->handle($payload); + } + + $previous = pcntl_signal_get_handler(\SIGALRM); + + pcntl_signal(\SIGALRM, static function () use ($payload, $seconds): never { + throw new \RuntimeException(sprintf( + 'Job [%s] (id %s) exceeded its declared timeout of %ds.', + $payload->jobClass(), + $payload->jobId() !== '' ? $payload->jobId() : '?', + $seconds, + )); + }); + + pcntl_alarm($seconds); + + try { + return $job->handle($payload); + } finally { + // Cancel first, THEN restore: an alarm that fires during teardown + // would surface as a timeout for whatever job runs next. + pcntl_alarm(0); + pcntl_signal(\SIGALRM, $previous); + } + } + + /** Route a rejected payload through the error pipeline so it is not silent. */ + private function report(\Throwable $e, JobPayload $payload): void + { + $this->errorPipeline->consume(ErrorContext::fromThrowable( + $e, + requestPath: 'job:' . $payload->jobClass(), + requestMethod: 'WORKER', + )); + } + + /** + * Trap SIGTERM/SIGINT so a shutdown finishes the current job first. + * + * Without this, the signal every process supervisor and container runtime + * sends to stop a worker (SIGTERM, then SIGKILL after a grace period) killed + * PHP outright — including in the window between handle() returning and + * ack() removing the message. A job that had already run its side effects + * came back on the next boot and ran them AGAIN. Trapping the signal turns + * that from luck into a guarantee. + * + * ext-pcntl is optional and absent on some builds; without it the loop keeps + * its old behaviour rather than refusing to start. + */ + private function trapShutdownSignals(): void + { + if ($this->signalsInstalled || !\function_exists('pcntl_signal')) { + return; + } + + $this->signalsInstalled = true; + + // Without this, PHP only runs a signal handler where the script calls + // pcntl_signal_dispatch() — which the loop does between jobs, but which + // is no help to the per-job alarm in runWithTimeout(): a job that hangs + // never reaches the next dispatch point. Async delivery is what makes + // the timeout able to interrupt anything at all. + if (\function_exists('pcntl_async_signals')) { + pcntl_async_signals(true); + } + + $stop = function (): void { + $this->stop(); + }; + + foreach ([\SIGTERM, \SIGINT, \SIGQUIT] as $signal) { + pcntl_signal($signal, $stop); } } @@ -149,8 +348,22 @@ private function queuePort(): ?QueuePort private function process(JobPayload $payload): JobResult { // 1. Signature check — never run an unsigned/tampered payload. + // + // This used to return skipped(), which processWithPort then ACKED: a + // payload that failed authentication was deleted without a trace, so a + // misconfigured producer and an active attacker looked identical, and + // both looked like nothing at all. Throwing routes it to the dead-letter + // queue and through the error pipeline instead. if ($this->signingSecret !== '' && !$payload->isSignatureValid($this->signingSecret)) { - return JobResult::skipped('Invalid job signature — payload rejected.'); + throw new RejectedJobException(sprintf( + 'Rejected job [%s] (id %s) on queue [%s]: %s. ' + . 'The producing QueuePort adapter must stamp JobPayload::signatureFor() ' + . 'onto the envelope at push() time.', + $payload->jobClass(), + $payload->jobId() !== '' ? $payload->jobId() : '?', + $payload->queue(), + $payload->signature() === '' ? 'payload is unsigned' : 'signature does not verify', + )); } $jobClass = $this->pipeline->resolve($payload->jobClass()) @@ -170,7 +383,7 @@ private function process(JobPayload $payload): JobResult : ($this->core->has($jobClass) ? $this->core->make($jobClass) : new $jobClass()); try { - return $job->handle($payload); + return $this->runWithTimeout($job, $payload); } catch (\Throwable $e) { $this->errorPipeline->consume(ErrorContext::fromThrowable( $e, diff --git a/src/Kernel/Security/SecurityGateway.php b/src/Kernel/Security/SecurityGateway.php index 508c9d5..b9c8a13 100644 --- a/src/Kernel/Security/SecurityGateway.php +++ b/src/Kernel/Security/SecurityGateway.php @@ -31,16 +31,33 @@ public function __construct( public function inspect(Request $request): SecurityVerdict { - foreach ($this->layers as $layer) { + $last = \count($this->layers) - 1; + + foreach ($this->layers as $i => $layer) { $verdict = $layer->check($request); if ($verdict->isDenied()) { return $verdict; // Short-circuit — nothing else runs } - // If this layer resolved an identity, attach it to the request - if ($verdict->identity() !== null) { - $request = $request->withIdentity($verdict->identity()); + // If this layer resolved an identity, attach it so LATER layers can + // see it (a role check after authentication). + // + // Only when there IS a later layer. withIdentity() deep-clones all + // seven parameter bags, and on the last layer that clone exists for + // exactly one statement: the allow() below, which reads the identity + // straight off it. SecurityStage then clones a second time to put the + // identity on the request the pipeline actually carries. The typical + // stack — CSRF, then an Auth layer that resolves the identity last — + // therefore paid for a whole request copy nothing ever read. + $identity = $verdict->identity(); + + if ($identity !== null) { + if ($i === $last) { + return SecurityVerdict::allowWithIdentity($identity); + } + + $request = $request->withIdentity($identity); } } diff --git a/src/Kernel/Security/SecurityVerdict.php b/src/Kernel/Security/SecurityVerdict.php index f0bf20a..efda4a1 100644 --- a/src/Kernel/Security/SecurityVerdict.php +++ b/src/Kernel/Security/SecurityVerdict.php @@ -19,6 +19,20 @@ public static function allow(Request $request): self return new self(true, 200, '', $request->identity()); } + /** + * Allow, carrying an identity that is NOT yet attached to a request. + * + * allow() reads the identity back off a request, which forces a caller that + * has only just resolved one to clone the whole request first — seven + * parameter bags — purely so this constructor can read one property off the + * copy. SecurityStage attaches the identity to the request the pipeline + * carries anyway, so that first clone was always thrown away. + */ + public static function allowWithIdentity(?Identity $identity): self + { + return new self(true, 200, '', $identity); + } + public static function deny(int $statusCode, string $reason): self { return new self(false, $statusCode, $reason, null); diff --git a/tests/Unit/Kernel/Http/RequestTest.php b/tests/Unit/Kernel/Http/RequestTest.php index dba9156..12baff7 100644 --- a/tests/Unit/Kernel/Http/RequestTest.php +++ b/tests/Unit/Kernel/Http/RequestTest.php @@ -58,4 +58,49 @@ public function test_with_identity_is_immutable_and_readable(): void self::assertNull($original->identity()); self::assertSame('u42', $withId->identity()->userId); } + + public function test_with_attributes_sets_several_in_one_instance(): void + { + $original = Request::create('/x', 'GET'); + $modified = $original->withAttributes([ + 'route_entry' => ['solves' => 'invoice.generation'], + 'route_params' => ['id' => '7'], + 'target_service' => 'invoice.generation', + ]); + + self::assertNotSame($original, $modified); + self::assertNull($original->attribute('route_entry'), 'the original must be untouched'); + self::assertSame(['id' => '7'], $modified->attribute('route_params')); + self::assertSame('invoice.generation', $modified->attribute('target_service')); + } + + public function test_with_attributes_matches_chaining_with_attribute(): void + { + // It exists purely to avoid the intermediate deep clones; if it ever + // diverged in EFFECT from the chain it replaces, that would be a bug in + // the routing stage that now calls it. + $original = Request::create('/x', 'GET'); + + $chained = $original->withAttribute('a', 1)->withAttribute('b', ['c' => 2]); + $batched = $original->withAttributes(['a' => 1, 'b' => ['c' => 2]]); + + self::assertSame($chained->attributes->all(), $batched->attributes->all()); + } + + public function test_with_attributes_keeps_earlier_attributes(): void + { + $request = Request::create('/x', 'GET') + ->withAttribute('locale', 'fr') + ->withAttributes(['route_params' => []]); + + self::assertSame('fr', $request->attribute('locale')); + } + + public function test_with_no_attributes_returns_the_same_instance(): void + { + // Nothing to copy, so nothing is copied. + $original = Request::create('/x', 'GET'); + + self::assertSame($original, $original->withAttributes([])); + } } diff --git a/tools/docs/hkm-cli-usage.md b/tools/docs/hkm-cli-usage.md old mode 100644 new mode 100755 diff --git a/tools/src/commands/run.zig b/tools/src/commands/run.zig index 064cc3d..459b93a 100644 --- a/tools/src/commands/run.zig +++ b/tools/src/commands/run.zig @@ -50,12 +50,25 @@ const Options = struct { }; pub fn run(allocator: std.mem.Allocator, io: Io, env: *EnvMap, args: []const []const u8) !u8 { - // Bare `hkm run` (no arguments) prints help rather than serving silently. - if (args.len <= 2) { - printHelp(); - return 2; - } - + // NO ARG-COUNT GUARD HERE. + // + // There used to be `if (args.len <= 2) { printHelp(); return 2; }`, which + // contradicted this file's own contract three ways: the module docblock + // above ("with no argument the current directory is used", `hkm run` as the + // first example), the help text ("serve a project (defaults to ./)"), and + // the resolver below, which already handles an empty target by resolving + // ".". The default-to-cwd path was fully implemented and simply unreachable. + // + // It also made `hkm run --dev` fail in a way nobody could reason about. + // --dev is stripped BEFORE command parsing, so that invocation arrives here + // as exactly ["hkm", "run"] — length 2 — and printed usage. Adding any + // unrelated flag (`hkm run --dev --port=8000`) got past the count and then + // worked perfectly, which makes the failure look like it is about --dev, or + // about the directory, rather than about how many words were typed. + // + // Resolution now belongs entirely to resolveRoot(): in a project directory + // `hkm run` serves it, and anywhere else the error below names the actual + // problem instead of dumping a usage screen that does not mention it. var opts = (try parse(allocator, args)) orelse { // parse() returned null: the arguments were invalid. printHelp(); @@ -90,6 +103,18 @@ pub fn run(allocator: std.mem.Allocator, io: Io, env: *EnvMap, args: []const []c "'{s}' is neither a project folder (with proj.json) nor a registered name.", .{if (opts.target.len == 0) "." else opts.target}, )); + + // A bare `hkm run` outside a project is the one case where the user + // named nothing at all, so there is no spelling to check — point at the + // ways to name one instead. This is what the old arg-count guard was + // really reaching for, except it fired even INSIDE a project. + if (opts.target.len == 0) { + prompt.hintLine("run it from a project folder, or name one:"); + prompt.hint("hkm run ", "serve a specific project"); + prompt.hint("hkm run --pick", "choose from the registered projects"); + prompt.hint("hkm list", "show what is registered"); + } + return 1; }; From 73c086c5cb2e75ee620431c7c84fabc283b6286c Mon Sep 17 00:00:00 2001 From: hakeemRash Date: Sat, 29 Aug 2026 22:45:53 +0300 Subject: [PATCH 3/4] =?UTF-8?q?release:=20v1.6.0=20=E2=80=94=20make=20the?= =?UTF-8?q?=20queue=20an=20authenticated=20channel,=20and=20make=20BOOT=5F?= =?UTF-8?q?CACHE=20work?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Folds the outstanding Unreleased entries into 1.6.0 and adds this cycle's work. A MINOR bump, not a patch: it adds public API callers can now depend on — Kernel::withWorkerSecret(), Request::withAttributes(), SecurityVerdict::allowWithIdentity(), and a memoryLimitMb parameter on WorkerLoop::run(). The job-signature material changed shape, which would normally be breaking. It is not: verification was unreachable before this release (WorkerLoop's secret had no path in from the Kernel and was always ''), so no deployment has signed payloads in flight to migrate. --- CHANGELOG.md | 102 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 102 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 679d3e2..a6b643c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,7 +6,84 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [1.6.0] - 2026-08-29 + +### Added +- **`Kernel::withWorkerSecret()` — the queue can finally be an authenticated + channel.** `WorkerLoop` has always carried a signature check, but the kernel + had no way to give it a key: `$signingSecret` defaulted to `''`, was never + passed at construction, and there was no builder method. In every deployment + that has ever run, the check was dead code and the worker executed whatever it + was handed. A queue is an input channel — whoever can write to it is calling + into the application — so this closes a hole, not a nicety. Defaults to + `JOB_SIGNING_SECRET` and stays OFF when that is unset, preserving today's + behaviour. It deliberately does **not** fall back to `APP_KEY`: that would + switch verification on for every existing application at once and reject every + job already in flight, since no `QueuePort` adapter signs by default. Turning + it on is a two-sided change — roll it out producer-first, teaching the adapter + to stamp `JobPayload::signatureFor()` at `push()` time. +- **Graceful worker shutdown, a memory ceiling, and per-job timeouts.** There + was no `pcntl` anywhere in the kernel, so SIGTERM — what every process + supervisor and container runtime sends to stop a worker — killed PHP outright, + including in the window between `handle()` returning and `ack()` removing the + message. A job that had already run its side effects came back on the next + boot and ran them **again**. The loop now traps SIGTERM/SIGINT/SIGQUIT, + finishes the job it is on, resolves its ack/release/fail, and exits. + `run()` takes a `memoryLimitMb` so a supervised worker exits between jobs + rather than being OOM-killed inside one, and a job's declared `timeout` is + enforced with `pcntl_alarm` — best effort, since SIGALRM is dispatched between + opcodes and cannot preempt a job blocked inside one long query. +- **`Request::withAttributes()`** — set several attributes in a single new + instance. Every `with*()` deep-clones all seven parameter bags, so a chain of + them pays that price once per link. `ResolveStage` attaching `route_entry`, + `route_params` and `target_service` is one logical step that cost three full + clones of a request nothing had read yet: **10.02 µs → 3.57 µs, 64% less**. +- **`SecurityVerdict::allowWithIdentity()`** — allow while carrying an identity + that is not yet attached to a request. `allow()` reads the identity back *off* + a request, forcing a layer that has just resolved one to clone the entire + request so the constructor can read a single property. + ### Fixed +- **`BOOT_CACHE` never hit for the essentials shape the docs recommend.** + `Kernel::build()` computed `buildHash()` twice — before and after + `resolveEssentialModules()`, which rewrites `essentials` from proj.json's + DOMAINS (`tenancy.routing`) into provider CLASSES. So the stamp was written + under one hash and read under another, and every request recompiled all ten + manifests **and** rewrote the stamp on top of the recompile it had failed to + skip — measurably *worse* than leaving the flag off. Measured on a three-route + application: **2604 µs → 39 µs per request under PHP-FPM.** The hash is now + taken once, from the raw builder inputs; the derived class list rides in the + stamp's payload, never its key. `BootStampTest` tests the stamp in isolation + and could not see this, so `KernelBootCacheTest` builds twice through the real + `Kernel::build()` and watches the manifest inode. +- **A job payload that failed verification was silently deleted.** The check + returned `skipped()`, which `processWithPort` then **acked** — removing the one + piece of evidence that something is writing to your queue. A misconfigured + producer and an active attacker were indistinguishable, and both looked like + nothing happening at all. An unverifiable payload now raises + `RejectedJobException`, goes through the `ErrorPipeline`, and is dead-lettered + via `fail()`. It is never retried: a signature that does not verify will not + verify on the second attempt. +- **`retry` and `timeout` in `module.json` compiled to nothing.** + `CompileJobManifestStage` read `handler`, `queue`, `module` and `solves` and + dropped the other two, so every job in every application shared one hardcoded + exponential strategy and ran unbounded — while its manifest said otherwise. A + declaration that compiles to nothing is worse than no declaration: it reads as + a guarantee. Both are compiled through now and honoured per job, including the + `"retry": 5` shorthand; an unknown strategy falls back rather than failing the + boot, and `"max": 0` is raised to 1 (a job that can never run is never what it + meant). +- **`hkm run` ignored its own documented default of `./`.** An `args.len <= 2` + guard printed usage and exited 2 before the resolver ever ran, contradicting + the command's module docblock, its help text, and the `resolveRoot()` call + below it — which already handled an empty target. It bit `hkm run --dev` + hardest: `--dev` is stripped before command parsing, so that invocation + arrived as exactly `["hkm", "run"]` and failed, while adding any unrelated flag + (`--port=8000`) got past the count and worked perfectly — making the failure + look like it was about `--dev`, or about the directory, rather than about how + many words were typed. Resolution now belongs entirely to `resolveRoot()`, and + a bare `hkm run` outside a project names the actual problem instead of dumping + a usage screen that does not mention it. - **The Homebrew bump job failed a release that had already published.** Its PR fallback pushed the bump branch, then called `gh pr create` — which the API refuses unless *Settings → Actions → General → "Allow GitHub Actions to create @@ -22,6 +99,31 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 … from untrusted tap". `brew trust alfacode-team/hkm` is now part of the documented sequence, in the README and in the formula's own header. +### Changed +- **A job signature now covers the whole envelope, not just `data`.** Signing + `data` alone left `jobClass` — the field that decides WHICH CODE RUNS — + unauthenticated. Capturing one legitimately signed envelope and swapping its + class for any other `JobContract` was enough; nothing about that required + forging a signature, only reusing one. The material is now + `jobId | jobClass | queue | maxAttempts | canonical(data)`. `attempts` is + deliberately excluded: the driver increments it on every `release()`, so + covering it would invalidate a job on its first retry — `maxAttempts` is signed + instead, so the retry budget cannot be widened in transit. The payload is + canonicalised (associative keys sorted at every depth, list order preserved) + because a driver round-tripping the envelope through JSON is under no + obligation to keep key order, and an unstable input makes an HMAC reject its + own legitimate messages. **No migration is required**: verification was + unreachable before this release, so no deployment has signed payloads in + flight. `JobPayload::signatureFor()` is the one implementation both producer + and verifier use. +- **`SecurityGateway` no longer clones the request on its final layer.** The + clone existed so `SecurityVerdict::allow()` could read the identity back off + it, and `SecurityStage` then cloned a second time to put that identity on the + request the pipeline actually carries — so for the documented CSRF-then-Auth + stack, where the last layer is the one that authenticates, a whole request copy + was built and read once. A later layer still sees an earlier layer's identity; + only the final layer takes the shortcut. **4.17 µs → 0.87 µs, 79% less.** + ## [1.5.0] - 2026-08-28 ### Added From b5ee45688fea75e7b14a31b9d8e10e7bf37c33dc Mon Sep 17 00:00:00 2001 From: hakeemRash Date: Sat, 29 Aug 2026 22:46:47 +0300 Subject: [PATCH 4/4] fix(http): pin the submodule to a commit that exists on the remote MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pointer was d1b1366, which lived only in a local checkout on a detached HEAD. Nothing could resolve it: `git submodule update --init` fails for a clone and for CI, and `Kernel\Http\Request` never loads — while ResolveStage now calls withAttributes(), which exists in no pushed commit at all. Releasing against that pin would have shipped a kernel whose routing calls a method the published http package does not have. c4fe527 is the same change cherry-picked onto origin/main (AlfaCode-Team/http#7) and pushed, so the pin resolves. --- modules/http | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/modules/http b/modules/http index d1b1366..c4fe527 160000 --- a/modules/http +++ b/modules/http @@ -1 +1 @@ -Subproject commit d1b1366d89c7514d262e8403931d5d71723521d0 +Subproject commit c4fe527a4317a347cf10dd630f21a9a33564234e