/
rezvich
/
SoapBatchProcessor
Обзор
Документация
Войти
/
rezvich
/
SoapBatchProcessor
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
Program.cs
339 строк
12 KB
rezvich
Подправил токен
04 авг 2026, 12:20
04 авг 2026, 12:20
632636d
Код
Авторство
О чём код?
using Polly; using Polly.Contrib.WaitAndRetry; using SoapBatchProcessor.Data; using SoapBatchProcessor.Infrastructure; using SoapBatchProcessor.Models.SOAP; using SoapBatchProcessor.Parsing; using System.Diagnostics; using System.Net.Http.Headers; using System.Text; using System.Text.Json; using System.Threading.Channels; var configJson = await File.ReadAllTextAsync("appsettings.json"); var root = JsonSerializer.Deserialize<JsonElement>(configJson); var settings = JsonSerializer.Deserialize<AppSettings>(root.GetProperty("AppSettings")); if (settings is null) throw new InvalidOperationException("AppSettings not found."); XmlNs.ApiVersion = settings.ApiVersion; Console.OutputEncoding = Encoding.UTF8; Console.WriteLine("SoapBatchProcessor starting..."); var serviceUrl = settings.SoapServiceUrl; if (string.IsNullOrWhiteSpace(serviceUrl)) throw new InvalidOperationException("Service URL is not configured."); // === ИНИЦИАЛИЗАЦИЯ Token Service === var tokenHttp = new HttpClient { BaseAddress = new Uri(settings.TokenServiceBaseUrl), Timeout = TimeSpan.FromSeconds(settings.TokenServiceTimeoutSeconds <= 0 ? 15 : settings.TokenServiceTimeoutSeconds) }; var tokenProvider = new SoapBatchProcessor.Infrastructure.RemoteAccessTokenProvider( tokenHttp, string.IsNullOrWhiteSpace(settings.TokenServiceApiKey) ? null : settings.TokenServiceApiKey ); try { Console.WriteLine("[IAM via TokenService] Pre-fetching access token..."); var initialLease = await tokenProvider.GetAccessTokenLeaseAsync(); Console.WriteLine( $"[IAM via TokenService] Token OK, length = {initialLease.AccessToken.Length}, version = {initialLease.TokenVersion}"); } catch (Exception ex) { Console.Error.WriteLine("[IAM via TokenService] Failed to get initial token:"); Console.Error.WriteLine(ex); throw; } Console.WriteLine($"Using SOAP endpoint: {serviceUrl}"); var http = HttpClientFactoryEx.Create(settings); var inputRepo = new InputRepository(settings.SqlConnectionString, settings.InputCommand, settings.CommandTimeoutSeconds); var bulk = new BulkWriter(settings.SqlConnectionString, settings.BulkCopyBatchSize, settings.BulkCopyTimeout); var mpiRepo = new MpiResponseRepository(settings.SqlConnectionString, settings.CommandTimeoutSeconds); var reqRepo = new MpiRequestsRepository(settings.SqlConnectionString, settings.CommandTimeoutSeconds); var requests = await inputRepo.LoadAsync(settings.BatchSize); Console.WriteLine($"Loaded {requests.Count} requests"); var sw = Stopwatch.StartNew(); long sent = 0, parsed = 0, failed = 0, saved = 0; int progressStep = Math.Max(100, requests.Count / 20); // лог каждые ~5% или минимум 100 var channel = Channel.CreateBounded<(InputRow row, string requestXml, string responseXml, string externalRequestId)>( new BoundedChannelOptions(settings.MaxDegreeOfParallelism * 2) { SingleReader = true, SingleWriter = false, FullMode = BoundedChannelFullMode.Wait }); var delays = Backoff.DecorrelatedJitterBackoffV2( medianFirstRetryDelay: TimeSpan.FromSeconds(2), retryCount: settings.MaxRetryCount, fastFirst: true); var timeoutPolicy = Policy .TimeoutAsync(TimeSpan.FromMinutes(2)); var retryPolicy = Policy .Handle<Exception>() .WaitAndRetryAsync( Backoff.DecorrelatedJitterBackoffV2( medianFirstRetryDelay: TimeSpan.FromSeconds(2), retryCount: settings.MaxRetryCount, fastFirst: true), (ex, ts, attempt, _) => { Console.WriteLine($"Retry {attempt}/{settings.MaxRetryCount} in {ts.TotalSeconds:F1}s: {ex.GetType().Name}: {ex.Message}"); }); var policy = Policy.WrapAsync(retryPolicy, timeoutPolicy); // ПИСАТЕЛЬ: параллельно шлём SOAP и отдаём сырой XML var writer = Task.Run(async () => { try { await Parallel.ForEachAsync( requests, new ParallelOptions { MaxDegreeOfParallelism = settings.MaxDegreeOfParallelism }, async (row, ct) => { var srch = new SearchRequest { ExternalRequestId = row.ExternalRequestId, Show = row.Show, Type = row.SearchType, ENP = row.ENP, OIP = row.OIP, DtFrom = row.Date1, DtTo = row.Date2, DudlSer = row.DudlSer, DudlNum = row.DudlNum, DudlType = row.DudlType, Snils = row.Snils, BirthDay = row.BirthDay, Surname = row.Surname, FirstName = row.FirstName, Patronymic = row.Patronymic }; var (envelope, externalRequestId) = SoapEnvelopeBuilder.BuildGetPersonDataHistory(srch); string responseXml = ""; await policy.ExecuteAsync(async (token) => { var lease = await tokenProvider.GetAccessTokenLeaseAsync(token); using var resp = await SendSoapRequestAsync( http, serviceUrl, settings.SoapActionHistory, envelope, lease.AccessToken, token); if (resp.StatusCode != System.Net.HttpStatusCode.Unauthorized) { resp.EnsureSuccessStatusCode(); responseXml = await resp.Content.ReadAsStringAsync(token); return; } await LogUnauthorizedAsync(resp, lease.TokenVersion, "initial", token); var refreshedLease = await tokenProvider.RefreshAfterRejectionAsync( lease.TokenVersion, token); using var refreshedResp = await SendSoapRequestAsync( http, serviceUrl, settings.SoapActionHistory, envelope, refreshedLease.AccessToken, token); if (refreshedResp.StatusCode == System.Net.HttpStatusCode.Unauthorized) { await LogUnauthorizedAsync( refreshedResp, refreshedLease.TokenVersion, "after refresh", token); } refreshedResp.EnsureSuccessStatusCode(); responseXml = await refreshedResp.Content.ReadAsStringAsync(token); }, ct); await channel.Writer.WriteAsync((row, envelope, responseXml, externalRequestId), ct); Interlocked.Increment(ref sent); }); } catch (Exception ex) { Console.Error.WriteLine($"[WRITER FAIL] {ex.GetType().Name}: {ex.Message}"); throw; } finally { // Иначе reader бесконечно ждёт данные, если Parallel.ForEachAsync завершился с ошибкой. channel.Writer.TryComplete(); } }); // ЧИТАТЕЛЬ: парсит → вставляет шапку → копит для bulk дочерних таблиц var reader = Task.Run(async () => { var acc = new Accumulator(); await foreach (var item in channel.Reader.ReadAllAsync()) { SoapResponseDto dto; try { dto = HistoryResponseParser.Parse(item.externalRequestId, item.responseXml); //dto.Success = true; } catch (Exception ex) { dto = new SoapResponseDto { Success = false, ErrorMessage = ex.Message }; } if (!dto.Success) Interlocked.Increment(ref failed); Interlocked.Increment(ref parsed); // Шапочные поля dto.ExternalRequestId = item.externalRequestId; dto.ENP = item.row.ENP; dto.DtFrom = item.row.Date1; dto.DtTo = item.row.Date2; dto.RawXml = item.responseXml; dto.RequestXml = item.requestXml; dto.RequestId = item.row.RequestId; dto.ResponseId = await mpiRepo.InsertAsync(dto); Interlocked.Increment(ref saved); // await reqRepo.MarkProcessedAsync(item.row.RequestId, dto.Success, dto.ErrorMessage); acc.Add(dto); if (acc.Count >= settings.BulkCopyBatchSize) { try { var failedIds = await bulk.FlushAsync(acc); foreach (var r in acc.Responses) { if (failedIds.Contains(r.RequestId)) await reqRepo.MarkProcessedAsync(r.RequestId, false, "BulkCopy row fail"); else await reqRepo.MarkProcessedAsync(r.RequestId, r.Success, r.ErrorMessage); } } catch (Exception ex) // сюда теперь не должны доходить, но на всякий случай { Console.Error.WriteLine($"[BULK FAIL] {ex.Message}"); foreach (var r in acc.Responses) await reqRepo.MarkProcessedAsync(r.RequestId, false, "BulkCopy batch fail: " + ex.Message); } finally { acc.Clear(); } } // ПРОГРЕСС ЛОГ if (parsed % progressStep == 0) { var rps = parsed / Math.Max(0.001, sw.Elapsed.TotalSeconds); Console.WriteLine($"[PROGRESS] Parsed {parsed}/{requests.Count}, Sent {sent}, Saved {saved}, Failed {failed}, Elapsed {sw.Elapsed:c}, RPS {rps:F1}"); } } if (acc.Count > 0) { try { var failedIds = await bulk.FlushAsync(acc); foreach (var r in acc.Responses) { if (failedIds.Contains(r.RequestId)) await reqRepo.MarkProcessedAsync(r.RequestId, false, "BulkCopy row fail"); else await reqRepo.MarkProcessedAsync(r.RequestId, r.Success, r.ErrorMessage); } } catch (Exception ex) { Console.Error.WriteLine($"[BULK FAIL] (tail) {ex.Message}"); foreach (var r in acc.Responses) await reqRepo.MarkProcessedAsync(r.RequestId, false, "BulkCopy batch fail: " + ex.Message); } finally { acc.Clear(); } } }); await Task.WhenAll(writer, reader); sw.Stop(); var rpsFinal = parsed / Math.Max(0.001, sw.Elapsed.TotalSeconds); Console.WriteLine("===================================="); Console.WriteLine($"TOTAL requests: {requests.Count}"); Console.WriteLine($"SENT ok: {sent}"); Console.WriteLine($"PARSED ok: {parsed - failed}"); Console.WriteLine($"PARSED failed: {failed}"); Console.WriteLine($"RESPONSES saved: {saved}"); Console.WriteLine($"ELAPSED: {sw.Elapsed:c}"); Console.WriteLine($"THROUGHPUT (RPS): {rpsFinal:F1}"); Console.WriteLine("===================================="); Console.WriteLine("Done."); Console.ReadLine(); static async Task<HttpResponseMessage> SendSoapRequestAsync( HttpClient http, string serviceUrl, string? soapAction, string envelope, string accessToken, CancellationToken ct) { using var req = new HttpRequestMessage(HttpMethod.Post, serviceUrl); req.Headers.Authorization = new AuthenticationHeaderValue("Bearer", accessToken); if (!string.IsNullOrWhiteSpace(soapAction)) req.Headers.TryAddWithoutValidation("SOAPAction", soapAction); var content = new StringContent(envelope, Encoding.UTF8, "text/xml"); content.Headers.ContentType!.CharSet = "utf-8"; req.Content = content; return await http.SendAsync(req, ct); } static async Task LogUnauthorizedAsync( HttpResponseMessage response, long tokenVersion, string stage, CancellationToken ct) { var challenge = response.Headers.WwwAuthenticate.Count > 0 ? string.Join(", ", response.Headers.WwwAuthenticate) : "<none>"; var body = await response.Content.ReadAsStringAsync(ct); var bodyPreview = body.Length <= 500 ? body : body[..500] + "..."; bodyPreview = bodyPreview.ReplaceLineEndings(" "); Console.Error.WriteLine( $"[SOAP 401] stage={stage}, tokenVersion={tokenVersion}, WWW-Authenticate={challenge}, body={bodyPreview}"); }