Files
Deal/src/telegram-service/Deal.Telegram/Sessions/TenantSession.cs
T
2026-09-13 20:51:43 +03:00

1048 lines
40 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
using Deal.Contracts.Integrations.Models;
using Deal.Telegram.Telegram;
using Grpc.Core;
namespace Deal.Telegram.Sessions;
/// <summary>
/// Сессия тенанта: id тенанта, клиент Telegram и состояние входа.
/// </summary>
public sealed class TenantSession : IAsyncDisposable
{
private readonly ITelegramClientFactory _clientFactory;
private readonly SessionStore _sessionStore;
private readonly ILogger<TenantSession> _logger;
private readonly SemaphoreSlim _gate = new(1, 1);
// Таймаут одной попытки переподключения heartbeat по умолчанию (зависший connect не должен блокировать
// heartbeat остальных тенантов и остановку хоста — замечание code-review).
private static readonly TimeSpan DefaultReconnectAttemptTimeout = TimeSpan.FromSeconds(10);
private readonly TimeSpan _reconnectAttemptTimeout;
private ISessionClient? _client;
private int _clientApiId;
private string? _clientApiHash;
private bool _registered;
private bool _loggedOut;
private AuthPhase _phase = AuthPhase.Idle;
private string? _error;
private string? _account;
private string? _qrUrl;
private string? _phone;
private CancellationTokenSource? _qrCts;
private Task? _qrWaitTask;
private volatile bool _listenerActive;
/// <summary>
/// Realtime-listener сессии жив
/// </summary>
public AuthPhase Phase => _phase;
/// <summary>
/// Событие входящего сообщения аккаунта.
/// </summary>
public event Func<TelegramMessage, Task>? MessageReceived;
/// <summary>
/// Создаёт сессию тенанта
/// </summary>
/// <param name="tenantId">Id тенанта (принадлежность сессии).</param>
/// <param name="clientFactory">Фабрика клиентов Telegram (реальная или фейк в тестах).</param>
/// <param name="sessionStore">Файловое хранилище сессий (шифрование at-rest).</param>
/// <param name="logger">Логгер.</param>
/// <param name="reconnectAttemptTimeout">Таймаут попытки переподключения (null — 10 с по умолчанию; тесты).</param>
public TenantSession(
string tenantId,
ITelegramClientFactory clientFactory,
SessionStore sessionStore,
ILogger<TenantSession> logger,
TimeSpan? reconnectAttemptTimeout = null)
{
TenantId = tenantId;
_clientFactory = clientFactory;
_sessionStore = sessionStore;
_logger = logger;
_reconnectAttemptTimeout = reconnectAttemptTimeout ?? DefaultReconnectAttemptTimeout;
}
/// <summary>
/// Id тенанта, которому принадлежит сессия
/// </summary>
public string TenantId { get; }
/// <summary>
/// Вход по номеру телефона
/// </summary>
/// <param name="apiId">api_id приложения (из тела запроса ядра).</param>
/// <param name="apiHash">api_hash приложения.</param>
/// <param name="phone">Номер телефона (международный формат).</param>
/// <returns>Снимок состояния после операции (фаза "code").</returns>
/// <exception cref="SessionException">Нет ключей (INVALID_ARGUMENT) / ошибки Telegram.</exception>
public async Task<TenantSessionSnapshot> StartPhoneAsync(
int apiId,
string apiHash,
string phone,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ValidateApiKeys(apiId, apiHash);
_registered = true;
_loggedOut = false;
_error = null;
_phone = phone;
CancelQrFlow();
await EnsureClientAsync(apiId, apiHash, cancellationToken).ConfigureAwait(false);
if (_client!.IsAuthorized)
{
// Сессия уже авторизована (например, возобновлена на старте) — код не нужен.
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
return Snapshot();
}
try
{
await _client.ConnectAsync(cancellationToken).ConfigureAwait(false);
await _client.RequestCodeAsync(phone, cancellationToken).ConfigureAwait(false);
}
catch (SessionException exception)
{
_phase = AuthPhase.Idle;
_error = exception.Message;
throw;
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
SessionException wrapped = new(StatusCode.Unavailable, SessionErrorMessages.TelegramUnavailable, exception);
_phase = AuthPhase.Idle;
_error = wrapped.Message;
throw wrapped;
}
_phase = AuthPhase.Code;
return Snapshot();
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Начать QR-вход: фаза "qr" + первый URL либо "ready", если уже вошли.
/// </summary>
/// <param name="apiId">api_id приложения.</param>
/// <param name="apiHash">api_hash приложения.</param>
/// <returns>Снимок состояния: Qr с qrUrl или Ready (авторизация уже была).</returns>
/// <exception cref="SessionException">Нет ключей (INVALID_ARGUMENT) / ошибки Telegram до первого URL.</exception>
public async Task<TenantSessionSnapshot> StartQrAsync(
int apiId,
string apiHash,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ValidateApiKeys(apiId, apiHash);
_registered = true;
_loggedOut = false;
_error = null;
await EnsureClientAsync(apiId, apiHash, cancellationToken).ConfigureAwait(false);
if (_client!.IsAuthorized)
{
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
return Snapshot();
}
if (_phase == AuthPhase.Qr && _qrWaitTask is { IsCompleted: false })
{
return Snapshot();
}
CancelQrFlow();
_phase = AuthPhase.Qr;
_qrUrl = null;
var qrCts = new CancellationTokenSource();
_qrCts = qrCts;
var firstUrlTcs = new TaskCompletionSource<string?>(TaskCreationOptions.RunContinuationsAsynchronously);
_qrWaitTask = RunQrFlowAsync(qrCts.Token, firstUrlTcs);
string? firstUrl;
try
{
firstUrl = await WaitForFirstQrUrlAsync(_qrWaitTask, firstUrlTcs, cancellationToken).ConfigureAwait(false);
}
catch (SessionException exception)
{
_phase = AuthPhase.Idle;
_qrUrl = null;
_error = exception.Message;
throw;
}
catch (OperationCanceledException)
{
if (_phase == AuthPhase.Qr)
{
_phase = AuthPhase.Idle;
_qrUrl = null;
}
CancelQrFlow();
throw;
}
if (firstUrl is not null)
{
_qrUrl = firstUrl;
}
// Задача завершилась без URL — авторизация произошла мгновенно (фаза Ready выставлена задачей)
// либо URL пришёл (фаза Qr). Возвращаем актуальный снимок.
return Snapshot();
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Отправить SMS-код.
/// </summary>
/// <param name="code">Код из SMS/Telegram-сообщения.</param>
/// <returns>Снимок состояния: "password" при 2FA либо "ready" после авторизации.</returns>
public async Task<TenantSessionSnapshot> SendCodeAsync(string code, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
EnsureLoginStarted();
if (_phase != AuthPhase.Code || _client is null)
{
throw new SessionException(StatusCode.FailedPrecondition, SessionErrorMessages.CodeNotRequested);
}
_error = null;
string? nextStep;
try
{
nextStep = await _client.SubmitCodeAsync(code, cancellationToken).ConfigureAwait(false);
}
catch (SessionException exception)
{
_error = exception.Message;
throw;
}
if (nextStep == TelegramAuthPhases.Password)
{
_phase = AuthPhase.Password;
return Snapshot();
}
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
return Snapshot();
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Отправить облачный пароль 2FA.
/// </summary>
/// <param name="password">Облачный пароль.</param>
/// <returns>Снимок состояния после операции (фаза "ready").</returns>
public async Task<TenantSessionSnapshot> SendPasswordAsync(string password, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
EnsureLoginStarted();
if (_phase != AuthPhase.Password || _client is null)
{
throw new SessionException(StatusCode.FailedPrecondition, SessionErrorMessages.PasswordNotRequested);
}
_error = null;
try
{
await _client.SubmitPasswordAsync(password, cancellationToken).ConfigureAwait(false);
}
catch (SessionException exception)
{
_error = exception.Message;
throw;
}
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
return Snapshot();
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Отключить аккаунт и удалить сессию тенанта.
/// </summary>
public async Task<TenantSessionSnapshot?> LogoutAsync(CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
CancelQrFlow();
if (_client is not null)
{
try
{
await _client.LogOutAsync(cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
// Auth_LogOut недоступен (сеть) — локальный выход и удаление файла всё равно выполняются.
_logger.LogWarning(exception, "Logout {TenantId}: Auth_LogOut не выполнен — продолжаем локальный выход", TenantId);
}
try
{
await _sessionStore.DeleteAsync(TenantId, cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Logout {TenantId}: файл сессии не удалён", TenantId);
}
DetachClientMessages(_client);
await _client.DisposeAsync().ConfigureAwait(false);
_client = null;
}
_clientApiId = 0;
_clientApiHash = null;
_phase = AuthPhase.Idle;
_error = null;
_account = null;
_qrUrl = null;
_phone = null;
_loggedOut = true;
_registered = false;
return null;
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Снимок состояния для GetStatus; null — сессии тенанта нет
/// </summary>
public async Task<TenantSessionSnapshot?> GetSnapshotAsync(CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
return _loggedOut || !_registered ? null : Snapshot();
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Список диалогов аккаунта; фаза обязана быть ready.
/// </summary>
/// <param name="limit">Верхняя граница числа диалогов.</param>
/// <returns>Диалоги аккаунта (нейтральный вид).</returns>
public async Task<IReadOnlyList<TelegramDialog>> ListDialogsAsync(int limit, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.GetDialogsAsync(limit, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Последние сообщения диалога
/// </summary>
/// <param name="dialogId">Подписанный id диалога.</param>
/// <param name="limit">Сколько последних сообщений запросить.</param>
/// <returns>Сообщения диалога (от новых к старым, непустые тексты).</returns>
public async Task<IReadOnlyList<TelegramMessage>> GetMessagesAsync(
string dialogId,
int limit,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.GetMessagesAsync(dialogId, limit, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Исходное сообщение диалога по id (remote-просмотр).
/// </summary>
/// <param name="dialogId">Подписанный id диалога.</param>
/// <param name="msgId">Id сообщения в Telegram.</param>
/// <returns>Сообщение с текстом либо null (медиа без текста/не найдено).</returns>
public async Task<TelegramMessage?> GetMessageAsync(
string dialogId,
long msgId,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.GetMessageAsync(dialogId, msgId, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Помечает диалог прочитанным
/// </summary>
/// <param name="dialogId">Подписанный id диалога.</param>
public async Task MarkReadAsync(string dialogId, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
await client.MarkReadAsync(dialogId, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Глобальный поиск каналов/групп по ключу; фаза ready.
/// </summary>
/// <param name="query">Поисковый запрос (ключ задачи discovery).</param>
/// <param name="limit">Верхняя граница результата.</param>
/// <returns>Найденные источники.</returns>
public async Task<IReadOnlyList<TelegramDialog>> SearchAsync(
string query,
int limit,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.SearchAsync(query, limit, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Инфо об источнике для оценки кандидата; фаза ready.
/// </summary>
/// <param name="dialogId">Подписанный id источника.</param>
/// <returns>Инфо (по умолчанию — только id; сбои определения не бросаются, см. ISessionClient).</returns>
public async Task<TelegramSourceInfo> GetInfoAsync(string dialogId, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.GetInfoAsync(dialogId, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Выборка сообщений источника для оценки; фаза ready.
/// </summary>
/// <param name="dialogId">Подписанный id источника.</param>
/// <param name="limit">Размер выборки (limit ≤ 0 — пусто без сети).</param>
/// <returns>Результат чтения (ok/messages либо no_history).</returns>
public async Task<DiscoveryReadResult> ReadForEvalAsync(
string dialogId,
int limit,
CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
return await client.ReadForEvalAsync(dialogId, limit, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Вступить в канал/группу по username; фаза ready.
/// </summary>
/// <param name="username">Username (без «@»; нормализует DiscoveryOps).</param>
public async Task JoinAsync(string username, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
await client.JoinAsync(username, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Выйти из канала/группы; фаза ready.
/// </summary>
/// <param name="dialogId">Подписанный id диалога.</param>
public async Task LeaveAsync(string dialogId, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ISessionClient client = await EnsureReadyConnectedAsync(cancellationToken).ConfigureAwait(false);
await client.LeaveAsync(dialogId, cancellationToken).ConfigureAwait(false);
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Ставит признак живого realtime-listener.
/// </summary>
/// <param name="active">True — listener сессии подписан на события сообщений.</param>
public void SetListenerActive(bool active)
=> _listenerActive = active;
// Проверяет готовность сессии и соединения (фаза ready + клиент).
// cancellationToken: Отмена операции.
private async Task<ISessionClient> EnsureReadyConnectedAsync(CancellationToken cancellationToken)
{
ISessionClient client = RequireReadyClient();
if (client.IsConnected)
{
return client;
}
try
{
await client.ConnectAsync(cancellationToken).ConfigureAwait(false);
}
catch (SessionException)
{
throw;
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
throw new SessionException(StatusCode.Unavailable, SessionErrorMessages.TelegramUnavailable, exception);
}
return client;
}
// Ready-клиент сессии (иначе «Telegram не подключён», FAILED_PRECONDITION).
private ISessionClient RequireReadyClient()
{
if (_loggedOut || !_registered || _client is null || _phase != AuthPhase.Ready)
{
throw new SessionException(StatusCode.FailedPrecondition, SessionErrorMessages.NotConnected);
}
return _client;
}
// --- Проброс realtime-сообщений клиента на уровень службы ---
// Передаёт сообщение клиента подписчикам сессии (каждый в своей ошибко-изоляции).
// message: Входящее сообщение аккаунта.
private async Task ForwardClientMessageAsync(TelegramMessage message)
{
Func<TelegramMessage, Task>? subscribers = MessageReceived;
if (subscribers is null)
{
return;
}
foreach (Delegate subscriber in subscribers.GetInvocationList())
{
try
{
await ((Func<TelegramMessage, Task>)subscriber)(message).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Обработчик сообщения {TenantId} завершился с ошибкой", TenantId);
}
}
}
// Подписывает проброс сообщений нового клиента (вызывается после создания клиента).
// client: Новый клиент сессии.
private void AttachClientMessages(ISessionClient client)
=> client.MessageReceived += ForwardClientMessageAsync;
// Отписывает проброс сообщений клиента (перед Dispose клиента).
// client: Уходящий клиент сессии.
private void DetachClientMessages(ISessionClient client)
=> client.MessageReceived -= ForwardClientMessageAsync;
/// <summary>
/// Авто-возобновление на старте...
/// </summary>
/// <param name="stored">Содержимое файла сессии тенанта.</param>
/// <returns>True — сессия возобновлена (ready); false — не авторизована/сбой.</returns>
public async Task<bool> TryResumeAsync(StoredSession stored, CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
ValidateApiKeys(stored.ApiId, stored.ApiHash);
_registered = true;
_loggedOut = false;
_error = null;
if (_client is not null && _client.ApiId == stored.ApiId && _client.ApiHash == stored.ApiHash)
{
// Клиент уже создан под те же ключи (например, жив после входа в этой сессии).
}
else
{
if (_client is not null)
{
DetachClientMessages(_client);
await _client.DisposeAsync().ConfigureAwait(false);
}
_client = _clientFactory.Create(stored.ApiId, stored.ApiHash, stored.SessionBytes);
_clientApiId = stored.ApiId;
_clientApiHash = stored.ApiHash;
AttachClientMessages(_client);
}
try
{
await _client!.ConnectAsync(cancellationToken).ConfigureAwait(false);
if (!_client.IsAuthorized)
{
_phase = AuthPhase.Idle;
return false;
}
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
return true;
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "auto_resume {TenantId} пропущен", TenantId);
_phase = AuthPhase.Idle;
_error = exception is SessionException sessionException ? sessionException.Message : null;
return false;
}
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Сердцебиение: для фазы "ready" при обрыве соединения — повторный connect.
/// </summary>
public async Task TryReconnectAsync(CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
if (_loggedOut || _client is null || _phase != AuthPhase.Ready || _client.IsConnected)
{
return;
}
try
{
using CancellationTokenSource attemptTimeout =
CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
attemptTimeout.CancelAfter(_reconnectAttemptTimeout);
await _client.ConnectAsync(attemptTimeout.Token).ConfigureAwait(false);
}
catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested)
{
// Таймаут собственной попытки — не отмена хоста: предупреждение, повтор следующим циклом.
_logger.LogWarning(
"Heartbeat {TenantId}: таймаут переподключения ({TimeoutSeconds:0} с) — повторим следующим циклом",
TenantId,
_reconnectAttemptTimeout.TotalSeconds);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Heartbeat {TenantId}: переподключение не удалось — повторим следующим циклом", TenantId);
}
}
finally
{
_gate.Release();
}
}
/// <summary>
/// Остановка (хост гасится)
/// </summary>
public async Task FlushAndDisposeAsync(CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
CancelQrFlow();
if (_client is not null)
{
try
{
await PersistCurrentSessionAsync(cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Остановка {TenantId}: сессия не сохранена", TenantId);
}
DetachClientMessages(_client);
await _client.DisposeAsync().ConfigureAwait(false);
_client = null;
}
}
finally
{
_gate.Release();
}
}
/// <inheritdoc />
public async ValueTask DisposeAsync()
{
// Сброс без сетевых операций: авторизованные сессии уже сохранены на ключевых событиях
// (FlushAndDisposeAsync зовёт хост при остановке; здесь — финальная очистка клиента).
await _gate.WaitAsync().ConfigureAwait(false);
try
{
CancelQrFlow();
if (_client is not null)
{
DetachClientMessages(_client);
await _client.DisposeAsync().ConfigureAwait(false);
_client = null;
}
}
finally
{
_gate.Release();
}
}
// Проверяет, что сессия тенанта существует и не закрыта (иначе «Telegram не подключён»).
private void EnsureLoginStarted()
{
if (_loggedOut || !_registered)
{
throw new SessionException(StatusCode.FailedPrecondition, SessionErrorMessages.NotConnected);
}
}
private static void ValidateApiKeys(int apiId, string apiHash)
{
if (apiId <= 0 || string.IsNullOrWhiteSpace(apiHash))
{
throw new SessionException(StatusCode.InvalidArgument, SessionErrorMessages.NoApiKeys);
}
}
// Гарантирует клиент под запрошенные ключи: существующий клиент с теми же ключами переиспользуется;
// при смене ключей живая сессия сначала сохраняется (свои ключи), затем клиент пересоздаётся;
// сохранённая сессия с диска подсевается только при совпадении ключей приложения.
// apiId: api_id приложения.
// apiHash: api_hash приложения.
// cancellationToken: Отмена операции.
private async Task EnsureClientAsync(
int apiId,
string apiHash,
CancellationToken cancellationToken)
{
if (_client is not null && _clientApiId == apiId && _clientApiHash == apiHash)
{
return;
}
if (_client is not null)
{
DetachClientMessages(_client);
try
{
await PersistCurrentSessionAsync(cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Сессия {TenantId}: предыдущий клиент не сохранён при смене ключей", TenantId);
}
await _client.DisposeAsync().ConfigureAwait(false);
_client = null;
}
StoredSession? stored = await _sessionStore.LoadAsync(TenantId, cancellationToken).ConfigureAwait(false);
byte[]? seed = stored is not null && stored.ApiId == apiId && stored.ApiHash == apiHash
? stored.SessionBytes
: null;
_client = _clientFactory.Create(apiId, apiHash, seed);
_clientApiId = apiId;
_clientApiHash = apiHash;
AttachClientMessages(_client);
}
private async Task CompleteAuthorizationAsync(CancellationToken cancellationToken)
{
if (_client is null)
{
return;
}
try
{
_account = await _client.GetAccountAsync(cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
_logger.LogWarning(exception, "Финализация {TenantId}: account не получен — готовность сохраняется", TenantId);
}
_phase = AuthPhase.Ready;
_qrUrl = null;
_error = null;
await PersistCurrentSessionAsync(cancellationToken).ConfigureAwait(false);
}
// Шифрует и сохраняет текущие байты сессии в файл (если клиент их уже сформировал).
// cancellationToken: Отмена операции.
private async Task PersistCurrentSessionAsync(CancellationToken cancellationToken)
{
byte[]? sessionBytes = _client?.SessionBytes;
if (_client is null || sessionBytes is not { Length: > 0 })
{
return;
}
await _sessionStore.SaveAsync(
TenantId,
new StoredSession
{
ApiId = _client.ApiId,
ApiHash = _client.ApiHash,
SessionBytes = sessionBytes,
},
cancellationToken).ConfigureAwait(false);
}
// Отменяет активный QR-вход (без ожидания фоновой задачи: её guard увидит смену фазы).
private void CancelQrFlow()
{
if (_qrCts is not null)
{
_qrCts.Cancel();
_qrCts.Dispose();
_qrCts = null;
}
_qrWaitTask = null;
}
// Фоновая задача QR-входа: ждёт сканирования; URL обновляет колбэк, после авторизации под замком
// финализирует сессию. Ошибка до первого URL пробрасывается (StartQrAsync держит замок и сам
// сбросит фазу); после выдачи URL состояние обновляется здесь под замком.
// cancellationToken: Токен отмены QR (StartPhone/Logout/остановка).
// firstUrlTcs: Завершается первым URL (StartQrAsync ждёт его под замком).
private async Task RunQrFlowAsync(CancellationToken cancellationToken, TaskCompletionSource<string?> firstUrlTcs)
{
ISessionClient client = _client!;
bool urlAlreadyDelivered;
try
{
await client.StartQrAsync(
url =>
{
if (_phase == AuthPhase.Qr)
{
_qrUrl = url;
}
firstUrlTcs.TrySetResult(url);
},
cancellationToken).ConfigureAwait(false);
// Сканирование принято — авторизация завершена: финализировать под замком.
await _gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
if (!_loggedOut && ReferenceEquals(_client, client) && _phase == AuthPhase.Qr)
{
await CompleteAuthorizationAsync(cancellationToken).ConfigureAwait(false);
}
}
finally
{
_gate.Release();
}
}
catch (OperationCanceledException)
{
urlAlreadyDelivered = firstUrlTcs.Task.IsCompletedSuccessfully;
if (urlAlreadyDelivered)
{
await ResetQrUnderGateAsync().ConfigureAwait(false);
return;
}
firstUrlTcs.TrySetCanceled();
throw;
}
catch (SessionException exception)
{
urlAlreadyDelivered = firstUrlTcs.Task.IsCompletedSuccessfully;
if (urlAlreadyDelivered)
{
await FailQrUnderGateAsync(exception).ConfigureAwait(false);
return;
}
firstUrlTcs.TrySetException(exception);
throw;
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
SessionException wrapped = new(StatusCode.Unavailable, SessionErrorMessages.TelegramUnavailable, exception);
urlAlreadyDelivered = firstUrlTcs.Task.IsCompletedSuccessfully;
if (urlAlreadyDelivered)
{
await FailQrUnderGateAsync(wrapped).ConfigureAwait(false);
return;
}
firstUrlTcs.TrySetException(wrapped);
throw wrapped;
}
}
// Ждёт первый URL QR (или завершение/ошибку фоновой задачи). Вызывается из StartQrAsync под замком:
// при ошибке/отмене до первого URL задача уже завершена — замок освобождается при unwind.
// qrTask: Фоновая задача QR.
// firstUrlTcs: TCS первого URL.
// cancellationToken: Отмена операции.
// Возвращает: Первый URL либо null (задача завершилась без URL).
private static async Task<string?> WaitForFirstQrUrlAsync(
Task qrTask,
TaskCompletionSource<string?> firstUrlTcs,
CancellationToken cancellationToken)
{
Task<string?> urlTask = firstUrlTcs.Task;
Task completed = await Task.WhenAny(urlTask, qrTask).WaitAsync(cancellationToken).ConfigureAwait(false);
if (completed == urlTask)
{
return await urlTask.ConfigureAwait(false);
}
// Задача завершилась без URL (ошибка/отмена/авторизация без URL) — проброс результата задачи.
await qrTask.ConfigureAwait(false);
return null;
}
// Сброс QR-состояния после отмены (URL уже был выдан, RPC вернулся): фаза idle, если QR всё ещё
// владеет сессией (StartPhone/Logout уже сменили фазу/клиента — не трогаем).
private async Task ResetQrUnderGateAsync()
{
await _gate.WaitAsync().ConfigureAwait(false);
try
{
if (!_loggedOut && _phase == AuthPhase.Qr)
{
_phase = AuthPhase.Idle;
_qrUrl = null;
}
}
finally
{
_gate.Release();
}
}
private async Task FailQrUnderGateAsync(SessionException exception)
{
await _gate.WaitAsync().ConfigureAwait(false);
try
{
if (!_loggedOut && _phase == AuthPhase.Qr)
{
_phase = AuthPhase.Idle;
_qrUrl = null;
_error = exception.Message;
}
}
finally
{
_gate.Release();
}
}
private TenantSessionSnapshot Snapshot()
=> new(
_phase,
connected: _client?.IsConnected ?? false,
listener: _listenerActive,
account: _phase == AuthPhase.Ready ? _account : null,
error: _error,
qrUrl: _phase == AuthPhase.Qr ? _qrUrl : null);
}