Spaces:
Running
Running
| from __future__ import annotations | |
| from collections import defaultdict | |
| from datetime import datetime, timedelta, timezone | |
| from typing import Any | |
| from zoneinfo import ZoneInfo, ZoneInfoNotFoundError | |
| from app.analytics.errors import AnalyticsValidationError | |
| from app.analytics.models import AnalyticsSyncRun | |
| from app.analytics.repository import AnalyticsRepository | |
| from app.analytics.schemas import ( | |
| AnalyticsCapabilities, | |
| AnalyticsFreshness, | |
| AnalyticsMetric, | |
| AnalyticsOverview, | |
| AnalyticsPostList, | |
| AnalyticsQuery, | |
| AnalyticsSyncRequest, | |
| AnalyticsSyncRunView, | |
| AnalyticsTimeseries, | |
| ) | |
| from app.core.logger import get_logger | |
| from app.social.providers.registry import ProviderRegistry | |
| from app.social.services.analytics_service import AnalyticsService as SocialAnalyticsService | |
| from app.social.services.audit_service import SocialAuditService | |
| logger = get_logger(__name__) | |
| COMMON_METRICS = ("views", "impressions", "likes", "comments", "shares", "engagement_rate") | |
| PROVIDER_METRICS: dict[str, tuple[str, ...]] = { | |
| "facebook": ("views", "impressions", "likes", "comments", "shares"), | |
| "instagram": ("views", "likes", "comments", "shares"), | |
| "tiktok": ("views", "likes", "comments", "shares"), | |
| "x": ("impressions", "likes", "comments", "shares"), | |
| "youtube": ("views", "likes", "comments"), | |
| "linkedin": ("impressions", "likes", "comments", "shares", "engagement_rate"), | |
| } | |
| class AnalyticsDomainService: | |
| """Tenant-scoped analytics aggregation over provider-reported snapshots.""" | |
| def __init__( | |
| self, | |
| repository: AnalyticsRepository, | |
| providers: ProviderRegistry, | |
| social_analytics: SocialAnalyticsService, | |
| accounts: object, | |
| audit: SocialAuditService, | |
| projects: object | None = None, | |
| ) -> None: | |
| self.repository = repository | |
| self.providers = providers | |
| self.social_analytics = social_analytics | |
| self.accounts = accounts | |
| self.audit = audit | |
| self.projects = projects | |
| self.ready = False | |
| async def initialize(self, social_ready: bool) -> None: | |
| self.ready = social_ready | |
| def ensure_ready(self) -> None: | |
| if not self.ready: | |
| from app.social.domain.errors import SocialProviderUnavailableError | |
| raise SocialProviderUnavailableError( | |
| "Analytics database schema is unavailable. Apply the analytics migration." | |
| ) | |
| def _range(query: AnalyticsQuery) -> tuple[datetime, datetime, str]: | |
| end = query.date_to or datetime.now(timezone.utc) | |
| start = query.date_from or end - timedelta(days=30) | |
| if start.tzinfo is None or end.tzinfo is None: | |
| raise AnalyticsValidationError("Analytics dates must include a UTC offset.") | |
| if end <= start: | |
| raise AnalyticsValidationError("date_to must be after date_from.") | |
| if end - start > timedelta(days=366): | |
| raise AnalyticsValidationError("Analytics date ranges are limited to 366 days.") | |
| try: | |
| ZoneInfo(query.timezone) | |
| except ZoneInfoNotFoundError as exc: | |
| raise AnalyticsValidationError("timezone must be a valid IANA timezone.") from exc | |
| return start.astimezone(timezone.utc), end.astimezone(timezone.utc), query.timezone | |
| def _metric_view(row: Any) -> AnalyticsMetric: | |
| return AnalyticsMetric( | |
| provider=row.provider, | |
| post_id=row.social_post_id, | |
| project_id=row.project_id, | |
| metric_date=row.metric_date, | |
| views=row.views, | |
| impressions=row.impressions, | |
| likes=row.likes, | |
| comments=row.comments, | |
| shares=row.shares, | |
| engagement_rate=row.engagement_rate, | |
| collected_at=row.collected_at, | |
| source=row.source, | |
| ) | |
| async def capabilities(self) -> list[AnalyticsCapabilities]: | |
| result: list[AnalyticsCapabilities] = [] | |
| for adapter in self.providers.list(): | |
| capabilities = adapter.capabilities | |
| available = bool(capabilities.analytics) | |
| result.append( | |
| AnalyticsCapabilities( | |
| provider=capabilities.provider.value, | |
| implementation_status=capabilities.implementation_status, | |
| status="available" if available else "unsupported", | |
| metrics=( | |
| list(PROVIDER_METRICS.get(capabilities.provider.value, ())) | |
| if available | |
| else [] | |
| ), | |
| required_scopes=list(capabilities.analytics_required_scopes), | |
| ) | |
| ) | |
| return result | |
| async def create_sync( | |
| self, workspace_id: str, user_id: str, payload: AnalyticsSyncRequest | |
| ) -> AnalyticsSyncRunView: | |
| query = AnalyticsQuery( | |
| date_from=payload.date_from, | |
| date_to=payload.date_to, | |
| timezone=payload.timezone, | |
| provider=payload.provider, | |
| project_id=payload.project_id, | |
| ) | |
| start, end, tz = self._range(query) | |
| if not payload.idempotency_key: | |
| raise AnalyticsValidationError( | |
| "Idempotency-Key is required for analytics synchronization." | |
| ) | |
| if payload.provider: | |
| self.providers.get(payload.provider) | |
| if payload.project_id and self.projects is not None: | |
| await self.projects.get( | |
| workspace_id=workspace_id, | |
| user_id=user_id, | |
| project_id=payload.project_id, | |
| ) | |
| run = await self.repository.create_sync( | |
| AnalyticsSyncRun( | |
| workspace_id=workspace_id, | |
| project_id=payload.project_id, | |
| provider=payload.provider, | |
| date_from=start, | |
| date_to=end, | |
| timezone=tz, | |
| idempotency_key=payload.idempotency_key, | |
| requested_by=user_id, | |
| ) | |
| ) | |
| await self.audit.record( | |
| workspace_id=workspace_id, | |
| event_type="analytics.sync_requested", | |
| metadata={"sync_run_id": run.id, "provider": payload.provider}, | |
| ) | |
| return self._sync_view(run) | |
| async def sync_once(self, run: AnalyticsSyncRun) -> AnalyticsSyncRunView: | |
| targets = await self.repository.published_targets( | |
| run.workspace_id, | |
| project_id=run.project_id, | |
| provider=run.provider, | |
| date_from=run.date_from, | |
| date_to=run.date_to, | |
| ) | |
| await self.audit.record( | |
| workspace_id=run.workspace_id, | |
| event_type="analytics.sync_started", | |
| metadata={"sync_run_id": run.id}, | |
| ) | |
| errors = 0 | |
| count = 0 | |
| by_post: dict[str, tuple[Any, list[Any]]] = {} | |
| for post, target in targets: | |
| by_post.setdefault(post.id, (post, []))[1].append(target) | |
| for post, post_targets in by_post.values(): | |
| try: | |
| latest = await self.repository.get_sync(run.workspace_id, run.id) | |
| if latest.status == "cancelled": | |
| return self._sync_view(latest) | |
| result = await self.social_analytics.post(run.workspace_id, post.id) | |
| for target in post_targets: | |
| matches = [ | |
| item | |
| for item in result.get("metrics", []) | |
| if isinstance(item, dict) and item.get("provider") == target.provider | |
| ] | |
| if not matches: | |
| errors += 1 | |
| continue | |
| metric = matches[-1] | |
| await self.repository.upsert_post_metric( | |
| workspace_id=run.workspace_id, | |
| post=post, | |
| target=target, | |
| metric=metric, | |
| collected_at=datetime.now(timezone.utc), | |
| ) | |
| count += 1 | |
| if result.get("unavailable"): | |
| errors += 1 | |
| except Exception as exc: | |
| errors += 1 | |
| logger.warning( | |
| "analytics target synchronization failed", | |
| extra={"sync_run_id": run.id, "post_id": post.id}, | |
| ) | |
| if not getattr(exc, "code", None): | |
| logger.debug("analytics sync exception", exc_info=True) | |
| status = ( | |
| "failed" if targets and count == 0 and errors else "partial" if errors else "succeeded" | |
| ) | |
| if not targets: | |
| status = "succeeded" | |
| updated = await self.repository.update_sync(run.id, status=status, metrics_count=count) | |
| await self.audit.record( | |
| workspace_id=run.workspace_id, | |
| event_type=( | |
| "analytics.sync_completed" if status != "failed" else "analytics.sync_failed" | |
| ), | |
| metadata={"sync_run_id": run.id, "status": status, "metrics_count": count}, | |
| ) | |
| return self._sync_view(updated) | |
| async def overview(self, workspace_id: str, query: AnalyticsQuery) -> AnalyticsOverview: | |
| start, end, tz = self._range(query) | |
| rows = await self.repository.metric_rows( | |
| workspace_id, | |
| { | |
| "date_from": start, | |
| "date_to": end, | |
| "provider": query.provider, | |
| "project_id": query.project_id, | |
| "social_account_id": query.social_account_id, | |
| "post_id": query.post_id, | |
| "search": query.search, | |
| }, | |
| ) | |
| totals = self._totals(rows) | |
| grouped: dict[str, list[Any]] = defaultdict(list) | |
| for row in rows: | |
| grouped[row.provider].append(row) | |
| platforms = [ | |
| {"provider": provider, "posts": len(items), **self._totals(items)} | |
| for provider, items in sorted(grouped.items()) | |
| ] | |
| top = sorted( | |
| rows, key=lambda row: self._sort_value(row, query.sort), reverse=query.descending | |
| )[:10] | |
| return AnalyticsOverview( | |
| date_from=start, | |
| date_to=end, | |
| timezone=tz, | |
| timezone_source="requested" if query.timezone != "UTC" else "utc_fallback", | |
| totals=totals, | |
| platforms=platforms, | |
| top_posts=[self._metric_view(row) for row in top], | |
| freshness=await self._freshness(workspace_id, query), | |
| ) | |
| async def timeseries(self, workspace_id: str, query: AnalyticsQuery) -> AnalyticsTimeseries: | |
| start, end, tz_name = self._range(query) | |
| rows = await self.repository.metric_rows( | |
| workspace_id, | |
| { | |
| "date_from": start, | |
| "date_to": end, | |
| "provider": query.provider, | |
| "project_id": query.project_id, | |
| "social_account_id": query.social_account_id, | |
| "post_id": query.post_id, | |
| "search": query.search, | |
| }, | |
| ) | |
| zone = ZoneInfo(tz_name) | |
| points: dict[str, dict[str, Any]] = {} | |
| for row in rows: | |
| local = row.metric_date.astimezone(zone) | |
| if query.granularity == "month": | |
| key = local.strftime("%Y-%m-01") | |
| elif query.granularity == "week": | |
| monday = local.date() - timedelta(days=local.weekday()) | |
| key = monday.isoformat() | |
| else: | |
| key = local.date().isoformat() | |
| item = points.setdefault(key, {"bucket": key, "value": 0, "posts": 0}) | |
| value = getattr(row, query.metric or "views") | |
| if value is not None: | |
| item["value"] += value | |
| item["posts"] += 1 | |
| return AnalyticsTimeseries( | |
| date_from=start, | |
| date_to=end, | |
| timezone=tz_name, | |
| granularity=query.granularity, | |
| metric=query.metric or "views", | |
| points=[points[key] for key in sorted(points)], | |
| freshness=await self._freshness(workspace_id, query), | |
| ) | |
| async def posts(self, workspace_id: str, query: AnalyticsQuery) -> AnalyticsPostList: | |
| start, end, _ = self._range(query) | |
| rows = await self.repository.metric_rows( | |
| workspace_id, | |
| { | |
| "date_from": start, | |
| "date_to": end, | |
| "provider": query.provider, | |
| "project_id": query.project_id, | |
| "social_account_id": query.social_account_id, | |
| "post_id": query.post_id, | |
| "search": query.search, | |
| }, | |
| ) | |
| rows.sort(key=lambda row: self._sort_value(row, query.sort), reverse=query.descending) | |
| return AnalyticsPostList( | |
| items=[ | |
| self._metric_view(row) for row in rows[query.offset : query.offset + query.limit] | |
| ], | |
| offset=query.offset, | |
| limit=query.limit, | |
| freshness=await self._freshness(workspace_id, query), | |
| ) | |
| async def post( | |
| self, workspace_id: str, post_id: str, query: AnalyticsQuery | |
| ) -> AnalyticsPostList: | |
| query.post_id = post_id | |
| return await self.posts(workspace_id, query) | |
| async def _freshness(self, workspace_id: str, query: AnalyticsQuery) -> AnalyticsFreshness: | |
| run = await self.repository.latest_sync( | |
| workspace_id, provider=query.provider, project_id=query.project_id | |
| ) | |
| if run is None: | |
| return AnalyticsFreshness( | |
| status="not_synchronized", reason="No analytics sync has completed." | |
| ) | |
| status = ( | |
| "fresh" | |
| if run.completed_at | |
| and run.completed_at >= datetime.now(timezone.utc) - timedelta(hours=24) | |
| else "stale" | |
| ) | |
| if run.status == "partial": | |
| status = "partial" | |
| return AnalyticsFreshness( | |
| status=status, last_collected_at=run.completed_at, last_sync_id=run.id | |
| ) | |
| def _totals(rows: list[Any]) -> dict[str, int | float | None]: | |
| result: dict[str, int | float | None] = {} | |
| for name in COMMON_METRICS: | |
| values = [getattr(row, name) for row in rows if getattr(row, name) is not None] | |
| result[name] = ( | |
| (sum(values) / len(values) if name == "engagement_rate" else sum(values)) | |
| if values | |
| else None | |
| ) | |
| return result | |
| def _sort_value(row: Any, field: str) -> float: | |
| if field == "date": | |
| return row.metric_date.timestamp() | |
| value = getattr(row, field, None) | |
| return float(value) if isinstance(value, (int, float)) else 0.0 | |
| def _sync_view(run: AnalyticsSyncRun) -> AnalyticsSyncRunView: | |
| return AnalyticsSyncRunView.model_validate(run, from_attributes=True) | |