Compare commits

...

3 Commits

Author SHA1 Message Date
Fringg b377c52fd9 feat(tasks): include target_meta/reward_meta in TaskListItem response
Cabinet admin needs target_meta.tariff_id to render the tariff badge in
the tasks list. Adds target_meta and reward_meta fields to the compact
TaskListItem schema (defaults to {} for backward compatibility).
2026-05-06 09:03:41 +03:00
Fringg 7bbf9fbc63 feat: tasks/rewards system with 9 task types and 2 reward types
Adds a comprehensive tasks system where users earn rewards for completing
configurable goals: purchase_tariff, subscribe_channel, traffic_used,
referrals_invited, purchase_period, spend_amount, multi_tariff,
gift_purchased, gifts_count.

Reward types: balance kopeks or subscription_days (with optional fallback
to Tariff.bonus_days_per_purchase).

Highlights:
- Multi-level chains via parent_task_id (cycle detection up to 10 hops)
- Promo group + audience (telegram/email/both) filters
- Trial users blocked from progress/claim via _user_eligible
- Calendar-month windowing for traffic_used (UTC, per-user aggregate)
- BigInteger for monetary fields, UNIQUE(user_id, task_id) constraint
- FOR UPDATE locks on progress/user/subscription for atomic claim
- chosen_reward_type gated by task.allow_user_choice
- Atomic upsert via PG ON CONFLICT DO NOTHING (SQLite fallback via savepoint)
- Triggers wired into: cabinet submit_purchase, bot tariff_purchase (custom +
  preset), subscription_renewal, 5 auto-purchase paths, daily_subscription,
  channel_member ChatMemberUpdated, gift purchase, referral first_topup,
  spend_amount in create_transaction commit=True path, traffic sync
  (multi-tariff + single-tariff)
- Admin CRUD endpoints for tasks + partner channels (tasks:read/edit perms)
- Adds Tariff.bonus_days_per_purchase column

Migration: 0075_create_tasks.py
2026-05-06 08:33:46 +03:00
Fringg 31e3ccd24c docs: add Antilopay/Etoplatezhi/Jupiter/Donut/Lava + Apple IAP to provider list
- Bump provider counter 18 → 24+ in Features and Documentation sections
- Add 7 new rows to providers table:
  - Antilopay (RSA signing)
  - Etoplatezhi
  - Jupiter (FPGate P2P) — partner via @k_juppiter
  - Donut (Donut P2P) — partner via @donut_payment
  - Lava Business
  - Apple In-App Purchase
