diff --git a/backend/app/api/v1/dashboard.py b/backend/app/api/v1/dashboard.py index 1586ba8..f8b99af 100644 --- a/backend/app/api/v1/dashboard.py +++ b/backend/app/api/v1/dashboard.py @@ -7,11 +7,12 @@ from datetime import datetime, timezone from uuid import UUID -from fastapi import APIRouter, Depends, HTTPException, Query +from fastapi import APIRouter, Depends, HTTPException, Query, status from pydantic import BaseModel from sqlalchemy.ext.asyncio import AsyncSession from app.dependencies.async_db import get_async_db +from app.dependencies.auth import get_current_seller_async from app.models.alert import ( Alert, AlertActionType, @@ -19,22 +20,31 @@ AlertUrgencyLevel, ) from app.models.recommendation import RecommendationAction, RecommendationConfidence +from app.models.seller import Seller from app.repositories.alert_repository import AlertRepository from app.repositories.competitor_repository import CompetitorRepository from app.repositories.query_utils import PaginationParams, SortDirection from app.repositories.recommendation_repository import RecommendationRepository +from app.schemas.dashboard.activity_schema import ActivityItem from app.schemas.dashboard.alert_schema import AlertListResponse, AlertSortField from app.schemas.dashboard.competitor_schema import CompetitorListResponse, CompetitorSortField +from app.schemas.dashboard.market_trend_schema import MarketTrendPoint from app.schemas.dashboard.opportunity_schema import OpportunityListResponse, OpportunitySortField +from app.schemas.dashboard.opportunity_metrics_schema import OpportunityMetricRecord +from app.schemas.dashboard.price_trend_schema import CompetitorPriceTrendPoint from app.schemas.dashboard.recommendation_schema import ( RecommendationListResponse, RecommendationSortField, ) from app.schemas.dashboard.summary_schema import DashboardSummaryResponse +from app.services.dashboard.activity_service import ActivityDashboardService from app.services.dashboard.alert_service import AlertDashboardService from app.services.dashboard.competitor_service import CompetitorDashboardService from app.services.dashboard.dashboard_service import DashboardService +from app.services.dashboard.market_trend_service import MarketTrendDashboardService +from app.services.dashboard.opportunity_metrics_service import OpportunityMetricsService from app.services.dashboard.opportunity_service import OpportunityDashboardService +from app.services.dashboard.price_trend_service import CompetitorPriceTrendService from app.services.dashboard.recommendation_service import RecommendationDashboardService router = APIRouter(prefix="/dashboard", tags=["dashboard"]) @@ -50,10 +60,12 @@ class AlertResolveResponse(BaseModel): async def get_dashboard_summary( seller_id: UUID | None = Query(None, description="Filter by seller"), db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> DashboardSummaryResponse: start = time.monotonic() service = DashboardService(db, logger=logger) - response = await service.get_summary(seller_id=str(seller_id) if seller_id else None) + scoped_seller_id = _require_seller_scope(current_seller, seller_id) + response = await service.get_summary(seller_id=scoped_seller_id) _log_endpoint("dashboard_summary", start, seller_id=seller_id) return response @@ -75,14 +87,16 @@ async def list_competitors( page: int | None = Query(None, ge=1), page_size: int | None = Query(None, ge=1, le=200), db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> CompetitorListResponse: start = time.monotonic() pagination = PaginationParams.from_request( limit=limit, offset=offset, page=page, page_size=page_size, max_page_size=200 ) service = CompetitorDashboardService(CompetitorRepository(db), logger=logger) + scoped_seller_id = _require_seller_scope(current_seller, seller_id) response = await service.list_competitors( - seller_id=str(seller_id) if seller_id else None, + seller_id=scoped_seller_id, product_id=str(product_id) if product_id else None, stock_status=stock_status, search=search, @@ -114,14 +128,16 @@ async def list_recommendations( page: int | None = Query(None, ge=1), page_size: int | None = Query(None, ge=1, le=200), db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> RecommendationListResponse: start = time.monotonic() pagination = PaginationParams.from_request( limit=limit, offset=offset, page=page, page_size=page_size, max_page_size=200 ) service = RecommendationDashboardService(RecommendationRepository(db), logger=logger) + scoped_seller_id = _require_seller_scope(current_seller, seller_id) response = await service.list_recommendations( - seller_id=str(seller_id) if seller_id else None, + seller_id=scoped_seller_id, product_id=str(product_id) if product_id else None, action_type=action_type.value if action_type else None, confidence=confidence.value if confidence else None, @@ -159,14 +175,16 @@ async def list_alerts( page: int | None = Query(None, ge=1), page_size: int | None = Query(None, ge=1, le=200), db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> AlertListResponse: start = time.monotonic() pagination = PaginationParams.from_request( limit=limit, offset=offset, page=page, page_size=page_size, max_page_size=200 ) service = AlertDashboardService(AlertRepository(db), logger=logger) + scoped_seller_id = _require_seller_scope(current_seller, seller_id) response = await service.list_alerts( - seller_id=str(seller_id) if seller_id else None, + seller_id=scoped_seller_id, product_id=str(product_id) if product_id else None, urgency_level=urgency_level.value if urgency_level else None, action_type=action_type.value if action_type else None, @@ -197,14 +215,16 @@ async def list_market_opportunities( page: int | None = Query(None, ge=1), page_size: int | None = Query(None, ge=1, le=200), db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> OpportunityListResponse: start = time.monotonic() pagination = PaginationParams.from_request( limit=limit, offset=offset, page=page, page_size=page_size, max_page_size=200 ) service = OpportunityDashboardService(RecommendationRepository(db), logger=logger) + scoped_seller_id = _require_seller_scope(current_seller, seller_id) response = await service.list_opportunities( - seller_id=str(seller_id) if seller_id else None, + seller_id=scoped_seller_id, min_revenue=min_revenue, date_from=date_from, date_to=date_to, @@ -220,9 +240,10 @@ async def list_market_opportunities( async def resolve_alert( alert_id: UUID, db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), ) -> AlertResolveResponse: alert = await db.get(Alert, alert_id) - if not alert: + if not alert or str(alert.seller_id) != str(current_seller.id): raise HTTPException(status_code=404, detail="Alert not found.") if alert.read_at is None: alert.read_at = datetime.now(timezone.utc) @@ -230,7 +251,80 @@ async def resolve_alert( return AlertResolveResponse(id=str(alert.id), status="resolved") +@router.get("/activity", response_model=list[ActivityItem]) +async def get_activity_feed( + seller_id: UUID | None = Query(None, description="Filter by seller"), + limit: int = Query(12, ge=1, le=50), + db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), +) -> list[ActivityItem]: + start = time.monotonic() + scoped_seller_id = _require_seller_scope(current_seller, seller_id) + service = ActivityDashboardService(db, logger=logger) + response = await service.get_activity_feed(seller_id=scoped_seller_id, limit=limit) + _log_endpoint("dashboard_activity", start, seller_id=seller_id) + return response + + +@router.get("/market-trends", response_model=list[MarketTrendPoint]) +async def get_market_trends( + seller_id: UUID | None = Query(None, description="Filter by seller"), + days: int = Query(14, ge=1, le=90), + db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), +) -> list[MarketTrendPoint]: + start = time.monotonic() + scoped_seller_id = _require_seller_scope(current_seller, seller_id) + service = MarketTrendDashboardService(db, logger=logger) + response = await service.get_market_trends(seller_id=scoped_seller_id, days=days) + _log_endpoint("dashboard_market_trends", start, seller_id=seller_id) + return response + + +@router.get("/opportunity-metrics", response_model=list[OpportunityMetricRecord]) +async def get_opportunity_metrics( + seller_id: UUID | None = Query(None, description="Filter by seller"), + db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), +) -> list[OpportunityMetricRecord]: + start = time.monotonic() + scoped_seller_id = _require_seller_scope(current_seller, seller_id) + service = OpportunityMetricsService(db, logger=logger) + response = await service.get_metrics(seller_id=scoped_seller_id) + _log_endpoint("dashboard_opportunity_metrics", start, seller_id=seller_id) + return response + + +@router.get("/competitor-price-trend", response_model=list[CompetitorPriceTrendPoint]) +async def get_competitor_price_trend( + seller_id: UUID | None = Query(None, description="Filter by seller"), + product_id: UUID | None = Query(None, description="Filter by product"), + days: int = Query(14, ge=1, le=90), + db: AsyncSession = Depends(get_async_db), + current_seller: Seller = Depends(get_current_seller_async), +) -> list[CompetitorPriceTrendPoint]: + start = time.monotonic() + scoped_seller_id = _require_seller_scope(current_seller, seller_id) + service = CompetitorPriceTrendService(db, logger=logger) + response = await service.get_price_trend( + seller_id=scoped_seller_id, + product_id=str(product_id) if product_id else None, + days=days, + ) + _log_endpoint("dashboard_competitor_price_trend", start, seller_id=seller_id) + return response + + def _log_endpoint(name: str, start: float, **fields: Any) -> None: duration_ms = int((time.monotonic() - start) * 1000) payload = " ".join(f"{key}={value}" for key, value in fields.items() if value is not None) logger.info("%s duration_ms=%s %s", name, duration_ms, payload) + + +def _require_seller_scope(current_seller: Seller, seller_id: UUID | None) -> str: + if seller_id and seller_id != current_seller.id: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Seller scope is not authorized.", + ) + return str(current_seller.id) diff --git a/backend/app/core/config.py b/backend/app/core/config.py index e8ed925..aae069d 100644 --- a/backend/app/core/config.py +++ b/backend/app/core/config.py @@ -5,6 +5,8 @@ from dotenv import load_dotenv +print("DATABASE_URL =", os.getenv("DATABASE_URL")) + def _split_csv(value: str) -> list[str]: return [item.strip() for item in value.split(",") if item.strip()] diff --git a/backend/app/main.py b/backend/app/main.py index 02423d4..0cc4034 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -23,6 +23,10 @@ from app.scraper.browser import BrowserManager from app.websocket.socket_manager import init_socketio +import os + +print("DEBUG DATABASE_URL:", os.environ.get("DATABASE_URL")) +print("ALL ENV KEYS:", list(os.environ.keys())) @asynccontextmanager diff --git a/backend/app/schemas/dashboard/__init__.py b/backend/app/schemas/dashboard/__init__.py index 69566d8..7c8518c 100644 --- a/backend/app/schemas/dashboard/__init__.py +++ b/backend/app/schemas/dashboard/__init__.py @@ -1,10 +1,14 @@ from __future__ import annotations from app.schemas.dashboard.alert_schema import AlertListResponse, AlertRecord +from app.schemas.dashboard.activity_schema import ActivityItem from app.schemas.dashboard.competitor_schema import CompetitorListResponse, CompetitorRecord from app.schemas.dashboard.error_schema import DashboardErrorResponse, ErrorDetail +from app.schemas.dashboard.market_trend_schema import MarketTrendPoint from app.schemas.dashboard.opportunity_schema import OpportunityListResponse, OpportunityRecord +from app.schemas.dashboard.opportunity_metrics_schema import OpportunityMetricRecord from app.schemas.dashboard.pagination_schema import PaginatedResponse, PaginationMeta +from app.schemas.dashboard.price_trend_schema import CompetitorPriceTrendPoint from app.schemas.dashboard.recommendation_schema import ( RecommendationListResponse, RecommendationRecord, @@ -14,14 +18,18 @@ __all__ = [ "AlertListResponse", "AlertRecord", + "ActivityItem", "CompetitorListResponse", "CompetitorRecord", "DashboardErrorResponse", "ErrorDetail", + "MarketTrendPoint", "OpportunityListResponse", "OpportunityRecord", + "OpportunityMetricRecord", "PaginatedResponse", "PaginationMeta", + "CompetitorPriceTrendPoint", "RecommendationListResponse", "RecommendationRecord", "DashboardSummaryResponse", diff --git a/backend/app/schemas/dashboard/activity_schema.py b/backend/app/schemas/dashboard/activity_schema.py new file mode 100644 index 0000000..0f1d1c4 --- /dev/null +++ b/backend/app/schemas/dashboard/activity_schema.py @@ -0,0 +1,43 @@ +from __future__ import annotations + +from datetime import datetime +from enum import Enum + +from pydantic import BaseModel, ConfigDict, Field + + +class ActivityCategory(str, Enum): + recommendation = "recommendation" + alert = "alert" + competitor = "competitor" + system = "system" + + +class ActivitySeverity(str, Enum): + info = "info" + warning = "warning" + critical = "critical" + + +class ActivityItem(BaseModel): + id: str + title: str + description: str + timestamp: datetime + category: ActivityCategory + severity: ActivitySeverity + + model_config = ConfigDict( + json_schema_extra={ + "examples": [ + { + "id": "rec-9f3a4f02", + "title": "Recommendation generated", + "description": "Lower price recommendation for Dove Body Lotion 200ml.", + "timestamp": "2026-05-26T22:05:00Z", + "category": "recommendation", + "severity": "info", + } + ] + } + ) diff --git a/backend/app/schemas/dashboard/market_trend_schema.py b/backend/app/schemas/dashboard/market_trend_schema.py new file mode 100644 index 0000000..bdd987a --- /dev/null +++ b/backend/app/schemas/dashboard/market_trend_schema.py @@ -0,0 +1,27 @@ +from __future__ import annotations + +from datetime import datetime + +from pydantic import BaseModel, ConfigDict, Field + + +class MarketTrendPoint(BaseModel): + timestamp: datetime + demand_index: float = Field(..., ge=0, le=100) + urgency_index: float = Field(..., ge=0, le=100) + opportunity_index: float = Field(..., ge=0, le=100) + price_index: float = Field(..., ge=-100, le=100) + + model_config = ConfigDict( + json_schema_extra={ + "examples": [ + { + "timestamp": "2026-05-26T00:00:00Z", + "demand_index": 62.5, + "urgency_index": 58.1, + "opportunity_index": 71.4, + "price_index": -3.2, + } + ] + } + ) diff --git a/backend/app/schemas/dashboard/opportunity_metrics_schema.py b/backend/app/schemas/dashboard/opportunity_metrics_schema.py new file mode 100644 index 0000000..4eb2125 --- /dev/null +++ b/backend/app/schemas/dashboard/opportunity_metrics_schema.py @@ -0,0 +1,25 @@ +from __future__ import annotations + +from pydantic import BaseModel, ConfigDict + + +class OpportunityMetricRecord(BaseModel): + id: str + label: str + value: float | str + delta: float | str | None + trend: str + + model_config = ConfigDict( + json_schema_extra={ + "examples": [ + { + "id": "revenue-opportunity", + "label": "Revenue opportunity", + "value": 4200.0, + "delta": 650.0, + "trend": "up", + } + ] + } + ) diff --git a/backend/app/schemas/dashboard/price_trend_schema.py b/backend/app/schemas/dashboard/price_trend_schema.py new file mode 100644 index 0000000..d5e515a --- /dev/null +++ b/backend/app/schemas/dashboard/price_trend_schema.py @@ -0,0 +1,25 @@ +from __future__ import annotations + +from datetime import datetime + +from pydantic import BaseModel, ConfigDict, Field + + +class CompetitorPriceTrendPoint(BaseModel): + timestamp: datetime + competitor_price: float = Field(..., ge=0) + seller_price: float = Field(..., ge=0) + market_average: float = Field(..., ge=0) + + model_config = ConfigDict( + json_schema_extra={ + "examples": [ + { + "timestamp": "2026-05-26T00:00:00Z", + "competitor_price": 410.0, + "seller_price": 450.0, + "market_average": 432.5, + } + ] + } + ) diff --git a/backend/app/schemas/dashboard/summary_schema.py b/backend/app/schemas/dashboard/summary_schema.py index 0e7969f..bd85eab 100644 --- a/backend/app/schemas/dashboard/summary_schema.py +++ b/backend/app/schemas/dashboard/summary_schema.py @@ -18,6 +18,7 @@ class DashboardSummaryResponse(BaseModel): opportunity_count: int = Field(..., ge=0) high_risk_stock_count: int = Field(..., ge=0) average_urgency_score: float = Field(..., ge=0, le=100) + latest_urgency_score: float = Field(..., ge=0, le=100) average_competitor_price_gap: float = Field(..., ge=-100, le=100) market_trend_summary: str realtime_activity: RealtimeActivityCounters @@ -32,6 +33,7 @@ class DashboardSummaryResponse(BaseModel): "opportunity_count": 22, "high_risk_stock_count": 4, "average_urgency_score": 58.3, + "latest_urgency_score": 74.2, "average_competitor_price_gap": 6.2, "market_trend_summary": "Moderate volatility with rising competitor pressure.", "realtime_activity": { diff --git a/backend/app/services/dashboard/__init__.py b/backend/app/services/dashboard/__init__.py index 1f45498..6a809ee 100644 --- a/backend/app/services/dashboard/__init__.py +++ b/backend/app/services/dashboard/__init__.py @@ -1,15 +1,23 @@ from __future__ import annotations from app.services.dashboard.alert_service import AlertDashboardService +from app.services.dashboard.activity_service import ActivityDashboardService from app.services.dashboard.competitor_service import CompetitorDashboardService from app.services.dashboard.dashboard_service import DashboardService +from app.services.dashboard.market_trend_service import MarketTrendDashboardService +from app.services.dashboard.opportunity_metrics_service import OpportunityMetricsService from app.services.dashboard.opportunity_service import OpportunityDashboardService +from app.services.dashboard.price_trend_service import CompetitorPriceTrendService from app.services.dashboard.recommendation_service import RecommendationDashboardService __all__ = [ "AlertDashboardService", + "ActivityDashboardService", "CompetitorDashboardService", "DashboardService", + "MarketTrendDashboardService", + "OpportunityMetricsService", "OpportunityDashboardService", + "CompetitorPriceTrendService", "RecommendationDashboardService", ] diff --git a/backend/app/services/dashboard/activity_service.py b/backend/app/services/dashboard/activity_service.py new file mode 100644 index 0000000..91ec3bd --- /dev/null +++ b/backend/app/services/dashboard/activity_service.py @@ -0,0 +1,203 @@ +from __future__ import annotations + +import logging +import time +from typing import Any + +from sqlalchemy import Float, cast, select +from sqlalchemy.exc import SQLAlchemyError +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.alert import Alert +from app.models.competitor import Competitor +from app.models.product import Product +from app.models.recommendation import Recommendation +from app.schemas.dashboard.activity_schema import ActivityCategory, ActivityItem, ActivitySeverity +from app.services.dashboard.exceptions import DashboardDatabaseError + + +class ActivityDashboardService: + def __init__(self, session: AsyncSession, *, logger: logging.Logger | None = None) -> None: + self._session = session + self._logger = logger or logging.getLogger(__name__) + + async def get_activity_feed( + self, *, seller_id: str, limit: int = 12 + ) -> list[ActivityItem]: + start = time.monotonic() + per_stream = max(limit, 6) + try: + recommendations = await self._fetch_recommendations(seller_id, per_stream) + alerts = await self._fetch_alerts(seller_id, per_stream) + competitors = await self._fetch_competitors(seller_id, per_stream) + except SQLAlchemyError as exc: + raise DashboardDatabaseError("Failed to load activity feed.") from exc + + items = recommendations + alerts + competitors + items.sort(key=lambda item: item.timestamp, reverse=True) + trimmed = items[:limit] + duration_ms = int((time.monotonic() - start) * 1000) + self._logger.info( + "dashboard_activity_feed count=%s duration_ms=%s seller_id=%s", + len(trimmed), + duration_ms, + seller_id, + ) + return trimmed + + async def _fetch_recommendations(self, seller_id: str, limit: int) -> list[ActivityItem]: + analytics = Recommendation.recommendation_context["analytics"] + urgency_expr = cast(analytics["urgency_score"].astext, Float).label("urgency_score") + stmt = ( + select( + Recommendation.id.label("recommendation_id"), + Recommendation.action.label("action_type"), + Recommendation.reasoning_bn.label("reasoning"), + Recommendation.created_at.label("created_at"), + Product.title.label("product_title"), + urgency_expr, + ) + .join(Product, Recommendation.product_id == Product.id) + .where(Product.seller_id == seller_id) + .order_by(Recommendation.created_at.desc(), Recommendation.id.desc()) + .limit(limit) + ) + result = await self._session.execute(stmt) + rows = result.mappings().all() + items: list[ActivityItem] = [] + for row in rows: + urgency = _normalize_score(row.get("urgency_score")) + severity = _severity_from_score(urgency) + action_label = _format_action(row.get("action_type")) + product_title = row.get("product_title") or "product" + description = ( + row.get("reasoning") or f"{action_label} recommendation for {product_title}." + ) + items.append( + ActivityItem( + id=f"rec-{row['recommendation_id']}", + title="Recommendation generated", + description=description, + timestamp=row.get("created_at"), + category=ActivityCategory.recommendation, + severity=severity, + ) + ) + return items + + async def _fetch_alerts(self, seller_id: str, limit: int) -> list[ActivityItem]: + stmt = ( + select( + Alert.id.label("alert_id"), + Alert.alert_type.label("alert_type"), + Alert.urgency_level.label("urgency_level"), + Alert.action_type.label("action_type"), + Alert.created_at.label("created_at"), + Product.title.label("product_title"), + ) + .join(Product, Alert.product_id == Product.id) + .where(Alert.seller_id == seller_id) + .order_by(Alert.created_at.desc(), Alert.id.desc()) + .limit(limit) + ) + result = await self._session.execute(stmt) + rows = result.mappings().all() + items: list[ActivityItem] = [] + for row in rows: + product_title = row.get("product_title") or "product" + action_label = _format_label(row.get("action_type")) + alert_label = _format_label(row.get("alert_type")) + items.append( + ActivityItem( + id=f"alert-{row['alert_id']}", + title=f"{action_label.title()} alert", + description=f"{alert_label} alert for {product_title}.", + timestamp=row.get("created_at"), + category=ActivityCategory.alert, + severity=_severity_from_urgency(row.get("urgency_level")), + ) + ) + return items + + async def _fetch_competitors(self, seller_id: str, limit: int) -> list[ActivityItem]: + stmt = ( + select( + Competitor.id.label("competitor_id"), + Competitor.competitor_title.label("competitor_title"), + Competitor.stock_status.label("stock_status"), + Competitor.updated_at.label("updated_at"), + Product.title.label("product_title"), + ) + .join(Product, Competitor.product_id == Product.id) + .where(Product.seller_id == seller_id) + .order_by(Competitor.updated_at.desc(), Competitor.id.desc()) + .limit(limit) + ) + result = await self._session.execute(stmt) + rows = result.mappings().all() + items: list[ActivityItem] = [] + for row in rows: + competitor_title = row.get("competitor_title") or "Competitor" + product_title = row.get("product_title") or "product" + severity = ( + ActivitySeverity.warning + if (row.get("stock_status") or "").lower() in {"out_of_stock", "low_stock"} + else ActivitySeverity.info + ) + items.append( + ActivityItem( + id=f"comp-{row['competitor_id']}", + title="Competitor update", + description=f"{competitor_title} adjusted {product_title}.", + timestamp=row.get("updated_at"), + category=ActivityCategory.competitor, + severity=severity, + ) + ) + return items + + +def _normalize_score(value: Any) -> float: + if value is None: + return 0.0 + try: + numeric = float(value) + except (TypeError, ValueError): + return 0.0 + if 0 <= numeric <= 1: + return numeric * 100.0 + return numeric + + +def _severity_from_score(score: float) -> ActivitySeverity: + if score >= 80: + return ActivitySeverity.critical + if score >= 60: + return ActivitySeverity.warning + return ActivitySeverity.info + + +def _severity_from_urgency(value: Any) -> ActivitySeverity: + level = str(value or "").upper() + if level == "CRITICAL": + return ActivitySeverity.critical + if level in {"HIGH", "MEDIUM"}: + return ActivitySeverity.warning + return ActivitySeverity.info + + +def _format_label(value: Any) -> str: + if value is None: + return "alert" + return str(value).replace("_", " ").lower() + + +def _format_action(value: Any) -> str: + normalized = str(value or "").upper() + if normalized == "LOWER": + return "Lower price" + if normalized == "RAISE": + return "Raise price" + if normalized == "REORDER": + return "Restock" + return "Hold" diff --git a/backend/app/services/dashboard/dashboard_service.py b/backend/app/services/dashboard/dashboard_service.py index 5be2e09..b3da65f 100644 --- a/backend/app/services/dashboard/dashboard_service.py +++ b/backend/app/services/dashboard/dashboard_service.py @@ -95,6 +95,9 @@ async def get_summary(self, *, seller_id: str | None = None) -> DashboardSummary average_urgency_score = await self._session.scalar( _avg_query(urgency_score_expr, seller_id=seller_id) ) + latest_urgency_score = await self._session.scalar( + _latest_score_query(urgency_score_expr, seller_id=seller_id) + ) average_price_gap = await self._session.scalar( _avg_query(price_gap_expr, seller_id=seller_id) ) @@ -135,6 +138,7 @@ async def get_summary(self, *, seller_id: str | None = None) -> DashboardSummary raise DashboardDatabaseError("Failed to load dashboard summary.") from exc normalized_urgency = _normalize_score(average_urgency_score) + latest_urgency = _normalize_score(latest_urgency_score) normalized_gap = float(average_price_gap or 0.0) trend_summary = _build_market_trend_summary(normalized_urgency, normalized_gap) @@ -152,6 +156,7 @@ async def get_summary(self, *, seller_id: str | None = None) -> DashboardSummary opportunity_count=opportunity_count or 0, high_risk_stock_count=high_risk_stock_count or 0, average_urgency_score=round(normalized_urgency, 2), + latest_urgency_score=round(latest_urgency, 2), average_competitor_price_gap=round(normalized_gap, 2), market_trend_summary=trend_summary, realtime_activity=RealtimeActivityCounters( @@ -204,6 +209,16 @@ def _avg_query(expression, *, seller_id: str | None): return stmt +def _latest_score_query(expression, *, seller_id: str | None): + stmt = select(expression).select_from(Recommendation) + if seller_id: + stmt = stmt.join(Product, Recommendation.product_id == Product.id).where( + Product.seller_id == seller_id + ) + stmt = stmt.order_by(Recommendation.created_at.desc(), Recommendation.id.desc()).limit(1) + return stmt + + def _build_market_trend_summary(avg_urgency: float, avg_gap: float) -> str: if avg_urgency >= 75: return "High volatility with aggressive competitor pressure." diff --git a/backend/app/services/dashboard/market_trend_service.py b/backend/app/services/dashboard/market_trend_service.py new file mode 100644 index 0000000..70e61c5 --- /dev/null +++ b/backend/app/services/dashboard/market_trend_service.py @@ -0,0 +1,102 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +import logging +import time +from typing import Any + +from sqlalchemy import Float, case, cast, func, select +from sqlalchemy.exc import SQLAlchemyError +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.product import Product +from app.models.recommendation import Recommendation +from app.schemas.dashboard.market_trend_schema import MarketTrendPoint +from app.services.dashboard.exceptions import DashboardDatabaseError + + +class MarketTrendDashboardService: + def __init__(self, session: AsyncSession, *, logger: logging.Logger | None = None) -> None: + self._session = session + self._logger = logger or logging.getLogger(__name__) + + async def get_market_trends( + self, *, seller_id: str, days: int = 14 + ) -> list[MarketTrendPoint]: + start = time.monotonic() + window_start = datetime.now(timezone.utc) - timedelta(days=days) + analytics = Recommendation.recommendation_context["analytics"] + urgency_expr = cast(analytics["urgency_score"].astext, Float) + market_comp_expr = cast(analytics["market_competitiveness_score"].astext, Float) + pricing_pressure_expr = cast(analytics["pricing_pressure_score"].astext, Float) + price_gap_expr = cast(analytics["price_gap_percent"].astext, Float) + out_of_stock_expr = cast(analytics["out_of_stock_competitor_count"].astext, Float) + in_stock_expr = cast(analytics["in_stock_competitor_count"].astext, Float) + total_expr = out_of_stock_expr + in_stock_expr + demand_ratio_expr = case( + (total_expr > 0, out_of_stock_expr / total_expr), + else_=0.0, + ) + bucket = func.date_trunc("day", Recommendation.created_at).label("bucket") + stmt = ( + select( + bucket, + func.avg(urgency_expr).label("urgency_score"), + func.avg(market_comp_expr).label("market_competitiveness_score"), + func.avg(pricing_pressure_expr).label("pricing_pressure_score"), + func.avg(price_gap_expr).label("price_gap"), + func.avg(demand_ratio_expr).label("demand_ratio"), + ) + .join(Product, Recommendation.product_id == Product.id) + .where(Product.seller_id == seller_id) + .where(Recommendation.created_at >= window_start) + .group_by(bucket) + .order_by(bucket) + ) + try: + result = await self._session.execute(stmt) + rows = result.mappings().all() + except SQLAlchemyError as exc: + raise DashboardDatabaseError("Failed to load market trends.") from exc + + points: list[MarketTrendPoint] = [] + for row in rows: + opportunity_raw = row.get("market_competitiveness_score") + if opportunity_raw is None: + opportunity_raw = row.get("pricing_pressure_score") + points.append( + MarketTrendPoint( + timestamp=row["bucket"], + demand_index=_normalize_index(row.get("demand_ratio"), scale=100.0), + urgency_index=_normalize_index(row.get("urgency_score"), scale=100.0), + opportunity_index=_normalize_index(opportunity_raw, scale=100.0), + price_index=_to_float(row.get("price_gap")) or 0.0, + ) + ) + + duration_ms = int((time.monotonic() - start) * 1000) + self._logger.info( + "dashboard_market_trends count=%s duration_ms=%s seller_id=%s", + len(points), + duration_ms, + seller_id, + ) + return points + + +def _normalize_index(value: Any, *, scale: float) -> float: + numeric = _to_float(value) + if numeric is None: + return 0.0 + if 0 <= numeric <= 1: + return numeric * scale + return numeric + + +def _to_float(value: Any) -> float | None: + if value is None: + return None + try: + return float(value) + except (TypeError, ValueError): + return None diff --git a/backend/app/services/dashboard/opportunity_metrics_service.py b/backend/app/services/dashboard/opportunity_metrics_service.py new file mode 100644 index 0000000..54ae1cc --- /dev/null +++ b/backend/app/services/dashboard/opportunity_metrics_service.py @@ -0,0 +1,180 @@ +from __future__ import annotations + +import logging +import time +from typing import Any, Mapping + +from sqlalchemy import select +from sqlalchemy.exc import SQLAlchemyError +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.product import Product +from app.models.recommendation import Recommendation, RecommendationAction, RecommendationConfidence +from app.schemas.dashboard.opportunity_metrics_schema import OpportunityMetricRecord +from app.services.dashboard.exceptions import DashboardDatabaseError + + +class OpportunityMetricsService: + def __init__(self, session: AsyncSession, *, logger: logging.Logger | None = None) -> None: + self._session = session + self._logger = logger or logging.getLogger(__name__) + + async def get_metrics(self, *, seller_id: str) -> list[OpportunityMetricRecord]: + start = time.monotonic() + stmt = ( + select(Recommendation, Product.title.label("product_title")) + .join(Product, Recommendation.product_id == Product.id) + .where(Product.seller_id == seller_id) + .order_by(Recommendation.created_at.desc(), Recommendation.id.desc()) + .limit(2) + ) + try: + result = await self._session.execute(stmt) + rows = result.all() + except SQLAlchemyError as exc: + raise DashboardDatabaseError("Failed to load opportunity metrics.") from exc + + latest = rows[0] if rows else None + previous = rows[1] if len(rows) > 1 else None + + latest_rec = latest[0] if latest else None + prev_rec = previous[0] if previous else None + + _, latest_analytics = _split_context_snapshot( + latest_rec.recommendation_context if latest_rec else None + ) + _, prev_analytics = _split_context_snapshot( + prev_rec.recommendation_context if prev_rec else None + ) + + revenue_now = _to_float(latest_rec.revenue_opportunity) if latest_rec else 0.0 + revenue_prev = _to_float(prev_rec.revenue_opportunity) if prev_rec else None + confidence_now = _confidence_score(latest_rec.confidence) if latest_rec else 0.0 + confidence_prev = _confidence_score(prev_rec.confidence) if prev_rec else None + opportunity_now = _normalize_index( + latest_analytics.get("market_competitiveness_score") + or latest_analytics.get("opportunity_score") + ) + opportunity_prev = _normalize_index( + prev_analytics.get("market_competitiveness_score") + or prev_analytics.get("opportunity_score") + ) + + metrics = [ + OpportunityMetricRecord( + id="revenue-opportunity", + label="Revenue opportunity", + value=revenue_now or 0.0, + delta=_delta(revenue_now, revenue_prev), + trend=_trend(revenue_now, revenue_prev), + ), + OpportunityMetricRecord( + id="confidence", + label="Confidence", + value=confidence_now, + delta=_delta(confidence_now, confidence_prev), + trend=_trend(confidence_now, confidence_prev), + ), + OpportunityMetricRecord( + id="action-recommendation", + label="Action recommendation", + value=_action_label(latest_rec.action if latest_rec else None), + delta=_action_delta( + latest_rec.action if latest_rec else None, + prev_rec.action if prev_rec else None, + ), + trend="flat", + ), + OpportunityMetricRecord( + id="opportunity-score", + label="Opportunity score", + value=opportunity_now, + delta=_delta(opportunity_now, opportunity_prev), + trend=_trend(opportunity_now, opportunity_prev), + ), + ] + + duration_ms = int((time.monotonic() - start) * 1000) + self._logger.info( + "dashboard_opportunity_metrics duration_ms=%s seller_id=%s", + duration_ms, + seller_id, + ) + return metrics + + +def _split_context_snapshot(payload: Mapping[str, Any] | None) -> tuple[Mapping[str, Any], Mapping[str, Any]]: + if not isinstance(payload, Mapping): + return {}, {} + if "context" in payload or "analytics" in payload: + context = payload.get("context") + analytics = payload.get("analytics") + return ( + context if isinstance(context, Mapping) else {}, + analytics if isinstance(analytics, Mapping) else {}, + ) + return payload, {} + + +def _to_float(value: Any) -> float: + if value is None: + return 0.0 + try: + return float(value) + except (TypeError, ValueError): + return 0.0 + + +def _confidence_score(value: RecommendationConfidence | None) -> float: + if value is RecommendationConfidence.HIGH: + return 90.0 + if value is RecommendationConfidence.MEDIUM: + return 70.0 + if value is RecommendationConfidence.LOW: + return 50.0 + return 0.0 + + +def _normalize_index(value: Any) -> float: + numeric = _to_float(value) + if 0 <= numeric <= 1: + return numeric * 100.0 + return numeric + + +def _delta(current: float, previous: float | None) -> float | None: + if previous is None: + return None + return round(current - previous, 2) + + +def _trend(current: float, previous: float | None) -> str: + if previous is None: + return "flat" + if current > previous: + return "up" + if current < previous: + return "down" + return "flat" + + +def _action_label(action: RecommendationAction | None) -> str: + if action is RecommendationAction.LOWER: + return "Lower price" + if action is RecommendationAction.RAISE: + return "Raise price" + if action is RecommendationAction.REORDER: + return "Restock" + if action is RecommendationAction.HOLD: + return "Hold" + return "Hold" + + +def _action_delta( + current: RecommendationAction | None, previous: RecommendationAction | None +) -> str: + if not previous: + return "Latest" + if current == previous: + return "Stable" + return "Updated" diff --git a/backend/app/services/dashboard/price_trend_service.py b/backend/app/services/dashboard/price_trend_service.py new file mode 100644 index 0000000..f90d5e7 --- /dev/null +++ b/backend/app/services/dashboard/price_trend_service.py @@ -0,0 +1,98 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +import logging +import time +from typing import Any +from uuid import UUID + +from sqlalchemy import func, select +from sqlalchemy.exc import SQLAlchemyError +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.competitor import Competitor +from app.models.price_history import PriceHistory +from app.models.product import Product +from app.schemas.dashboard.price_trend_schema import CompetitorPriceTrendPoint +from app.services.dashboard.exceptions import DashboardDatabaseError + + +class CompetitorPriceTrendService: + def __init__(self, session: AsyncSession, *, logger: logging.Logger | None = None) -> None: + self._session = session + self._logger = logger or logging.getLogger(__name__) + + async def get_price_trend( + self, + *, + seller_id: str, + product_id: str | None = None, + days: int = 14, + ) -> list[CompetitorPriceTrendPoint]: + start = time.monotonic() + window_start = datetime.now(timezone.utc) - timedelta(days=days) + bucket = func.date_trunc("day", PriceHistory.observed_at).label("bucket") + + stmt = ( + select( + bucket, + func.min(PriceHistory.observed_price).label("competitor_price"), + func.avg(PriceHistory.observed_price).label("market_average"), + func.avg(Product.selling_price).label("seller_price"), + ) + .join(Competitor, PriceHistory.competitor_id == Competitor.id) + .join(Product, Competitor.product_id == Product.id) + .where(Product.seller_id == seller_id) + .where(PriceHistory.observed_at >= window_start) + .group_by(bucket) + .order_by(bucket) + ) + + product_uuid = _to_uuid(product_id) + if product_id and not product_uuid: + return [] + if product_uuid: + stmt = stmt.where(Product.id == product_uuid) + + try: + result = await self._session.execute(stmt) + rows = result.mappings().all() + except SQLAlchemyError as exc: + raise DashboardDatabaseError("Failed to load competitor price trends.") from exc + + points = [ + CompetitorPriceTrendPoint( + timestamp=row["bucket"], + competitor_price=_to_float(row.get("competitor_price")) or 0.0, + seller_price=_to_float(row.get("seller_price")) or 0.0, + market_average=_to_float(row.get("market_average")) or 0.0, + ) + for row in rows + ] + duration_ms = int((time.monotonic() - start) * 1000) + self._logger.info( + "dashboard_price_trend count=%s duration_ms=%s seller_id=%s product_id=%s", + len(points), + duration_ms, + seller_id, + product_id, + ) + return points + + +def _to_float(value: Any) -> float | None: + if value is None: + return None + try: + return float(value) + except (TypeError, ValueError): + return None + + +def _to_uuid(value: str | None) -> UUID | None: + if not value: + return None + try: + return UUID(value) + except ValueError: + return None diff --git a/backend/app/services/realtime/realtime_service.py b/backend/app/services/realtime/realtime_service.py index 72f2fe2..9ab1013 100644 --- a/backend/app/services/realtime/realtime_service.py +++ b/backend/app/services/realtime/realtime_service.py @@ -285,7 +285,9 @@ def _default_rooms( def _default_dashboard_rooms( self, seller_id: UUID | None, product_id: UUID | None ) -> set[str]: - rooms = {DASHBOARD_ROOM, ACTIVITY_ROOM} + rooms: set[str] = set() + if seller_id is None: + rooms.update({DASHBOARD_ROOM, ACTIVITY_ROOM}) if seller_id: rooms.add(seller_room(str(seller_id))) if product_id: @@ -295,7 +297,9 @@ def _default_dashboard_rooms( def _default_alert_rooms( self, seller_id: UUID | None, product_id: UUID | None ) -> set[str]: - rooms = {ALERTS_ROOM, ACTIVITY_ROOM, DASHBOARD_ROOM} + rooms: set[str] = set() + if seller_id is None: + rooms.update({ALERTS_ROOM, ACTIVITY_ROOM, DASHBOARD_ROOM}) if seller_id: rooms.add(seller_room(str(seller_id))) if product_id: diff --git a/backend/app/services/recommendation_service.py b/backend/app/services/recommendation_service.py index d9c6457..1b31a8f 100644 --- a/backend/app/services/recommendation_service.py +++ b/backend/app/services/recommendation_service.py @@ -25,13 +25,28 @@ from app.models.product import Product from app.models.recommendation import Recommendation, RecommendationAction, RecommendationConfidence from app.models.scrape_run import ScrapeRun +from app.recommendations.opportunity_calculator import OpportunityCalculator, OpportunityInput from app.schemas.recommendation import ( RecommendationAnalyticsResponse, RecommendationGenerateResponse, RecommendationResponse, ) -from app.services.dashboard.events import DashboardEvent, DashboardEventEmitter -from app.websocket.socket_manager import socket_manager +from app.schemas.dashboard.summary_schema import DashboardSummaryResponse +from app.services.dashboard.events import DashboardEventEmitter +from app.services.dashboard.dashboard_service import DashboardService +from app.services.realtime import RealtimeService +from app.websocket.socket_events import ( + DashboardSummaryPayload, + DashboardSummarySnapshot, + MarketOpportunityPayload, + MarketOpportunitySnapshot, + RecommendationCreatedPayload, + RecommendationSnapshot, + UrgencyScoreUpdatedPayload, +) +from app.websocket.socket_manager import get_socket_emitter, socket_manager +from app.websocket.socket_rooms import seller_room +from app.db.async_session import AsyncSessionLocal class RecommendationServiceError(Exception): @@ -97,6 +112,8 @@ def __init__( self._context_builder = RecommendationContextBuilder() self._ai_generator = RecommendationGenerator() self._event_emitter = event_emitter or socket_manager + self._realtime: RealtimeService | None = None + self._opportunity_calculator = OpportunityCalculator() self._stock_risk_threshold = ( stock_risk_threshold if stock_risk_threshold is not None @@ -156,6 +173,7 @@ async def generate_recommendation( await self._emit_recommendation_created( seller_id=str(seller_id), recommendation=recommendation, + product=product, analytics=analytics, ) response = self.build_response(recommendation, analytics) @@ -386,27 +404,162 @@ def _log_event(self, level: int, event: str, **fields: Any) -> None: payload = " ".join(f"{key}={value}" for key, value in safe_fields.items()) self._logger.log(level, f"{event} {payload}".strip()) + def _resolve_realtime(self) -> RealtimeService | None: + if self._realtime is not None: + return self._realtime + try: + emitter = get_socket_emitter() + except RuntimeError: + self._logger.debug("socket_emitter_unavailable") + return None + self._realtime = RealtimeService(emitter, logging.getLogger("app.websocket.realtime")) + return self._realtime + + async def _emit_urgency_score_updated( + self, + *, + realtime: RealtimeService, + seller_id: UUID, + recommendation: Recommendation, + urgency_score: float, + ) -> None: + previous_score = self._load_previous_urgency_score( + seller_id=seller_id, + recommendation_id=recommendation.id, + product_id=recommendation.product_id, + ) + delta = int(round(urgency_score - previous_score)) if previous_score is not None else None + payload = UrgencyScoreUpdatedPayload( + urgency_score=int(round(urgency_score)), + previous_score=int(round(previous_score)) if previous_score is not None else None, + delta=delta, + updated_at=recommendation.created_at or datetime.now(timezone.utc), + reason="Recommendation generated.", + ) + await realtime.emit_urgency_score_updated( + payload, seller_id=seller_id, rooms=[seller_room(str(seller_id))] + ) + + async def _emit_dashboard_summary_updated( + self, + *, + realtime: RealtimeService, + seller_id: UUID, + ) -> None: + async with AsyncSessionLocal() as session: + service = DashboardService(session, logger=self._logger) + summary = await service.get_summary(seller_id=str(seller_id)) + payload = DashboardSummaryPayload( + summary=_build_dashboard_summary_snapshot(summary) + ) + await realtime.emit_dashboard_summary_updated( + payload, seller_id=seller_id, rooms=[seller_room(str(seller_id))] + ) + + async def _emit_market_opportunity_detected( + self, + *, + realtime: RealtimeService, + seller_id: UUID, + recommendation: Recommendation, + product: Product, + analytics: Mapping[str, Any], + ) -> None: + if recommendation.revenue_opportunity <= 0: + return + snapshot = _build_market_opportunity_snapshot( + recommendation=recommendation, + product=product, + analytics=analytics, + calculator=self._opportunity_calculator, + ) + if not snapshot: + return + payload = MarketOpportunityPayload(opportunity=snapshot) + await realtime.emit_market_opportunity_detected( + payload, + seller_id=seller_id, + product_id=recommendation.product_id, + rooms=[seller_room(str(seller_id))], + ) + + def _load_previous_urgency_score( + self, + *, + seller_id: UUID, + recommendation_id: UUID, + product_id: UUID, + ) -> float | None: + statement = ( + select(Recommendation.recommendation_context) + .join(Product, Recommendation.product_id == Product.id) + .where( + Recommendation.product_id == product_id, + Recommendation.id != recommendation_id, + Product.seller_id == seller_id, + ) + .order_by(Recommendation.created_at.desc(), Recommendation.id.desc()) + .limit(1) + ) + payload = self._db.execute(statement).scalar_one_or_none() + if not payload: + return None + _, analytics = self._split_context_snapshot(payload) + return _normalize_score(analytics.get("urgency_score")) + async def _emit_recommendation_created( self, *, seller_id: str, recommendation: Recommendation, + product: Product, analytics: Mapping[str, Any], ) -> None: - payload = { - "recommendation_id": str(recommendation.id), - "product_id": str(recommendation.product_id), - "seller_id": seller_id, - "action": recommendation.action.value, - "confidence": recommendation.confidence.value, - "revenue_opportunity": float(recommendation.revenue_opportunity), - "urgency_score": analytics.get("urgency_score"), - "stock_risk_score": analytics.get("stock_risk_score"), - "created_at": recommendation.created_at.isoformat() - if recommendation.created_at - else None, - } - await self._event_emitter.emit(DashboardEvent.recommendation_created, payload) + realtime = self._resolve_realtime() + if not realtime: + self._logger.debug("realtime_emit_skipped reason=socket_unavailable") + return + seller_uuid = UUID(seller_id) + urgency_score = _normalize_score(analytics.get("urgency_score")) + recommendation_snapshot = RecommendationSnapshot( + id=recommendation.id, + title=_format_recommendation_title(recommendation.action), + summary=recommendation.reasoning_bn, + action=_map_action(recommendation.action), + status="new", + priority=_map_priority(urgency_score), + confidence=_map_confidence(recommendation.confidence), + expected_lift=_expected_lift(urgency_score), + product=product.title, + competitor=None, + created_at=recommendation.created_at or datetime.now(timezone.utc), + ) + payload = RecommendationCreatedPayload( + recommendation=recommendation_snapshot, + urgency_score=int(round(urgency_score)), + confidence=recommendation.confidence.value, + opportunity_score=float(recommendation.revenue_opportunity), + ) + await realtime.emit_recommendation_created( + payload, + seller_id=seller_uuid, + product_id=recommendation.product_id, + rooms=[seller_room(str(seller_uuid))], + ) + await self._emit_urgency_score_updated( + realtime=realtime, + seller_id=seller_uuid, + recommendation=recommendation, + urgency_score=urgency_score, + ) + await self._emit_dashboard_summary_updated(realtime=realtime, seller_id=seller_uuid) + await self._emit_market_opportunity_detected( + realtime=realtime, + seller_id=seller_uuid, + recommendation=recommendation, + product=product, + analytics=analytics, + ) async def _emit_stock_risk_event( self, @@ -420,14 +573,12 @@ async def _emit_stock_risk_event( return if stock_risk_score < self._stock_risk_threshold: return - payload = { - "product_id": str(product.id), - "seller_id": seller_id, - "stock_risk_score": stock_risk_score, - "estimated_days_of_stock_left": analytics.get("estimated_days_of_stock_left"), - "stock_quantity": product.stock_quantity, - } - await self._event_emitter.emit(DashboardEvent.stock_risk_detected, payload) + self._logger.info( + "stock_risk_detected seller_id=%s product_id=%s score=%s", + seller_id, + product.id, + stock_risk_score, + ) def _build_analytics_response( self, @@ -507,3 +658,210 @@ def _get_env_float(key: str, default: float) -> float: return float(value) except ValueError: return default + + +def _normalize_score(value: Any) -> float: + if value is None: + return 0.0 + try: + numeric = float(value) + except (TypeError, ValueError): + return 0.0 + if 0 <= numeric <= 1: + return numeric * 100.0 + return numeric + + +def _map_action(action: RecommendationAction) -> str: + if action is RecommendationAction.LOWER: + return "undercut" + if action is RecommendationAction.RAISE: + return "raise_price" + if action is RecommendationAction.REORDER: + return "restock" + return "hold" + + +def _format_recommendation_title(action: RecommendationAction) -> str: + if action is RecommendationAction.LOWER: + return "Lower price recommendation" + if action is RecommendationAction.RAISE: + return "Raise price recommendation" + if action is RecommendationAction.REORDER: + return "Restock recommendation" + return "Hold price recommendation" + + +def _map_priority(urgency_score: float) -> str: + if urgency_score >= 80: + return "critical" + if urgency_score >= 60: + return "high" + if urgency_score >= 40: + return "medium" + return "low" + + +def _map_confidence(confidence: RecommendationConfidence) -> float: + if confidence is RecommendationConfidence.HIGH: + return 0.9 + if confidence is RecommendationConfidence.MEDIUM: + return 0.7 + if confidence is RecommendationConfidence.LOW: + return 0.5 + return 0.6 + + +def _expected_lift(urgency_score: float) -> float: + return min(max(urgency_score / 100.0, 0.0), 1.0) + + +def _build_dashboard_summary_snapshot(summary: DashboardSummaryResponse) -> DashboardSummarySnapshot: + insights = [] + if summary.market_trend_summary: + severity = "high" if summary.average_urgency_score >= 75 else "medium" + if summary.average_urgency_score < 55: + severity = "low" + insights.append( + { + "id": "market-trend", + "title": "Market trend", + "summary": summary.market_trend_summary, + "severity": severity, + } + ) + activity = summary.realtime_activity + insights.append( + { + "id": "realtime-activity", + "title": "Last hour activity", + "summary": ( + f"Recommendations {activity.recommendations_last_hour}, alerts " + f"{activity.alerts_last_hour}, competitor updates " + f"{activity.competitor_updates_last_hour}." + ), + "severity": "low", + } + ) + system_status = "degraded" if summary.high_risk_stock_count > 0 or summary.active_alerts > 0 else "healthy" + return DashboardSummarySnapshot( + kpis=[], + insights=insights, + active_alerts=summary.active_alerts, + live_competitors=summary.total_monitored_products, + ai_recommendations=summary.realtime_activity.recommendations_last_hour, + revenue_opportunities=summary.opportunity_count, + urgency_score=int(round(summary.latest_urgency_score)), + system_status=system_status, + last_synced_at=summary.generated_at, + ) + + +def _to_float(value: Any) -> float | None: + if value is None: + return None + try: + return float(value) + except (TypeError, ValueError): + return None + + +def _to_int(value: Any) -> int | None: + if value is None: + return None + try: + return int(value) + except (TypeError, ValueError): + return None + + +def _build_market_opportunity_snapshot( + *, + recommendation: Recommendation, + product: Product, + analytics: Mapping[str, Any], + calculator: OpportunityCalculator, +) -> MarketOpportunitySnapshot | None: + context, analytics_snapshot = RecommendationService._split_context_snapshot( + recommendation.recommendation_context + ) + merged_analytics = {**analytics_snapshot, **dict(analytics)} + seller_price = _to_float(context.get("seller_price")) or float(product.selling_price) + recommended_price = _to_float(context.get("recommended_price")) + stock_quantity = _to_int(context.get("seller_stock")) or int(product.stock_quantity or 0) + margin_percent = _to_float(merged_analytics.get("estimated_margin_percent")) + if margin_percent is None: + margin_percent = _compute_margin_percent( + float(product.selling_price), float(product.cost_price) + ) + opportunity_score = _to_float(merged_analytics.get("market_competitiveness_score")) or 0.0 + demand_indicator = _to_float(merged_analytics.get("demand_indicator")) or 0.0 + competitor_oos_ratio = _compute_oos_ratio(merged_analytics) + inventory_turnover_days = _to_float(context.get("inventory_turnover_days")) + + try: + result = calculator.calculate( + OpportunityInput( + seller_price=seller_price, + recommended_price=recommended_price, + current_stock=stock_quantity, + margin_percent=margin_percent, + opportunity_score=opportunity_score, + demand_indicator=demand_indicator, + competitor_oos_ratio=competitor_oos_ratio, + inventory_turnover_days=inventory_turnover_days, + ) + ) + except Exception: + return None + + description = _build_opportunity_description(result) + time_window = ( + recommendation.created_at.date().isoformat() + if recommendation.created_at + else "Rolling" + ) + category = _resolve_opportunity_category( + competitor_oos_ratio, result.projected_demand_shift + ) + urgency_score = int(round(_normalize_score(merged_analytics.get("urgency_score")))) + + return MarketOpportunitySnapshot( + id=recommendation.id, + title=product.title, + description=description, + potential_revenue=float(recommendation.revenue_opportunity), + probability=competitor_oos_ratio, + urgency_score=urgency_score, + time_window=time_window, + category=category, + ) + + +def _build_opportunity_description(result) -> str: + units = round(result.inventory_capture_units, 2) + profit = round(result.estimated_profit_impact, 2) + return f"Capture {units} units, est. profit impact {profit}" + + +def _resolve_opportunity_category(competitor_oos_ratio: float, demand_shift: float) -> str: + if competitor_oos_ratio >= 0.5: + return "inventory" + if demand_shift >= 15: + return "promotion" + return "pricing" + + +def _compute_oos_ratio(analytics: Mapping[str, Any]) -> float: + out_of_stock = _to_int(analytics.get("out_of_stock_competitor_count")) or 0 + in_stock = _to_int(analytics.get("in_stock_competitor_count")) or 0 + total = out_of_stock + in_stock + if total == 0: + return 0.0 + return out_of_stock / total + + +def _compute_margin_percent(selling_price: float, cost_price: float) -> float: + if selling_price <= 0: + return 0.0 + return ((selling_price - cost_price) / selling_price) * 100.0 diff --git a/backend/app/services/scrape_service.py b/backend/app/services/scrape_service.py index 1a602fd..42b99fb 100644 --- a/backend/app/services/scrape_service.py +++ b/backend/app/services/scrape_service.py @@ -9,6 +9,8 @@ from sqlalchemy import select from sqlalchemy.orm import Session +from app.models.competitor import Competitor +from app.models.price_history import PriceHistory from app.models.product import Product from app.models.scrape_run import ScrapeRun from app.models.seller import Seller @@ -19,8 +21,12 @@ from app.schemas.scrape import ScrapeRunResponse from app.schemas.competitor import CompetitorSnapshot from app.services.competitor_service import CompetitorService -from app.services.dashboard.events import DashboardEvent, DashboardEventEmitter -from app.websocket.socket_manager import socket_manager +from app.services.dashboard.events import DashboardEventEmitter +from app.services.realtime import RealtimeService +from app.websocket.socket_events import CompetitorSnapshot as RealtimeCompetitorSnapshot +from app.websocket.socket_events import CompetitorUpdatedPayload +from app.websocket.socket_manager import get_socket_emitter +from app.websocket.socket_rooms import seller_room class ScrapeProductNotFoundError(Exception): @@ -48,7 +54,8 @@ def __init__( ) -> None: self._db = db self._competitor_service = CompetitorService(db) - self._event_emitter = event_emitter or socket_manager + self._event_emitter = event_emitter + self._realtime: RealtimeService | None = None async def run_manual_scrape( self, @@ -94,7 +101,7 @@ async def run_manual_scrape( matched_results = len(filtered) filtered_results = max(0, competitors_found - matched_results) - self._competitor_service.store_competitors( + competitor_rows = self._competitor_service.store_competitors( scrape_run=scrape_run, product=product, marketplace=scrape_response.marketplace, @@ -113,6 +120,7 @@ async def run_manual_scrape( await self._emit_competitor_updated( scrape_run=scrape_run, product=product, + competitors=competitor_rows, matched_results=matched_results, marketplace=scrape_response.marketplace, scraped_at=scrape_response.scraped_at, @@ -227,16 +235,107 @@ async def _emit_competitor_updated( *, scrape_run: ScrapeRun, product: Product, + competitors: list[Competitor], matched_results: int, marketplace: str, scraped_at: datetime | None, ) -> None: - payload = { - "scrape_run_id": str(scrape_run.id), - "product_id": str(product.id), - "seller_id": str(product.seller_id), - "marketplace": marketplace, - "matched_results": matched_results, - "scraped_at": scraped_at.isoformat() if scraped_at else None, - } - await self._event_emitter.emit(DashboardEvent.competitor_updated, payload) + realtime = self._resolve_realtime() + if not realtime: + return + if not competitors: + return + primary = self._select_primary_competitor(competitors) + if not primary or primary.competitor_price is None: + return + seller_price = float(product.selling_price) + competitor_price = float(primary.competitor_price) + price_gap = seller_price - competitor_price + trend = self._resolve_price_trend(primary.id) + snapshot = RealtimeCompetitorSnapshot( + id=primary.id, + product=product.title, + competitor=primary.competitor_title, + competitor_price=competitor_price, + seller_price=seller_price, + price_gap=price_gap, + stock_status=primary.stock_status or "in_stock", + urgency_score=_compute_competitor_urgency( + price_gap_percent=_compute_price_gap_percent(seller_price, competitor_price), + stock_status=primary.stock_status, + ), + recommendation_action=_recommendation_action_label(seller_price, competitor_price), + trend=trend, + last_updated=primary.updated_at or scraped_at or datetime.now(timezone.utc), + status="active", + ) + payload = CompetitorUpdatedPayload(competitor=snapshot) + await realtime.emit_competitor_updated( + payload, + seller_id=product.seller_id, + product_id=product.id, + rooms=[seller_room(str(product.seller_id))], + ) + + def _resolve_realtime(self) -> RealtimeService | None: + if self._realtime is not None: + return self._realtime + try: + emitter = get_socket_emitter() + except RuntimeError: + return None + self._realtime = RealtimeService(emitter) + return self._realtime + + @staticmethod + def _select_primary_competitor(competitors: list[Competitor]) -> Competitor | None: + if not competitors: + return None + return max( + competitors, + key=lambda item: item.similarity_score or 0.0, + ) + + def _resolve_price_trend(self, competitor_id: UUID) -> str: + stmt = ( + select(PriceHistory.observed_price) + .where(PriceHistory.competitor_id == competitor_id) + .order_by(PriceHistory.observed_at.desc()) + .limit(2) + ) + prices = self._db.execute(stmt).scalars().all() + if len(prices) < 2: + return "flat" + latest, previous = prices[0], prices[1] + if latest is None or previous is None: + return "flat" + if latest > previous: + return "up" + if latest < previous: + return "down" + return "flat" + + +def _compute_price_gap_percent(seller_price: float, competitor_price: float) -> float: + if competitor_price <= 0: + return 0.0 + return ((seller_price - competitor_price) / competitor_price) * 100.0 + + +def _compute_competitor_urgency( + *, price_gap_percent: float, stock_status: str | None +) -> int: + base = min(abs(price_gap_percent) * 2, 100) + if stock_status == "out_of_stock": + return int(max(base, 85)) + if stock_status == "low_stock": + return int(max(base, 65)) + return int(round(base)) + + +def _recommendation_action_label(seller_price: float, competitor_price: float) -> str: + if seller_price > competitor_price: + return "Lower price" + if seller_price < competitor_price: + return "Raise price" + return "Hold" diff --git a/backend/app/websocket/connection_manager.py b/backend/app/websocket/connection_manager.py index cf2ddcc..9ab5ac9 100644 --- a/backend/app/websocket/connection_manager.py +++ b/backend/app/websocket/connection_manager.py @@ -56,8 +56,10 @@ async def register_connection( async with self._lock: restored_rooms: set[str] = set() previous = self._sessions.pop(session_id, None) - if previous: + if previous and (seller_id is None or previous.seller_id == seller_id): restored_rooms = set(previous.rooms) + if seller_id is None: + seller_id = previous.seller_id self._connections[sid] = ConnectionState( sid=sid, session_id=session_id, @@ -95,6 +97,13 @@ async def get_rooms(self, sid: str) -> set[str]: return set() return set(state.rooms) + async def get_seller_id(self, sid: str) -> str | None: + async with self._lock: + state = self._connections.get(sid) + if not state: + return None + return state.seller_id + async def mark_disconnected(self, sid: str) -> SessionState | None: async with self._lock: state = self._connections.pop(sid, None) diff --git a/backend/app/websocket/socket_manager.py b/backend/app/websocket/socket_manager.py index f977733..6b8840b 100644 --- a/backend/app/websocket/socket_manager.py +++ b/backend/app/websocket/socket_manager.py @@ -6,14 +6,20 @@ import logging import time from typing import Any, Literal -from uuid import uuid4 +from uuid import UUID, uuid4 import socketio from fastapi import FastAPI +from jwt import ExpiredSignatureError, InvalidTokenError from pydantic import ValidationError +from sqlalchemy import select from app.core.config import settings +from app.core.security import decode_access_token from app.core.socketio_config import build_socketio_resources +from app.db.async_session import AsyncSessionLocal +from app.models.product import Product +from app.models.seller import Seller from app.websocket.connection_manager import ConnectionManager from app.websocket.socket_emitter import SocketEmitter from app.websocket.socket_events import ( @@ -30,9 +36,11 @@ ACTIVITY_ROOM, ALERTS_ROOM, DASHBOARD_ROOM, + RoomScope, RoomSubscriptionRequest, - resolve_rooms, + resolve_room, seller_room, + product_room, ) @@ -54,8 +62,10 @@ def register_handlers(self) -> None: @self._server.event async def connect(sid: str, _environ: dict, auth: dict | None) -> bool: auth = auth or {} - session_id = auth.get("session_id") or str(uuid4()) - seller_id = auth.get("seller_id") + auth_result = await self._authenticate(auth) + if not auth_result: + return False + seller_id, session_id = auth_result restored_rooms = await self._connections.register_connection( sid, session_id, seller_id=seller_id ) @@ -86,6 +96,14 @@ async def disconnect(sid: str) -> None: @self._server.on("subscribe_rooms") async def subscribe_rooms(sid: str, payload: dict | None) -> None: await self._connections.mark_seen(sid) + seller_id = await self._connections.get_seller_id(sid) + if not seller_id: + await self._emit_error( + sid, + code="unauthorized", + message="Authentication required to subscribe to rooms.", + ) + return try: request = RoomSubscriptionRequest.model_validate(payload or {}) except ValidationError as exc: @@ -96,7 +114,16 @@ async def subscribe_rooms(sid: str, payload: dict | None) -> None: details={"error": str(exc)}, ) return - rooms = resolve_rooms(request) + rooms, rejected = await self._resolve_authorized_rooms( + request, seller_id=seller_id, include_seller_room=True + ) + if rejected: + await self._emit_error( + sid, + code="unauthorized_room", + message="One or more requested rooms are not authorized.", + details={"rooms": ",".join(rejected)}, + ) for room in rooms: await self._server.enter_room(sid, room) active_rooms = await self._connections.add_rooms(sid, rooms) @@ -106,6 +133,14 @@ async def subscribe_rooms(sid: str, payload: dict | None) -> None: @self._server.on("unsubscribe_rooms") async def unsubscribe_rooms(sid: str, payload: dict | None) -> None: await self._connections.mark_seen(sid) + seller_id = await self._connections.get_seller_id(sid) + if not seller_id: + await self._emit_error( + sid, + code="unauthorized", + message="Authentication required to unsubscribe from rooms.", + ) + return try: request = RoomSubscriptionRequest.model_validate(payload or {}) except ValidationError as exc: @@ -116,7 +151,16 @@ async def unsubscribe_rooms(sid: str, payload: dict | None) -> None: details={"error": str(exc)}, ) return - rooms = resolve_rooms(request) + rooms, rejected = await self._resolve_authorized_rooms( + request, seller_id=seller_id, include_seller_room=False + ) + if rejected: + await self._emit_error( + sid, + code="unauthorized_room", + message="One or more requested rooms are not authorized.", + details={"rooms": ",".join(rejected)}, + ) for room in rooms: await self._server.leave_room(sid, room) active_rooms = await self._connections.remove_rooms(sid, rooms) @@ -126,14 +170,8 @@ async def unsubscribe_rooms(sid: str, payload: dict | None) -> None: @self._server.on("subscribe_war_room") async def subscribe_war_room(sid: str, payload: dict | None) -> None: await self._connections.mark_seen(sid) - seller_id = None - if payload: - seller_id = payload.get("sellerId") or payload.get("seller_id") - rooms = { - DASHBOARD_ROOM, - ALERTS_ROOM, - ACTIVITY_ROOM, - } + seller_id = await self._connections.get_seller_id(sid) + rooms = {DASHBOARD_ROOM, ALERTS_ROOM, ACTIVITY_ROOM} if seller_id: rooms.add(seller_room(seller_id)) for room in rooms: @@ -164,6 +202,84 @@ async def client_ping(sid: str, payload: dict | None) -> None: ) await self._emitter.emit_event(event, to=sid) + async def _authenticate(self, auth: dict[str, Any]) -> tuple[str, str] | None: + token = auth.get("token") + if not token: + self._logger.warning("socket_auth_missing") + return None + try: + payload = decode_access_token(token) + except ExpiredSignatureError: + self._logger.warning("socket_auth_expired") + return None + except InvalidTokenError: + self._logger.warning("socket_auth_invalid") + return None + subject = payload.get("sub") + if not subject: + self._logger.warning("socket_auth_subject_missing") + return None + try: + seller_uuid = UUID(subject) + except ValueError: + self._logger.warning("socket_auth_subject_invalid") + return None + async with AsyncSessionLocal() as session: + seller = await session.get(Seller, seller_uuid) + if not seller: + self._logger.warning("socket_auth_seller_not_found seller_id=%s", seller_uuid) + return None + session_id = auth.get("session_id") or str(uuid4()) + return str(seller_uuid), session_id + + async def _resolve_authorized_rooms( + self, + request: RoomSubscriptionRequest, + *, + seller_id: str, + include_seller_room: bool, + ) -> tuple[set[str], list[str]]: + rooms: set[str] = set() + rejected: list[str] = [] + if include_seller_room: + rooms.add(seller_room(seller_id)) + for target in request.rooms: + if target.scope is RoomScope.SELLER: + if target.id and target.id != seller_id: + rejected.append(resolve_room(target) or "seller") + continue + rooms.add(seller_room(seller_id)) + continue + if target.scope is RoomScope.PRODUCT: + if not target.id: + continue + if await self._owns_product(seller_id, target.id): + rooms.add(product_room(target.id)) + else: + rejected.append(resolve_room(target) or f"product:{target.id}") + continue + if target.scope in {RoomScope.DASHBOARD, RoomScope.ALERTS, RoomScope.ACTIVITY}: + resolved = resolve_room(target) + if resolved: + rooms.add(resolved) + continue + if target.scope is RoomScope.CUSTOM: + rejected.append(resolve_room(target) or "custom") + return rooms, rejected + + async def _owns_product(self, seller_id: str, product_id: str) -> bool: + try: + product_uuid = UUID(product_id) + except ValueError: + return False + async with AsyncSessionLocal() as session: + stmt = select(Product.id).where( + Product.id == product_uuid, + Product.seller_id == seller_id, + ) + result = await session.execute(stmt) + return result.scalar_one_or_none() is not None + async def start(self) -> None: if self._heartbeat_task: return diff --git a/backend/dockerfile b/backend/dockerfile new file mode 100644 index 0000000..fc6e641 --- /dev/null +++ b/backend/dockerfile @@ -0,0 +1,11 @@ +FROM python:3.12-slim + +WORKDIR /app + +COPY requirements.txt . + +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["uvicorn","app.main:app","--host","0.0.0.0","--port","8000"] \ No newline at end of file