feat(platform): route scheduler notifications through platform messenger
PR Checks / test-and-build (pull_request) Successful in 7m9s

This commit is contained in:
2026-05-21 12:30:35 +03:00
parent 5dbec1a0a4
commit 2a707e4825
49 changed files with 2158 additions and 846 deletions
@@ -0,0 +1,356 @@
using System.Globalization;
using Dapper;
using GmRelay.Shared.Domain;
using GmRelay.Shared.Features.Notifications;
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Logging;
using Npgsql;
namespace GmRelay.Shared.Features.Confirmation.HandleRsvp;
public sealed record HandleRsvpCommand(
Guid SessionId,
PlatformUser User,
string Status,
string InteractionId,
PlatformGroup Group,
PlatformMessageRef ConfirmationMessage);
internal sealed record RsvpCounts(int Total, int Confirmed, int Declined);
internal sealed record RsvpSessionContext(
Guid GroupId,
string Title,
DateTime ScheduledAt,
string Status);
internal sealed record ParticipantRsvpRow(
string Platform,
string ExternalUserId,
string DisplayName,
string? ExternalUsername,
string RsvpStatus,
string RegistrationStatus,
bool IsGm);
internal sealed record RsvpRecipientRow(
string Platform,
string ExternalUserId,
string DisplayName,
string? ExternalUsername);
public sealed class HandleRsvpHandler(
NpgsqlDataSource dataSource,
IPlatformMessenger messenger,
ILogger<HandleRsvpHandler> logger)
{
public async Task HandleAsync(HandleRsvpCommand command, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
await using var transaction = await connection.BeginTransactionAsync(ct);
var participantExists = await connection.ExecuteScalarAsync<bool>(
"""
SELECT EXISTS (
SELECT 1
FROM session_participants sp
JOIN players p ON p.id = sp.player_id
WHERE sp.session_id = @SessionId
AND COALESCE(p.platform, 'Telegram') = @Platform
AND COALESCE(p.external_user_id, p.telegram_id::TEXT) = @ExternalUserId
AND sp.is_gm = false
AND sp.registration_status = @Active
)
""",
new
{
command.SessionId,
Platform = command.User.Platform.ToString(),
command.User.ExternalUserId,
Active = ParticipantRegistrationStatus.Active
},
transaction);
if (!participantExists)
{
await messenger.AnswerInteractionAsync(
new PlatformInteractionReply(
command.InteractionId,
"Вы не являетесь участником этой сессии."),
ct);
return;
}
var updated = await connection.ExecuteAsync(
"""
UPDATE session_participants
SET rsvp_status = @Status,
responded_at = now()
WHERE session_id = @SessionId
AND player_id = (
SELECT id
FROM players
WHERE COALESCE(platform, 'Telegram') = @Platform
AND COALESCE(external_user_id, telegram_id::TEXT) = @ExternalUserId
LIMIT 1
)
AND registration_status = @Active
AND rsvp_status != @Status
""",
new
{
command.SessionId,
command.Status,
Platform = command.User.Platform.ToString(),
command.User.ExternalUserId,
Active = ParticipantRegistrationStatus.Active
},
transaction);
if (updated == 0)
{
var alreadyText = command.Status == RsvpStatus.Confirmed
? "Вы уже подтвердили участие."
: "Вы уже отказались от участия.";
await messenger.AnswerInteractionAsync(
new PlatformInteractionReply(command.InteractionId, alreadyText),
ct);
return;
}
var session = await connection.QuerySingleAsync<RsvpSessionContext>(
"""
SELECT s.group_id AS GroupId,
s.title,
s.scheduled_at AS ScheduledAt,
s.status AS Status
FROM sessions s
WHERE s.id = @SessionId
""",
new { command.SessionId },
transaction);
if (command.Status == RsvpStatus.Declined)
{
var decision = RsvpFlowRules.Evaluate(
command.Status,
session.Status,
totalParticipants: 0,
confirmedParticipants: 0);
if (decision.ShouldRevertSessionToConfirmationSent)
{
await connection.ExecuteAsync(
"""
UPDATE sessions
SET status = @ConfirmationSent, updated_at = now()
WHERE id = @SessionId AND status = @Confirmed
""",
new
{
command.SessionId,
ConfirmationSent = SessionStatus.ConfirmationSent,
Confirmed = SessionStatus.Confirmed
},
transaction);
}
var gmRecipients = (await GetGmRecipientsAsync(connection, session.GroupId, transaction))
.ToList();
await transaction.CommitAsync(ct);
if (gmRecipients.Count > 0)
{
await messenger.SendRsvpOutcomeAsync(
new PlatformRsvpOutcomeNotification(
PlatformRsvpOutcomeKind.GmPlayerDeclined,
Group: null,
gmRecipients,
command.SessionId,
session.Title,
session.ScheduledAt,
ActorDisplayName: command.User.DisplayName),
ct);
}
await messenger.AnswerInteractionAsync(
new PlatformInteractionReply(command.InteractionId, decision.CallbackText),
ct);
}
else
{
var counts = await connection.QuerySingleAsync<RsvpCounts>(
"""
SELECT
count(*) AS Total,
count(*) FILTER (WHERE rsvp_status = @Confirmed) AS Confirmed,
count(*) FILTER (WHERE rsvp_status = @Declined) AS Declined
FROM session_participants
WHERE session_id = @SessionId AND is_gm = false
AND registration_status = @Active
""",
new
{
command.SessionId,
Confirmed = RsvpStatus.Confirmed,
Declined = RsvpStatus.Declined,
Active = ParticipantRegistrationStatus.Active
},
transaction);
var decision = RsvpFlowRules.Evaluate(command.Status, session.Status, counts.Total, counts.Confirmed);
if (decision.ShouldMarkSessionConfirmed)
{
await connection.ExecuteAsync(
"""
UPDATE sessions
SET status = @Confirmed, updated_at = now()
WHERE id = @SessionId
""",
new { command.SessionId, Confirmed = SessionStatus.Confirmed },
transaction);
}
var gmRecipients = decision.ShouldNotifyGm
? (await GetGmRecipientsAsync(connection, session.GroupId, transaction)).ToList()
: [];
await transaction.CommitAsync(ct);
if (decision.ShouldNotifyGroup)
{
await messenger.SendRsvpOutcomeAsync(
new PlatformRsvpOutcomeNotification(
PlatformRsvpOutcomeKind.GroupAllConfirmed,
command.Group,
[],
command.SessionId,
session.Title,
session.ScheduledAt),
ct);
}
if (decision.ShouldNotifyGm && gmRecipients.Count > 0)
{
await messenger.SendRsvpOutcomeAsync(
new PlatformRsvpOutcomeNotification(
PlatformRsvpOutcomeKind.GmAllConfirmed,
Group: null,
gmRecipients,
command.SessionId,
session.Title,
session.ScheduledAt),
ct);
}
await messenger.AnswerInteractionAsync(
new PlatformInteractionReply(command.InteractionId, decision.CallbackText),
ct);
}
await UpdateConfirmationMessage(command, session, ct);
}
private async Task UpdateConfirmationMessage(
HandleRsvpCommand command,
RsvpSessionContext session,
CancellationToken ct)
{
try
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var participants = (await connection.QueryAsync<ParticipantRsvpRow>(
"""
SELECT COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername,
sp.rsvp_status AS RsvpStatus,
sp.registration_status AS RegistrationStatus,
sp.is_gm AS IsGm
FROM session_participants sp
JOIN players p ON p.id = sp.player_id
WHERE sp.session_id = @SessionId
AND sp.is_gm = false
AND sp.registration_status = @Active
ORDER BY sp.responded_at NULLS LAST
""",
new { command.SessionId, Active = ParticipantRegistrationStatus.Active }))
.Select(ToParticipant)
.ToList();
var disableActions = participants.Count > 0 &&
participants.All(participant => participant.RsvpStatus == RsvpStatus.Confirmed);
await messenger.UpdateConfirmationRequestAsync(
new PlatformRsvpMessageUpdate(
new PlatformConfirmationRequest(
command.Group,
command.SessionId,
session.Title,
session.ScheduledAt,
participants,
command.ConfirmationMessage),
disableActions),
ct);
}
catch (Exception ex)
{
logger.LogWarning(ex, "Failed to update confirmation message for session {SessionId}", command.SessionId);
}
}
private static async Task<IEnumerable<PlatformUser>> GetGmRecipientsAsync(
NpgsqlConnection connection,
Guid groupId,
NpgsqlTransaction transaction)
{
var rows = await connection.QueryAsync<RsvpRecipientRow>(
"""
SELECT DISTINCT
COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername
FROM group_managers gm
JOIN players p ON p.id = gm.player_id
WHERE gm.group_id = @GroupId
UNION
SELECT DISTINCT
COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername
FROM game_groups g
JOIN players p ON p.telegram_id = g.gm_telegram_id
WHERE g.id = @GroupId
AND g.gm_telegram_id IS NOT NULL
""",
new { GroupId = groupId },
transaction);
return rows.Select(row => new PlatformUser(
ParsePlatform(row.Platform),
row.ExternalUserId,
row.DisplayName,
row.ExternalUsername));
}
private static PlatformSessionParticipant ToParticipant(ParticipantRsvpRow row) =>
new(
new PlatformUser(
ParsePlatform(row.Platform),
row.ExternalUserId,
row.DisplayName,
row.ExternalUsername),
row.RsvpStatus,
row.RegistrationStatus,
row.IsGm);
private static PlatformKind ParsePlatform(string platform) =>
Enum.Parse<PlatformKind>(platform, ignoreCase: true);
}
@@ -0,0 +1,42 @@
using GmRelay.Shared.Domain;
namespace GmRelay.Shared.Features.Confirmation.HandleRsvp;
public sealed record RsvpFlowDecision(
string CallbackText,
bool ShouldAlertGm,
bool ShouldRevertSessionToConfirmationSent,
bool ShouldMarkSessionConfirmed,
bool ShouldNotifyGroup,
bool ShouldNotifyGm);
public static class RsvpFlowRules
{
public static RsvpFlowDecision Evaluate(
string requestedStatus,
string currentSessionStatus,
int totalParticipants,
int confirmedParticipants)
{
if (requestedStatus == RsvpStatus.Declined)
{
return new RsvpFlowDecision(
CallbackText: "Вы отказались от участия.",
ShouldAlertGm: true,
ShouldRevertSessionToConfirmationSent: currentSessionStatus == SessionStatus.Confirmed,
ShouldMarkSessionConfirmed: false,
ShouldNotifyGroup: false,
ShouldNotifyGm: false);
}
var everyoneConfirmed = confirmedParticipants == totalParticipants;
return new RsvpFlowDecision(
CallbackText: "Вы подтвердили участие!",
ShouldAlertGm: false,
ShouldRevertSessionToConfirmationSent: false,
ShouldMarkSessionConfirmed: everyoneConfirmed,
ShouldNotifyGroup: everyoneConfirmed,
ShouldNotifyGm: everyoneConfirmed);
}
}
@@ -0,0 +1,6 @@
namespace GmRelay.Shared.Features.Confirmation.SendConfirmation;
public interface ISendConfirmationHandler
{
Task HandleAsync(Guid sessionId, CancellationToken ct);
}
@@ -0,0 +1,217 @@
using System.Globalization;
using Dapper;
using GmRelay.Shared.Domain;
using GmRelay.Shared.Features.Notifications;
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Logging;
using Npgsql;
namespace GmRelay.Shared.Features.Confirmation.SendConfirmation;
internal sealed record ConfirmationSessionRow(
Guid Id,
string Title,
DateTime ScheduledAt,
Guid GroupId,
string Platform,
string ExternalGroupId,
string DisplayName,
string? ExternalChannelId,
int? ThreadId,
string NotificationMode);
internal sealed record ConfirmationParticipantRow(
string Platform,
string ExternalUserId,
string DisplayName,
string? ExternalUsername,
string RsvpStatus,
string RegistrationStatus,
bool IsGm);
public sealed class SendConfirmationHandler(
NpgsqlDataSource dataSource,
IPlatformMessenger messenger,
PlatformDirectNotificationSender directSender,
ILogger<SendConfirmationHandler> logger) : ISendConfirmationHandler
{
public async Task HandleAsync(Guid sessionId, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var session = await connection.QuerySingleOrDefaultAsync<ConfirmationSessionRow>(
"""
SELECT s.id,
s.title,
s.scheduled_at AS ScheduledAt,
s.group_id AS GroupId,
COALESCE(g.platform, 'Telegram') AS Platform,
COALESCE(g.external_group_id, g.telegram_chat_id::TEXT) AS ExternalGroupId,
g.name AS DisplayName,
COALESCE(g.external_channel_id, g.telegram_chat_id::TEXT) AS ExternalChannelId,
s.thread_id AS ThreadId,
s.notification_mode AS NotificationMode
FROM sessions s
JOIN game_groups g ON g.id = s.group_id
WHERE s.id = @SessionId AND s.status = @Planned
""",
new { SessionId = sessionId, Planned = SessionStatus.Planned });
if (session is null)
{
logger.LogWarning("Session {SessionId} not found or not in Planned status", sessionId);
return;
}
var participants = (await connection.QueryAsync<ConfirmationParticipantRow>(
"""
SELECT COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername,
sp.rsvp_status AS RsvpStatus,
sp.registration_status AS RegistrationStatus,
sp.is_gm AS IsGm
FROM session_participants sp
JOIN players p ON p.id = sp.player_id
WHERE sp.session_id = @SessionId
AND sp.is_gm = false
AND sp.registration_status = @Active
ORDER BY sp.created_at ASC
""",
new { SessionId = sessionId, Active = ParticipantRegistrationStatus.Active }))
.Select(ToParticipant)
.ToList();
if (participants.Count == 0)
{
logger.LogWarning("Session {SessionId} has no non-GM participants", sessionId);
return;
}
var group = CreateGroup(session);
var message = await messenger.SendConfirmationRequestAsync(
new PlatformConfirmationRequest(
group,
session.Id,
session.Title,
session.ScheduledAt,
participants),
ct);
await connection.ExecuteAsync(
"""
UPDATE sessions
SET status = @Status,
confirmation_message_id = @MessageId,
confirmation_sent_at = now(),
updated_at = now()
WHERE id = @SessionId
AND confirmation_sent_at IS NULL
""",
new
{
SessionId = sessionId,
Status = SessionStatus.ConfirmationSent,
MessageId = TryGetTelegramMessageId(message)
});
await PersistPlatformMessageAsync(
connection,
message,
session.GroupId,
session.Id,
batchId: null,
purpose: "confirmation");
var mode = SessionNotificationModeExtensions.FromDatabaseValue(session.NotificationMode);
if (mode.ShouldSendDirectMessages())
{
await directSender.SendAsync(
PlatformDirectSessionNotificationKind.ConfirmationRequest,
participants.Select(p => p.User),
session.Id,
session.Title,
session.ScheduledAt,
joinLink: null,
actorDisplayName: null,
reason: null,
ct);
}
logger.LogInformation(
"Confirmation sent for session {SessionId} ({Title}), platform={Platform}, message_id={MessageId}",
sessionId,
session.Title,
message.Platform,
message.ExternalMessageId);
}
private static PlatformSessionParticipant ToParticipant(ConfirmationParticipantRow row) =>
new(
new PlatformUser(
ParsePlatform(row.Platform),
row.ExternalUserId,
row.DisplayName,
row.ExternalUsername),
row.RsvpStatus,
row.RegistrationStatus,
row.IsGm);
private static PlatformGroup CreateGroup(ConfirmationSessionRow row) =>
new(
ParsePlatform(row.Platform),
row.ExternalGroupId,
row.DisplayName,
row.ExternalChannelId,
row.ThreadId?.ToString(CultureInfo.InvariantCulture));
private static PlatformKind ParsePlatform(string platform) =>
Enum.Parse<PlatformKind>(platform, ignoreCase: true);
private static int? TryGetTelegramMessageId(PlatformMessageRef message) =>
message.Platform == PlatformKind.Telegram &&
int.TryParse(message.ExternalMessageId, NumberStyles.Integer, CultureInfo.InvariantCulture, out var messageId)
? messageId
: null;
private static Task PersistPlatformMessageAsync(
NpgsqlConnection connection,
PlatformMessageRef message,
Guid groupId,
Guid? sessionId,
Guid? batchId,
string purpose) =>
connection.ExecuteAsync(
"""
INSERT INTO platform_messages (
platform,
group_id,
batch_id,
session_id,
external_channel_id,
external_thread_id,
external_message_id,
purpose)
VALUES (
@Platform,
@GroupId,
@BatchId,
@SessionId,
@ExternalChannelId,
@ExternalThreadId,
@ExternalMessageId,
@Purpose)
""",
new
{
Platform = message.Platform.ToString(),
GroupId = groupId,
BatchId = batchId,
SessionId = sessionId,
ExternalChannelId = message.ExternalGroupId,
message.ExternalThreadId,
message.ExternalMessageId,
Purpose = purpose
});
}
@@ -0,0 +1,50 @@
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Logging;
namespace GmRelay.Shared.Features.Notifications;
public sealed class PlatformDirectNotificationSender(
IPlatformMessenger messenger,
ILogger<PlatformDirectNotificationSender> logger)
{
public async Task SendAsync(
PlatformDirectSessionNotificationKind kind,
IEnumerable<PlatformUser> recipients,
Guid sessionId,
string title,
DateTime scheduledAt,
string? joinLink,
string? actorDisplayName,
string? reason,
CancellationToken ct)
{
foreach (var recipient in recipients)
{
try
{
await messenger.SendDirectSessionNotificationAsync(
new PlatformDirectSessionNotification(
kind,
recipient,
sessionId,
title,
scheduledAt,
joinLink,
actorDisplayName,
reason),
ct);
}
catch (Exception ex)
{
logger.LogWarning(
ex,
"Failed to send {NotificationKind} notification for session {SessionId} to {Platform} user {ExternalUserId} ({DisplayName})",
kind,
sessionId,
recipient.Platform,
recipient.ExternalUserId,
recipient.DisplayName);
}
}
}
}
@@ -0,0 +1,6 @@
namespace GmRelay.Shared.Features.Reminders.SendJoinLink;
public interface ISendJoinLinkHandler
{
Task HandleAsync(Guid sessionId, CancellationToken ct);
}
@@ -0,0 +1,228 @@
using System.Globalization;
using Dapper;
using GmRelay.Shared.Domain;
using GmRelay.Shared.Features.Notifications;
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Logging;
using Npgsql;
namespace GmRelay.Shared.Features.Reminders.SendJoinLink;
internal sealed record JoinLinkSessionRow(
Guid Id,
Guid GroupId,
string Title,
string JoinLink,
DateTime ScheduledAt,
string Platform,
string ExternalGroupId,
string DisplayName,
string? ExternalChannelId,
int? ThreadId,
string NotificationMode);
internal sealed record JoinLinkPlayerRow(
string Platform,
string ExternalUserId,
string DisplayName,
string? ExternalUsername,
string RsvpStatus,
string RegistrationStatus,
bool IsGm);
public sealed class SendJoinLinkHandler(
NpgsqlDataSource dataSource,
IPlatformMessenger messenger,
PlatformDirectNotificationSender directSender,
ILogger<SendJoinLinkHandler> logger) : ISendJoinLinkHandler
{
public async Task HandleAsync(Guid sessionId, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var session = await connection.QuerySingleOrDefaultAsync<JoinLinkSessionRow>(
"""
SELECT s.id,
s.group_id AS GroupId,
s.title,
s.join_link AS JoinLink,
s.scheduled_at AS ScheduledAt,
COALESCE(g.platform, 'Telegram') AS Platform,
COALESCE(g.external_group_id, g.telegram_chat_id::TEXT) AS ExternalGroupId,
g.name AS DisplayName,
COALESCE(g.external_channel_id, g.telegram_chat_id::TEXT) AS ExternalChannelId,
s.thread_id AS ThreadId,
s.notification_mode AS NotificationMode
FROM sessions s
JOIN game_groups g ON g.id = s.group_id
WHERE s.id = @SessionId
AND s.status = @Confirmed
AND (
(COALESCE(g.platform, 'Telegram') = 'Telegram' AND s.link_message_id IS NULL)
OR (
COALESCE(g.platform, 'Telegram') <> 'Telegram'
AND NOT EXISTS (
SELECT 1
FROM platform_messages pm
WHERE pm.session_id = s.id
AND pm.platform = COALESCE(g.platform, 'Telegram')
AND pm.purpose = 'join_link'
)
)
)
""",
new { SessionId = sessionId, Confirmed = SessionStatus.Confirmed });
if (session is null)
{
logger.LogWarning("Session {SessionId} not eligible for join link", sessionId);
return;
}
var players = (await connection.QueryAsync<JoinLinkPlayerRow>(
"""
SELECT COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername,
sp.rsvp_status AS RsvpStatus,
sp.registration_status AS RegistrationStatus,
sp.is_gm AS IsGm
FROM session_participants sp
JOIN players p ON p.id = sp.player_id
WHERE sp.session_id = @SessionId
AND sp.rsvp_status = @Confirmed
AND sp.registration_status = @Active
ORDER BY sp.created_at ASC
""",
new
{
SessionId = sessionId,
Confirmed = RsvpStatus.Confirmed,
Active = ParticipantRegistrationStatus.Active
}))
.Select(ToParticipant)
.ToList();
var group = CreateGroup(session);
var message = await messenger.SendJoinLinkNotificationAsync(
new PlatformJoinLinkNotification(
group,
session.Id,
session.Title,
session.ScheduledAt,
session.JoinLink,
players),
ct);
await connection.ExecuteAsync(
"""
UPDATE sessions
SET link_message_id = @MessageId, updated_at = now()
WHERE id = @SessionId
""",
new
{
SessionId = sessionId,
MessageId = TryGetTelegramMessageId(message)
});
await PersistPlatformMessageAsync(
connection,
message,
session.GroupId,
session.Id,
batchId: null,
purpose: "join_link");
var mode = SessionNotificationModeExtensions.FromDatabaseValue(session.NotificationMode);
if (mode.ShouldSendDirectMessages())
{
await directSender.SendAsync(
PlatformDirectSessionNotificationKind.JoinLink,
players.Select(p => p.User),
session.Id,
session.Title,
session.ScheduledAt,
session.JoinLink,
actorDisplayName: null,
reason: null,
ct);
}
logger.LogInformation(
"Join link sent for session {SessionId} ({Title}), platform={Platform}, message_id={MessageId}",
sessionId,
session.Title,
message.Platform,
message.ExternalMessageId);
}
private static PlatformSessionParticipant ToParticipant(JoinLinkPlayerRow row) =>
new(
new PlatformUser(
ParsePlatform(row.Platform),
row.ExternalUserId,
row.DisplayName,
row.ExternalUsername),
row.RsvpStatus,
row.RegistrationStatus,
row.IsGm);
private static PlatformGroup CreateGroup(JoinLinkSessionRow row) =>
new(
ParsePlatform(row.Platform),
row.ExternalGroupId,
row.DisplayName,
row.ExternalChannelId,
row.ThreadId?.ToString(CultureInfo.InvariantCulture));
private static PlatformKind ParsePlatform(string platform) =>
Enum.Parse<PlatformKind>(platform, ignoreCase: true);
private static int? TryGetTelegramMessageId(PlatformMessageRef message) =>
message.Platform == PlatformKind.Telegram &&
int.TryParse(message.ExternalMessageId, NumberStyles.Integer, CultureInfo.InvariantCulture, out var messageId)
? messageId
: null;
private static Task PersistPlatformMessageAsync(
NpgsqlConnection connection,
PlatformMessageRef message,
Guid groupId,
Guid? sessionId,
Guid? batchId,
string purpose) =>
connection.ExecuteAsync(
"""
INSERT INTO platform_messages (
platform,
group_id,
batch_id,
session_id,
external_channel_id,
external_thread_id,
external_message_id,
purpose)
VALUES (
@Platform,
@GroupId,
@BatchId,
@SessionId,
@ExternalChannelId,
@ExternalThreadId,
@ExternalMessageId,
@Purpose)
""",
new
{
Platform = message.Platform.ToString(),
GroupId = groupId,
BatchId = batchId,
SessionId = sessionId,
ExternalChannelId = message.ExternalGroupId,
message.ExternalThreadId,
message.ExternalMessageId,
Purpose = purpose
});
}
@@ -0,0 +1,6 @@
namespace GmRelay.Shared.Features.Reminders.SendOneHourReminder;
public interface ISendOneHourReminderHandler
{
Task HandleAsync(Guid sessionId, CancellationToken ct);
}
@@ -0,0 +1,117 @@
using Dapper;
using GmRelay.Shared.Domain;
using GmRelay.Shared.Features.Notifications;
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Logging;
using Npgsql;
namespace GmRelay.Shared.Features.Reminders.SendOneHourReminder;
internal sealed record OneHourReminderSessionRow(
Guid Id,
string Title,
string JoinLink,
DateTime ScheduledAt,
string NotificationMode);
internal sealed record OneHourReminderRecipientRow(
string Platform,
string ExternalUserId,
string DisplayName,
string? ExternalUsername);
public sealed class SendOneHourReminderHandler(
NpgsqlDataSource dataSource,
PlatformDirectNotificationSender directSender,
ILogger<SendOneHourReminderHandler> logger) : ISendOneHourReminderHandler
{
public async Task HandleAsync(Guid sessionId, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var session = await connection.QuerySingleOrDefaultAsync<OneHourReminderSessionRow>(
"""
SELECT id,
title,
join_link AS JoinLink,
scheduled_at AS ScheduledAt,
notification_mode AS NotificationMode
FROM sessions
WHERE id = @SessionId
AND status IN (@Confirmed, @ConfirmationSent)
AND one_hour_reminder_processed_at IS NULL
""",
new
{
SessionId = sessionId,
Confirmed = SessionStatus.Confirmed,
ConfirmationSent = SessionStatus.ConfirmationSent
});
if (session is null)
{
logger.LogWarning("Session {SessionId} not eligible for one-hour reminder", sessionId);
return;
}
var recipients = (await connection.QueryAsync<OneHourReminderRecipientRow>(
"""
SELECT COALESCE(p.platform, 'Telegram') AS Platform,
COALESCE(p.external_user_id, p.telegram_id::TEXT) AS ExternalUserId,
p.display_name AS DisplayName,
COALESCE(p.external_username, p.telegram_username) AS ExternalUsername
FROM session_participants sp
JOIN players p ON p.id = sp.player_id
WHERE sp.session_id = @SessionId
AND sp.is_gm = false
AND sp.registration_status = @Active
AND sp.rsvp_status != @Declined
""",
new
{
SessionId = sessionId,
Active = ParticipantRegistrationStatus.Active,
Declined = RsvpStatus.Declined
}))
.Select(row => new PlatformUser(
ParsePlatform(row.Platform),
row.ExternalUserId,
row.DisplayName,
row.ExternalUsername))
.ToList();
var mode = SessionNotificationModeExtensions.FromDatabaseValue(session.NotificationMode);
if (mode.ShouldSendDirectMessages() && recipients.Count > 0)
{
await directSender.SendAsync(
PlatformDirectSessionNotificationKind.OneHourReminder,
recipients,
session.Id,
session.Title,
session.ScheduledAt,
session.JoinLink,
actorDisplayName: null,
reason: null,
ct);
}
await connection.ExecuteAsync(
"""
UPDATE sessions
SET one_hour_reminder_processed_at = now(),
updated_at = now()
WHERE id = @SessionId
AND one_hour_reminder_processed_at IS NULL
""",
new { SessionId = sessionId });
logger.LogInformation(
"One-hour reminder processed for session {SessionId} ({Title}) with mode {NotificationMode}",
sessionId,
session.Title,
session.NotificationMode);
}
private static PlatformKind ParsePlatform(string platform) =>
Enum.Parse<PlatformKind>(platform, ignoreCase: true);
}
+1
View File
@@ -11,6 +11,7 @@
<ItemGroup>
<PackageReference Include="Dapper" Version="2.1.72" />
<PackageReference Include="Dapper.AOT" Version="1.0.48" PrivateAssets="all" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" Version="10.0.5" />
<PackageReference Include="Microsoft.Extensions.Logging.Abstractions" Version="10.0.5" />
<PackageReference Include="Npgsql" Version="10.0.2" />
</ItemGroup>
@@ -0,0 +1,109 @@
using Dapper;
using GmRelay.Shared.Domain;
using Npgsql;
namespace GmRelay.Shared.Infrastructure.Scheduling;
public interface ISessionTriggerStore
{
Task<IReadOnlyList<Guid>> GetSessionsNeedingConfirmationAsync(DateTimeOffset now, CancellationToken ct);
Task<IReadOnlyList<Guid>> GetSessionsNeedingOneHourReminderAsync(DateTimeOffset now, CancellationToken ct);
Task<IReadOnlyList<Guid>> GetSessionsNeedingJoinLinkAsync(DateTimeOffset now, CancellationToken ct);
}
public sealed class DbSessionTriggerStore(
NpgsqlDataSource dataSource,
PlatformSchedulerOptions options) : ISessionTriggerStore
{
private static readonly TimeSpan ConfirmationLeadTime = TimeSpan.FromHours(24);
private static readonly TimeSpan OneHourReminderLeadTime = TimeSpan.FromHours(1);
private static readonly TimeSpan JoinLinkLeadTime = TimeSpan.FromMinutes(5);
public async Task<IReadOnlyList<Guid>> GetSessionsNeedingConfirmationAsync(DateTimeOffset now, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var results = await connection.QueryAsync<Guid>(
"""
SELECT s.id
FROM sessions s
JOIN game_groups g ON g.id = s.group_id
WHERE g.platform = @Platform
AND s.status = @Planned
AND s.scheduled_at - @LeadTime <= @Now
AND s.confirmation_sent_at IS NULL
""",
new
{
Platform = options.Platform.ToString(),
Planned = SessionStatus.Planned,
LeadTime = ConfirmationLeadTime,
Now = now.UtcDateTime
});
return results.ToList();
}
public async Task<IReadOnlyList<Guid>> GetSessionsNeedingOneHourReminderAsync(DateTimeOffset now, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var results = await connection.QueryAsync<Guid>(
"""
SELECT s.id
FROM sessions s
JOIN game_groups g ON g.id = s.group_id
WHERE g.platform = @Platform
AND s.status IN (@Confirmed, @ConfirmationSent)
AND s.scheduled_at - @LeadTime <= @Now
AND s.one_hour_reminder_processed_at IS NULL
""",
new
{
Platform = options.Platform.ToString(),
Confirmed = SessionStatus.Confirmed,
ConfirmationSent = SessionStatus.ConfirmationSent,
LeadTime = OneHourReminderLeadTime,
Now = now.UtcDateTime
});
return results.ToList();
}
public async Task<IReadOnlyList<Guid>> GetSessionsNeedingJoinLinkAsync(DateTimeOffset now, CancellationToken ct)
{
await using var connection = await dataSource.OpenConnectionAsync(ct);
var results = await connection.QueryAsync<Guid>(
"""
SELECT s.id
FROM sessions s
JOIN game_groups g ON g.id = s.group_id
WHERE g.platform = @Platform
AND s.status = @Confirmed
AND s.scheduled_at - @LeadTime <= @Now
AND (
(g.platform = 'Telegram' AND s.link_message_id IS NULL)
OR (
g.platform <> 'Telegram'
AND NOT EXISTS (
SELECT 1
FROM platform_messages pm
WHERE pm.session_id = s.id
AND pm.platform = g.platform
AND pm.purpose = 'join_link'
)
)
)
""",
new
{
Platform = options.Platform.ToString(),
Confirmed = SessionStatus.Confirmed,
LeadTime = JoinLinkLeadTime,
Now = now.UtcDateTime
});
return results.ToList();
}
}
@@ -0,0 +1,5 @@
using GmRelay.Shared.Platform;
namespace GmRelay.Shared.Infrastructure.Scheduling;
public sealed record PlatformSchedulerOptions(PlatformKind Platform);
@@ -0,0 +1,139 @@
using GmRelay.Shared.Features.Confirmation.SendConfirmation;
using GmRelay.Shared.Features.Reminders.SendJoinLink;
using GmRelay.Shared.Features.Reminders.SendOneHourReminder;
using GmRelay.Shared.Platform;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace GmRelay.Shared.Infrastructure.Scheduling;
/// <summary>
/// Stateless scheduler: wakes every 60 seconds, queries PostgreSQL for actionable sessions.
/// All state is kept in the database so worker restarts do not lose scheduled work.
/// </summary>
public sealed class SessionSchedulerService(
ISessionTriggerStore triggerStore,
ISendConfirmationHandler confirmationHandler,
ISendOneHourReminderHandler oneHourReminderHandler,
ISendJoinLinkHandler joinLinkHandler,
ISystemClock clock,
ILogger<SessionSchedulerService> logger) : BackgroundService
{
private static readonly TimeSpan TickInterval = TimeSpan.FromMinutes(1);
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
logger.LogInformation("Session scheduler started (interval: {Interval})", TickInterval);
using var timer = new PeriodicTimer(TickInterval);
do
{
try
{
await TickAsync(stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
break;
}
catch (Exception ex)
{
logger.LogError(ex, "Scheduler tick failed, will retry next tick");
}
}
while (await timer.WaitForNextTickAsync(stoppingToken));
logger.LogInformation("Session scheduler stopped");
}
public async Task TickAsync(CancellationToken ct)
{
var now = clock.UtcNow;
await ProcessConfirmationTriggers(now, ct);
await ProcessOneHourReminderTriggers(now, ct);
await ProcessJoinLinkTriggers(now, ct);
}
private async Task ProcessConfirmationTriggers(DateTimeOffset now, CancellationToken ct)
{
IReadOnlyList<Guid> sessionIds;
try
{
sessionIds = await triggerStore.GetSessionsNeedingConfirmationAsync(now, ct);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to query confirmation triggers");
return;
}
foreach (var sessionId in sessionIds)
{
try
{
await confirmationHandler.HandleAsync(sessionId, ct);
logger.LogInformation("Confirmation sent for session {SessionId}", sessionId);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to send confirmation for session {SessionId}", sessionId);
}
}
}
private async Task ProcessOneHourReminderTriggers(DateTimeOffset now, CancellationToken ct)
{
IReadOnlyList<Guid> sessionIds;
try
{
sessionIds = await triggerStore.GetSessionsNeedingOneHourReminderAsync(now, ct);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to query one-hour reminder triggers");
return;
}
foreach (var sessionId in sessionIds)
{
try
{
await oneHourReminderHandler.HandleAsync(sessionId, ct);
logger.LogInformation("One-hour reminder processed for session {SessionId}", sessionId);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to process one-hour reminder for session {SessionId}", sessionId);
}
}
}
private async Task ProcessJoinLinkTriggers(DateTimeOffset now, CancellationToken ct)
{
IReadOnlyList<Guid> sessionIds;
try
{
sessionIds = await triggerStore.GetSessionsNeedingJoinLinkAsync(now, ct);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to query join-link triggers");
return;
}
foreach (var sessionId in sessionIds)
{
try
{
await joinLinkHandler.HandleAsync(sessionId, ct);
logger.LogInformation("Join link sent for session {SessionId}", sessionId);
}
catch (Exception ex)
{
logger.LogError(ex, "Failed to send join link for session {SessionId}", sessionId);
}
}
}
}
@@ -13,4 +13,22 @@ public interface IPlatformMessenger
Task AnswerInteractionAsync(PlatformInteractionReply reply, CancellationToken ct);
Task SendCalendarFileAsync(PlatformCalendarFile file, CancellationToken ct);
Task<PlatformMessageRef> SendConfirmationRequestAsync(PlatformConfirmationRequest request, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support confirmation requests.");
Task UpdateConfirmationRequestAsync(PlatformRsvpMessageUpdate update, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support confirmation request updates.");
Task<PlatformMessageRef> SendJoinLinkNotificationAsync(PlatformJoinLinkNotification notification, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support join-link notifications.");
Task SendDirectSessionNotificationAsync(PlatformDirectSessionNotification notification, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support direct session notifications.");
Task SendRsvpOutcomeAsync(PlatformRsvpOutcomeNotification notification, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support RSVP outcome notifications.");
Task UpdateRescheduleVoteAsync(PlatformRescheduleVoteUpdate update, CancellationToken ct) =>
throw new NotSupportedException("This platform messenger does not support reschedule vote updates.");
}
@@ -1,3 +1,4 @@
using GmRelay.Shared.Features.Sessions.RescheduleSession;
using GmRelay.Shared.Rendering;
namespace GmRelay.Shared.Platform;
@@ -34,3 +35,81 @@ public sealed record PlatformCalendarFile(
byte[] Content,
string CaptionHtml,
IReadOnlyList<PlatformMessageAction> Actions);
public sealed record PlatformSessionParticipant(
PlatformUser User,
string RsvpStatus,
string RegistrationStatus,
bool IsGm = false);
public sealed record PlatformConfirmationRequest(
PlatformGroup Group,
Guid SessionId,
string Title,
DateTime ScheduledAt,
IReadOnlyList<PlatformSessionParticipant> Participants,
PlatformMessageRef? ExistingMessage = null);
public sealed record PlatformJoinLinkNotification(
PlatformGroup Group,
Guid SessionId,
string Title,
DateTime ScheduledAt,
string JoinLink,
IReadOnlyList<PlatformSessionParticipant> ConfirmedPlayers,
PlatformMessageRef? ExistingMessage = null);
public enum PlatformDirectSessionNotificationKind
{
ConfirmationRequest = 0,
OneHourReminder = 1,
JoinLink = 2,
RsvpAllConfirmed = 3,
RsvpDeclined = 4,
RescheduleApproved = 5,
RescheduleRejected = 6
}
public sealed record PlatformDirectSessionNotification(
PlatformDirectSessionNotificationKind Kind,
PlatformUser Recipient,
Guid SessionId,
string Title,
DateTime ScheduledAt,
string? JoinLink = null,
string? ActorDisplayName = null,
string? Reason = null);
public sealed record PlatformRsvpMessageUpdate(
PlatformConfirmationRequest Request,
bool DisableActions);
public enum PlatformRsvpOutcomeKind
{
GroupAllConfirmed = 0,
GmAllConfirmed = 1,
GmPlayerDeclined = 2
}
public sealed record PlatformRsvpOutcomeNotification(
PlatformRsvpOutcomeKind Kind,
PlatformGroup? Group,
IReadOnlyList<PlatformUser> Recipients,
Guid SessionId,
string Title,
DateTime ScheduledAt,
string? ActorDisplayName = null);
public sealed record PlatformRescheduleVoteUpdate(
PlatformGroup Group,
PlatformMessageRef ExistingMessage,
Guid ProposalId,
Guid SessionId,
string Title,
DateTime CurrentScheduledAt,
DateTimeOffset VotingDeadlineAt,
RescheduleVoteDecision Decision,
RescheduleOptionDto? SelectedOption,
IReadOnlyList<RescheduleOptionDto> Options,
IReadOnlyList<RescheduleOptionVoteDto> Votes,
IReadOnlyList<VoteParticipantDto> Participants);
+52
View File
@@ -14,6 +14,19 @@
"resolved": "1.0.48",
"contentHash": "rsLM3yKr4g+YKKox9lhc8D+kz67P7Q9+xdyn1LmCsoYr1kYpJSm+Nt6slo5UrfUrcTiGJ57zUlyO8XUdV7G7iA=="
},
"Microsoft.Extensions.Hosting.Abstractions": {
"type": "Direct",
"requested": "[10.0.5, )",
"resolved": "10.0.5",
"contentHash": "+Wb7KAMVZTomwJkQrjuPTe5KBzGod7N8XeG+ScxRlkPOB4sZLG4ccVwjV4Phk5BCJt7uIMnGHVoN6ZMVploX+g==",
"dependencies": {
"Microsoft.Extensions.Configuration.Abstractions": "10.0.5",
"Microsoft.Extensions.DependencyInjection.Abstractions": "10.0.5",
"Microsoft.Extensions.Diagnostics.Abstractions": "10.0.5",
"Microsoft.Extensions.FileProviders.Abstractions": "10.0.5",
"Microsoft.Extensions.Logging.Abstractions": "10.0.5"
}
},
"Microsoft.Extensions.Logging.Abstractions": {
"type": "Direct",
"requested": "[10.0.5, )",
@@ -38,10 +51,49 @@
"resolved": "5.6.7",
"contentHash": "WIE9RJswdSc2j+rLz2gW6U+gMUjMHzY2j7C/CL8/R2olXNM/+twarfMnWqm+rZodDBvaYDApJyxM8mVYf9FGrQ=="
},
"Microsoft.Extensions.Configuration.Abstractions": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "P09QpTHjqHmCLQOTC+WyLkoRNxek4NIvfWt+TnU0etoDUSRxcltyd6+j/ouRbMdLR0j44GqGO+lhI2M4fAHG4g==",
"dependencies": {
"Microsoft.Extensions.Primitives": "10.0.5"
}
},
"Microsoft.Extensions.DependencyInjection.Abstractions": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "iVMtq9eRvzyhx8949EGT0OCYJfXi737SbRVzWXE5GrOgGj5AaZ9eUuxA/BSUfmOMALKn/g8KfFaNQw0eiB3lyA=="
},
"Microsoft.Extensions.Diagnostics.Abstractions": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "/nYGrpa9/0BZofrVpBbbj+Ns8ZesiPE0V/KxsuHgDgHQopIzN54nRaQGSuvPw16/kI9sW1Zox5yyAPqvf0Jz6A==",
"dependencies": {
"Microsoft.Extensions.DependencyInjection.Abstractions": "10.0.5",
"Microsoft.Extensions.Options": "10.0.5"
}
},
"Microsoft.Extensions.FileProviders.Abstractions": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "nCBmCx0Xemlu65ZiWMcXbvfvtznKxf4/YYKF9R28QkqdI9lTikedGqzJ28/xmdGGsxUnsP5/3TQGpiPwVjK0dA==",
"dependencies": {
"Microsoft.Extensions.Primitives": "10.0.5"
}
},
"Microsoft.Extensions.Options": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "MDaQMdUplw0AIRhWWmbLA7yQEXaLIHb+9CTroTiNS8OlI0LMXS4LCxtopqauiqGCWlRgJ+xyraVD8t6veRAFbw==",
"dependencies": {
"Microsoft.Extensions.DependencyInjection.Abstractions": "10.0.5",
"Microsoft.Extensions.Primitives": "10.0.5"
}
},
"Microsoft.Extensions.Primitives": {
"type": "Transitive",
"resolved": "10.0.5",
"contentHash": "/HUHJ0tw/LQvD0DZrz50eQy/3z7PfX7WWEaXnjKTV9/TNdcgFlNTZGo49QhS7PTmhDqMyHRMqAXSBxLh0vso4g=="
}
}
}