using GB5Shared.DTO.Framework.Login; using GB5Shared.ExternalAI; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace KmsBLL.Sync { /// /// KMS's concrete IKmsAISyncClient — resolves its own AIEngineOptions from the /// "KMS:AIExtract" config section directly (see GB5Shared.ExternalAI.AIEngineOptions's doc /// comment for why this isn't bound via generic IOptions<AIEngineOptions> DI). Targets /// gbEAI's "kms_discovery" agent route (POST /kms_discovery), matching /// AIExtractNodeExecutor.cs's proven multipart-form + JWT pattern. /// public class KmsAISyncClient : IKmsAISyncClient { private const string SectionName = "KMS:AIExtract"; private const string RoutePath = "/kms_discovery"; private readonly IAIEngineSyncClient _transport; private readonly AIEngineOptions _options; private readonly ILogger _logger; public KmsAISyncClient( IAIEngineSyncClient transport, IConfiguration configuration, ILogger logger) { _transport = transport ?? throw new ArgumentNullException(nameof(transport)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); _options = configuration.GetSection(SectionName).Get() ?? new AIEngineOptions(); if (string.IsNullOrWhiteSpace(_options.SigningKeyVaultPath)) _options.SigningKeyVaultPath = "kms/ai-enterprise-jwt-secret"; if (string.IsNullOrWhiteSpace(_options.Issuer)) _options.Issuer = "GB5-KMS"; // gbEAI's extraction pipeline is fully synchronous (no Celery) and can legitimately // take minutes for a real document — a dedicated, longer default than the shared // AIEngineOptions.TimeoutSeconds default (120s), overridable per-config too. if (_options.TimeoutSeconds <= 120) _options.TimeoutSeconds = 600; } public async Task TriggerExtractionAsync( int jobId, int domainId, string? sourceText, LoginDTO login, CancellationToken ct = default) { var input = new Dictionary { ["job_id"] = jobId, ["domain_id"] = domainId, }; if (!string.IsNullOrWhiteSpace(sourceText)) input["source_text"] = sourceText; var formFields = new Dictionary { ["input"] = JsonSerializer.Serialize(input), ["capability"] = "KMS_EXTRACT", ["session_id"] = jobId.ToString(), }; var response = await _transport.SendMultipartAsync( _options, RoutePath, formFields, file: null, login, ct).ConfigureAwait(false); if (!response.IsSuccess) { _logger.LogWarning( "KmsAISyncClient | TriggerExtraction | Failed | JobId={JobId} | {Error}", jobId, response.ErrorMessage); return new KmsSyncResult { Success = false, Message = response.ErrorMessage }; } return ParseEnvelope(response.Body); } public async Task EmbedArtifactAsync( int artifactId, int domainId, byte artifactType, string content, LoginDTO login, CancellationToken ct = default) { var input = new Dictionary { ["artifact_id"] = artifactId, ["domain_id"] = domainId, ["artifact_type"] = artifactType, ["content"] = content, }; var formFields = new Dictionary { ["input"] = JsonSerializer.Serialize(input), ["capability"] = "KMS_EMBED_ARTIFACT", }; var response = await _transport.SendMultipartAsync( _options, RoutePath, formFields, file: null, login, ct).ConfigureAwait(false); if (!response.IsSuccess) { _logger.LogWarning( "KmsAISyncClient | EmbedArtifact | Failed | ArtifactId={ArtifactId} | {Error}", artifactId, response.ErrorMessage); return new KmsSyncResult { Success = false, Message = response.ErrorMessage }; } return ParseEnvelope(response.Body); } public async Task RetireArtifactAsync( int artifactId, string? qdrantPointId, LoginDTO login, CancellationToken ct = default) { var input = new Dictionary { ["artifact_id"] = artifactId, }; if (!string.IsNullOrWhiteSpace(qdrantPointId)) input["qdrant_point_id"] = qdrantPointId; var formFields = new Dictionary { ["input"] = JsonSerializer.Serialize(input), ["capability"] = "KMS_RETIRE_ARTIFACT", }; var response = await _transport.SendMultipartAsync( _options, RoutePath, formFields, file: null, login, ct).ConfigureAwait(false); if (!response.IsSuccess) { _logger.LogWarning( "KmsAISyncClient | RetireArtifact | Failed | ArtifactId={ArtifactId} | {Error}", artifactId, response.ErrorMessage); return new KmsSyncResult { Success = false, Message = response.ErrorMessage }; } return ParseEnvelope(response.Body); } // gbEAI's engine envelope is {"executionId": , "result": {...}}; _embed_artifact/ // _retire_artifact's own result shape is {"status": "...", "point_id": "...", "artifact_id": ...}. private static KmsSyncResult ParseEnvelope(string body) { try { using var doc = JsonDocument.Parse(body); var root = doc.RootElement; var result = root.ValueKind == JsonValueKind.Object && root.TryGetProperty("result", out var resultEl) && resultEl.ValueKind == JsonValueKind.Object ? resultEl : root; string? status = result.TryGetProperty("status", out var s) && s.ValueKind == JsonValueKind.String ? s.GetString() : null; string? pointId = result.TryGetProperty("point_id", out var p) && p.ValueKind == JsonValueKind.String ? p.GetString() : null; // "index_failed"/"remove_failed" are HTTP 200s from gbEAI's own perspective (the // call succeeded, the downstream Qdrant op didn't) — surface as Success=false so // the caller applies its own partial-failure handling. bool ok = status is null || (status != "index_failed" && status != "remove_failed"); return new KmsSyncResult { Success = ok, Status = status, QdrantPointId = pointId }; } catch (JsonException) { return new KmsSyncResult { Success = true }; } } } }