Skip to content
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
6 changes: 5 additions & 1 deletion backend/app/Console/Commands/SyncScheduledPlugins.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
use App\Jobs\RunPluginVersionIngestion;
use App\Models\Plugin;
use App\Services\Ingestion\IngestionService;
use App\Support\CorrelationId;
use Illuminate\Console\Command;

class SyncScheduledPlugins extends Command
Expand All @@ -27,7 +28,10 @@ public function handle(IngestionService $service): int
return;
}

$run = $service->start($pluginVersion, 'scheduled');
// Each scheduled run is its own logical operation, not part
// of an inbound request, so it gets its own fresh
// correlation ID rather than sharing one across the batch.
$run = $service->start($pluginVersion, 'scheduled', correlationId: CorrelationId::generate());

if ($run->wasRecentlyCreated) {
RunPluginVersionIngestion::dispatch($pluginVersion->id, $run->id);
Expand Down
27 changes: 25 additions & 2 deletions backend/app/Http/Controllers/Api/V1/GitHubWebhookController.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
use App\Jobs\RunPluginVersionIngestion;
use App\Models\Plugin;
use App\Services\Ingestion\IngestionService;
use App\Support\CorrelationId;
use App\Support\Telemetry;
use Illuminate\Http\JsonResponse;
use Illuminate\Http\Request;
use Illuminate\Support\Str;
Expand All @@ -15,7 +17,14 @@ class GitHubWebhookController extends Controller
{
public function __invoke(Request $request, IngestionService $service): JsonResponse
{
$correlationId = CorrelationId::fromRequest($request);

if (! $this->hasValidSignature($request)) {
Telemetry::event('webhook.github.rejected', [
'correlation_id' => $correlationId,
'reason' => 'invalid_signature',
]);

return response()->json([
'message' => 'Invalid GitHub webhook signature.',
'code' => 'webhook.invalid_signature',
Expand All @@ -26,6 +35,13 @@ public function __invoke(Request $request, IngestionService $service): JsonRespo
$branch = Str::after($request->string('ref')->toString(), 'refs/heads/');

if ($repository === '' || $branch === '') {
Telemetry::event('webhook.github.accepted', [
'correlation_id' => $correlationId,
'repository' => $repository !== '' ? $repository : null,
'branch' => $branch !== '' ? $branch : null,
'queued' => 0,
]);

return response()->json([
'message' => 'Ignored GitHub webhook payload.',
'code' => 'webhook.ignored',
Expand All @@ -39,21 +55,28 @@ public function __invoke(Request $request, IngestionService $service): JsonRespo
->with(['versions' => fn ($query) => $query->orderByDesc('is_latest')->orderByDesc('id')])
->get()
->filter(fn (Plugin $plugin): bool => $this->matchesPlugin($plugin, $repository, $branch))
->each(function (Plugin $plugin) use ($service, &$queued): void {
->each(function (Plugin $plugin) use ($service, $correlationId, &$queued): void {
$pluginVersion = $plugin->versions->first();

if ($pluginVersion === null) {
return;
}

$run = $service->start($pluginVersion, 'webhook');
$run = $service->start($pluginVersion, 'webhook', correlationId: $correlationId);

if ($run->wasRecentlyCreated) {
RunPluginVersionIngestion::dispatch($pluginVersion->id, $run->id);
$queued++;
}
});

Telemetry::event('webhook.github.accepted', [
'correlation_id' => $correlationId,
'repository' => $repository,
'branch' => $branch,
'queued' => $queued,
]);

return response()->json([
'message' => 'GitHub webhook processed.',
'code' => 'webhook.processed',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,7 @@ private function payload(IngestionRun $run): array
'id' => $run->id,
'plugin_version_id' => $run->plugin_version_id,
'source' => $run->source,
'correlation_id' => $run->correlation_id,
'status' => $run->status,
'stats' => $run->stats,
'log' => $run->log,
Expand Down
7 changes: 5 additions & 2 deletions backend/app/Http/Controllers/Api/V1/PluginSyncController.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,13 @@
use App\Jobs\RunPluginVersionIngestion;
use App\Models\Plugin;
use App\Services\Ingestion\IngestionService;
use App\Support\CorrelationId;
use Illuminate\Http\JsonResponse;
use Illuminate\Http\Request;

class PluginSyncController extends Controller
{
public function __invoke(Plugin $plugin, IngestionService $service): JsonResponse
public function __invoke(Request $request, Plugin $plugin, IngestionService $service): JsonResponse
{
$plugin->loadMissing('versions');

Expand All @@ -20,7 +22,7 @@ public function __invoke(Plugin $plugin, IngestionService $service): JsonRespons
->first()
?? $plugin->versions()->latest('id')->firstOrFail();

$run = $service->start($pluginVersion, 'manual');
$run = $service->start($pluginVersion, 'manual', correlationId: CorrelationId::fromRequest($request));

if ($run->wasRecentlyCreated) {
RunPluginVersionIngestion::dispatch($pluginVersion->id, $run->id);
Expand All @@ -31,6 +33,7 @@ public function __invoke(Plugin $plugin, IngestionService $service): JsonRespons
'id' => $run->id,
'plugin_version_id' => $run->plugin_version_id,
'source' => $run->source,
'correlation_id' => $run->correlation_id,
'status' => $run->status,
'stats' => $run->stats,
'log' => $run->log,
Expand Down
31 changes: 31 additions & 0 deletions backend/app/Http/Middleware/AssignCorrelationId.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
<?php

namespace App\Http\Middleware;

use App\Support\CorrelationId;
use Closure;
use Illuminate\Http\Request;
use Symfony\Component\HttpFoundation\Response;

/**
* Assigns every API request a correlation ID: a valid client-supplied
* `X-Request-Id` is reused, otherwise one is generated. The ID is exposed
* to downstream code via the request attribute bag and echoed back on the
* response so a caller can correlate their request with server-side logs
* and ingestion runs.
*/
class AssignCorrelationId
{
public function handle(Request $request, Closure $next): Response
{
$incoming = $request->header(CorrelationId::HEADER);
$correlationId = CorrelationId::isValid($incoming) ? $incoming : CorrelationId::generate();

$request->attributes->set(CorrelationId::REQUEST_ATTRIBUTE, $correlationId);

$response = $next($request);
$response->headers->set(CorrelationId::HEADER, $correlationId);

return $response;
}
}
1 change: 1 addition & 0 deletions backend/app/Models/IngestionRun.php
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ class IngestionRun extends Model
protected $fillable = [
'plugin_version_id',
'source',
'correlation_id',
'status',
'stats',
'log',
Expand Down
75 changes: 71 additions & 4 deletions backend/app/Services/Ingestion/IngestionService.php
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@
use App\Models\PluginVersion;
use App\Services\Discovery\DiscoveryCache;
use App\Services\Markdown\MarkdownParser;
use App\Support\CorrelationId;
use App\Support\Telemetry;
use Illuminate\Support\Facades\Auth;
use Illuminate\Support\Facades\DB;
use Illuminate\Support\Str;
use Throwable;
Expand All @@ -27,9 +30,20 @@ public function __construct(
private readonly DiscoveryCache $cache,
) {}

public function start(PluginVersion $pluginVersion, string $source = 'manual', bool $force = false): IngestionRun
/**
* Start (or reuse) an ingestion run.
*
* The correlation ID identifies the caller's request, not the run: if
* an active run is reused, the run keeps the correlation ID of whoever
* created it, and this call's own ID is only recorded on the
* `ingestion.enqueued` telemetry event, referencing the reused run.
* Run history is never rewritten.
*/
public function start(PluginVersion $pluginVersion, string $source = 'manual', bool $force = false, ?string $correlationId = null): IngestionRun
{
$run = DB::transaction(function () use ($pluginVersion, $source, $force): IngestionRun {
$requestCorrelationId = $correlationId ?? CorrelationId::generate();

$run = DB::transaction(function () use ($pluginVersion, $source, $force, $requestCorrelationId): IngestionRun {
if (! $force) {
$existingRun = IngestionRun::query()
->where('plugin_version_id', $pluginVersion->id)
Expand All @@ -46,6 +60,7 @@ public function start(PluginVersion $pluginVersion, string $source = 'manual', b
return IngestionRun::query()->create([
'plugin_version_id' => $pluginVersion->id,
'source' => $source,
'correlation_id' => $requestCorrelationId,
'status' => 'queued',
'stats' => [
'docs_parsed' => 0,
Expand All @@ -60,20 +75,41 @@ public function start(PluginVersion $pluginVersion, string $source = 'manual', b
IngestionRunStatusChanged::dispatch($run);
}

Telemetry::event('ingestion.enqueued', [
'correlation_id' => $requestCorrelationId,
'ingestion_run_id' => $run->id,
'plugin_version_id' => $pluginVersion->id,
'plugin_id' => $pluginVersion->plugin_id,
'community_id' => $pluginVersion->plugin?->community_id,
'source' => $source,
'reused' => ! $run->wasRecentlyCreated,
'user_id' => Auth::id(),
]);

return $run;
}

public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): IngestionRun
{
$pluginVersion->loadMissing('plugin');
$run ??= $this->start($pluginVersion);
$correlationId = $run->correlation_id ?? CorrelationId::generate();
$run->update([
'status' => 'running',
'started_at' => now(),
'finished_at' => null,
]);
IngestionRunStatusChanged::dispatch($run->refresh());

Telemetry::event('ingestion.started', [
'correlation_id' => $correlationId,
'ingestion_run_id' => $run->id,
'plugin_version_id' => $pluginVersion->id,
'plugin_id' => $pluginVersion->plugin_id,
'community_id' => $pluginVersion->plugin?->community_id,
'source' => $run->source,
]);

$stats = [
'docs_parsed' => 0,
'commands_extracted' => 0,
Expand All @@ -88,7 +124,8 @@ public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): In
'level' => 'error',
'code' => $exception->failureCode(),
'message' => $exception->getMessage(),
]]);
'correlation_id' => $correlationId,
]], $correlationId);
}

foreach ($files as $file) {
Expand All @@ -98,6 +135,7 @@ public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): In
'code' => 'document_not_modified',
'path' => $file->path,
'message' => "Skipped unchanged markdown file [{$file->path}].",
'correlation_id' => $correlationId,
];

continue;
Expand All @@ -117,6 +155,7 @@ public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): In
'code' => 'markdown_warning',
'path' => $parsedDocument->path,
'message' => $warning,
'correlation_id' => $correlationId,
];
}
} catch (Throwable $exception) {
Expand All @@ -126,6 +165,7 @@ public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): In
'code' => 'document_parse_failed',
'path' => $file->path,
'message' => $exception->getMessage(),
'correlation_id' => $correlationId,
];
}
}
Expand All @@ -139,6 +179,18 @@ public function run(PluginVersion $pluginVersion, ?IngestionRun $run = null): In
'finished_at' => now(),
]);

Telemetry::event('ingestion.completed', [
'correlation_id' => $correlationId,
'ingestion_run_id' => $run->id,
'plugin_version_id' => $pluginVersion->id,
'plugin_id' => $pluginVersion->plugin_id,
'community_id' => $pluginVersion->plugin?->community_id,
'status' => $status,
'docs_parsed' => $stats['docs_parsed'],
'commands_extracted' => $stats['commands_extracted'],
'warnings' => $stats['warnings'],
]);

if ($stats['docs_parsed'] > 0) {
$pluginVersion->commands()->with(['category', 'pluginVersion.plugin.community'])->get()->searchable();
$this->cache->invalidate();
Expand Down Expand Up @@ -226,7 +278,7 @@ private function categoryId(ParsedCommand $command): ?int
* @param array<string, int> $stats
* @param list<array<string, mixed>> $log
*/
private function fail(IngestionRun $run, array $stats, array $log): IngestionRun
private function fail(IngestionRun $run, array $stats, array $log, ?string $correlationId = null): IngestionRun
{
$run->update([
'status' => 'failed',
Expand All @@ -238,6 +290,21 @@ private function fail(IngestionRun $run, array $stats, array $log): IngestionRun
$run = $run->refresh();
IngestionRunStatusChanged::dispatch($run);

// The failure code is consumed from the `code` key already
// persisted in the log entry, per ADR-27/BRAIN-007: never
// recomputed, never renamed.
$run->loadMissing('pluginVersion.plugin');
$pluginVersion = $run->pluginVersion;

Telemetry::event('ingestion.failed', [
'correlation_id' => $correlationId ?? $run->correlation_id,
'ingestion_run_id' => $run->id,
'plugin_version_id' => $run->plugin_version_id,
'plugin_id' => $pluginVersion?->plugin_id,
'community_id' => $pluginVersion?->plugin?->community_id,
'code' => $log[0]['code'] ?? null,
]);

return $run;
}
}
62 changes: 62 additions & 0 deletions backend/app/Support/CorrelationId.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
<?php

namespace App\Support;

use Illuminate\Http\Request;
use Illuminate\Support\Str;

/**
* Generates and validates correlation IDs used to trace a single logical
* operation (HTTP request, webhook delivery, scheduled sync iteration)
* across ingestion runs, queued jobs, and structured log events.
*
* A client-supplied `X-Request-Id` is only trusted when it matches a safe,
* bounded character set. This prevents log injection (no control
* characters, no newlines) and unbounded input; anything else is replaced
* with a freshly generated ID rather than echoed back.
*/
final class CorrelationId
{
// The "D" modifier forces "$" to match only at the absolute end of the
// subject, not before a trailing newline. Without it, PCRE's default
// behavior lets a value like "id\n" pass validation and be echoed back
// raw in the response header and persisted on the run row.
private const PATTERN = '/^[A-Za-z0-9._:-]{1,128}$/D';

public const HEADER = 'X-Request-Id';

public const REQUEST_ATTRIBUTE = 'correlation_id';

/**
* Produce a new, safe-by-construction correlation ID.
*/
public static function generate(): string
{
return (string) Str::uuid();
}

/**
* Determine whether a client-supplied value is safe to reuse as-is.
*/
public static function isValid(?string $value): bool
{
return $value !== null && preg_match(self::PATTERN, $value) === 1;
}

/**
* Resolve the correlation ID assigned to the current request by
* `App\Http\Middleware\AssignCorrelationId`, generating one as a
* fallback for call sites reached outside that middleware (e.g. Artisan
* commands, unit tests exercising services directly).
*/
public static function fromRequest(Request $request): string
{
$assigned = $request->attributes->get(self::REQUEST_ATTRIBUTE);

if (is_string($assigned) && $assigned !== '') {
return $assigned;
}

return self::generate();
}
}
Loading
Loading