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 };
}
}
}
}