- Add 2 partner cards (Jupiter, Donut) following the Platega/PayPear pattern
2026-05-04 21:20:44 +03:00
25 changed files with 2508 additions and 7 deletions
+34 -2
View File
@@ -55,7 +55,7 @@ Bedolaga — полнофункциональная платформа для п
### 💳 Платежи
- 🏦 **18 платёжных провайдеров** одновременно
- 🏦 **24+ платёжных провайдера** одновременно
- 💰 Единый баланс: пополнение любым способом → покупка с баланса
- ⚡ Автопокупка подписки после пополнения
- 💾 Рекуррентные платежи (сохранённые карты)
@@ -127,6 +127,12 @@ Bedolaga — полнофункциональная платформа для п
| 🤝 | **[RollyPay](https://rollypay.io/?utm_source=bedolaga&utm_medium=community&utm_campaign=integration)** 🔸 | СБП, карты, крипто | RUB → USDT |
| 🤝 | **[AuraPay](https://aurapay.tech/)** 🔸 | Карты, СБП | RUB |
| 🤝 | **[Overpay](https://overpay.pro/)** 🔸 | Карты, СБП | RUB |
| 🦌 | **Antilopay** | Карты, СБП, SberPay (RSA подпись) | RUB |
| 💳 | **Etoplatezhi** | Карты, СБП | RUB |
| 🪐 | **[Jupiter](https://t.me/k_juppiter)** 🔸 | СБП через QR (FPGate P2P v2.1) | RUB |
| 🍩 | **[Donut](https://t.me/donut_payment)** 🔸 | Карты, СБП по телефону, СБП QR (P2P) | RUB |
| 🌋 | **Lava Business** | Карты, СБП (gate.lava.ru) | RUB |
| 🍎 | **Apple In-App Purchase** | Покупки через iOS App Store | USD |
| 📲 | **Tribute** | Telegram-платежи | RUB |
</div>
@@ -210,6 +216,32 @@ Bedolaga — официальный партнёр платёжного шлюз
📩 Менеджер: [@A_OverPay](https://t.me/A_OverPay) | 🌐 [overpay.pro](https://overpay.pro/)
</td>
</tr>
<tr>
<td align="center">
**🤝 Официальный партнёр Jupiter (FPGate P2P)**
Bedolaga — официальный партнёр платёжного шлюза **Jupiter** (FPGate P2P v2.1).<br>
Эквайринг СБП через QR-код банковского приложения, HMAC-SHA256 подпись.<br>
Высокая проходимость, callback-driven архитектура, защита от replay-атак.<br>
Подключение по кодовому слову **`БЕДОЛАГА`** — **спец. условия**
📩 Менеджер: [@k_juppiter](https://t.me/k_juppiter)
</td>
<td align="center">
**🤝 Официальный партнёр Donut**
Bedolaga — официальный партнёр платёжной системы **Donut** (Donut P2P).<br>
P2P-оплата картой, СБП по номеру телефона и СБП QR — три метода через единый API.<br>
HMAC-SHA256 подпись, sticky terminal-status guard, защита от amount tampering.<br>
Подключение по кодовому слову **`БЕДОЛАГА`** — **спец. условия**
📩 Менеджер: [@donut_payment](https://t.me/donut_payment)
</td>
</tr>
</table>
@@ -275,7 +307,7 @@ docker compose up -d
| | Раздел | Описание |
|:---:|:---|:---|
| 🚀 | [Быстрый старт](https://docs.bedolagam.ru/getting-started/quickstart) | Развёртывание за 5 минут |
| 💳 | [Настройка платежей](https://docs.bedolagam.ru/bot/payments) | 18 провайдеров, webhook, фискализация |
| 💳 | [Настройка платежей](https://docs.bedolagam.ru/bot/payments) | 24+ провайдера, webhook, фискализация, Apple IAP |
| 📦 | [Подписки и тарифы](https://docs.bedolagam.ru/bot/subscriptions) | Конфигурация планов и трафика |
| 👥 | [Реферальная программа](https://docs.bedolagam.ru/bot/referral-program) | Партнёрка и вывод средств |
| 🖥 | [Cabinet](https://docs.bedolagam.ru/cabinet/overview) | Настройка веб-кабинета |
+4
View File
@@ -34,6 +34,7 @@ from .admin_servers import router as admin_servers_router
from .admin_settings import router as admin_settings_router
from .admin_stats import router as admin_stats_router
from .admin_tariffs import router as admin_tariffs_router
from .admin_tasks import router as admin_tasks_router
from .admin_tickets import router as admin_tickets_router
from .admin_traffic import router as admin_traffic_router
from .admin_updates import router as admin_updates_router
@@ -64,6 +65,7 @@ from .ticket_notifications import (
router as ticket_notifications_router,
)
from .tickets import router as tickets_router
from .user_tasks import router as user_tasks_router
from .websocket import router as websocket_router
from .wheel import router as wheel_router
from .withdrawal import router as withdrawal_router
@@ -109,6 +111,7 @@ router.include_router(landing_router)
router.include_router(media_router)
router.include_router(news_router)
router.include_router(info_pages_router)
router.include_router(user_tasks_router)
# Wheel routes
router.include_router(wheel_router)
@@ -158,6 +161,7 @@ router.include_router(admin_news_tags_router)
router.include_router(admin_news_media_router)
router.include_router(admin_news_router)
router.include_router(admin_info_pages_router)
router.include_router(admin_tasks_router)
# WebSocket route
router.include_router(websocket_router)
+7
View File
@@ -273,6 +273,8 @@ async def get_tariff(
external_squad_uuid=tariff.external_squad_uuid,
# Показывать в подарках
show_in_gift=tariff.show_in_gift,
# Бонусные дни Tasks
bonus_days_per_purchase=getattr(tariff, 'bonus_days_per_purchase', 0) or 0,
created_at=tariff.created_at,
updated_at=tariff.updated_at,
)
@@ -331,6 +333,8 @@ async def create_new_tariff(
external_squad_uuid=request.external_squad_uuid,
# Показывать в подарках
show_in_gift=request.show_in_gift,
# Бонусные дни Tasks
bonus_days_per_purchase=request.bonus_days_per_purchase,
)
logger.info('Admin created tariff', admin_id=admin.id, tariff_id=tariff.id, tariff_name=tariff.name)
@@ -430,6 +434,9 @@ async def update_existing_tariff(
# Показывать в подарках
if request.show_in_gift is not None:
updates['show_in_gift'] = request.show_in_gift
# Бонусные дни Tasks
if request.bonus_days_per_purchase is not None:
updates['bonus_days_per_purchase'] = request.bonus_days_per_purchase
if updates:
await update_tariff(db, tariff, **updates)
+207
View File
@@ -0,0 +1,207 @@
"""Admin endpoints для системы заданий с наградами и партнёрских каналов."""
from __future__ import annotations
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.ext.asyncio import AsyncSession
from app.cabinet.dependencies import get_cabinet_db, require_permission
from app.cabinet.schemas.tasks import (
TaskCreateRequest,
TaskListItem,
TaskPartnerChannelCreateRequest,
TaskPartnerChannelResponse,
TaskPartnerChannelUpdateRequest,
TaskResponse,
TaskUpdateRequest,
)
from app.database.crud import tasks as tasks_crud
from app.database.models import User
router = APIRouter(prefix='/admin', tags=['Cabinet Admin Tasks'])
# ===========================================================================
# Tasks
# ===========================================================================
@router.get('/tasks', response_model=list[TaskListItem])
async def admin_list_tasks(
include_inactive: bool = True,
parent_task_id: int | None = None,
admin: User = Depends(require_permission('tasks:read')),
db: AsyncSession = Depends(get_cabinet_db),
):
tasks = await tasks_crud.list_tasks(db, include_inactive=include_inactive, parent_task_id=parent_task_id)
return [TaskListItem.model_validate(t) for t in tasks]
@router.get('/tasks/{task_id}', response_model=TaskResponse)
async def admin_get_task(
task_id: int,
admin: User = Depends(require_permission('tasks:read')),
db: AsyncSession = Depends(get_cabinet_db),
):
task = await tasks_crud.get_task_by_id(db, task_id)
if task is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='task_not_found')
return TaskResponse.model_validate(task)
@router.post('/tasks', response_model=TaskResponse, status_code=status.HTTP_201_CREATED)
async def admin_create_task(
request: TaskCreateRequest,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
if request.parent_task_id is not None:
parent = await tasks_crud.get_task_by_id(db, request.parent_task_id)
if parent is None:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='parent_task_not_found')
task = await tasks_crud.create_task(
db,
title=request.title,
description=request.description,
task_type=request.task_type,
reward_type=request.reward_type,
target_value=request.target_value,
reward_value=request.reward_value,
target_meta=request.target_meta,
reward_meta=request.reward_meta,
icon=request.icon,
is_active=request.is_active,
sort_order=request.sort_order,
allow_user_choice=request.allow_user_choice,
user_audience=request.user_audience,
promo_group_id=request.promo_group_id,
parent_task_id=request.parent_task_id,
level=request.level,
starts_at=request.starts_at,
ends_at=request.ends_at,
)
return TaskResponse.model_validate(task)
async def _parent_chain_has_cycle(db: AsyncSession, *, task_id: int, parent_id: int, max_depth: int = 10) -> bool:
"""Идёт вверх по цепочке parent — проверяет, не возвращается ли в task_id."""
current = parent_id
visited: set[int] = set()
for _ in range(max_depth):
if current == task_id:
return True
if current in visited:
return False
visited.add(current)
parent = await tasks_crud.get_task_by_id(db, current)
if parent is None or parent.parent_task_id is None:
return False
current = parent.parent_task_id
return False # max_depth достигнут — дальше не считаем циклом
@router.put('/tasks/{task_id}', response_model=TaskResponse)
async def admin_update_task(
task_id: int,
request: TaskUpdateRequest,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
task = await tasks_crud.get_task_by_id(db, task_id)
if task is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='task_not_found')
if request.parent_task_id is not None and request.parent_task_id == task_id:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='parent_task_cannot_be_self')
if request.parent_task_id is not None:
if await _parent_chain_has_cycle(db, task_id=task_id, parent_id=request.parent_task_id):
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='parent_chain_cycle_detected')
fields = request.model_dump(exclude_unset=True)
if not fields:
return TaskResponse.model_validate(task)
updated = await tasks_crud.update_task(db, task, **fields)
return TaskResponse.model_validate(updated)
@router.delete('/tasks/{task_id}', status_code=status.HTTP_204_NO_CONTENT)
async def admin_delete_task(
task_id: int,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
task = await tasks_crud.get_task_by_id(db, task_id)
if task is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='task_not_found')
await tasks_crud.delete_task(db, task)
# ===========================================================================
# Partner channels
# ===========================================================================
@router.get('/task-partner-channels', response_model=list[TaskPartnerChannelResponse])
async def admin_list_partner_channels(
include_inactive: bool = True,
admin: User = Depends(require_permission('tasks:read')),
db: AsyncSession = Depends(get_cabinet_db),
):
channels = await tasks_crud.list_partner_channels(db, include_inactive=include_inactive)
return [TaskPartnerChannelResponse.model_validate(c) for c in channels]
@router.post(
'/task-partner-channels',
response_model=TaskPartnerChannelResponse,
status_code=status.HTTP_201_CREATED,
)
async def admin_create_partner_channel(
request: TaskPartnerChannelCreateRequest,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
existing = await tasks_crud.get_partner_channel_by_channel_id(db, request.channel_id)
if existing is not None:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='channel_id_already_exists')
channel = await tasks_crud.create_partner_channel(
db,
channel_id=request.channel_id,
title=request.title,
channel_link=request.channel_link,
description=request.description,
is_active=request.is_active,
sort_order=request.sort_order,
)
return TaskPartnerChannelResponse.model_validate(channel)
@router.put('/task-partner-channels/{channel_pk}', response_model=TaskPartnerChannelResponse)
async def admin_update_partner_channel(
channel_pk: int,
request: TaskPartnerChannelUpdateRequest,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
channel = await tasks_crud.get_partner_channel_by_id(db, channel_pk)
if channel is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='channel_not_found')
fields = request.model_dump(exclude_unset=True)
updated = await tasks_crud.update_partner_channel(db, channel, **fields)
return TaskPartnerChannelResponse.model_validate(updated)
@router.delete('/task-partner-channels/{channel_pk}', status_code=status.HTTP_204_NO_CONTENT)
async def admin_delete_partner_channel(
channel_pk: int,
admin: User = Depends(require_permission('tasks:edit')),
db: AsyncSession = Depends(get_cabinet_db),
):
channel = await tasks_crud.get_partner_channel_by_id(db, channel_pk)
if channel is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='channel_not_found')
await tasks_crud.delete_partner_channel(db, channel)
+30
View File
@@ -517,6 +517,36 @@ async def create_gift_purchase(
description=tx_description,
)
# Tasks: триггерим прогресс по подаркам
try:
from app.database.models import TaskType as _TaskType
from app.services.tasks_service import record_event as _record_event
await _record_event(
db,
user_id=user.id,
event_type=_TaskType.GIFT_PURCHASED,
payload={'purchase_id': purchase.id},
)
await _record_event(
db,
user_id=user.id,
event_type=_TaskType.GIFTS_COUNT,
payload={'purchase_id': purchase.id},
)
# record_event делает только flush(); коммитим явно. Для has_recipient=True далее
# fulfill_purchase сделает свой commit, для has_recipient=False — это единственный
# шанс закоммитить task-прогресс перед return.
await db.commit()
except Exception as task_err:
# Сессия может быть в poisoned state — откатываем, чтобы fulfill_purchase ниже
# мог продолжить работу с сессией.
try:
await db.rollback()
except Exception:
pass
logger.warning('Tasks: ошибка GIFT триггеров', user_id=user.id, error=task_err)
# Capture token before fulfill_purchase — session state may change after rollback inside fulfill
purchase_token = purchase.token
+136
View File
@@ -0,0 +1,136 @@
"""User-side endpoints для системы заданий с наградами."""
from __future__ import annotations
import structlog
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.ext.asyncio import AsyncSession
from app.cabinet.dependencies import get_cabinet_db, get_current_cabinet_user
from app.cabinet.schemas.tasks import (
ClaimRewardRequest,
ClaimRewardResponse,
UserTaskProgressResponse,
UserTasksAvailabilityResponse,
UserTasksListResponse,
)
from app.database.models import User
from app.services import tasks_service
logger = structlog.get_logger(__name__)
router = APIRouter(prefix='/tasks', tags=['Cabinet Tasks'])
@router.get('/availability', response_model=UserTasksAvailabilityResponse)
async def get_tasks_availability(
user: User = Depends(get_current_cabinet_user),
db: AsyncSession = Depends(get_cabinet_db),
):
"""Краткая информация для условного показа вкладки «Задания» в меню."""
visible = await tasks_service.get_available_tasks_for_user(db, user)
has_available = len(visible) > 0
unclaimed = await tasks_service.count_completed_unclaimed(db, user_id=user.id)
return UserTasksAvailabilityResponse(
has_available_tasks=has_available,
unclaimed_count=unclaimed,
)
@router.get('', response_model=UserTasksListResponse)
async def list_my_tasks(
user: User = Depends(get_current_cabinet_user),
db: AsyncSession = Depends(get_cabinet_db),
):
"""Список доступных заданий пользователя с их прогрессом."""
visible = await tasks_service.get_available_tasks_for_user(db, user)
items: list[UserTaskProgressResponse] = []
unclaimed_count = 0
for task, progress in visible:
current_value = progress.current_value if progress else 0
is_completed = progress.completed_at is not None if progress else False
is_claimed = progress.claimed_at is not None if progress else False
if is_completed and not is_claimed:
unclaimed_count += 1
percent = (
int(min(current_value, task.target_value) / max(task.target_value, 1) * 100)
if task.target_value
else 0
)
items.append(
UserTaskProgressResponse(
task_id=task.id,
title=task.title or {},
description=task.description or {},
icon=task.icon,
task_type=task.task_type,
target_value=task.target_value,
target_meta=task.target_meta or {},
reward_type=task.reward_type,
reward_value=task.reward_value,
reward_meta=task.reward_meta or {},
allow_user_choice=task.allow_user_choice,
level=task.level,
parent_task_id=task.parent_task_id,
current_value=current_value,
percent=percent,
is_completed=is_completed,
is_claimed=is_claimed,
completed_at=progress.completed_at if progress else None,
claimed_at=progress.claimed_at if progress else None,
reward_granted_meta=progress.reward_granted_meta if progress else None,
)
)
return UserTasksListResponse(
items=items,
has_unclaimed=unclaimed_count > 0,
unclaimed_count=unclaimed_count,
)
@router.post('/{task_id}/claim', response_model=ClaimRewardResponse)
async def claim_task_reward(
task_id: int,
request: ClaimRewardRequest,
user: User = Depends(get_current_cabinet_user),
db: AsyncSession = Depends(get_cabinet_db),
):
"""Получить награду за выполненное задание."""
try:
granted = await tasks_service.claim_reward(
db,
user_id=user.id,
task_id=task_id,
chosen_subscription_id=request.chosen_subscription_id,
chosen_reward_type=request.chosen_reward_type,
)
except ValueError as exc:
msg = str(exc)
# Маппим внутренние коды на HTTP-статусы
not_found = {'progress_not_found', 'task_not_found', 'user_not_found'}
bad_request = {
'not_completed',
'already_claimed',
'user_not_eligible',
'user_choice_not_allowed',
'no_paid_subscription',
'no_subscription_with_target_tariff',
'chosen_subscription_invalid',
'need_choose_subscription',
'invalid_reward_amount',
'invalid_reward_days',
}
if msg in not_found:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=msg) from exc
if msg in bad_request or msg.startswith('unknown_reward_type'):
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=msg) from exc
logger.exception('claim_reward unexpected error', error=msg)
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail='internal_error'
) from exc
return ClaimRewardResponse(success=True, reward=granted)
+6
View File
@@ -117,6 +117,8 @@ class TariffDetailResponse(BaseModel):
external_squad_uuid: str | None = None
# Показывать в подарках
show_in_gift: bool = True
# Бонусные дни для системы Tasks (при награде subscription_days)
bonus_days_per_purchase: int = 0
created_at: datetime
updated_at: datetime | None = None
@@ -175,6 +177,8 @@ class TariffCreateRequest(BaseModel):
external_squad_uuid: str | None = Field(None, pattern=UUID_PATTERN)
# Показывать в подарках
show_in_gift: bool = True
# Бонусные дни для Tasks (subscription_days reward)
bonus_days_per_purchase: int = Field(0, ge=0)
class TariffUpdateRequest(BaseModel):
@@ -216,6 +220,8 @@ class TariffUpdateRequest(BaseModel):
external_squad_uuid: str | None = Field(None, pattern=UUID_PATTERN)
# Показывать в подарках
show_in_gift: bool | None = None
# Бонусные дни для Tasks (subscription_days reward)
bonus_days_per_purchase: int | None = Field(None, ge=0)
class TariffSortOrderRequest(BaseModel):
+316
View File
@@ -0,0 +1,316 @@
"""Pydantic schemas для системы заданий с наградами."""
from __future__ import annotations
from datetime import datetime
from typing import Any, Literal
from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
# ---------------------------------------------------------------------------
# Task partner channels (admin)
# ---------------------------------------------------------------------------
class TaskPartnerChannelBase(BaseModel):
channel_id: str = Field(min_length=1, max_length=100)
title: str = Field(min_length=1, max_length=255)
channel_link: str | None = Field(default=None, max_length=500)
description: str | None = None
is_active: bool = True
sort_order: int = 0
class TaskPartnerChannelCreateRequest(TaskPartnerChannelBase):
pass
class TaskPartnerChannelUpdateRequest(BaseModel):
title: str | None = Field(default=None, min_length=1, max_length=255)
channel_link: str | None = Field(default=None, max_length=500)
description: str | None = None
is_active: bool | None = None
sort_order: int | None = None
class TaskPartnerChannelResponse(TaskPartnerChannelBase):
id: int
created_at: datetime
updated_at: datetime | None = None
model_config = ConfigDict(from_attributes=True)
# ---------------------------------------------------------------------------
# Tasks (admin)
# ---------------------------------------------------------------------------
TASK_TYPES: tuple[str, ...] = (
'purchase_tariff',
'subscribe_channel',
'traffic_used',
'referrals_invited',
'purchase_period',
'spend_amount',
'multi_tariff',
'gift_purchased',
'gifts_count',
)
REWARD_TYPES: tuple[str, ...] = ('balance', 'subscription_days')
USER_AUDIENCES: tuple[str, ...] = ('telegram', 'email', 'both')
class TaskCreateRequest(BaseModel):
"""Создание шаблона задания."""
title: dict[str, str] = Field(..., description='i18n: { "ru": "...", "en": "..." }')
description: dict[str, str] = Field(default_factory=dict)
icon: str | None = None
is_active: bool = True
sort_order: int = 0
task_type: str = Field(...)
target_value: int = Field(default=1, ge=1)
target_meta: dict[str, Any] = Field(default_factory=dict)
reward_type: str = Field(...)
reward_value: int = Field(default=0, ge=0)
reward_meta: dict[str, Any] = Field(default_factory=dict)
allow_user_choice: bool = False
user_audience: str = Field(default='both')
promo_group_id: int | None = None
parent_task_id: int | None = None
level: int = Field(default=1, ge=1)
starts_at: datetime | None = None
ends_at: datetime | None = None
@field_validator('task_type')
@classmethod
def _validate_task_type(cls, v: str) -> str:
if v not in TASK_TYPES:
raise ValueError(f'invalid task_type: {v}')
return v
@field_validator('reward_type')
@classmethod
def _validate_reward_type(cls, v: str) -> str:
if v not in REWARD_TYPES:
raise ValueError(f'invalid reward_type: {v}')
return v
@field_validator('user_audience')
@classmethod
def _validate_user_audience(cls, v: str) -> str:
if v not in USER_AUDIENCES:
raise ValueError(f'invalid user_audience: {v}')
return v
@field_validator('title')
@classmethod
def _validate_title(cls, v: dict[str, str]) -> dict[str, str]:
if not v or not any(value.strip() for value in v.values() if isinstance(value, str)):
raise ValueError('title must contain at least one non-empty translation')
return v
@model_validator(mode='after')
def _validate_meta_per_type(self) -> TaskCreateRequest:
"""Per-type validation: target_meta required keys, reward_value sanity."""
# PURCHASE_TARIFF требует tariff_id
if self.task_type == 'purchase_tariff' and 'tariff_id' not in (self.target_meta or {}):
raise ValueError('PURCHASE_TARIFF requires target_meta.tariff_id')
# SUBSCRIBE_CHANNEL требует channel_id (строкой, как в TaskPartnerChannel.channel_id)
if self.task_type == 'subscribe_channel':
channel_id = (self.target_meta or {}).get('channel_id')
if channel_id is None:
raise ValueError('SUBSCRIBE_CHANNEL requires target_meta.channel_id')
if not isinstance(channel_id, str) or not channel_id.strip():
raise ValueError('SUBSCRIBE_CHANNEL target_meta.channel_id must be a non-empty string')
# PURCHASE_PERIOD требует period_days
if self.task_type == 'purchase_period' and 'period_days' not in (self.target_meta or {}):
raise ValueError('PURCHASE_PERIOD requires target_meta.period_days')
# BALANCE reward требует reward_value > 0
if self.reward_type == 'balance' and self.reward_value <= 0:
raise ValueError('BALANCE reward requires reward_value > 0')
# SUBSCRIPTION_DAYS reward: либо reward_value > 0, либо tariff_id указан
if self.reward_type == 'subscription_days':
tariff_id = (self.reward_meta or {}).get('tariff_id')
if self.reward_value <= 0 and tariff_id is None:
raise ValueError(
'SUBSCRIPTION_DAYS reward requires reward_value > 0 or '
'reward_meta.tariff_id (to use Tariff.bonus_days_per_purchase)'
)
return self
class TaskUpdateRequest(BaseModel):
"""Частичное обновление задания (все поля опциональны)."""
title: dict[str, str] | None = None
description: dict[str, str] | None = None
icon: str | None = None
is_active: bool | None = None
sort_order: int | None = None
task_type: str | None = None
target_value: int | None = Field(default=None, ge=1)
target_meta: dict[str, Any] | None = None
reward_type: str | None = None
reward_value: int | None = Field(default=None, ge=0)
reward_meta: dict[str, Any] | None = None
allow_user_choice: bool | None = None
user_audience: str | None = None
promo_group_id: int | None = None
parent_task_id: int | None = None
level: int | None = Field(default=None, ge=1)
starts_at: datetime | None = None
ends_at: datetime | None = None
@field_validator('task_type')
@classmethod
def _validate_task_type(cls, v: str | None) -> str | None:
if v is not None and v not in TASK_TYPES:
raise ValueError(f'invalid task_type: {v}')
return v
@field_validator('reward_type')
@classmethod
def _validate_reward_type(cls, v: str | None) -> str | None:
if v is not None and v not in REWARD_TYPES:
raise ValueError(f'invalid reward_type: {v}')
return v
@field_validator('user_audience')
@classmethod
def _validate_user_audience(cls, v: str | None) -> str | None:
if v is not None and v not in USER_AUDIENCES:
raise ValueError(f'invalid user_audience: {v}')
return v
class TaskResponse(BaseModel):
"""Полное представление задания (для админа)."""
id: int
title: dict[str, str]
description: dict[str, str]
icon: str | None = None
is_active: bool
sort_order: int
task_type: str
target_value: int
target_meta: dict[str, Any]
reward_type: str
reward_value: int
reward_meta: dict[str, Any]
allow_user_choice: bool
user_audience: str
promo_group_id: int | None = None
parent_task_id: int | None = None
level: int
starts_at: datetime | None = None
ends_at: datetime | None = None
created_at: datetime
updated_at: datetime | None = None
model_config = ConfigDict(from_attributes=True)
class TaskListItem(BaseModel):
"""Компактное представление для списка."""
id: int
title: dict[str, str]
icon: str | None = None
is_active: bool
sort_order: int
task_type: str
target_value: int
target_meta: dict[str, Any] = Field(default_factory=dict)
reward_type: str
reward_value: int
reward_meta: dict[str, Any] = Field(default_factory=dict)
user_audience: str
promo_group_id: int | None = None
parent_task_id: int | None = None
level: int
updated_at: datetime | None = None
model_config = ConfigDict(from_attributes=True)
# ---------------------------------------------------------------------------
# User-side schemas
# ---------------------------------------------------------------------------
class UserTaskProgressResponse(BaseModel):
"""Прогресс пользователя по конкретному заданию."""
task_id: int
title: dict[str, str]
description: dict[str, str]
icon: str | None = None
task_type: str
target_value: int
target_meta: dict[str, Any] = Field(default_factory=dict)
reward_type: str
reward_value: int
reward_meta: dict[str, Any] = Field(default_factory=dict)
allow_user_choice: bool
level: int
parent_task_id: int | None = None
current_value: int
percent: int
is_completed: bool
is_claimed: bool
completed_at: datetime | None = None
claimed_at: datetime | None = None
reward_granted_meta: dict[str, Any] | None = None
class UserTasksListResponse(BaseModel):
"""Список заданий пользователя."""
items: list[UserTaskProgressResponse]
has_unclaimed: bool
unclaimed_count: int
class UserTasksAvailabilityResponse(BaseModel):
"""Краткая инфа для условного показа вкладки."""
has_available_tasks: bool
unclaimed_count: int
class ClaimRewardRequest(BaseModel):
"""Запрос на получение награды."""
chosen_subscription_id: int | None = Field(
default=None,
description='Для multi-tariff / subscription_days reward — какой подписке начислить дни',
)
chosen_reward_type: Literal['balance', 'subscription_days'] | None = Field(
default=None,
description='Если allow_user_choice=true, юзер может выбрать тип награды',
)
class ClaimRewardResponse(BaseModel):
"""Результат claim награды."""
success: bool
reward: dict[str, Any]
+9
View File
@@ -200,6 +200,8 @@ async def create_tariff(
traffic_reset_mode: str | None = None, # DAY, WEEK, MONTH, MONTH_ROLLING, NO_RESET, None = глобальная настройка
# Внешний сквад RemnaWave
external_squad_uuid: str | None = None,
# Бонусные дни для Tasks (subscription_days reward)
bonus_days_per_purchase: int = 0,
) -> Tariff:
"""Создает новый тариф."""
normalized_prices = _normalize_period_prices(period_prices)
@@ -240,6 +242,8 @@ async def create_tariff(
traffic_reset_mode=traffic_reset_mode,
# Внешний сквад
external_squad_uuid=external_squad_uuid,
# Бонусные дни Tasks
bonus_days_per_purchase=max(0, bonus_days_per_purchase),
)
db.add(tariff)
@@ -309,6 +313,8 @@ async def update_tariff(
traffic_reset_mode: str | None = ..., # ... = не передан, None = сбросить к глобальной настройке
# Внешний сквад RemnaWave
external_squad_uuid: str | None = ..., # ... = не передан, None = убрать внешний сквад
# Бонусные дни для Tasks
bonus_days_per_purchase: int | None = None,
) -> Tariff:
"""Обновляет существующий тариф."""
if name is not None:
@@ -378,6 +384,9 @@ async def update_tariff(
# Внешний сквад
if external_squad_uuid is not ...:
tariff.external_squad_uuid = external_squad_uuid
# Бонусные дни Tasks
if bonus_days_per_purchase is not None:
tariff.bonus_days_per_purchase = max(0, bonus_days_per_purchase)
# Обновляем промогруппы если указаны
if promo_group_ids is not None:
+388
View File
@@ -0,0 +1,388 @@
"""CRUD операции для системы заданий с наградами."""
from __future__ import annotations
from datetime import UTC, datetime
from typing import Any
import structlog
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
from app.database.models import (
Task,
TaskPartnerChannel,
TaskRewardType,
TaskType,
TaskUserAudience,
UserTaskProgress,
)
logger = structlog.get_logger(__name__)
# ===========================================================================
# Task CRUD
# ===========================================================================
async def create_task(
db: AsyncSession,
*,
title: dict[str, str],
description: dict[str, str],
task_type: TaskType | str,
reward_type: TaskRewardType | str,
target_value: int = 1,
reward_value: int = 0,
target_meta: dict[str, Any] | None = None,
reward_meta: dict[str, Any] | None = None,
icon: str | None = None,
is_active: bool = True,
sort_order: int = 0,
allow_user_choice: bool = False,
user_audience: TaskUserAudience | str = TaskUserAudience.BOTH,
promo_group_id: int | None = None,
parent_task_id: int | None = None,
level: int = 1,
starts_at: datetime | None = None,
ends_at: datetime | None = None,
) -> Task:
"""Создаёт новый шаблон задания."""
task = Task(
title=title,
description=description,
icon=icon,
is_active=is_active,
sort_order=sort_order,
task_type=task_type.value if isinstance(task_type, TaskType) else task_type,
target_value=target_value,
target_meta=target_meta or {},
reward_type=reward_type.value if isinstance(reward_type, TaskRewardType) else reward_type,
reward_value=reward_value,
reward_meta=reward_meta or {},
allow_user_choice=allow_user_choice,
user_audience=user_audience.value if isinstance(user_audience, TaskUserAudience) else user_audience,
promo_group_id=promo_group_id,
parent_task_id=parent_task_id,
level=level,
starts_at=starts_at,
ends_at=ends_at,
)
db.add(task)
await db.commit()
await db.refresh(task)
logger.info('Создано задание', task_id=task.id, type=task.task_type, level=task.level)
return task
async def get_task_by_id(db: AsyncSession, task_id: int) -> Task | None:
result = await db.execute(
select(Task).options(selectinload(Task.promo_group)).where(Task.id == task_id)
)
return result.scalar_one_or_none()
async def list_tasks(
db: AsyncSession,
*,
include_inactive: bool = False,
parent_task_id: int | None = None,
) -> list[Task]:
"""Список заданий (для админа). По умолчанию исключает неактивные."""
stmt = (
select(Task)
.options(selectinload(Task.promo_group))
.order_by(Task.sort_order, Task.id)
)
if not include_inactive:
stmt = stmt.where(Task.is_active == True)
if parent_task_id is not None:
stmt = stmt.where(Task.parent_task_id == parent_task_id)
result = await db.execute(stmt)
return list(result.scalars().all())
async def update_task(db: AsyncSession, task: Task, **fields: Any) -> Task:
"""Обновляет поля задания. Enum-поля принимаются как enum или строка."""
for key, value in fields.items():
if key == 'task_type' and isinstance(value, TaskType):
value = value.value
if key == 'reward_type' and isinstance(value, TaskRewardType):
value = value.value
if key == 'user_audience' and isinstance(value, TaskUserAudience):
value = value.value
setattr(task, key, value)
task.updated_at = datetime.now(UTC)
await db.commit()
await db.refresh(task)
logger.info('Обновлено задание', task_id=task.id)
return task
async def delete_task(db: AsyncSession, task: Task) -> None:
await db.delete(task)
await db.commit()
logger.info('Удалено задание', task_id=task.id)
async def list_active_tasks_for_user(
db: AsyncSession,
*,
user_audience: TaskUserAudience | str,
promo_group_id: int | None,
now: datetime | None = None,
) -> list[Task]:
"""Возвращает активные задания, доступные конкретному пользователю.
Учитывает:
- is_active = True
- starts_at <= now <= ends_at (если заданы)
- user_audience: задание для 'both' или совпадающего типа аудитории
- promo_group_id: задание без промогруппы либо совпадающей с user
"""
audience_value = (
user_audience.value if isinstance(user_audience, TaskUserAudience) else user_audience
)
now = now or datetime.now(UTC)
stmt = (
select(Task)
.where(Task.is_active == True)
.where((Task.starts_at == None) | (Task.starts_at <= now))
.where((Task.ends_at == None) | (Task.ends_at >= now))
.where(Task.user_audience.in_(['both', audience_value]))
.order_by(Task.level, Task.sort_order, Task.id)
)
if promo_group_id is not None:
# Задание без promo_group_id — для всех; либо ровно та же группа
stmt = stmt.where((Task.promo_group_id == None) | (Task.promo_group_id == promo_group_id))
else:
stmt = stmt.where(Task.promo_group_id == None)
result = await db.execute(stmt)
return list(result.scalars().all())
# ===========================================================================
# UserTaskProgress CRUD
# ===========================================================================
async def get_progress(db: AsyncSession, *, user_id: int, task_id: int) -> UserTaskProgress | None:
result = await db.execute(
select(UserTaskProgress)
.options(selectinload(UserTaskProgress.task))
.where(UserTaskProgress.user_id == user_id, UserTaskProgress.task_id == task_id)
)
return result.scalar_one_or_none()
async def get_progress_for_update(
db: AsyncSession, *, user_id: int, task_id: int
) -> UserTaskProgress | None:
"""FOR UPDATE lock — для атомарного claim."""
result = await db.execute(
select(UserTaskProgress)
.where(UserTaskProgress.user_id == user_id, UserTaskProgress.task_id == task_id)
.with_for_update()
.execution_options(populate_existing=True)
)
return result.scalar_one_or_none()
async def get_progress_by_id_for_update(db: AsyncSession, progress_id: int) -> UserTaskProgress | None:
result = await db.execute(
select(UserTaskProgress)
.where(UserTaskProgress.id == progress_id)
.with_for_update()
.execution_options(populate_existing=True)
)
return result.scalar_one_or_none()
async def get_or_create_progress(
db: AsyncSession,
*,
user_id: int,
task_id: int,
period_started_at: datetime | None = None,
baseline_value: int = 0,
) -> tuple[UserTaskProgress, bool]:
"""Получает или создаёт запись прогресса. Возвращает (progress, created).
На PostgreSQL использует ``INSERT ... ON CONFLICT DO NOTHING`` (атомарный upsert),
защищая от race на uq_user_task при параллельных record_event.
На SQLite (dev/test mode) использует savepoint + try/IntegrityError atomic upsert
тоже доступен в SQLite dialect, но проще и надёжнее savepoint pattern.
"""
from sqlalchemy.exc import IntegrityError
from app.database.database import IS_SQLITE
if IS_SQLITE:
existing = await get_progress(db, user_id=user_id, task_id=task_id)
if existing is not None:
return existing, False
# ВАЖНО: ``db.add(progress)`` ДОЛЖЕН быть внутри ``begin_nested()``, иначе
# SQLAlchemy в ``_take_snapshot`` сделает flush до открытия savepoint, и
# IntegrityError повредит outer transaction вместо savepoint. См. эталонный
# паттерн в ``app/database/crud/promocode.py``.
try:
async with db.begin_nested():
progress = UserTaskProgress(
user_id=user_id,
task_id=task_id,
current_value=0,
baseline_value=baseline_value,
period_started_at=period_started_at,
)
db.add(progress)
await db.flush()
return progress, True
except IntegrityError:
# Параллельный insert проскочил впереди — fetch'нем существующую запись.
# Savepoint уже откатился, _new очищен через _restore_snapshot.
existing = await get_progress(db, user_id=user_id, task_id=task_id)
if existing is None:
raise RuntimeError('UserTaskProgress race: row disappeared after IntegrityError')
return existing, False
from sqlalchemy.dialects.postgresql import insert as pg_insert
stmt = (
pg_insert(UserTaskProgress)
.values(
user_id=user_id,
task_id=task_id,
current_value=0,
baseline_value=baseline_value,
period_started_at=period_started_at,
)
.on_conflict_do_nothing(index_elements=['user_id', 'task_id'])
.returning(UserTaskProgress.id)
)
result = await db.execute(stmt)
row = result.first()
created = row is not None
if created:
await db.flush()
progress = await get_progress(db, user_id=user_id, task_id=task_id)
if progress is None:
# Не должно случаться: либо мы только что вставили, либо запись уже была.
raise RuntimeError('UserTaskProgress disappeared after upsert')
return progress, created
async def list_user_progress(
db: AsyncSession, *, user_id: int, task_ids: list[int] | None = None
) -> list[UserTaskProgress]:
stmt = (
select(UserTaskProgress)
.options(selectinload(UserTaskProgress.task))
.where(UserTaskProgress.user_id == user_id)
)
if task_ids is not None:
stmt = stmt.where(UserTaskProgress.task_id.in_(task_ids))
result = await db.execute(stmt)
return list(result.scalars().all())
async def mark_progress_completed(
db: AsyncSession, progress: UserTaskProgress
) -> UserTaskProgress:
if progress.completed_at is None:
progress.completed_at = datetime.now(UTC)
progress.updated_at = datetime.now(UTC)
await db.flush()
return progress
async def mark_progress_claimed(
db: AsyncSession,
progress: UserTaskProgress,
*,
reward_granted_meta: dict[str, Any] | None = None,
) -> UserTaskProgress:
if progress.claimed_at is None:
progress.claimed_at = datetime.now(UTC)
if reward_granted_meta is not None:
progress.reward_granted_meta = reward_granted_meta
progress.updated_at = datetime.now(UTC)
await db.flush()
return progress
# ===========================================================================
# TaskPartnerChannel CRUD
# ===========================================================================
async def list_partner_channels(db: AsyncSession, *, include_inactive: bool = False) -> list[TaskPartnerChannel]:
stmt = select(TaskPartnerChannel).order_by(TaskPartnerChannel.sort_order, TaskPartnerChannel.id)
if not include_inactive:
stmt = stmt.where(TaskPartnerChannel.is_active == True)
result = await db.execute(stmt)
return list(result.scalars().all())
async def get_partner_channel_by_id(db: AsyncSession, channel_pk: int) -> TaskPartnerChannel | None:
result = await db.execute(
select(TaskPartnerChannel).where(TaskPartnerChannel.id == channel_pk)
)
return result.scalar_one_or_none()
async def get_partner_channel_by_channel_id(
db: AsyncSession, channel_id: str
) -> TaskPartnerChannel | None:
result = await db.execute(
select(TaskPartnerChannel).where(TaskPartnerChannel.channel_id == channel_id)
)
return result.scalar_one_or_none()
async def create_partner_channel(
db: AsyncSession,
*,
channel_id: str,
title: str,
channel_link: str | None = None,
description: str | None = None,
is_active: bool = True,
sort_order: int = 0,
) -> TaskPartnerChannel:
channel = TaskPartnerChannel(
channel_id=channel_id,
title=title,
channel_link=channel_link,
description=description,
is_active=is_active,
sort_order=sort_order,
)
db.add(channel)
await db.commit()
await db.refresh(channel)
logger.info('Создан партнёрский канал', channel_id=channel_id)
return channel
async def update_partner_channel(
db: AsyncSession, channel: TaskPartnerChannel, **fields: Any
) -> TaskPartnerChannel:
for key, value in fields.items():
setattr(channel, key, value)
channel.updated_at = datetime.now(UTC)
await db.commit()
await db.refresh(channel)
return channel
async def delete_partner_channel(db: AsyncSession, channel: TaskPartnerChannel) -> None:
await db.delete(channel)
await db.commit()
+36
View File
@@ -132,6 +132,28 @@ async def create_transaction(
except Exception as exc:
logger.debug('Не удалось записать событие конкурса для пользователя', user_id=user_id, exc=exc)
# Tasks: SPEND_AMOUNT — учитываем потраченное на подписку. Триггер также есть
# в emit_transaction_side_effects (для commit=False callers); тут — для commit=True.
try:
from app.database.models import TaskType
from app.services.tasks_service import record_event
await record_event(
db,
user_id=user_id,
event_type=TaskType.SPEND_AMOUNT,
payload={'amount_kopeks': abs(amount_kopeks), 'is_trial': False},
)
# commit=True путь: фиксируем task-прогресс в той же транзакции
await db.commit()
except Exception as exc:
# Откатываем poisoned state, чтобы caller не получил PendingRollbackError.
try:
await db.rollback()
except Exception:
pass
logger.warning('Tasks: SPEND_AMOUNT прогресс не обновлён', user_id=user_id, exc=exc)
return transaction
@@ -193,6 +215,20 @@ async def emit_transaction_side_effects(
except Exception as exc:
logger.debug('Не удалось записать событие конкурса для пользователя', user_id=user_id, exc=exc)
# Tasks: SPEND_AMOUNT — учитываем потраченное на подписку
try:
from app.database.models import TaskType
from app.services.tasks_service import record_event
await record_event(
db,
user_id=user_id,
event_type=TaskType.SPEND_AMOUNT,
payload={'amount_kopeks': abs(amount_kopeks), 'is_trial': False},
)
except Exception as exc:
logger.debug('Tasks: не удалось обновить SPEND_AMOUNT прогресс', user_id=user_id, exc=exc)
async def get_transaction_by_id(db: AsyncSession, transaction_id: int) -> Transaction | None:
result = await db.execute(
+4
View File
@@ -511,6 +511,9 @@ async def add_user_balance(
if create_transaction:
from app.database.crud.transaction import create_transaction as create_trans
# Пропагируем commit=False вниз, иначе вложенный create_trans выпустит
# преждевременный db.commit() и сбросит FOR UPDATE-локи (включая progress lock
# из claim_reward).
await create_trans(
db=db,
user_id=user.id,
@@ -518,6 +521,7 @@ async def add_user_balance(
amount_kopeks=amount_kopeks,
description=description,
payment_method=payment_method,
commit=commit,
)
if commit:
+173
View File
@@ -1649,6 +1649,10 @@ class Tariff(Base):
# Видимость в разделе подарков
show_in_gift = Column(Boolean, default=True, server_default='true', nullable=False)
# Бонусные дни — сколько дней начислять при выдаче награды subscription_days,
# если задание ссылается на этот тариф (используется системой Tasks)
bonus_days_per_purchase = Column(Integer, default=0, nullable=False, server_default='0')
# Режим сброса трафика: DAY, WEEK, MONTH, MONTH_ROLLING, NO_RESET (по умолчанию берётся из конфига)
traffic_reset_mode = Column(String(20), nullable=True, default=None) # None = использовать глобальную настройку
@@ -4070,3 +4074,172 @@ class InfoPage(Base):
replaces_tab = Column(String(20), nullable=True) # 'faq', 'rules', 'privacy', 'offer', or null
created_at = Column(AwareDateTime(), server_default=func.now())
updated_at = Column(AwareDateTime(), server_default=func.now(), onupdate=func.now())
class TaskType(Enum):
"""Тип задания (что нужно выполнить пользователю)."""
PURCHASE_TARIFF = 'purchase_tariff' # купить конкретный тариф
SUBSCRIBE_CHANNEL = 'subscribe_channel' # подписаться на партнёрский канал
TRAFFIC_USED = 'traffic_used' # использовать N ГБ трафика за месяц на подписке
REFERRALS_INVITED = 'referrals_invited' # пригласить N рефералов
PURCHASE_PERIOD = 'purchase_period' # купить любой тариф минимум на N дней
SPEND_AMOUNT = 'spend_amount' # совокупно потратить N копеек
MULTI_TARIFF = 'multi_tariff' # иметь N+ активных тарифов в multi-tariff режиме
GIFT_PURCHASED = 'gift_purchased' # купить хотя бы 1 подписку в подарок
GIFTS_COUNT = 'gifts_count' # купить N подписок в подарок (накопительно)
class TaskRewardType(Enum):
"""Тип награды за задание."""
BALANCE = 'balance' # деньги на баланс (в копейках)
SUBSCRIPTION_DAYS = 'subscription_days' # бонусные дни подписки
class TaskUserAudience(Enum):
"""Аудитория задания (фильтр по типу пользователя)."""
TELEGRAM = 'telegram' # только Telegram-юзеры
EMAIL = 'email' # только email-юзеры (cabinet)
BOTH = 'both' # все
class TaskPartnerChannel(Base):
"""Партнёрский канал, на который можно требовать подписку в задании.
Отдельный список от ``RequiredChannel`` (обязательная подписка), чтобы каналы
для заданий не пересекались с системой обязательной подписки.
"""
__tablename__ = 'task_partner_channels'
id = Column(Integer, primary_key=True, index=True)
channel_id = Column(String(100), unique=True, nullable=False, index=True) # формат -100xxxxxxxx
title = Column(String(255), nullable=False)
channel_link = Column(String(500), nullable=True) # https://t.me/...
description = Column(Text, nullable=True)
is_active = Column(Boolean, nullable=False, default=True, server_default='true')
sort_order = Column(Integer, nullable=False, default=0, server_default='0')
created_at = Column(AwareDateTime(), server_default=func.now())
updated_at = Column(AwareDateTime(), server_default=func.now(), onupdate=func.now())
class Task(Base):
"""Шаблон задания с наградой.
Создаётся админом, выдаётся пользователям. Прогресс трекается в ``UserTaskProgress``.
"""
__tablename__ = 'tasks'
id = Column(Integer, primary_key=True, index=True)
# Идентификация и видимость
title = Column(JSONB, nullable=False, server_default='{}') # i18n
description = Column(JSONB, nullable=False, server_default='{}') # i18n
icon = Column(String(50), nullable=True)
is_active = Column(Boolean, nullable=False, default=True, server_default='true')
sort_order = Column(Integer, nullable=False, default=0, server_default='0')
# Тип задания и его параметры
task_type = Column(String(32), nullable=False, index=True) # значение TaskType.value
target_value = Column(BigInteger, nullable=False, default=1) # цель: 5 рефералов / 100 ГБ / N копеек
# Дополнительные параметры в зависимости от типа:
# - PURCHASE_TARIFF: {"tariff_id": 12}
# - SUBSCRIBE_CHANNEL: {"channel_id": "-1001234"}
# - PURCHASE_PERIOD: {"period_days": 30}
# - TRAFFIC_USED: {} (target_value = ГБ)
# - SPEND_AMOUNT: {} (target_value = копейки)
# - REFERRALS_INVITED / MULTI_TARIFF / GIFT_PURCHASED / GIFTS_COUNT: {} (target_value = шт.)
target_meta = Column(JSON, nullable=False, default=dict, server_default='{}')
# Награда
reward_type = Column(String(32), nullable=False) # значение TaskRewardType.value
reward_value = Column(BigInteger, nullable=False, default=0) # копейки или дни
# Для SUBSCRIPTION_DAYS: { "tariff_id": 12 } — если задано, дни начисляются на этот тариф
# (если у юзера в multi-tariff несколько подписок — он выберет какую продлевать).
reward_meta = Column(JSON, nullable=False, default=dict, server_default='{}')
# Может ли user сам выбрать тип награды (если админ задал альтернативу в reward_meta.alt)
allow_user_choice = Column(Boolean, nullable=False, default=False, server_default='false')
# Фильтры аудитории (значение TaskUserAudience.value: 'telegram' / 'email' / 'both')
user_audience = Column(String(16), nullable=False, default='both', server_default='both')
promo_group_id = Column(Integer, ForeignKey('promo_groups.id', ondelete='SET NULL'), nullable=True, index=True)
# Цепочка уровней (последовательное открытие)
parent_task_id = Column(Integer, ForeignKey('tasks.id', ondelete='SET NULL'), nullable=True, index=True)
level = Column(Integer, nullable=False, default=1, server_default='1')
# Период действия задания (опционально)
starts_at = Column(AwareDateTime(), nullable=True)
ends_at = Column(AwareDateTime(), nullable=True)
created_at = Column(AwareDateTime(), server_default=func.now())
updated_at = Column(AwareDateTime(), server_default=func.now(), onupdate=func.now())
# Relationships
promo_group = relationship('PromoGroup', backref='tasks')
parent_task = relationship('Task', remote_side=[id], backref='child_tasks')
progress_records = relationship('UserTaskProgress', back_populates='task', cascade='all, delete-orphan')
def __repr__(self) -> str: # pragma: no cover
return f'<Task(id={self.id}, type={self.task_type}, level={self.level})>'
class UserTaskProgress(Base):
"""Прогресс пользователя по конкретному заданию."""
__tablename__ = 'user_task_progress'
__table_args__ = (UniqueConstraint('user_id', 'task_id', name='uq_user_task'),)
id = Column(Integer, primary_key=True, index=True)
user_id = Column(Integer, ForeignKey('users.id', ondelete='CASCADE'), nullable=False, index=True)
task_id = Column(Integer, ForeignKey('tasks.id', ondelete='CASCADE'), nullable=False, index=True)
# Текущий прогресс к target_value (BigInteger чтобы поддерживать SPEND_AMOUNT в копейках за всё время)
current_value = Column(BigInteger, nullable=False, default=0, server_default='0')
# Снапшоты для типов с периодом (TRAFFIC_USED — за месяц на подписке)
# period_started_at — начало текущего окна (для traffic_used = первое число месяца)
# baseline_value — снапшот значения на начало периода (для traffic_used = traffic_used_gb на старте)
period_started_at = Column(AwareDateTime(), nullable=True)
baseline_value = Column(BigInteger, nullable=False, default=0, server_default='0')
# Статусы
completed_at = Column(AwareDateTime(), nullable=True) # когда выполнено
claimed_at = Column(AwareDateTime(), nullable=True) # когда награда получена
# Метаданные о выданной награде:
# { "type": "balance" | "subscription_days",
# "value": 100000,
# "subscription_id": 42, # для multi-tariff: какой подписке начислили дни
# "transaction_id": 1234, # если создана транзакция
# "old_end_date": "...",
# "new_end_date": "..." }
reward_granted_meta = Column(JSON, nullable=True)
created_at = Column(AwareDateTime(), server_default=func.now())
updated_at = Column(AwareDateTime(), server_default=func.now(), onupdate=func.now())
# Relationships
user = relationship('User', backref='task_progress')
task = relationship('Task', back_populates='progress_records')
@property
def is_completed(self) -> bool:
return self.completed_at is not None
@property
def is_claimed(self) -> bool:
return self.claimed_at is not None
@property
def percent(self) -> int:
if not self.task or self.task.target_value <= 0:
return 0
ratio = max(0, min(self.current_value, self.task.target_value)) / self.task.target_value
return int(ratio * 100)
def __repr__(self) -> str: # pragma: no cover
return f'<UserTaskProgress(user={self.user_id}, task={self.task_id}, {self.current_value}/{self.task.target_value if self.task else "?"})>'
+55 -4
View File
@@ -17,14 +17,16 @@ from aiogram.types import ChatMemberUpdated
from app.config import settings
from app.database.crud.subscription import deactivate_subscription, reactivate_subscription
from app.database.crud.tasks import get_partner_channel_by_channel_id
from app.database.crud.user import get_user_by_telegram_id
from app.database.database import AsyncSessionLocal
from app.database.models import SubscriptionStatus, UserStatus
from app.database.models import SubscriptionStatus, TaskType, UserStatus
from app.keyboards.inline import get_channel_sub_keyboard
from app.localization.loader import DEFAULT_LANGUAGE
from app.localization.texts import get_texts
from app.services.channel_subscription_service import channel_subscription_service
from app.services.subscription_service import SubscriptionService
from app.services.tasks_service import record_event
logger = structlog.get_logger(__name__)
@@ -38,14 +40,58 @@ async def _is_required_channel(channel_id: str) -> bool:
return channel_id in required_ids
async def _is_task_partner_channel(channel_id: str) -> bool:
"""Check if the channel_id is a partner channel used for SUBSCRIBE_CHANNEL tasks."""
async with AsyncSessionLocal() as db:
partner = await get_partner_channel_by_channel_id(db, channel_id)
return partner is not None and partner.is_active
async def _trigger_subscribe_channel_task(telegram_id: int, channel_id: str) -> None:
"""Записывает событие SUBSCRIBE_CHANNEL для системы заданий.
Вызывается, когда юзер подписался на канал, у которого есть TaskPartnerChannel.
Безопасно не пробрасывает исключения, чтобы не сломать обработку события.
"""
try:
async with AsyncSessionLocal() as db:
db_user = await get_user_by_telegram_id(db, telegram_id)
if db_user is None:
return
await record_event(
db,
user_id=db_user.id,
event_type=TaskType.SUBSCRIBE_CHANNEL,
payload={'channel_id': channel_id},
)
await db.commit()
except Exception as exc:
logger.error(
'Failed to record SUBSCRIBE_CHANNEL task event',
telegram_id=telegram_id,
channel_id=channel_id,
error=exc,
)
@router.chat_member(ChatMemberUpdatedFilter(member_status_changed=IS_NOT_MEMBER >> IS_MEMBER))
async def on_user_joined_channel(event: ChatMemberUpdated, bot: Bot) -> None:
"""User subscribed to a channel -- update cache and reactivate VPN if applicable."""
user = event.new_chat_member.user
channel_id = str(event.chat.id) # Normalize int to str (DB stores string)
# FILTER: Only process events for required channels
if not await _is_required_channel(channel_id):
is_required = await _is_required_channel(channel_id)
is_partner = await _is_task_partner_channel(channel_id)
# FILTER: Only process events for required or task partner channels
if not is_required and not is_partner:
return
# Если канал партнёрский (для заданий) — триггерим прогресс задания
if is_partner:
await _trigger_subscribe_channel_task(user.id, channel_id)
if not is_required:
return
await channel_subscription_service.on_user_joined(user.id, channel_id)
@@ -125,7 +171,12 @@ async def on_user_joined_channel(event: ChatMemberUpdated, bot: Bot) -> None:
@router.chat_member(ChatMemberUpdatedFilter(member_status_changed=IS_MEMBER >> IS_NOT_MEMBER))
async def on_user_left_channel(event: ChatMemberUpdated, bot: Bot) -> None:
"""User unsubscribed from a channel -- update cache and deactivate VPN if applicable."""
"""User unsubscribed from a channel -- update cache and deactivate VPN if applicable.
Партнёрские каналы (TaskPartnerChannel) тут не учитываем: задания SUBSCRIBE_CHANNEL
в режиме absolute completed-once. Отписка не должна откатывать выполненное задание
(юзер уже claim'нул награду).
"""
user = event.old_chat_member.user
channel_id = str(event.chat.id) # Normalize int to str (DB stores string)
@@ -1111,6 +1111,22 @@ async def handle_custom_confirm(
description=f'Покупка тарифа {tariff.name} на {custom_days} дней',
)
# Tasks: триггерим прогресс по платным покупкам подписок
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=db_user.id,
tariff_id=getattr(tariff, 'id', None),
period_days=custom_days,
amount_kopeks=total_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров tariff_purchase (custom)', error=task_err)
# Отправляем уведомление админу
try:
admin_notification_service = AdminNotificationService(callback.bot)
@@ -1670,6 +1686,22 @@ async def confirm_tariff_purchase(
except Exception as e:
logger.error('Ошибка создания транзакции', error=e)
# Tasks: триггерим прогресс по платным покупкам подписок
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=db_user.id,
tariff_id=getattr(tariff, 'id', None),
period_days=period,
amount_kopeks=final_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров tariff_purchase (preset)', error=task_err)
# Отправляем уведомление админу
try:
admin_notification_service = AdminNotificationService(callback.bot)
@@ -209,6 +209,39 @@ class DailySubscriptionService:
await db.commit()
await db.refresh(user)
# Tasks: триггерим прогресс по платным purchase-tasks (суточное списание).
# SPEND_AMOUNT не сработает через emit_transaction_side_effects (он не вызывается
# тут), поэтому делаем это вручную через trigger_paid_purchase_tasks +
# отдельный SPEND_AMOUNT.
try:
from app.database.models import TaskType
from app.services.tasks_service import record_event, trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(tariff, 'id', None),
period_days=1,
amount_kopeks=daily_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
# SPEND_AMOUNT отдельно — суточный путь не идёт через стандартный
# create_transaction(commit=True), значит SPEND_AMOUNT нужно явно.
await record_event(
db,
user_id=user.id,
event_type=TaskType.SPEND_AMOUNT,
payload={'amount_kopeks': daily_price, 'is_trial': False},
)
await db.commit()
except Exception as task_err:
try:
await db.rollback()
except Exception:
pass
logger.warning('Tasks: ошибка триггеров daily charge', error=task_err)
user_id_display = user.telegram_id or user.email or f'#{user.id}'
logger.info(
'✅ Суточное списание: подписка сумма коп., пользователь',
+1
View File
@@ -80,6 +80,7 @@ PERMISSION_REGISTRY: dict[str, list[str]] = {
'bulk_actions': ['read', 'execute'],
'info_pages': ['read', 'create', 'edit', 'delete'],
'news': ['read', 'create', 'edit', 'delete'],
'tasks': ['read', 'edit'],
}
+1
View File
@@ -71,6 +71,7 @@ _PRESET_ROLES: list[dict] = [
'landings:create',
'landings:edit',
'landings:delete',
'tasks:*',
],
'color': '#F59E0B',
'icon': 'crown',
+29
View File
@@ -262,6 +262,11 @@ async def process_referral_registration(db: AsyncSession, new_user_id: int, refe
except Exception as exc:
logger.debug('Не удалось записать конкурсную регистрацию', exc=exc)
# NB: REFERRALS_INVITED task progress больше НЕ триггерится здесь (registration).
# Триггер перенесён в process_referral_topup → блок first_topup, чтобы предотвратить
# multi-account farming (создал 100 пустых аккаунтов → claimed награду без покупок).
# Засчитываем реферала только когда он сделал первое qualifying пополнение.
if bot:
commission_percent = get_effective_referral_commission_percent(referrer)
referral_notification = (
@@ -410,6 +415,30 @@ async def process_referral_topup(db: AsyncSession, user_id: int, topup_amount_ko
user.has_made_first_topup = True
await db.commit()
# Tasks: засчитываем REFERRALS_INVITED у реферера только сейчас — реферал сделал
# первое qualifying пополнение, значит это «настоящий» приглашённый юзер.
# Триггер при регистрации удалён, чтобы предотвратить multi-account farming.
try:
from app.database.models import TaskType
from app.services.tasks_service import record_event
await record_event(
db,
user_id=referrer.id,
event_type=TaskType.REFERRALS_INVITED,
payload={'referred_user_id': user.id},
)
await db.commit()
except Exception as task_err:
try:
await db.rollback()
except Exception:
pass
logger.warning(
'Tasks: ошибка REFERRALS_INVITED триггера в process_referral_topup',
exc=task_err,
)
try:
await db.execute(
delete(ReferralEarning).where(
+28 -1
View File
@@ -1935,8 +1935,24 @@ class RemnaWaveService:
# Update traffic
used_traffic_bytes = panel_user.get('usedTrafficBytes', 0) or 0
traffic_used_gb = used_traffic_bytes / (1024**3)
if abs(subscription.traffic_used_gb - traffic_used_gb) > 0.01:
traffic_changed = abs(subscription.traffic_used_gb - traffic_used_gb) > 0.01
if traffic_changed:
subscription.traffic_used_gb = traffic_used_gb
# Tasks: TRAFFIC_USED обновляем per-user аггрегат (sum по
# всем платным подпискам). update_traffic_progress сам
# подсчитает суммарный traffic_used_gb после обновления.
try:
from app.services.tasks_service import update_traffic_progress
await update_traffic_progress(
db, user_id=subscription.user_id
)
except Exception as task_err:
logger.warning(
'Tasks: ошибка TRAFFIC_USED триггера',
user_id=subscription.user_id,
error=task_err,
)
# traffic_limit_gb: bot is source of truth, do not overwrite from panel
@@ -2173,6 +2189,17 @@ class RemnaWaveService:
if abs(subscription.traffic_used_gb - traffic_used_gb) > 0.01:
subscription.traffic_used_gb = traffic_used_gb
logger.debug('Обновлен использованный трафик', traffic_used_gb=traffic_used_gb)
# Tasks: TRAFFIC_USED обновляем per-user аггрегат (sum по платным подпискам)
try:
from app.services.tasks_service import update_traffic_progress
await update_traffic_progress(db, user_id=subscription.user_id)
except Exception as task_err:
logger.warning(
'Tasks: ошибка TRAFFIC_USED триггера (single-tariff)',
user_id=subscription.user_id,
error=task_err,
)
# traffic_limit_gb, device_limit: bot is source of truth, do not overwrite from panel
@@ -564,6 +564,22 @@ async def _auto_extend_subscription(
exc_info=True,
)
# Tasks: триггерим прогресс по платным purchase-tasks
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(prepared, 'tariff_id', None) or getattr(prepared.subscription, 'tariff_id', None),
period_days=prepared.period_days,
amount_kopeks=prepared.price_kopeks,
subscription_id=getattr(prepared.subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров auto-purchase (renewal)', error=task_err)
subscription_service = SubscriptionService()
# Сброс трафика: при смене тарифа — по RESET_TRAFFIC_ON_TARIFF_SWITCH, при оплате — по RESET_TRAFFIC_ON_PAYMENT
if is_tariff_change:
@@ -946,6 +962,22 @@ async def _auto_purchase_tariff(
)
transaction = None
# Tasks: триггерим прогресс по платным purchase-tasks
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(tariff, 'id', None),
period_days=period_days,
amount_kopeks=final_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров auto-purchase (tariff)', error=task_err)
# Обновляем Remnawave
# При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
@@ -1302,6 +1334,22 @@ async def _auto_purchase_daily_tariff(
format_user_id=_format_user_id(user),
error=error,
)
# Tasks: триггерим прогресс по платным purchase-tasks (суточная подписка)
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(tariff, 'id', None),
period_days=1,
amount_kopeks=final_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров auto-purchase (daily)', error=task_err)
transaction = None
# Обновляем Remnawave
@@ -2361,6 +2409,22 @@ async def try_auto_extend_expired_after_topup(
exc_info=True,
)
# Tasks: триггерим прогресс по платным purchase-tasks (renewal expired)
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(updated_subscription, 'tariff_id', None) or getattr(subscription, 'tariff_id', None),
period_days=period_days,
amount_kopeks=renewal_cost,
subscription_id=getattr(updated_subscription, 'id', None) or getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров auto-purchase (expired-renewal)', error=task_err)
# Update RemnaWave
try:
await subscription_service.update_remnawave_user(
@@ -2692,6 +2756,22 @@ async def try_resume_disabled_daily_after_topup(
error=error,
)
# Tasks: триггерим прогресс по платным purchase-tasks (daily resume)
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(subscription, 'tariff_id', None),
period_days=1,
amount_kopeks=daily_price,
subscription_id=getattr(subscription, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров auto-purchase (daily-resume)', error=task_err)
# Update charge time and end_date (+24h)
old_end_date = subscription.end_date
try:
@@ -1226,6 +1226,26 @@ class MiniAppSubscriptionPurchaseService:
)
message = f'{message}\n\n{note}'
# Tasks: триггерим прогресс по платным покупкам подписок (PURCHASE_TARIFF /
# PURCHASE_PERIOD / MULTI_TARIFF). Триал-конверсия тоже считается как платная.
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
tariff_id = getattr(pricing.selection, 'tariff_id', None) or getattr(
subscription, 'tariff_id', None
)
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=tariff_id,
period_days=pricing.selection.period.days,
amount_kopeks=pricing.final_total,
subscription_id=subscription.id,
is_trial=False,
)
except Exception as task_err: # pragma: no cover - defensive
logger.warning('Tasks: ошибка триггеров submit_purchase', error=task_err)
return {
'subscription': subscription,
'transaction': transaction,
@@ -566,6 +566,23 @@ class SubscriptionRenewalService:
error=error,
)
# Tasks: триггерим прогресс по платным purchase-tasks (renewal — это покупка
# того же тарифа на новый период, считается как platная покупка).
try:
from app.services.tasks_service import trigger_paid_purchase_tasks
await trigger_paid_purchase_tasks(
db,
user_id=user.id,
tariff_id=getattr(subscription_after, 'tariff_id', None),
period_days=period_days,
amount_kopeks=final_total,
subscription_id=getattr(subscription_after, 'id', None),
is_trial=False,
)
except Exception as task_err:
logger.warning('Tasks: ошибка триггеров renewal', error=task_err)
await db.refresh(user)
await db.refresh(subscription_after)
+756
View File
@@ -0,0 +1,756 @@
"""Сервис системы заданий с наградами.
Архитектура:
- Внешние действия (покупка тарифа, реферал, трафик и т.д.) триггерят
``record_event(...)`` через event_emitter или прямые вызовы.
- ``record_event`` для каждого активного задания пользователя считает прогресс
под конкретный тип задания (см. ``_apply_event_to_progress``).
- При достижении ``target_value`` прогресс помечается ``completed_at``.
- Пользователь вызывает ``claim_reward(...)`` чтобы получить награду.
Многоуровневые задания: ``parent_task_id`` указывает на предыдущий уровень.
Уровень N+1 виден пользователю только после успешного ``claim`` уровня N.
Триал-юзеры (subscription.is_trial=True без платных подписок) не получают
прогресс ``_user_eligible`` возвращает False.
"""
from __future__ import annotations
from datetime import UTC, datetime
from typing import Any
import structlog
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
from app.database.crud import tasks as tasks_crud
from app.database.models import (
Subscription,
Tariff,
Task,
TaskRewardType,
TaskType,
TaskUserAudience,
User,
UserTaskProgress,
)
logger = structlog.get_logger(__name__)
# ===========================================================================
# Eligibility helpers
# ===========================================================================
async def _user_eligible(db: AsyncSession, user: User) -> bool:
"""Eligibility пользователя для системы заданий.
По требованию: триал-юзеры (есть только триал-подписки) НЕ получают доступ
к выполнению заданий (`выполнение заданий триал юзерам недоступно`).
Юзеры без подписок вообще допускаем (могут выполнять реферал/гифт/spend задания).
Юзеры с хотя бы одной платной подпиской допускаем.
Per-task-type фильтр (например, ``purchase_tariff`` payload помечается ``is_trial``)
дополнительно блокирует засчитывание trial-конверсий ретроактивно.
"""
if not user.subscriptions:
return True
has_paid = any(not getattr(sub, 'is_trial', False) for sub in user.subscriptions)
return has_paid
def _user_audience_for(user: User) -> str:
"""Определяет аудиторию пользователя для фильтрации заданий."""
has_telegram = user.telegram_id is not None
has_email = bool(getattr(user, 'email', None))
if has_telegram and has_email:
return 'both'
if has_telegram:
return TaskUserAudience.TELEGRAM.value
if has_email:
return TaskUserAudience.EMAIL.value
return 'both'
def _user_promo_group_id(user: User) -> int | None:
pg = getattr(user, 'promo_group', None)
if pg is None:
return None
return pg.id
def _is_within_period(task: Task, now: datetime | None = None) -> bool:
now = now or datetime.now(UTC)
if task.starts_at is not None and task.starts_at > now:
return False
if task.ends_at is not None and task.ends_at < now:
return False
return True
def _audience_matches(task: Task, user_audience: str) -> bool:
if task.user_audience == TaskUserAudience.BOTH.value:
return True
return task.user_audience == user_audience
def _promo_group_matches(task: Task, user_promo_group_id: int | None) -> bool:
if task.promo_group_id is None:
return True
return task.promo_group_id == user_promo_group_id
async def _parent_completed_and_claimed(db: AsyncSession, *, user_id: int, parent_task_id: int) -> bool:
"""Проверяет, что предыдущий уровень в цепочке выполнен и награда получена."""
progress = await tasks_crud.get_progress(db, user_id=user_id, task_id=parent_task_id)
if progress is None:
return False
return progress.claimed_at is not None
async def get_available_tasks_for_user(db: AsyncSession, user: User) -> list[tuple[Task, UserTaskProgress | None]]:
"""Возвращает список доступных заданий для пользователя с текущим прогрессом.
Отфильтровывает скрытые (по audience/promo_group/period/parent).
"""
if not await _user_eligible(db, user):
return []
audience = _user_audience_for(user)
promo_id = _user_promo_group_id(user)
candidates = await tasks_crud.list_active_tasks_for_user(db, user_audience=audience, promo_group_id=promo_id)
progress_map: dict[int, UserTaskProgress] = {}
if candidates:
progress_records = await tasks_crud.list_user_progress(db, user_id=user.id, task_ids=[t.id for t in candidates])
progress_map = {p.task_id: p for p in progress_records}
visible: list[tuple[Task, UserTaskProgress | None]] = []
for task in candidates:
if not _is_within_period(task):
continue
if not _audience_matches(task, audience):
continue
if not _promo_group_matches(task, promo_id):
continue
if task.parent_task_id is not None:
parent_done = await _parent_completed_and_claimed(db, user_id=user.id, parent_task_id=task.parent_task_id)
if not parent_done:
continue
visible.append((task, progress_map.get(task.id)))
return visible
# ===========================================================================
# Event recording (вызывается из бизнес-логики при событиях)
# ===========================================================================
async def record_event(
db: AsyncSession,
*,
user_id: int,
event_type: TaskType | str,
payload: dict[str, Any] | None = None,
) -> list[UserTaskProgress]:
"""Регистрирует событие и обновляет прогресс по соответствующим заданиям.
Возвращает список прогрессов, которые были изменены (для последующих
уведомлений / триггеров).
"""
payload = payload or {}
event_value = event_type.value if isinstance(event_type, TaskType) else event_type
# Загружаем юзера со связями
user_q = await db.execute(
select(User)
.options(
selectinload(User.subscriptions),
selectinload(User.promo_group),
)
.where(User.id == user_id)
)
user = user_q.scalar_one_or_none()
if user is None:
return []
if not await _user_eligible(db, user):
return []
audience = _user_audience_for(user)
promo_id = _user_promo_group_id(user)
# Берём только те задания, тип которых совпадает с событием.
# Это важно для перфоманса — иначе пришлось бы пробегать все активные.
candidates = await tasks_crud.list_active_tasks_for_user(db, user_audience=audience, promo_group_id=promo_id)
matching = [t for t in candidates if t.task_type == event_value and _is_within_period(t)]
if not matching:
return []
changed: list[UserTaskProgress] = []
for task in matching:
# Multi-level: пропустить, если родитель не зачищен
if task.parent_task_id is not None:
parent_done = await _parent_completed_and_claimed(db, user_id=user.id, parent_task_id=task.parent_task_id)
if not parent_done:
continue
progress = await _apply_event_to_progress(db, user, task, payload)
if progress is not None:
changed.append(progress)
# Пишем изменения в текущую транзакцию через flush; commit делает caller —
# это сохраняет атомарность операции, в рамках которой триггерится событие
# (покупка, реферал, etc).
if changed:
await db.flush()
return changed
async def _apply_event_to_progress(
db: AsyncSession, user: User, task: Task, payload: dict[str, Any]
) -> UserTaskProgress | None:
"""Применяет событие к прогрессу конкретного задания.
Возвращает обновлённый ``UserTaskProgress`` либо None если событие не подошло
(например, тариф в задании не совпал с купленным).
Использует FOR UPDATE-lock на progress row для серилизации параллельных обновлений
(защита от lost-update / double-credit на gонке).
"""
progress, created = await tasks_crud.get_or_create_progress(db, user_id=user.id, task_id=task.id)
if not created:
# Берём lock на существующую строку, чтобы исключить race condition
locked = await tasks_crud.get_progress_by_id_for_update(db, progress.id)
if locked is not None:
progress = locked
if progress.completed_at is not None:
# Уже выполнено, повторно не обновляем
return None
value, target_meta_match, mode = _compute_increment(task, payload)
if not target_meta_match:
return None
if value == 0:
return None
if mode == 'absolute':
# Установить абсолютное значение (но не больше target_value)
new_value = min(value, task.target_value)
# Не уменьшаем прогресс назад (если admin отозвал тариф — оставим что есть)
if new_value <= progress.current_value:
return None
else:
new_value = min(progress.current_value + value, task.target_value)
if new_value == progress.current_value:
return None
progress.current_value = new_value
progress.updated_at = datetime.now(UTC)
if progress.current_value >= task.target_value:
await tasks_crud.mark_progress_completed(db, progress)
return progress
def _compute_increment(task: Task, payload: dict[str, Any]) -> tuple[int, bool, str]:
"""Возвращает ``(value, target_meta_match, mode)``.
``mode``:
- ``'increment'`` прибавить ``value`` к ``current_value`` (накопительные счётчики)
- ``'absolute'`` установить ``value`` как абсолютное (например, кол-во активных подписок)
``target_meta_match`` соответствует ли событие требованиям задания
(например, конкретный ``tariff_id`` для PURCHASE_TARIFF). Если False событие игнорируется.
Триал-юзер фильтруется отдельно через ``_user_eligible``; здесь дополнительно блокируем
события, помеченные ``payload['is_trial']=True``, чтобы конверсия trialpaid не давала
ретроактивно зачитанный прогресс.
"""
ttype = task.task_type
meta = task.target_meta or {}
if ttype == TaskType.PURCHASE_TARIFF.value:
if payload.get('is_trial'):
return 0, False, 'increment'
target_tariff_id = meta.get('tariff_id')
if target_tariff_id is None:
return 0, False, 'increment'
if int(payload.get('tariff_id') or 0) != int(target_tariff_id):
return 0, False, 'increment'
return 1, True, 'increment'
if ttype == TaskType.SUBSCRIBE_CHANNEL.value:
target_channel_id = meta.get('channel_id')
if target_channel_id is None:
return 0, False, 'absolute'
if str(payload.get('channel_id') or '') != str(target_channel_id):
return 0, False, 'absolute'
# Подписался — задание выполнено сразу (1 канал = 1 шаг к цели)
return task.target_value, True, 'absolute'
if ttype == TaskType.TRAFFIC_USED.value:
# TRAFFIC_USED обрабатывается ТОЛЬКО через update_traffic_progress (calendar-month
# windowing + per-user aggregate). Direct record_event(...TRAFFIC_USED) — no-op.
return 0, False, 'absolute'
if ttype == TaskType.REFERRALS_INVITED.value:
return 1, True, 'increment'
if ttype == TaskType.PURCHASE_PERIOD.value:
if payload.get('is_trial'):
return 0, False, 'increment'
target_period = int(meta.get('period_days') or 0)
period = int(payload.get('period_days') or 0)
if period < target_period:
return 0, False, 'increment'
return 1, True, 'increment'
if ttype == TaskType.SPEND_AMOUNT.value:
if payload.get('is_trial'):
return 0, False, 'increment'
amount_kopeks = int(payload.get('amount_kopeks') or 0)
if amount_kopeks <= 0:
return 0, True, 'increment'
return amount_kopeks, True, 'increment'
if ttype == TaskType.MULTI_TARIFF.value:
# Абсолютное число активных платных подписок пользователя в текущий момент.
count = int(payload.get('active_paid_subscriptions') or 0)
if count <= 0:
return 0, True, 'absolute'
return count, True, 'absolute'
if ttype == TaskType.GIFT_PURCHASED.value:
# Достаточно одного подарка — устанавливаем target_value сразу
return task.target_value, True, 'absolute'
if ttype == TaskType.GIFTS_COUNT.value:
return 1, True, 'increment'
return 0, False, 'increment'
return 0, False
# ===========================================================================
# Claim reward
# ===========================================================================
async def claim_reward(
db: AsyncSession,
*,
user_id: int,
task_id: int,
chosen_subscription_id: int | None = None,
chosen_reward_type: TaskRewardType | str | None = None,
) -> dict[str, Any]:
"""Выдаёт награду пользователю за выполненное задание.
``chosen_subscription_id`` для multi-tariff: какой подписке начислить дни.
``chosen_reward_type`` если задание ``allow_user_choice=True``, юзер
может выбрать тип награды.
Возвращает meta-словарь о выданной награде. Бросает ValueError при
некорректном вызове (не выполнено, уже выдано, и т.д.).
"""
progress = await tasks_crud.get_progress_for_update(db, user_id=user_id, task_id=task_id)
if progress is None:
raise ValueError('progress_not_found')
if progress.completed_at is None:
raise ValueError('not_completed')
if progress.claimed_at is not None:
raise ValueError('already_claimed')
task_q = await db.execute(select(Task).where(Task.id == task_id))
task = task_q.scalar_one_or_none()
if task is None:
raise ValueError('task_not_found')
user_q = await db.execute(select(User).options(selectinload(User.subscriptions)).where(User.id == user_id))
user = user_q.scalar_one_or_none()
if user is None:
raise ValueError('user_not_found')
# Триал-юзер не может клеймить
if not await _user_eligible(db, user):
raise ValueError('user_not_eligible')
# chosen_reward_type разрешён только если task.allow_user_choice=True. Без этой проверки
# юзер мог бы подменить balance-награду на subscription_days (или наоборот) и получить
# значение, не предусмотренное админом. См. claim_reward в tasks_service.
chosen_value = (
chosen_reward_type.value if isinstance(chosen_reward_type, TaskRewardType) else chosen_reward_type
)
if chosen_value is not None and chosen_value != task.reward_type and not task.allow_user_choice:
raise ValueError('user_choice_not_allowed')
reward_type = chosen_value or task.reward_type
granted_meta: dict[str, Any] = {'type': reward_type}
if reward_type == TaskRewardType.BALANCE.value:
granted_meta.update(await _grant_balance_reward(db, user=user, task=task))
elif reward_type == TaskRewardType.SUBSCRIPTION_DAYS.value:
granted_meta.update(
await _grant_subscription_days_reward(
db,
user=user,
task=task,
chosen_subscription_id=chosen_subscription_id,
)
)
else:
raise ValueError(f'unknown_reward_type:{reward_type}')
await tasks_crud.mark_progress_claimed(db, progress, reward_granted_meta=granted_meta)
await db.commit()
logger.info(
'Награда за задание выдана',
user_id=user_id,
task_id=task_id,
reward_type=reward_type,
meta=granted_meta,
)
return granted_meta
async def _grant_balance_reward(db: AsyncSession, *, user: User, task: Task) -> dict[str, Any]:
"""Начисляет деньги на баланс пользователя через ``add_user_balance``.
Не делает commit caller (``claim_reward``) сам закрывает транзакцию.
"""
from app.database.crud.user import add_user_balance
from app.database.models import PaymentMethod, TransactionType
amount_kopeks = int(task.reward_value or 0)
if amount_kopeks <= 0:
raise ValueError('invalid_reward_amount')
old_balance = user.balance_kopeks
# add_user_balance берёт FOR UPDATE на user, создаёт транзакцию.
# transaction_type=DEPOSIT (единственный безопасный non-monetary тип в проекте — у
# TransactionType нет MANUAL). payment_method=MANUAL исключает транзакцию из revenue
# графиков (см. REAL_PAYMENT_METHODS).
success = await add_user_balance(
db,
user,
amount_kopeks,
description=f'Награда за задание #{task.id}',
create_transaction=True,
transaction_type=TransactionType.DEPOSIT,
payment_method=PaymentMethod.MANUAL,
commit=False,
)
if not success:
raise ValueError('balance_grant_failed')
# После lock_for_update в add_user_balance объект user обновлён в session.
await db.refresh(user)
return {
'value': amount_kopeks,
'old_balance': old_balance,
'new_balance': user.balance_kopeks,
}
async def _grant_subscription_days_reward(
db: AsyncSession,
*,
user: User,
task: Task,
chosen_subscription_id: int | None,
) -> dict[str, Any]:
"""Начисляет бонусные дни на платную подписку.
Логика выбора подписки:
1. Если у задания ``reward_meta['tariff_id']`` задан ищем платную подписку
пользователя с этим тарифом.
2. Иначе если у юзера несколько платных подписок и ``chosen_subscription_id``
не задан ошибка ``need_choose_subscription``.
3. Иначе единственная платная подписка.
Количество дней берётся из ``task.reward_value`` (приоритетно), либо из
``Tariff.bonus_days_per_purchase`` если ``reward_value=0`` и задан target tariff.
"""
from datetime import timedelta
from app.database.models import SubscriptionStatus
# Только активные платные подписки — продлевать expired/disabled нельзя
active_paid_subs = [
s
for s in user.subscriptions
if not getattr(s, 'is_trial', False) and getattr(s, 'status', None) == SubscriptionStatus.ACTIVE.value
]
if not active_paid_subs:
raise ValueError('no_paid_subscription')
target_tariff_id = (task.reward_meta or {}).get('tariff_id')
candidate_subs: list[Subscription]
if target_tariff_id is not None:
candidate_subs = [s for s in active_paid_subs if s.tariff_id == int(target_tariff_id)]
if not candidate_subs:
raise ValueError('no_subscription_with_target_tariff')
else:
candidate_subs = active_paid_subs
if chosen_subscription_id is not None:
chosen = next((s for s in candidate_subs if s.id == int(chosen_subscription_id)), None)
if chosen is None:
raise ValueError('chosen_subscription_invalid')
target_sub = chosen
elif len(candidate_subs) == 1:
target_sub = candidate_subs[0]
else:
raise ValueError('need_choose_subscription')
# Берём FOR UPDATE на выбранную подписку — защита от lost-update между
# параллельными claim'ами на одной подписке (например, два task с разными reward).
locked_sub_q = await db.execute(
select(Subscription)
.where(Subscription.id == target_sub.id)
.with_for_update()
.execution_options(populate_existing=True)
)
locked_sub = locked_sub_q.scalar_one_or_none()
if locked_sub is None:
raise ValueError('chosen_subscription_invalid')
target_sub = locked_sub
# Количество дней: task.reward_value > 0 имеет приоритет;
# иначе fallback на Tariff.bonus_days_per_purchase (target → выбранной подписки).
days = int(task.reward_value or 0)
if days <= 0 and target_tariff_id is not None:
tariff_q = await db.execute(select(Tariff).where(Tariff.id == int(target_tariff_id)))
tariff = tariff_q.scalar_one_or_none()
if tariff is not None:
days = int(getattr(tariff, 'bonus_days_per_purchase', 0) or 0)
if days <= 0 and target_sub.tariff_id is not None:
tariff_q = await db.execute(select(Tariff).where(Tariff.id == target_sub.tariff_id))
tariff = tariff_q.scalar_one_or_none()
if tariff is not None:
days = int(getattr(tariff, 'bonus_days_per_purchase', 0) or 0)
if days <= 0:
raise ValueError('invalid_reward_days')
old_end_date = target_sub.end_date
target_sub.end_date = (
target_sub.end_date + timedelta(days=days)
if target_sub.end_date is not None
else datetime.now(UTC) + timedelta(days=days)
)
target_sub.updated_at = datetime.now(UTC)
return {
'value': days,
'subscription_id': target_sub.id,
'tariff_id': target_sub.tariff_id,
'old_end_date': old_end_date.isoformat() if old_end_date else None,
'new_end_date': target_sub.end_date.isoformat() if target_sub.end_date else None,
}
# ===========================================================================
# Helpers (для интеграции из других сервисов)
# ===========================================================================
def _current_month_start(now: datetime | None = None) -> datetime:
now = now or datetime.now(UTC)
return now.replace(day=1, hour=0, minute=0, second=0, microsecond=0)
async def update_traffic_progress(
db: AsyncSession,
*,
user_id: int,
traffic_used_gb_total: float | None = None,
) -> list[UserTaskProgress]:
"""Обновляет прогресс по заданиям типа TRAFFIC_USED (за календарный месяц на пользователя).
``traffic_used_gb_total`` суммарный трафик по всем платным подпискам пользователя.
Если ``None`` считаем сами через ``user.subscriptions``.
Логика per-user (НЕ per-subscription):
- При пересечении календарного месяца (UTC): фиксируем ``baseline_value =
текущий cumulative`` и ``current_value = 0`` (новый период начался с нуля).
- На каждом обновлении: ``current_value = min(target, max(0, current - baseline))``.
- Если cumulative уменьшился (panel reset / удаление подписки): сдвигаем
``baseline`` вниз, чтобы прогресс не уехал в минус и продолжал считаться корректно.
"""
user_q = await db.execute(
select(User).options(selectinload(User.subscriptions), selectinload(User.promo_group)).where(User.id == user_id)
)
user = user_q.scalar_one_or_none()
if user is None:
return []
if not await _user_eligible(db, user):
return []
# Считаем суммарный использованный трафик по платным подпискам
if traffic_used_gb_total is None:
traffic_used_gb_total = sum(
float(getattr(sub, 'traffic_used_gb', 0) or 0)
for sub in user.subscriptions
if not getattr(sub, 'is_trial', False)
)
audience = _user_audience_for(user)
promo_id = _user_promo_group_id(user)
candidates = await tasks_crud.list_active_tasks_for_user(db, user_audience=audience, promo_group_id=promo_id)
matching = [t for t in candidates if t.task_type == TaskType.TRAFFIC_USED.value and _is_within_period(t)]
if not matching:
return []
month_start = _current_month_start()
used_gb_int = int(traffic_used_gb_total) # GB округляем вниз
changed: list[UserTaskProgress] = []
for task in matching:
if task.parent_task_id is not None:
parent_done = await _parent_completed_and_claimed(db, user_id=user.id, parent_task_id=task.parent_task_id)
if not parent_done:
continue
progress, created = await tasks_crud.get_or_create_progress(db, user_id=user.id, task_id=task.id)
if not created:
locked = await tasks_crud.get_progress_by_id_for_update(db, progress.id)
if locked is not None:
progress = locked
if progress.completed_at is not None:
continue
# Сброс при пересечении календарного месяца
if progress.period_started_at is None or progress.period_started_at < month_start:
progress.period_started_at = month_start
progress.baseline_value = used_gb_int
progress.current_value = 0
progress.updated_at = datetime.now(UTC)
delta = used_gb_int - progress.baseline_value
if delta < 0:
# Cumulative ушёл вниз (удалили подписку, panel reset). Сдвигаем baseline,
# сохраняя current_value (юзер уже потратил эти GB в этом месяце).
progress.baseline_value = used_gb_int
progress.updated_at = datetime.now(UTC)
continue
new_value = min(delta, task.target_value)
if new_value != progress.current_value:
progress.current_value = new_value
progress.updated_at = datetime.now(UTC)
changed.append(progress)
if progress.current_value >= task.target_value:
await tasks_crud.mark_progress_completed(db, progress)
if changed:
await db.flush()
return changed
async def trigger_paid_purchase_tasks(
db: AsyncSession,
*,
user_id: int,
tariff_id: int | None,
period_days: int,
amount_kopeks: int,
subscription_id: int | None,
is_trial: bool = False,
) -> None:
"""Триггерит PURCHASE_TARIFF / PURCHASE_PERIOD / MULTI_TARIFF после платной покупки.
Используется из bot handlers и cabinet submit_purchase. Делает явный db.commit()
record_event сам коммит не делает.
Безопасно все исключения логируются как warning, но не пробрасываются.
"""
try:
from app.database.crud.subscription import get_active_subscriptions_by_user_id
from app.database.models import SubscriptionStatus
common_payload = {
'is_trial': is_trial,
'tariff_id': tariff_id,
'period_days': period_days,
'amount_kopeks': amount_kopeks,
'subscription_id': subscription_id,
}
await record_event(
db, user_id=user_id, event_type=TaskType.PURCHASE_TARIFF, payload=common_payload
)
await record_event(
db, user_id=user_id, event_type=TaskType.PURCHASE_PERIOD, payload=common_payload
)
try:
user_subs = await get_active_subscriptions_by_user_id(db, user_id)
active_paid = sum(
1
for s in user_subs
if not getattr(s, 'is_trial', False)
and getattr(s, 'status', None) == SubscriptionStatus.ACTIVE.value
)
except Exception:
active_paid = 0
if active_paid > 0:
await record_event(
db,
user_id=user_id,
event_type=TaskType.MULTI_TARIFF,
payload={'active_paid_subscriptions': active_paid, **common_payload},
)
await db.commit()
except Exception as exc:
# Сессия может быть в poisoned state (например, IntegrityError race на uq_user_task) —
# явно откатываем, чтобы caller мог продолжить работу с сессией.
try:
await db.rollback()
except Exception:
pass
logger.warning(
'Tasks: ошибка trigger_paid_purchase_tasks',
user_id=user_id,
error=exc,
)
async def has_available_tasks(db: AsyncSession, user: User) -> bool:
"""Быстрый чек: есть ли у пользователя хотя бы одно доступное задание.
Используется фронтом для условного показа вкладки «Задания».
"""
visible = await get_available_tasks_for_user(db, user)
return len(visible) > 0
async def count_completed_unclaimed(db: AsyncSession, *, user_id: int) -> int:
"""Количество выполненных, но не полученных наград (для бейджа на иконке)."""
from sqlalchemy import func
result = await db.execute(
select(func.count())
.select_from(UserTaskProgress)
.where(
UserTaskProgress.user_id == user_id,
UserTaskProgress.completed_at.isnot(None),
UserTaskProgress.claimed_at.is_(None),
)
)
return int(result.scalar() or 0)
@@ -0,0 +1,106 @@
"""create tasks system: tasks, user_task_progress, task_partner_channels + tariffs.bonus_days_per_purchase
Revision ID: 0075
Revises: 0074
Create Date: 2026-05-06
"""
from typing import Sequence, Union
import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql
revision: str = '0075'
down_revision: Union[str, None] = '0074'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# 1. Tariff: bonus_days_per_purchase
op.add_column(
'tariffs',
sa.Column('bonus_days_per_purchase', sa.Integer(), nullable=False, server_default='0'),
)
# 2. Партнёрские каналы для заданий
op.create_table(
'task_partner_channels',
sa.Column('id', sa.Integer(), primary_key=True, autoincrement=True),
sa.Column('channel_id', sa.String(100), unique=True, nullable=False, index=True),
sa.Column('title', sa.String(255), nullable=False),
sa.Column('channel_link', sa.String(500), nullable=True),
sa.Column('description', sa.Text(), nullable=True),
sa.Column('is_active', sa.Boolean(), nullable=False, server_default=sa.text('true')),
sa.Column('sort_order', sa.Integer(), nullable=False, server_default='0'),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
)
# 3. Шаблоны заданий
op.create_table(
'tasks',
sa.Column('id', sa.Integer(), primary_key=True, autoincrement=True),
sa.Column('title', postgresql.JSONB(astext_type=sa.Text()), nullable=False, server_default='{}'),
sa.Column('description', postgresql.JSONB(astext_type=sa.Text()), nullable=False, server_default='{}'),
sa.Column('icon', sa.String(50), nullable=True),
sa.Column('is_active', sa.Boolean(), nullable=False, server_default=sa.text('true')),
sa.Column('sort_order', sa.Integer(), nullable=False, server_default='0'),
sa.Column('task_type', sa.String(32), nullable=False, index=True),
sa.Column('target_value', sa.BigInteger(), nullable=False, server_default='1'),
sa.Column('target_meta', sa.JSON(), nullable=False, server_default='{}'),
sa.Column('reward_type', sa.String(32), nullable=False),
sa.Column('reward_value', sa.BigInteger(), nullable=False, server_default='0'),
sa.Column('reward_meta', sa.JSON(), nullable=False, server_default='{}'),
sa.Column('allow_user_choice', sa.Boolean(), nullable=False, server_default=sa.text('false')),
sa.Column('user_audience', sa.String(16), nullable=False, server_default='both'),
sa.Column(
'promo_group_id',
sa.Integer(),
sa.ForeignKey('promo_groups.id', ondelete='SET NULL'),
nullable=True,
index=True,
),
sa.Column(
'parent_task_id',
sa.Integer(),
sa.ForeignKey('tasks.id', ondelete='SET NULL'),
nullable=True,
index=True,
),
sa.Column('level', sa.Integer(), nullable=False, server_default='1'),
sa.Column('starts_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('ends_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
)
# 4. Прогресс пользователей
op.create_table(
'user_task_progress',
sa.Column('id', sa.Integer(), primary_key=True, autoincrement=True),
sa.Column(
'user_id', sa.Integer(), sa.ForeignKey('users.id', ondelete='CASCADE'), nullable=False, index=True
),
sa.Column(
'task_id', sa.Integer(), sa.ForeignKey('tasks.id', ondelete='CASCADE'), nullable=False, index=True
),
sa.Column('current_value', sa.BigInteger(), nullable=False, server_default='0'),
sa.Column('period_started_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('baseline_value', sa.BigInteger(), nullable=False, server_default='0'),
sa.Column('completed_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('claimed_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('reward_granted_meta', sa.JSON(), nullable=True),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
sa.UniqueConstraint('user_id', 'task_id', name='uq_user_task'),
)
def downgrade() -> None:
op.drop_table('user_task_progress')
op.drop_table('tasks')
op.drop_table('task_partner_channels')
op.drop_column('tariffs', 'bonus_days_per_purchase')