using System; using System.Net.Http; using System.Net.Http.Headers; using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using FrameworkBLL.GOP.Worker.AI; using GB5Shared.Auth.Jwt; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.GOP; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; namespace FrameworkBLL.GOP.Worker.NodeExecutors { // ============================================================ // AIExtractNodeExecutor — executes an AIExtract node type: converts an // unstructured or semi-structured source (PDF/text/scanned document/voice // transcript) into structured JSON by calling the external Enterprise AI // engine (AI-Enterprise-v1.0, a separate Python/FastAPI service — NOT // part of this repo, confirmed live and reachable from GB5DEMO at // http://217.217.249.121:8006 as of 2026-08-22). // // Output flows into the exact same Qualifier→Mapper→(ApiCall|IceMap) // chain every other GOP flow already uses — this node is purely a JSON // producer, no different from a Qualifier's Facts output. // // Classification gate: no real EAI governance module exists yet (that's // a separate, explicitly deferred phase), so this round carries the // "safe to send to an external AI" decision as an explicit, temporary // per-step config value on RequestSchemaJson: {"Classification": // "SafeForExternalAI","Slug":"extractor"}. Anything else — missing, // malformed, or a different Classification value — throws BEFORE any // outbound call is attempted. This is the ONLY safety gate in the whole // chain: the external engine itself has zero PII/data-classification // concept of its own (confirmed by reading its source). // // Auth: mints a short-lived JWT (tenantId/userId claims, HS256) via the // same shared IJwtAccessTokenIssuer/IJwtSigningKeyResolver every other // GB5 M2M integration uses (DXP/Partner/Entitlement) — the signing key // itself is a secret owned by the external engine's operator, not GB5, // and must be provisioned into Vault out-of-band (see AIExtractOptions). // // Route: always the deterministic per-agent route (POST /, e.g. // /extractor) — never the engine's guess-based /api/ai/execute, so GOP // is never at the mercy of that engine's content-shape heuristics. // ============================================================ public class AIExtractNodeExecutor : INodeExecutor { private readonly IHttpClientFactory _HttpClientFactory; private readonly IJwtAccessTokenIssuer _JwtIssuer; private readonly AIExtractOptions _Options; private readonly ILogger _Logger; public AIExtractNodeExecutor( IHttpClientFactory httpClientFactory, IJwtAccessTokenIssuer jwtIssuer, IOptions options, ILogger logger) { _HttpClientFactory = httpClientFactory; _JwtIssuer = jwtIssuer; _Options = options.Value; _Logger = logger; } public string NodeType => "AIExtract"; public async Task ExecuteAsync( GopFlowSnapshotStepDTO step, GopExecutionHeaderDTO header, string? inputPayloadJson, LoginDTO loginDTO, CancellationToken ct) { using var activity = GB5ActivitySources.AICapability.StartActivity("gop.ai_extract"); activity?.SetTag("gop.execution_id", header.ExecutionId); activity?.SetTag("gop.node_code", step.NodeCode); var (classification, slug) = ParseStepConfig(step.RequestSchemaJson); activity?.SetTag("ai.classification", classification ?? "(none)"); activity?.SetTag("ai.slug", slug ?? "(none)"); if (!string.Equals(classification, "SafeForExternalAI", StringComparison.Ordinal)) { activity?.SetStatus(System.Diagnostics.ActivityStatusCode.Error, "classification gate refused"); throw new InvalidOperationException( $"AIExtract step '{step.NodeCode}' refused: RequestSchemaJson must set " + "{\"Classification\":\"SafeForExternalAI\"} before any payload leaves GB5's " + "process boundary. No outbound call was made."); } if (string.IsNullOrWhiteSpace(slug)) { throw new InvalidOperationException( $"AIExtract step '{step.NodeCode}' has no target Slug configured in RequestSchemaJson " + "(e.g. \"extractor\", \"document_identifier\", \"parse_validate\")."); } var claims = new AIExtractAccessTokenClaims { TenantId = loginDTO.ClientId, UserId = loginDTO.UserId }; var issued = await _JwtIssuer.IssueAsync( claims, _Options.SigningKeyVaultPath, _Options.Issuer, _Options.Audience, _Options.AccessTokenMinutes, ct).ConfigureAwait(false); var client = _HttpClientFactory.CreateClient("ai-enterprise"); using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(ct); timeoutCts.CancelAfter(TimeSpan.FromSeconds(step.TimeoutSeconds > 0 ? step.TimeoutSeconds : 120)); var requestUrl = $"{_Options.EngineBaseUrl.TrimEnd('/')}/{slug.TrimStart('/')}"; using var request = new HttpRequestMessage(HttpMethod.Post, requestUrl); request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", issued.AccessToken); using var form = new MultipartFormDataContent(); // The upstream payload IS the "unstructured content" this node converts — carried as // the "file" multipart part every one of the engine's extraction routes requires. // GOP payloads are JSON strings end-to-end; no binary-content-carrying step config // exists yet, so this treats the raw payload text as the file body. A real PDF/scanned // -document source will need a dedicated upstream binding to carry actual file bytes — // out of scope for this first version, flagged here rather than silently assumed away. var fileBytes = Encoding.UTF8.GetBytes(inputPayloadJson ?? string.Empty); var fileContent = new ByteArrayContent(fileBytes); fileContent.Headers.ContentType = new MediaTypeHeaderValue("application/octet-stream"); form.Add(fileContent, "file", "payload.txt"); form.Add(new StringContent(header.ExecutionId.ToString()), "session_id"); // Per-flow prompt granularity (see MEAIPROMPTTEMPLATE.VariantCode / the approved // plan's "Prompt variants" section) — gbEAI's prompt_resolver falls back to the // capability-wide default row when no variant-specific prompt exists, so this is // additive-only and safe even for flows with no per-flow prompt configured yet. if (!string.IsNullOrWhiteSpace(header.FlowCode)) form.Add(new StringContent(header.FlowCode), "variant_code"); request.Content = form; _Logger.LogDebug( "AIExtractNodeExecutor: execution {ExecutionId} → POST {Url}", header.ExecutionId, requestUrl); using var response = await client.SendAsync(request, timeoutCts.Token).ConfigureAwait(false); var body = await response.Content.ReadAsStringAsync(ct).ConfigureAwait(false); activity?.SetTag("ai.http_status", (int)response.StatusCode); if (!response.IsSuccessStatusCode) { activity?.SetStatus(System.Diagnostics.ActivityStatusCode.Error, $"HTTP {(int)response.StatusCode}"); throw new HttpRequestException( $"AIExtract step '{step.NodeCode}' returned HTTP {(int)response.StatusCode}: {body}", null, response.StatusCode); } // Engine's own envelope: {"executionId": , "result": {...}}. Unwrap to "result" // so downstream Qualifier/Mapper steps see just the extracted data, not the wrapper. string? extractedJson = body; try { using var doc = JsonDocument.Parse(body); if (doc.RootElement.ValueKind == JsonValueKind.Object && doc.RootElement.TryGetProperty("result", out var resultEl)) { extractedJson = resultEl.GetRawText(); if (doc.RootElement.TryGetProperty("executionId", out var execIdEl)) activity?.SetTag("ai.engine_execution_id", execIdEl.ToString()); } } catch (JsonException) { // Not the expected envelope shape — pass the raw body through unchanged. } _Logger.LogInformation( "AIExtractNodeExecutor: execution {ExecutionId} step {NodeCode} — HTTP {Status}, slug={Slug}", header.ExecutionId, step.NodeCode, (int)response.StatusCode, slug); return extractedJson; } // Reads the temporary per-step config carried on RequestSchemaJson — see the class-level // doc comment for why this lives here instead of a real EAI governance table. private static (string? Classification, string? Slug) ParseStepConfig(string? requestSchemaJson) { if (string.IsNullOrWhiteSpace(requestSchemaJson)) return (null, null); try { using var doc = JsonDocument.Parse(requestSchemaJson); if (doc.RootElement.ValueKind != JsonValueKind.Object) return (null, null); string? classification = doc.RootElement.TryGetProperty("Classification", out var c) && c.ValueKind == JsonValueKind.String ? c.GetString() : null; string? slug = doc.RootElement.TryGetProperty("Slug", out var s) && s.ValueKind == JsonValueKind.String ? s.GetString() : null; return (classification, slug); } catch (JsonException) { return (null, null); } } } }