Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
63242775ee | ||
|
|
e81f1ebf30 | ||
|
|
babfbf8006 |
@@ -1,19 +0,0 @@
|
||||
# Dev-edge «Дейла» (compose.dev.yml, сервис frontend): SPA + /api на core.
|
||||
# Отличие от prod-Caddyfile: HTTP без TLS и без плейсхолдер-домена — для локального просмотра UI.
|
||||
# Статика — собранный SPA (Vite) в /srv, неизвестные пути отдают index.html (история браузера).
|
||||
|
||||
:80 {
|
||||
# API core: /api/* уходит на core:5080 без перезаписи (контракт /api неизменен).
|
||||
# SSE (/api/events), файлы и QR-SVG проходят reverse_proxy потоково.
|
||||
handle /api/* {
|
||||
reverse_proxy core:5080
|
||||
}
|
||||
|
||||
handle {
|
||||
# SPA/ассеты в dev не кэшируем: пересборка фронта должна подхватываться по F5.
|
||||
header Cache-Control "no-cache"
|
||||
root * /srv
|
||||
try_files {path} /index.html
|
||||
file_server
|
||||
}
|
||||
}
|
||||
@@ -254,18 +254,6 @@ services:
|
||||
timeout: 3s
|
||||
retries: 10
|
||||
|
||||
# Фронтенд (SPA) — сборка образа (Vite) и отдача через Caddy; /api → core:5080. UI — http://localhost:8080.
|
||||
frontend:
|
||||
build:
|
||||
context: ..
|
||||
dockerfile: src/frontend/Dockerfile
|
||||
container_name: deal-frontend
|
||||
ports:
|
||||
- "8080:80"
|
||||
depends_on:
|
||||
core:
|
||||
condition: service_healthy
|
||||
|
||||
# Prometheus (профиль observability, этап 12/пакет A) — сбор /metrics всех 4 процессов (:9464)
|
||||
# внутри dev-сети. Подъём: docker compose -f deploy/compose.dev.yml --profile observability up -d.
|
||||
# Конфиг — общий deploy/observability/prometheus.yml (те же имена сервисов и таргеты). UI — 9090.
|
||||
@@ -344,56 +332,6 @@ services:
|
||||
- /:/host/root:ro
|
||||
pid: host
|
||||
|
||||
# Loki — хранилище логов (профиль observability), UI/API — :3100.
|
||||
loki:
|
||||
image: grafana/loki:3.4.2
|
||||
container_name: deal-loki
|
||||
profiles: ["observability"]
|
||||
command: -config.file=/etc/loki/loki.yml
|
||||
ports:
|
||||
- "3100:3100"
|
||||
volumes:
|
||||
- ./observability/loki.yml:/etc/loki/loki.yml:ro
|
||||
- deal_loki_data:/loki
|
||||
|
||||
# Promtail — сбор docker-логов deal-процессов в Loki (docker.sock, профиль observability).
|
||||
promtail:
|
||||
image: grafana/promtail:3.4.2
|
||||
container_name: deal-promtail
|
||||
profiles: ["observability"]
|
||||
command: -config.file=/etc/promtail/promtail.yml
|
||||
volumes:
|
||||
- ./observability/promtail.yml:/etc/promtail/promtail.yml:ro
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
- deal_promtail_data:/var/lib/promtail
|
||||
depends_on:
|
||||
loki:
|
||||
condition: service_started
|
||||
|
||||
# Grafana — UI логов/метрик/трейсов (профиль observability), локальный вход admin/admin.
|
||||
grafana:
|
||||
image: grafana/grafana:11.5.2
|
||||
container_name: deal-grafana
|
||||
profiles: ["observability"]
|
||||
environment:
|
||||
GF_SECURITY_ADMIN_USER: admin
|
||||
GF_SECURITY_ADMIN_PASSWORD: admin
|
||||
GF_USERS_ALLOW_SIGN_UP: "false"
|
||||
GF_AUTH_ANONYMOUS_ENABLED: "false"
|
||||
ports:
|
||||
- "3001:3000"
|
||||
volumes:
|
||||
- ./observability/grafana/provisioning:/etc/grafana/provisioning:ro
|
||||
- ./observability/grafana/dashboards:/var/lib/grafana/dashboards:ro
|
||||
- deal_grafana_data:/var/lib/grafana
|
||||
depends_on:
|
||||
loki:
|
||||
condition: service_started
|
||||
prometheus:
|
||||
condition: service_started
|
||||
tempo:
|
||||
condition: service_started
|
||||
|
||||
volumes:
|
||||
deal_pgdata:
|
||||
deal_minio_data:
|
||||
@@ -402,6 +340,3 @@ volumes:
|
||||
deal_api_data:
|
||||
deal_prometheus_data:
|
||||
deal_tempo_data:
|
||||
deal_loki_data:
|
||||
deal_promtail_data:
|
||||
deal_grafana_data:
|
||||
|
||||
@@ -22,7 +22,6 @@ public sealed class TgStatusService(
|
||||
TelegramKeysService keys)
|
||||
{
|
||||
private const string IdlePhase = "idle";
|
||||
private const string ReadyPhase = "ready";
|
||||
|
||||
// Опции JSON KV-значений статуса: camelCase (как пишет ингресс) + терпимость регистра.
|
||||
private static readonly JsonSerializerOptions KvJsonOptions = new()
|
||||
@@ -38,19 +37,13 @@ public sealed class TgStatusService(
|
||||
public async Task<TgStatusDto> GetAsync(CancellationToken ct)
|
||||
{
|
||||
TelegramAccountStatusDto live = await ReadLiveAsync(ct).ConfigureAwait(false);
|
||||
// «Подключён» для UI = авторизован (phase ready). Транспортный connected сервиса
|
||||
// означает лишь живость соединения и не гарантирует вход — в UI он даёт «зависание».
|
||||
bool authorized = string.Equals(live.Phase, ReadyPhase, StringComparison.Ordinal);
|
||||
// Живой account (имя из Telegram) приоритетнее KV: при QR-входе KV ещё не заполнен.
|
||||
string account = string.IsNullOrEmpty(live.Account)
|
||||
? await ReadAccountAsync(ct).ConfigureAwait(false)
|
||||
: live.Account;
|
||||
string account = await ReadAccountAsync(ct).ConfigureAwait(false);
|
||||
int monitored = (await dialogs.ListMonitoredIdsAsync(ct).ConfigureAwait(false)).Count;
|
||||
TgKeysSnapshot snapshot = await keys.GetAsync(ct).ConfigureAwait(false);
|
||||
|
||||
return new TgStatusDto(
|
||||
Phase: live.Phase,
|
||||
Connected: authorized,
|
||||
Connected: live.Connected,
|
||||
Listener: live.Listener,
|
||||
Account: account,
|
||||
Monitored: monitored,
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
using Deal.SharedKernel.Resilience;
|
||||
using Grpc.Core;
|
||||
|
||||
namespace Deal.Infrastructure.Integrations.Resilience;
|
||||
|
||||
/// <summary>
|
||||
/// Повтор транзиентных gRPC-сбоев клиентов автономных сервисов.
|
||||
/// </summary>
|
||||
public static class GrpcRetry
|
||||
{
|
||||
/// <summary>
|
||||
/// Число повторов после первой попытки.
|
||||
/// </summary>
|
||||
public const int RetryCount = 2;
|
||||
|
||||
/// <summary>
|
||||
/// Базовая задержка повтора (далее — экспоненциально с джиттером).
|
||||
/// </summary>
|
||||
public static readonly TimeSpan BaseDelay = TimeSpan.FromMilliseconds(200);
|
||||
|
||||
/// <summary>
|
||||
/// Выполняет gRPC-вызов с повтором транзиентных сбоев.
|
||||
/// </summary>
|
||||
/// <param name="operation">Вызов (принимает токен отмены).</param>
|
||||
/// <param name="cancellationToken">Токен отмены.</param>
|
||||
/// <returns>Ответ вызова.</returns>
|
||||
public static Task<TResult> ExecuteAsync<TResult>(
|
||||
Func<CancellationToken, Task<TResult>> operation,
|
||||
CancellationToken cancellationToken)
|
||||
=> ExecuteAsync(operation, DefaultDelayAsync, cancellationToken);
|
||||
|
||||
/// <summary>
|
||||
/// Выполняет gRPC-вызов с повтором и заданной паузой между попытками.
|
||||
/// </summary>
|
||||
/// <param name="operation">Вызов (принимает токен отмены).</param>
|
||||
/// <param name="delayAsync">Пауза между попытками (в тестах — мгновенная).</param>
|
||||
/// <param name="cancellationToken">Токен отмены.</param>
|
||||
/// <returns>Ответ вызова.</returns>
|
||||
public static Task<TResult> ExecuteAsync<TResult>(
|
||||
Func<CancellationToken, Task<TResult>> operation,
|
||||
Func<TimeSpan, CancellationToken, Task> delayAsync,
|
||||
CancellationToken cancellationToken)
|
||||
=> RetryExecutor.ExecuteAsync(
|
||||
operation,
|
||||
RetryCount,
|
||||
BaseDelay,
|
||||
IsTransient,
|
||||
delayAsync,
|
||||
cancellationToken);
|
||||
|
||||
/// <summary>
|
||||
/// Признак транзиентного сбоя транспорта (недоступность/дедлайн).
|
||||
/// </summary>
|
||||
/// <param name="exception">Исключение вызова.</param>
|
||||
/// <returns>True — сбой имеет смысл повторить.</returns>
|
||||
public static bool IsTransient(Exception exception)
|
||||
=> exception is RpcException rpc
|
||||
&& rpc.StatusCode is StatusCode.Unavailable or StatusCode.DeadlineExceeded;
|
||||
|
||||
// Экспоненциальная задержка с джиттером 0.5–1.5× (сглаживает синхронные ретраи воркеров).
|
||||
private static Task DefaultDelayAsync(TimeSpan delay, CancellationToken cancellationToken)
|
||||
{
|
||||
double factor = 0.5 + Random.Shared.NextDouble();
|
||||
return Task.Delay(TimeSpan.FromMilliseconds(delay.TotalMilliseconds * factor), cancellationToken);
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,7 @@ using Deal.Contracts.Integrations.Models;
|
||||
using Deal.Grpc.Ai;
|
||||
using Deal.Infrastructure.Integrations.Exceptions;
|
||||
using Deal.Infrastructure.Integrations.Models;
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Deal.Modules.Pipeline.Application.Services;
|
||||
using Deal.SharedKernel.Tenants.Abstractions;
|
||||
using Deal.SharedKernel.Tenants.Models;
|
||||
@@ -75,14 +76,16 @@ public sealed class GrpcAiClassifier : IAiClassifier
|
||||
string prompt = await _contextBuilder.BuildFilterPromptAsync(ct);
|
||||
ProviderConfig providerConfig = await _providerConfigBuilder.BuildAsync(ct);
|
||||
AiService.AiServiceClient client = _connection.CreateClient();
|
||||
FilterReply reply = await client.FilterAsync(
|
||||
new FilterRequest
|
||||
{
|
||||
Prompt = prompt,
|
||||
Text = SliceCodePoints(text, MaxFilterTextCodePoints), // python L193: text[:4000]
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, ct));
|
||||
FilterReply reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.FilterAsync(
|
||||
new FilterRequest
|
||||
{
|
||||
Prompt = prompt,
|
||||
Text = SliceCodePoints(text, MaxFilterTextCodePoints), // python L193: text[:4000]
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, token)).ResponseAsync,
|
||||
ct);
|
||||
await _usageRecorder.AddAsync(reply.Usage, providerConfig.ProviderId, providerConfig.Model, ct);
|
||||
|
||||
return new AiFilterResultDto(
|
||||
@@ -113,14 +116,16 @@ public sealed class GrpcAiClassifier : IAiClassifier
|
||||
ClassifyReply reply;
|
||||
try
|
||||
{
|
||||
reply = await client.ClassifyAsync(
|
||||
new ClassifyRequest
|
||||
{
|
||||
SystemPrompt = systemPrompt,
|
||||
UserContext = userContext,
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, ct));
|
||||
reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.ClassifyAsync(
|
||||
new ClassifyRequest
|
||||
{
|
||||
SystemPrompt = systemPrompt,
|
||||
UserContext = userContext,
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, token)).ResponseAsync,
|
||||
ct);
|
||||
}
|
||||
catch (RpcException exception)
|
||||
{
|
||||
|
||||
@@ -3,6 +3,7 @@ using Deal.Contracts.Integrations.Models;
|
||||
using Deal.Grpc.Ai;
|
||||
using Deal.Infrastructure.Integrations.Exceptions;
|
||||
using Deal.Infrastructure.Integrations.Models;
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Deal.SharedKernel.Tenants.Abstractions;
|
||||
using Deal.SharedKernel.Tenants.Models;
|
||||
using Grpc.Core;
|
||||
@@ -74,13 +75,15 @@ public sealed class GrpcAiTools : IAiTools
|
||||
{
|
||||
AiService.AiServiceClient client = _connection.CreateClient();
|
||||
ProviderConfig providerConfig = await _providerConfigBuilder.BuildAsync(ct);
|
||||
GenerateKeywordsReply reply = await client.GenerateKeywordsAsync(
|
||||
new GenerateKeywordsRequest
|
||||
{
|
||||
Description = SliceCodePoints(description ?? string.Empty, MaxDescriptionCodePoints),
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, ct));
|
||||
GenerateKeywordsReply reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.GenerateKeywordsAsync(
|
||||
new GenerateKeywordsRequest
|
||||
{
|
||||
Description = SliceCodePoints(description ?? string.Empty, MaxDescriptionCodePoints),
|
||||
ProviderConfig = providerConfig,
|
||||
},
|
||||
CallOptions(tenantId.Value, token)).ResponseAsync,
|
||||
ct);
|
||||
await _usageRecorder.AddAsync(reply.Usage, providerConfig.ProviderId, providerConfig.Model, ct);
|
||||
return new AiGenerateKeywordsResultDto(
|
||||
Ok: true,
|
||||
@@ -125,7 +128,9 @@ public sealed class GrpcAiTools : IAiTools
|
||||
}
|
||||
}
|
||||
|
||||
EvaluateFitReply reply = await client.EvaluateFitAsync(request, CallOptions(tenantId.Value, ct));
|
||||
EvaluateFitReply reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.EvaluateFitAsync(request, CallOptions(tenantId.Value, token)).ResponseAsync,
|
||||
ct);
|
||||
await _usageRecorder.AddAsync(reply.Usage, providerConfig.ProviderId, providerConfig.Model, ct);
|
||||
return new AiEvaluateFitResultDto(
|
||||
Fit: reply.Fit,
|
||||
|
||||
@@ -4,6 +4,7 @@ using Deal.Contracts.Integrations.Models;
|
||||
using Deal.Grpc.Ml;
|
||||
using Deal.Infrastructure.Integrations.Abstractions;
|
||||
using Deal.Infrastructure.Integrations.Models;
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Deal.Modules.Kanban.Application.Abstractions;
|
||||
using Deal.Modules.Kanban.Application.Models;
|
||||
using Deal.Modules.Settings.Application.Abstractions;
|
||||
@@ -124,9 +125,11 @@ public sealed class GrpcMlClient : IMlClient, IMlTrainClient
|
||||
try
|
||||
{
|
||||
MlService.MlServiceClient client = _connection.CreateClient();
|
||||
PredictReply reply = await client.PredictAsync(
|
||||
new PredictRequest { Text = text ?? string.Empty },
|
||||
CallOptions(tenantId.Value, TimeSpan.FromSeconds(PredictDeadlineSeconds), ct));
|
||||
PredictReply reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.PredictAsync(
|
||||
new PredictRequest { Text = text ?? string.Empty },
|
||||
CallOptions(tenantId.Value, TimeSpan.FromSeconds(PredictDeadlineSeconds), token)).ResponseAsync,
|
||||
ct);
|
||||
|
||||
await _usageRecorder.AddEstimatedAsync(text, TokenUsageSources.Local, TokenUsageSources.Ml, ct);
|
||||
return MapPredict(reply);
|
||||
@@ -220,9 +223,11 @@ public sealed class GrpcMlClient : IMlClient, IMlTrainClient
|
||||
try
|
||||
{
|
||||
MlService.MlServiceClient client = _connection.CreateClient();
|
||||
StatusReply reply = await client.StatusAsync(
|
||||
new StatusRequest(),
|
||||
CallOptions(tenantId.Value, TimeSpan.FromSeconds(StatusDeadlineSeconds), ct));
|
||||
StatusReply reply = await GrpcRetry.ExecuteAsync(
|
||||
token => client.StatusAsync(
|
||||
new StatusRequest(),
|
||||
CallOptions(tenantId.Value, TimeSpan.FromSeconds(StatusDeadlineSeconds), token)).ResponseAsync,
|
||||
ct);
|
||||
MlServiceStatusDto service = MapStatus(reply);
|
||||
_statusCache.Set(tenantId.Value, service, reachable: true);
|
||||
return _statusCache.TryGet(tenantId.Value, out MlStatusCache.Snapshot updated)
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
namespace Deal.SharedKernel.Resilience;
|
||||
|
||||
/// <summary>
|
||||
/// Повтор операции при транзиентном сбое.
|
||||
/// </summary>
|
||||
public static class RetryExecutor
|
||||
{
|
||||
/// <summary>
|
||||
/// Выполняет операцию, повторяя её при транзиентном сбое с задержкой.
|
||||
/// </summary>
|
||||
/// <param name="operation">Операция (принимает токен отмены).</param>
|
||||
/// <param name="retryCount">Число повторов после первой попытки.</param>
|
||||
/// <param name="baseDelay">Базовая задержка; для повтора N — baseDelay * 2^N.</param>
|
||||
/// <param name="shouldRetry">Предикат транзиентности сбоя.</param>
|
||||
/// <param name="delayAsync">Пауза между попытками (в тестах — мгновенная).</param>
|
||||
/// <param name="cancellationToken">Токен отмены.</param>
|
||||
/// <returns>Результат первой успешной попытки.</returns>
|
||||
public static async Task<TResult> ExecuteAsync<TResult>(
|
||||
Func<CancellationToken, Task<TResult>> operation,
|
||||
int retryCount,
|
||||
TimeSpan baseDelay,
|
||||
Func<Exception, bool> shouldRetry,
|
||||
Func<TimeSpan, CancellationToken, Task> delayAsync,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(operation);
|
||||
ArgumentOutOfRangeException.ThrowIfNegative(retryCount);
|
||||
ArgumentNullException.ThrowIfNull(shouldRetry);
|
||||
ArgumentNullException.ThrowIfNull(delayAsync);
|
||||
|
||||
for (int attempt = 0; ; attempt++)
|
||||
{
|
||||
try
|
||||
{
|
||||
return await operation(cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
catch (Exception exception) when (attempt < retryCount
|
||||
&& shouldRetry(exception)
|
||||
&& !cancellationToken.IsCancellationRequested)
|
||||
{
|
||||
await delayAsync(BackoffDelay(baseDelay, attempt), cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Задержка повтора с экспоненциальным ростом от базовой.
|
||||
private static TimeSpan BackoffDelay(TimeSpan baseDelay, int attempt)
|
||||
=> TimeSpan.FromMilliseconds(baseDelay.TotalMilliseconds * Math.Pow(2, attempt));
|
||||
}
|
||||
@@ -67,22 +67,21 @@ public sealed class TgStatusServiceTests
|
||||
/// Готовый аккаунт
|
||||
/// </summary>
|
||||
[Fact]
|
||||
public async Task GetAsync_ReadyGateway_PrefersLiveAccountOverKvAndComputesConnected()
|
||||
public async Task GetAsync_ReadyGateway_ComposesLiveFieldsWithKvAccountMonitoredAndKeys()
|
||||
{
|
||||
(TgStatusService service, TestTelegramStore store, TestSettingsStore settings, TestTelegramGateway gateway, _, TestGlobalSettingsStore globalSettings) = Create();
|
||||
store.Seed(Dialog("d_1", "Канал", "channel", Monitor: true));
|
||||
settings.Preload(SettingsKeys.TgAccount, "\"@realuser\"");
|
||||
PreloadKeys(globalSettings, "123456", "abcdefghijklmnop");
|
||||
// Транспортный Connected=false, но phase ready → UI должен считать аккаунт подключённым.
|
||||
gateway.Status = new TelegramAccountStatusDto(
|
||||
Phase: "ready", Connected: false, Listener: true, Account: "@liveuser", Error: null, QrUrl: null);
|
||||
Phase: "ready", Connected: true, Listener: true, Account: "gateway-account", Error: null, QrUrl: null);
|
||||
|
||||
TgStatusDto status = await service.GetAsync(CancellationToken.None);
|
||||
|
||||
Assert.Equal("ready", status.Phase);
|
||||
Assert.True(status.Connected); // connected = авторизация (phase ready), не транспорт
|
||||
Assert.True(status.Connected);
|
||||
Assert.True(status.Listener);
|
||||
Assert.Equal("@liveuser", status.Account); // живой account приоритетнее KV
|
||||
Assert.Equal("@realuser", status.Account); // KV tgAccount перекрывает справочное поле гейта (Ruling 8)
|
||||
Assert.Equal(1, status.Monitored);
|
||||
Assert.True(status.KeysSet);
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ using Deal.Infrastructure.Data;
|
||||
using Deal.Infrastructure.Integrations.Abstractions;
|
||||
using Deal.Infrastructure.Integrations.Models;
|
||||
using Deal.Infrastructure.Integrations.Options;
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Deal.Infrastructure.Integrations.Services;
|
||||
using Deal.Modules.Kanban.Application.Models;
|
||||
using Deal.Modules.Settings.Application.Models;
|
||||
@@ -91,7 +92,8 @@ public sealed class GrpcMlClientTests
|
||||
Assert.Null(result.Margin);
|
||||
Assert.Empty(result.Terms);
|
||||
Assert.Null(result.Type);
|
||||
Assert.Single(service.RequestTenantIds);
|
||||
// Недоступность транспорта повторяется — на сервер приходит первая попытка и повторы.
|
||||
Assert.Equal(GrpcRetry.RetryCount + 1, service.RequestTenantIds.Count);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -148,7 +150,8 @@ public sealed class GrpcMlClientTests
|
||||
Assert.False(down.Reachable);
|
||||
Assert.False(down.Service.Ready);
|
||||
Assert.False(down.Stats.Reachable);
|
||||
Assert.Equal(1, service.StatusCalls);
|
||||
// При недоступности транспорта идёт повтор — считаем все попытки.
|
||||
Assert.Equal(GrpcRetry.RetryCount + 1, service.StatusCalls);
|
||||
|
||||
// «Поднялся»: после TTL 15 с следующий StatusAsync обновляет кэш (ready=true, reachable=true).
|
||||
service.StatusUnavailable = false;
|
||||
@@ -165,7 +168,7 @@ public sealed class GrpcMlClientTests
|
||||
Assert.True(up.Reachable);
|
||||
Assert.True(up.Service.Ready);
|
||||
Assert.Equal(3, up.Service.Learned);
|
||||
Assert.Equal(2, service.StatusCalls);
|
||||
Assert.Equal(GrpcRetry.RetryCount + 2, service.StatusCalls);
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ using Deal.Grpc.Ai;
|
||||
using Deal.Infrastructure.Data;
|
||||
using Deal.Infrastructure.Integrations.Models;
|
||||
using Deal.Infrastructure.Integrations.Options;
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Deal.Infrastructure.Integrations.Services;
|
||||
using Deal.Modules.Kanban.Application.Models;
|
||||
using Deal.Modules.Pipeline.Application.Models;
|
||||
@@ -101,8 +102,9 @@ public sealed class PipelineWorkerGrpcAiTests
|
||||
Assert.Equal(1, result.AiFail);
|
||||
Assert.Equal(1, result.AiStored);
|
||||
Assert.Single(result.CreatedCards);
|
||||
Assert.Equal(1, service.FilterCalls);
|
||||
Assert.Equal(1, service.ClassifyCalls);
|
||||
// Недоступность транспорта повторяется — на сервер приходит первая попытка и повторы.
|
||||
Assert.Equal(GrpcRetry.RetryCount + 1, service.FilterCalls);
|
||||
Assert.Equal(GrpcRetry.RetryCount + 1, service.ClassifyCalls);
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
using Deal.Infrastructure.Integrations.Resilience;
|
||||
using Grpc.Core;
|
||||
|
||||
namespace Deal.Tests.Unit.Support;
|
||||
|
||||
/// <summary>
|
||||
/// Тесты <see cref="GrpcRetry"/> — повтор транзиентных gRPC-сбоев.
|
||||
/// </summary>
|
||||
public sealed class GrpcRetryTests
|
||||
{
|
||||
[Theory]
|
||||
[InlineData(StatusCode.Unavailable)]
|
||||
[InlineData(StatusCode.DeadlineExceeded)]
|
||||
public void IsTransient_TransportFailures_True(StatusCode statusCode)
|
||||
{
|
||||
Assert.True(GrpcRetry.IsTransient(new RpcException(new Status(statusCode, "сбой"))));
|
||||
}
|
||||
|
||||
[Theory]
|
||||
[InlineData(StatusCode.NotFound)]
|
||||
[InlineData(StatusCode.InvalidArgument)]
|
||||
[InlineData(StatusCode.Internal)]
|
||||
public void IsTransient_ApplicationFailures_False(StatusCode statusCode)
|
||||
{
|
||||
Assert.False(GrpcRetry.IsTransient(new RpcException(new Status(statusCode, "сбой"))));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void IsTransient_NonRpcException_False()
|
||||
{
|
||||
Assert.False(GrpcRetry.IsTransient(new InvalidOperationException("сбой")));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ExecuteAsync_TransientThenSuccess_Retries()
|
||||
{
|
||||
int calls = 0;
|
||||
|
||||
string result = await GrpcRetry.ExecuteAsync(
|
||||
_ =>
|
||||
{
|
||||
calls++;
|
||||
return calls < 2
|
||||
? Task.FromException<string>(new RpcException(new Status(StatusCode.Unavailable, "down")))
|
||||
: Task.FromResult("ok");
|
||||
},
|
||||
InstantDelayAsync,
|
||||
CancellationToken.None);
|
||||
|
||||
Assert.Equal("ok", result);
|
||||
Assert.Equal(2, calls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ExecuteAsync_TransientExhausted_ThrowsRpcException()
|
||||
{
|
||||
int calls = 0;
|
||||
|
||||
RpcException thrown = await Assert.ThrowsAsync<RpcException>(
|
||||
() => GrpcRetry.ExecuteAsync(
|
||||
_ =>
|
||||
{
|
||||
calls++;
|
||||
return Task.FromException<string>(new RpcException(new Status(StatusCode.Unavailable, "down")));
|
||||
},
|
||||
InstantDelayAsync,
|
||||
CancellationToken.None));
|
||||
|
||||
Assert.Equal(StatusCode.Unavailable, thrown.StatusCode);
|
||||
Assert.Equal(GrpcRetry.RetryCount + 1, calls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ExecuteAsync_ApplicationFailure_NotRetried()
|
||||
{
|
||||
int calls = 0;
|
||||
|
||||
await Assert.ThrowsAsync<RpcException>(
|
||||
() => GrpcRetry.ExecuteAsync(
|
||||
_ =>
|
||||
{
|
||||
calls++;
|
||||
return Task.FromException<string>(new RpcException(new Status(StatusCode.InvalidArgument, "bad")));
|
||||
},
|
||||
InstantDelayAsync,
|
||||
CancellationToken.None));
|
||||
|
||||
Assert.Equal(1, calls);
|
||||
}
|
||||
|
||||
private static Task InstantDelayAsync(TimeSpan delay, CancellationToken cancellationToken)
|
||||
=> Task.CompletedTask;
|
||||
}
|
||||
@@ -0,0 +1,158 @@
|
||||
using Deal.SharedKernel.Resilience;
|
||||
|
||||
namespace Deal.Tests.Unit.Support;
|
||||
|
||||
/// <summary>
|
||||
/// Тесты <see cref="RetryExecutor"/> — повтор транзиентных сбоев.
|
||||
/// </summary>
|
||||
public sealed class RetryExecutorTests
|
||||
{
|
||||
private const int RetryCount = 2;
|
||||
private static readonly TimeSpan BaseDelay = TimeSpan.FromMilliseconds(10);
|
||||
|
||||
[Fact]
|
||||
public async Task FirstAttemptSucceeds_NoRetry()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
|
||||
string result = await ExecuteAsync(
|
||||
() =>
|
||||
{
|
||||
calls++;
|
||||
return Task.FromResult("ok");
|
||||
},
|
||||
shouldRetry: _ => true,
|
||||
delays);
|
||||
|
||||
Assert.Equal("ok", result);
|
||||
Assert.Equal(1, calls);
|
||||
Assert.Empty(delays);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task TransientFailureThenSuccess_RetriesUntilSuccess()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
|
||||
string result = await ExecuteAsync(
|
||||
() =>
|
||||
{
|
||||
calls++;
|
||||
return calls < 3
|
||||
? throw new InvalidOperationException("транзиент")
|
||||
: Task.FromResult("ok");
|
||||
},
|
||||
shouldRetry: _ => true,
|
||||
delays);
|
||||
|
||||
Assert.Equal("ok", result);
|
||||
Assert.Equal(3, calls);
|
||||
Assert.Equal(2, delays.Count);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RetriesExhausted_ThrowsLastFailure()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
|
||||
InvalidOperationException thrown = await Assert.ThrowsAsync<InvalidOperationException>(
|
||||
() => ExecuteAsync<string>(
|
||||
async () =>
|
||||
{
|
||||
calls++;
|
||||
await Task.Yield();
|
||||
throw new InvalidOperationException($"сбой {calls}");
|
||||
},
|
||||
shouldRetry: _ => true,
|
||||
delays));
|
||||
|
||||
Assert.Equal(3, calls); // первая попытка + 2 повтора
|
||||
Assert.Equal("сбой 3", thrown.Message);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task NonTransientFailure_NotRetried()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
|
||||
await Assert.ThrowsAsync<InvalidOperationException>(
|
||||
() => ExecuteAsync<string>(
|
||||
async () =>
|
||||
{
|
||||
calls++;
|
||||
await Task.Yield();
|
||||
throw new InvalidOperationException("не транзиент");
|
||||
},
|
||||
shouldRetry: _ => false,
|
||||
delays));
|
||||
|
||||
Assert.Equal(1, calls);
|
||||
Assert.Empty(delays);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Cancellation_DoesNotRetry()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
using var cts = new CancellationTokenSource();
|
||||
cts.Cancel();
|
||||
|
||||
await Assert.ThrowsAsync<InvalidOperationException>(
|
||||
() => ExecuteAsync<string>(
|
||||
async () =>
|
||||
{
|
||||
calls++;
|
||||
await Task.Yield();
|
||||
throw new InvalidOperationException("сбой");
|
||||
},
|
||||
shouldRetry: _ => true,
|
||||
delays,
|
||||
cts.Token));
|
||||
|
||||
Assert.Equal(1, calls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Backoff_GrowsExponentially()
|
||||
{
|
||||
var delays = new List<TimeSpan>();
|
||||
int calls = 0;
|
||||
|
||||
await ExecuteAsync(
|
||||
() =>
|
||||
{
|
||||
calls++;
|
||||
return calls < 3
|
||||
? throw new InvalidOperationException("транзиент")
|
||||
: Task.FromResult(1);
|
||||
},
|
||||
shouldRetry: _ => true,
|
||||
delays);
|
||||
|
||||
Assert.Equal(2, delays.Count);
|
||||
Assert.Equal(BaseDelay, delays[0]);
|
||||
Assert.Equal(BaseDelay * 2, delays[1]);
|
||||
}
|
||||
|
||||
private static Task<TResult> ExecuteAsync<TResult>(
|
||||
Func<Task<TResult>> operation,
|
||||
Func<Exception, bool> shouldRetry,
|
||||
List<TimeSpan> delays,
|
||||
CancellationToken cancellationToken = default)
|
||||
=> RetryExecutor.ExecuteAsync(
|
||||
_ => operation(),
|
||||
RetryCount,
|
||||
BaseDelay,
|
||||
shouldRetry,
|
||||
(delay, _) =>
|
||||
{
|
||||
delays.Add(delay);
|
||||
return Task.CompletedTask;
|
||||
},
|
||||
cancellationToken);
|
||||
}
|
||||
@@ -1,14 +0,0 @@
|
||||
# Фронтенд Deal: сборка SPA (Vite) и отдача статики через Caddy.
|
||||
# Контекст сборки — корень репозитория (см. deploy/compose.dev.yml: build.context ..).
|
||||
# /api/* проксируется на core:5080 (см. deploy/caddy/Caddyfile.dev).
|
||||
|
||||
FROM node:24-alpine AS build
|
||||
WORKDIR /app
|
||||
COPY src/frontend/package.json src/frontend/package-lock.json ./
|
||||
RUN npm ci
|
||||
COPY src/frontend/ ./
|
||||
RUN npm run build
|
||||
|
||||
FROM caddy:2.9.1
|
||||
COPY --from=build /app/dist /srv
|
||||
COPY deploy/caddy/Caddyfile.dev /etc/caddy/Caddyfile
|
||||
@@ -2,9 +2,10 @@
|
||||
// Вкладка «Telegram» настроек: подключение аккаунта (QR / номер / код /
|
||||
// облачный пароль) и авто-мониторинг новых чатов. Ключи приложения
|
||||
// (api_id/api_hash) задаёт оператор глобально — у тенанта их нет.
|
||||
// Пока идёт QR-вход, вкладка сама опрашивает статус и обновляет картинку.
|
||||
import { onBeforeUnmount, onMounted, ref, watch } from 'vue'
|
||||
import { state, connectTgStep, tgStart, disconnectTg, refreshTgStatus } from '../../store.js'
|
||||
// Опрос статуса при QR-входе живёт в SettingsView (watch на state.tgState),
|
||||
// чтобы поведение при переключении вкладок осталось прежним.
|
||||
import { ref } from 'vue'
|
||||
import { state, connectTgStep, tgStart, disconnectTg } from '../../store.js'
|
||||
import Icon from '../Icon.vue'
|
||||
import UiToggle from '../ui/ToggleSwitch.vue'
|
||||
|
||||
@@ -13,58 +14,6 @@ const code = ref('')
|
||||
const tgPass = ref('')
|
||||
const qrTick = ref(0)
|
||||
|
||||
// QR-ссылка приходит после старта: перезагружаем картинку, когда она появилась/сменилась.
|
||||
watch(
|
||||
() => state.tgQrUrl,
|
||||
() => {
|
||||
if (state.tgState === 'qr') qrTick.value++
|
||||
},
|
||||
)
|
||||
|
||||
// Таймер QR: пока идёт вход — периодически обновляем статус и саму картинку
|
||||
// (токен Telegram меняется); при подключении показываем учётку и останавливаемся.
|
||||
const QR_POLL_MS = 5000
|
||||
let qrTimer = null
|
||||
|
||||
function stopQrTimer() {
|
||||
if (qrTimer) {
|
||||
clearInterval(qrTimer)
|
||||
qrTimer = null
|
||||
}
|
||||
}
|
||||
|
||||
function startQrTimer() {
|
||||
stopQrTimer()
|
||||
qrTimer = setInterval(async () => {
|
||||
await refreshTgStatus()
|
||||
// Продолжаем опрос, пока не авторизованы: переходной статус (phase ready до
|
||||
// подтверждения транспорта) не должен останавливать обновление.
|
||||
if (state.tgConnected) {
|
||||
stopQrTimer()
|
||||
return
|
||||
}
|
||||
if (state.tgState === 'idle') {
|
||||
stopQrTimer()
|
||||
return
|
||||
}
|
||||
qrTick.value++
|
||||
}, QR_POLL_MS)
|
||||
}
|
||||
|
||||
// Таймер стартует при входе в режим QR и сам останавливается при подключении/сбросе.
|
||||
watch(
|
||||
() => state.tgState,
|
||||
(s) => {
|
||||
if (s === 'qr') startQrTimer()
|
||||
},
|
||||
)
|
||||
|
||||
onMounted(async () => {
|
||||
await refreshTgStatus()
|
||||
if (state.tgState === 'qr') startQrTimer()
|
||||
})
|
||||
onBeforeUnmount(stopQrTimer)
|
||||
|
||||
async function submitCode() {
|
||||
await connectTgStep(code.value)
|
||||
if (state.tgState === 'done') code.value = ''
|
||||
|
||||
@@ -5,12 +5,11 @@ import { state, toast, errMsg, fmtMsgTime } from './core.js'
|
||||
|
||||
export function mapTgStatus(st) {
|
||||
if (!st) return
|
||||
const phase = st.phase || 'idle'
|
||||
// «Подключён» = авторизован (phase ready); транспортный connected ненадёжен для UI.
|
||||
state.tgConnected = phase === 'ready' || !!st.connected
|
||||
state.tgConnected = !!st.connected
|
||||
state.tgAccount = st.account || ''
|
||||
state.tgKeysSet = !!st.keysSet
|
||||
if (st.error) state.tgError = st.error
|
||||
const phase = st.phase || 'idle'
|
||||
if (phase === 'ready') {
|
||||
state.tgState = 'done'
|
||||
} else if (['phone', 'code', 'password', 'qr'].includes(phase)) {
|
||||
@@ -148,12 +147,10 @@ export async function refreshTgStatus() {
|
||||
export async function tgStart() {
|
||||
state.tgError = ''
|
||||
if (state.tgQrMode) {
|
||||
state.tgState = 'qr'
|
||||
try {
|
||||
// Сначала стартуем QR на бэке, затем показываем картинку: иначе <img> успевает запросить
|
||||
// /api/tg/qr-image до готовности QR и остаётся пустым.
|
||||
const r = await api.post('/api/tg/start-qr')
|
||||
state.tgQrUrl = r.qrUrl || ''
|
||||
state.tgState = 'qr'
|
||||
toast(t('settings.ssylka-dlya-vhoda-sgenerirovana-otkrojte'), { icon: 'send' })
|
||||
} catch (e) {
|
||||
state.tgError = errMsg(e)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
<script setup>
|
||||
import { t } from '@/i18n/index.js'
|
||||
import { state } from '../store.js'
|
||||
import { onBeforeUnmount, watch } from 'vue'
|
||||
import { state, refreshTgStatus } from '../store.js'
|
||||
import Icon from '../components/Icon.vue'
|
||||
import TelegramTab from '../components/settings/TelegramTab.vue'
|
||||
import AiTab from '../components/settings/AiTab.vue'
|
||||
@@ -25,6 +26,32 @@ const TABS = [
|
||||
{ id: 'appearance', name: t('settings.vneshnij-vid'), icon: 'palette' },
|
||||
{ id: 'profile', name: t('settings.profil'), icon: 'key' },
|
||||
]
|
||||
|
||||
// пока идёт QR-вход (вкладка Telegram) — опрашиваем статус, чтобы поймать
|
||||
// момент подтверждения. Watch живёт на уровне экрана настроек: при
|
||||
// переключении вкладок опрос не прерывается (как в исходном SettingsView).
|
||||
let qrTimer = null
|
||||
watch(
|
||||
() => state.tgState,
|
||||
(s) => {
|
||||
if (s === 'qr') {
|
||||
stopQrPoll()
|
||||
qrTimer = setInterval(async () => {
|
||||
await refreshTgStatus()
|
||||
if (state.tgState === 'done') stopQrPoll()
|
||||
}, 4000)
|
||||
} else {
|
||||
stopQrPoll()
|
||||
}
|
||||
},
|
||||
)
|
||||
function stopQrPoll() {
|
||||
if (qrTimer) {
|
||||
clearInterval(qrTimer)
|
||||
qrTimer = null
|
||||
}
|
||||
}
|
||||
onBeforeUnmount(stopQrPoll)
|
||||
</script>
|
||||
|
||||
<template>
|
||||
|
||||
Reference in New Issue
Block a user