feature: Add local subscription notification worker
This commit is contained in:
@@ -9,7 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.future import select
|
||||
from sqlalchemy.orm import selectinload
|
||||
|
||||
from db.models import Subscription
|
||||
from db.models import Subscription, SubscriptionNotification
|
||||
|
||||
INSTALL_SHARE_TOKEN_BYTES = 16
|
||||
|
||||
@@ -318,6 +318,45 @@ async def update_subscription_notification_time(
|
||||
)
|
||||
|
||||
|
||||
async def has_subscription_notification(
|
||||
session: AsyncSession,
|
||||
subscription_id: int,
|
||||
notification_key: str,
|
||||
) -> bool:
|
||||
stmt = (
|
||||
select(SubscriptionNotification.notification_id)
|
||||
.where(
|
||||
SubscriptionNotification.subscription_id == subscription_id,
|
||||
SubscriptionNotification.notification_key == notification_key,
|
||||
)
|
||||
.limit(1)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
return result.scalar_one_or_none() is not None
|
||||
|
||||
|
||||
async def record_subscription_notification(
|
||||
session: AsyncSession,
|
||||
subscription_id: int,
|
||||
notification_key: str,
|
||||
*,
|
||||
sent_at: Optional[datetime] = None,
|
||||
) -> None:
|
||||
if sent_at is None:
|
||||
sent_at = datetime.now(timezone.utc)
|
||||
existing = await has_subscription_notification(session, subscription_id, notification_key)
|
||||
if existing:
|
||||
return
|
||||
session.add(
|
||||
SubscriptionNotification(
|
||||
subscription_id=subscription_id,
|
||||
notification_key=notification_key,
|
||||
sent_at=sent_at,
|
||||
)
|
||||
)
|
||||
await update_subscription_notification_time(session, subscription_id, sent_at)
|
||||
|
||||
|
||||
async def find_subscription_for_notification_update(
|
||||
session: AsyncSession, user_id: int, subscription_end_date_to_match: datetime
|
||||
) -> Optional[Subscription]:
|
||||
|
||||
@@ -18,6 +18,7 @@ from ..models import (
|
||||
Payment,
|
||||
PromoCodeActivation,
|
||||
Subscription,
|
||||
SubscriptionNotification,
|
||||
SupportTicket,
|
||||
SupportTicketMessage,
|
||||
TariffChange,
|
||||
@@ -779,6 +780,11 @@ async def delete_user_and_relations(session: AsyncSession, user_id: int) -> bool
|
||||
await session.execute(
|
||||
delete(TrafficWarning).where(TrafficWarning.subscription_id.in_(subscription_ids))
|
||||
)
|
||||
await session.execute(
|
||||
delete(SubscriptionNotification).where(
|
||||
SubscriptionNotification.subscription_id.in_(subscription_ids)
|
||||
)
|
||||
)
|
||||
await session.execute(
|
||||
delete(SupportTicketMessage).where(SupportTicketMessage.ticket_id.in_(support_ticket_ids))
|
||||
)
|
||||
|
||||
@@ -1012,6 +1012,41 @@ def _migration_0030_add_hwid_pricing_metadata(connection: Connection) -> None:
|
||||
)
|
||||
|
||||
|
||||
def _migration_0031_add_subscription_notifications(connection: Connection) -> None:
|
||||
connection.execute(
|
||||
text(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS subscription_notifications (
|
||||
notification_id SERIAL PRIMARY KEY,
|
||||
subscription_id INTEGER NOT NULL REFERENCES subscriptions(subscription_id),
|
||||
notification_key VARCHAR(64) NOT NULL,
|
||||
sent_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
CONSTRAINT uq_subscription_notification_key UNIQUE (
|
||||
subscription_id,
|
||||
notification_key
|
||||
)
|
||||
)
|
||||
"""
|
||||
)
|
||||
)
|
||||
connection.execute(
|
||||
text(
|
||||
"""
|
||||
CREATE INDEX IF NOT EXISTS ix_subscription_notifications_subscription_id
|
||||
ON subscription_notifications (subscription_id)
|
||||
"""
|
||||
)
|
||||
)
|
||||
connection.execute(
|
||||
text(
|
||||
"""
|
||||
CREATE INDEX IF NOT EXISTS ix_subscription_notifications_notification_key
|
||||
ON subscription_notifications (notification_key)
|
||||
"""
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
MIGRATIONS: List[Migration] = [
|
||||
Migration(
|
||||
id="0001_add_channel_subscription_fields",
|
||||
@@ -1174,6 +1209,11 @@ MIGRATIONS: List[Migration] = [
|
||||
description="Persist quoted HWID top-up pricing windows and conversion audit",
|
||||
upgrade=_migration_0030_add_hwid_pricing_metadata,
|
||||
),
|
||||
Migration(
|
||||
id="0031_add_subscription_notifications",
|
||||
description="Track sent subscription notification stages",
|
||||
upgrade=_migration_0031_add_subscription_notifications,
|
||||
),
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -277,6 +277,26 @@ class TrafficWarning(Base):
|
||||
subscription = relationship("Subscription")
|
||||
|
||||
|
||||
class SubscriptionNotification(Base):
|
||||
__tablename__ = "subscription_notifications"
|
||||
__table_args__ = (
|
||||
UniqueConstraint(
|
||||
"subscription_id",
|
||||
"notification_key",
|
||||
name="uq_subscription_notification_key",
|
||||
),
|
||||
)
|
||||
|
||||
notification_id = Column(Integer, primary_key=True, autoincrement=True)
|
||||
subscription_id = Column(
|
||||
Integer, ForeignKey("subscriptions.subscription_id"), nullable=False, index=True
|
||||
)
|
||||
notification_key = Column(String(64), nullable=False, index=True)
|
||||
sent_at = Column(DateTime(timezone=True), server_default=func.now())
|
||||
|
||||
subscription = relationship("Subscription")
|
||||
|
||||
|
||||
class TariffChange(Base):
|
||||
__tablename__ = "tariff_changes"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user