From dbedf0cfa53e604e4b9bc3f26f29691b691ddeb4 Mon Sep 17 00:00:00 2001 From: "DESKTOP-RTLN3BA\\$punk" Date: Wed, 22 Jul 2026 22:16:28 -0700 Subject: [PATCH] feat(backend): integrate PostHog analytics for enhanced observability - Added PostHog configuration options to .env.example files for both Docker and Surfsense backend. - Introduced PostHog dependency in pyproject.toml. - Implemented analytics middleware to capture various events across the application, including user authentication, automation runs, and API requests. - Enhanced existing routes and services to emit analytics events, providing insights into user interactions and system performance. - Ensured graceful shutdown of analytics clients in worker processes and application lifecycles. --- docker/.env.example | 9 + surfsense_backend/.env.example | 8 + surfsense_backend/app/app.py | 52 +++ .../app/automations/runtime/executor.py | 30 ++ .../app/automations/services/automation.py | 47 +++ .../app/automations/services/trigger.py | 40 ++ .../app/capabilities/core/access/rest.py | 92 ++++- surfsense_backend/app/celery_app.py | 13 + surfsense_backend/app/config/__init__.py | 13 + .../app/observability/__init__.py | 2 +- .../app/observability/analytics.py | 202 ++++++++++ .../app/podcasts/tasks/render.py | 27 ++ .../app/routes/anonymous_chat_routes.py | 42 +++ .../app/routes/image_generation_routes.py | 15 + .../app/routes/incentive_tasks_routes.py | 15 + .../app/routes/new_chat_routes.py | 10 + .../routes/personal_access_tokens_routes.py | 10 + .../app/routes/public_chat_routes.py | 16 +- surfsense_backend/app/routes/rbac_routes.py | 24 ++ .../routes/search_source_connectors_routes.py | 58 +++ surfsense_backend/app/routes/stripe_routes.py | 24 ++ .../app/routes/workspaces_routes.py | 11 + .../app/tasks/celery_tasks/connector_tasks.py | 83 ++++- .../app/tasks/celery_tasks/document_tasks.py | 98 ++++- .../celery_tasks/video_presentation_tasks.py | 15 + .../streaming/flows/new_chat/orchestrator.py | 42 +++ .../flows/resume_chat/orchestrator.py | 39 ++ .../chat/streaming/flows/shared/analytics.py | 117 ++++++ surfsense_backend/app/users.py | 28 +- surfsense_backend/pyproject.toml | 1 + .../unit/observability/test_analytics.py | 149 ++++++++ surfsense_backend/uv.lock | 158 ++++++++ surfsense_mcp/mcp_server/core/client.py | 6 +- .../app/(home)/login/LocalLoginForm.tsx | 6 +- surfsense_web/app/(home)/register/page.tsx | 10 +- .../[workspace_id]/team/team-content.tsx | 23 +- .../app/invite/[invite_code]/page.tsx | 14 +- .../automations/automations-mutation.atoms.ts | 73 +--- .../hooks/use-connector-dialog.ts | 28 +- .../ui/dialogs/CreateWorkspaceDialog.tsx | 4 +- .../settings/earn-credits-content.tsx | 19 +- .../components/sources/DocumentUploadTab.tsx | 14 +- surfsense_web/content/docs/observability.mdx | 34 ++ surfsense_web/hooks/use-run-stream.ts | 4 +- surfsense_web/lib/apis/base-api.service.ts | 51 ++- .../lib/chat/stream-engine/engine.ts | 4 +- surfsense_web/lib/posthog/events.ts | 351 ++---------------- 47 files changed, 1618 insertions(+), 513 deletions(-) create mode 100644 surfsense_backend/app/observability/analytics.py create mode 100644 surfsense_backend/app/tasks/chat/streaming/flows/shared/analytics.py create mode 100644 surfsense_backend/tests/unit/observability/test_analytics.py diff --git a/docker/.env.example b/docker/.env.example index 78c0bd93a..fb7f5ff42 100644 --- a/docker/.env.example +++ b/docker/.env.example @@ -392,6 +392,15 @@ STT_SERVICE=local/base # OTEL_HTTP_PORT=4318 # OTEL_HEALTH_PORT=13133 +# PostHog product analytics (server-side). Opt-in like OTel: leave +# POSTHOG_API_KEY unset for zero telemetry. Passed to backend/worker/beat via +# env_file, so no compose changes are needed. Use the SAME project key as the +# frontend's NEXT_PUBLIC_POSTHOG_KEY so server events merge onto the persons the +# web app already identifies by user id. +# POSTHOG_API_KEY=phc_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx +# POSTHOG_HOST=https://us.i.posthog.com +# POSTHOG_AI_PRIVACY_MODE=true # false ships LLM prompt/completion bodies to PostHog + # ------------------------------------------------------------------------------ # Advanced (optional) # ------------------------------------------------------------------------------ diff --git a/surfsense_backend/.env.example b/surfsense_backend/.env.example index 266985e6c..1d60c98c3 100644 --- a/surfsense_backend/.env.example +++ b/surfsense_backend/.env.example @@ -555,6 +555,14 @@ LANGSMITH_PROJECT=surfsense # OTEL_METRIC_EXPORT_INTERVAL=300000 # ms; 5 minutes # OTEL_SDK_DISABLED=true # emergency kill-switch +# Observability - PostHog product analytics (server-side) +# Opt-in like OTel: leave POSTHOG_API_KEY unset for zero telemetry. Use the +# SAME project key as the frontend's NEXT_PUBLIC_POSTHOG_KEY so server events +# merge onto the persons the web app already identifies by user id. +# POSTHOG_API_KEY=phc_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx +# POSTHOG_HOST=https://us.i.posthog.com +# POSTHOG_AI_PRIVACY_MODE=true # false ships LLM prompt/completion bodies to PostHog + # Skills + subagents # SURFSENSE_ENABLE_SKILLS=false # SURFSENSE_ENABLE_SPECIALIZED_SUBAGENTS=false diff --git a/surfsense_backend/app/app.py b/surfsense_backend/app/app.py index dde2ce7fa..42efe38a7 100644 --- a/surfsense_backend/app/app.py +++ b/surfsense_backend/app/app.py @@ -50,6 +50,7 @@ from app.gateway.inbox_worker import ( start_gateway_inbox_worker, stop_gateway_inbox_worker, ) +from app.observability import analytics as ph_analytics from app.observability import metrics as ot_metrics from app.observability.bootstrap import init_otel, shutdown_otel from app.rate_limiter import get_real_client_ip, limiter @@ -690,6 +691,7 @@ async def lifespan(app: FastAPI): await stop_gateway_inbox_worker() _stop_openrouter_background_refresh() await close_checkpointer() + ph_analytics.shutdown() shutdown_otel() @@ -796,6 +798,56 @@ class RequestPerfMiddleware(BaseHTTPMiddleware): app.add_middleware(RequestPerfMiddleware) + +# --------------------------------------------------------------------------- +# PAT / MCP API attribution middleware +# --------------------------------------------------------------------------- +# Emits a PostHog ``pat_api_request`` event for any request authenticated by a +# Personal Access Token, so "documents added via MCP", "searches via MCP" etc. +# are queryable without instrumenting each route. Relies on ``get_auth_context`` +# stashing the resolved principal on ``request.state.auth_context``; requests +# that never resolve a PAT principal are silently skipped. No-op when PostHog +# is unconfigured. + + +class PatApiAnalyticsMiddleware(BaseHTTPMiddleware): + """Capture PAT-authenticated API usage (incl. MCP) after each response.""" + + async def dispatch( + self, request: StarletteRequest, call_next: RequestResponseEndpoint + ) -> StarletteResponse: + response = await call_next(request) + with contextlib.suppress(Exception): + ctx = getattr(request.state, "auth_context", None) + if ( + ctx is not None + and ctx.method == "pat" + and ph_analytics.is_enabled() + ): + # Use the route *template* (e.g. /documents/{id}) to keep the + # ``route`` property low-cardinality; fall back to the raw path. + route = request.scope.get("route") + route_path = getattr(route, "path", None) or request.url.path + client = ( + "mcp" + if request.headers.get("X-SurfSense-Client") == "mcp" + else "pat_script" + ) + ph_analytics.capture_for( + ctx, + "pat_api_request", + { + "route": route_path, + "method": request.method, + "status_code": response.status_code, + "client": client, + }, + ) + return response + + +app.add_middleware(PatApiAnalyticsMiddleware) + # Add SlowAPI middleware for automatic rate limiting # Uses Starlette BaseHTTPMiddleware (not the raw ASGI variant) to avoid # corrupting StreamingResponse — SlowAPIASGIMiddleware re-sends diff --git a/surfsense_backend/app/automations/runtime/executor.py b/surfsense_backend/app/automations/runtime/executor.py index d9544cd0c..cbf682d60 100644 --- a/surfsense_backend/app/automations/runtime/executor.py +++ b/surfsense_backend/app/automations/runtime/executor.py @@ -15,11 +15,38 @@ from app.automations.schemas.definition.envelope import ( ) from app.automations.schemas.definition.plan_step import PlanStep from app.automations.templating import build_run_context +from app.observability import analytics as ph_analytics from . import repository from .step import execute_step +def _capture_run_outcome(run: AutomationRun, status: str) -> None: + """Emit ``automation_run_completed`` — headless runs the frontend never sees. + + No-op when PostHog is unconfigured or the owning user is unknown. + """ + automation = run.automation + creator_id = getattr(automation, "created_by_user_id", None) + if not creator_id: + return + workspace_id = getattr(automation, "workspace_id", None) + ph_analytics.capture( + "automation_run_completed", + distinct_id=str(creator_id), + properties={ + "automation_id": run.automation_id, + "run_id": run.id, + "workspace_id": workspace_id, + "status": status, + "trigger_type": run.trigger.type.value if run.trigger else None, + }, + groups={"workspace": str(workspace_id)} + if workspace_id is not None + else None, + ) + + async def execute_run(session: AsyncSession, run_id: int) -> None: """Load run ``run_id`` and execute its snapshot plan to a terminal state.""" run = await repository.load_run(session, run_id) @@ -41,6 +68,7 @@ async def execute_run(session: AsyncSession, run_id: int) -> None: }, ) await session.commit() + _capture_run_outcome(run, "failed") return await repository.mark_running(session, run) @@ -66,6 +94,7 @@ async def execute_run(session: AsyncSession, run_id: int) -> None: await _run_on_failure(session, run, definition) await repository.mark_failed(session, run, result.get("error")) await session.commit() + _capture_run_outcome(run, "failed") return if result["status"] == "succeeded": @@ -73,6 +102,7 @@ async def execute_run(session: AsyncSession, run_id: int) -> None: await repository.mark_succeeded(session, run) await session.commit() + _capture_run_outcome(run, "succeeded") async def _run_on_failure( diff --git a/surfsense_backend/app/automations/services/automation.py b/surfsense_backend/app/automations/services/automation.py index a419327bf..a7bbd3af3 100644 --- a/surfsense_backend/app/automations/services/automation.py +++ b/surfsense_backend/app/automations/services/automation.py @@ -29,6 +29,7 @@ from app.automations.services.model_policy import ( from app.automations.triggers import get_trigger from app.automations.triggers.builtin.schedule import compute_next_fire_at from app.db import Permission, Workspace, get_async_session +from app.observability import analytics as ph_analytics from app.users import get_auth_context from app.utils.rbac import check_permission @@ -75,6 +76,18 @@ class AutomationService: self.session.add(automation) await self.session.commit() + + # Authoritative creation (migrated from automations-mutation.atoms.ts). + ph_analytics.capture_for( + self.auth, + "automation_created", + { + "automation_id": automation.id, + "workspace_id": automation.workspace_id, + "trigger_count": len(payload.triggers), + }, + groups={"workspace": str(automation.workspace_id)}, + ) return await self._get_with_triggers_or_raise(automation.id) async def list( @@ -150,6 +163,31 @@ class AutomationService: automation.version += 1 await self.session.commit() + + # Migrated from automations-mutation.atoms.ts: a status-only change is a + # distinct event; other field edits are ``automation_updated``. + if "status" in data: + ph_analytics.capture_for( + self.auth, + "automation_status_changed", + { + "automation_id": automation.id, + "workspace_id": automation.workspace_id, + "next_status": str(data["status"]), + }, + groups={"workspace": str(automation.workspace_id)}, + ) + if any(k in data for k in ("name", "description", "definition")): + ph_analytics.capture_for( + self.auth, + "automation_updated", + { + "automation_id": automation.id, + "workspace_id": automation.workspace_id, + "has_definition_change": "definition" in data, + }, + groups={"workspace": str(automation.workspace_id)}, + ) return await self._get_with_triggers_or_raise(automation_id) async def delete(self, automation_id: int) -> None: @@ -158,9 +196,18 @@ class AutomationService: await self._authorize( automation.workspace_id, Permission.AUTOMATIONS_DELETE.value ) + workspace_id = automation.workspace_id await self.session.delete(automation) await self.session.commit() + # Authoritative deletion (migrated from automations-mutation.atoms.ts). + ph_analytics.capture_for( + self.auth, + "automation_deleted", + {"automation_id": automation_id, "workspace_id": workspace_id}, + groups={"workspace": str(workspace_id)}, + ) + async def _get_or_raise(self, automation_id: int) -> Automation: automation = await self.session.get(Automation, automation_id) if automation is None: diff --git a/surfsense_backend/app/automations/services/trigger.py b/surfsense_backend/app/automations/services/trigger.py index 5175eb15e..866e66395 100644 --- a/surfsense_backend/app/automations/services/trigger.py +++ b/surfsense_backend/app/automations/services/trigger.py @@ -16,6 +16,7 @@ from app.automations.schemas.api import TriggerCreate, TriggerUpdate from app.automations.triggers import get_trigger from app.automations.triggers.builtin.schedule import compute_next_fire_at from app.db import Permission, get_async_session +from app.observability import analytics as ph_analytics from app.users import get_auth_context from app.utils.rbac import check_permission @@ -48,6 +49,18 @@ class TriggerService: self.session.add(trigger) await self.session.commit() await self.session.refresh(trigger) + + # Migrated from automations-mutation.atoms.ts. + ph_analytics.capture_for( + self.auth, + "automation_trigger_added", + { + "automation_id": automation_id, + "trigger_id": trigger.id, + "trigger_type": getattr(trigger.type, "value", str(trigger.type)), + "enabled": trigger.enabled, + }, + ) return trigger async def update( @@ -82,6 +95,26 @@ class TriggerService: await self.session.commit() await self.session.refresh(trigger) + + # Migrated from automations-mutation.atoms.ts. ``change`` mirrors the + # frontend's coarse categorisation. + _change = ( + "enabled" + if "enabled" in data and "params" not in data + else "params" + if "params" in data + else "other" + ) + ph_analytics.capture_for( + self.auth, + "automation_trigger_updated", + { + "automation_id": automation_id, + "trigger_id": trigger_id, + "change": _change, + "enabled": trigger.enabled, + }, + ) return trigger async def remove(self, *, automation_id: int, trigger_id: int) -> None: @@ -92,6 +125,13 @@ class TriggerService: await self.session.delete(trigger) await self.session.commit() + # Migrated from automations-mutation.atoms.ts. + ph_analytics.capture_for( + self.auth, + "automation_trigger_removed", + {"automation_id": automation_id, "trigger_id": trigger_id}, + ) + async def _authorize_automation( self, automation_id: int, permission: str ) -> Automation: diff --git a/surfsense_backend/app/capabilities/core/access/rest.py b/surfsense_backend/app/capabilities/core/access/rest.py index 2047b21ae..26ac46318 100644 --- a/surfsense_backend/app/capabilities/core/access/rest.py +++ b/surfsense_backend/app/capabilities/core/access/rest.py @@ -45,6 +45,7 @@ from app.capabilities.core.store import all_capabilities from app.capabilities.core.types import Capability, CapabilityContext from app.db import Run, async_session_maker, get_async_session from app.exceptions import ExternalServiceError, SurfSenseError +from app.observability import analytics as ph_analytics from app.services.web_crawl_credit_service import InsufficientCreditsError from app.users import get_auth_context from app.utils.rbac import check_workspace_access @@ -107,6 +108,41 @@ def _origin_for(auth: AuthContext) -> str: return "ui" if getattr(auth, "method", None) == "session" else "api" +def _capture_scraper_run( + *, + capability: str, + status: str, + user_id, + origin: str, + duration_ms: int | None = None, + item_count: int | None = None, + cost_micros: int | None = None, +) -> None: + """Emit ``scraper_run_completed`` — scrapers are the highest-value MCP surface. + + Covers both sync and async runs; the async background task never flows + through the PAT middleware, so this is the only place its outcome is + captured. No-op when PostHog is unconfigured. + """ + if not ph_analytics.is_enabled() or not user_id: + return + platform, _, verb = capability.partition(".") + ph_analytics.capture( + "scraper_run_completed", + distinct_id=str(user_id), + properties={ + "platform": platform, + "verb": verb, + "capability": capability, + "status": status, + "origin": origin, + "duration_ms": duration_ms, + "item_count": item_count, + "cost_micros": cost_micros, + }, + ) + + def _now_ms() -> int: return int(time.time() * 1000) @@ -185,6 +221,8 @@ async def _execute_async_run( unit, executor, payload, + user_id=None, + origin: str = "api", ) -> None: """Run a scrape in the background: stream progress, charge, finalize the row. @@ -218,6 +256,13 @@ async def _execute_async_run( progress=reporter.coarse, ) _publish_finished(run_id, "error", error=str(exc)) + _capture_scraper_run( + capability=capability, + status="error", + user_id=user_id, + origin=origin, + duration_ms=int((time.perf_counter() - started) * 1000), + ) return except Exception: logger.exception("async run %s failed with an upstream error", run_id) @@ -229,6 +274,13 @@ async def _execute_async_run( progress=reporter.coarse, ) _publish_finished(run_id, "error", error="upstream error") + _capture_scraper_run( + capability=capability, + status="error", + user_id=user_id, + origin=origin, + duration_ms=int((time.perf_counter() - started) * 1000), + ) return duration_ms = int((time.perf_counter() - started) * 1000) @@ -252,6 +304,15 @@ async def _execute_async_run( progress=reporter.coarse, ) _publish_finished(run_id, "success", item_count=serialized.item_count) + _capture_scraper_run( + capability=capability, + status="success", + user_id=user_id, + origin=origin, + duration_ms=duration_ms, + item_count=serialized.item_count, + cost_micros=cost_micros, + ) async def _finalize_async( @@ -358,6 +419,8 @@ def _register_verb(router: APIRouter, capability: Capability) -> None: unit=unit, executor=executor, payload=payload, + user_id=user_id, + origin=origin, ) ) run_event_bus.register_task(run_id, task) @@ -373,6 +436,7 @@ def _register_verb(router: APIRouter, capability: Capability) -> None: try: output = await executor(payload) except (SurfSenseError, HTTPException) as exc: + _sync_err_duration = int((time.perf_counter() - started) * 1000) await _record_rest_run( workspace_id=workspace_id, capability=name, @@ -381,11 +445,19 @@ def _register_verb(router: APIRouter, capability: Capability) -> None: input=input_dump, user_id=user_id, error=str(exc), - duration_ms=int((time.perf_counter() - started) * 1000), + duration_ms=_sync_err_duration, progress=reporter.coarse, ) + _capture_scraper_run( + capability=name, + status="error", + user_id=user_id, + origin=origin, + duration_ms=_sync_err_duration, + ) raise except Exception as exc: + _sync_err_duration = int((time.perf_counter() - started) * 1000) await _record_rest_run( workspace_id=workspace_id, capability=name, @@ -394,9 +466,16 @@ def _register_verb(router: APIRouter, capability: Capability) -> None: input=input_dump, user_id=user_id, error=str(exc), - duration_ms=int((time.perf_counter() - started) * 1000), + duration_ms=_sync_err_duration, progress=reporter.coarse, ) + _capture_scraper_run( + capability=name, + status="error", + user_id=user_id, + origin=origin, + duration_ms=_sync_err_duration, + ) raise ExternalServiceError( f"The '{name}' capability failed due to an upstream error.", code="CAPABILITY_UPSTREAM_ERROR", @@ -418,6 +497,15 @@ def _register_verb(router: APIRouter, capability: Capability) -> None: cost_micros=cost_micros, progress=reporter.coarse, ) + _capture_scraper_run( + capability=name, + status="success", + user_id=user_id, + origin=origin, + duration_ms=duration_ms, + item_count=serialized.item_count, + cost_micros=cost_micros, + ) if run_id is not None: response.headers["X-Run-Id"] = f"run_{run_id}" return output diff --git a/surfsense_backend/app/celery_app.py b/surfsense_backend/app/celery_app.py index 471d80197..2a9cddbd5 100644 --- a/surfsense_backend/app/celery_app.py +++ b/surfsense_backend/app/celery_app.py @@ -10,6 +10,7 @@ from celery.signals import ( task_postrun, task_prerun, worker_process_init, + worker_process_shutdown, ) from dotenv import load_dotenv @@ -123,6 +124,18 @@ def init_worker(**kwargs): initialize_image_gen_router() +@worker_process_shutdown.connect +def shutdown_worker(**kwargs): + """Flush queued PostHog events before a Celery worker process exits. + + The analytics client init is lazy (fork-safe), so there is nothing to + start here — only a flush to avoid dropping events captured by tasks. + """ + from app.observability import analytics as ph_analytics + + ph_analytics.shutdown() + + # Celery configuration, sourced from the central Config singleton CELERY_BROKER_URL = config.CELERY_BROKER_URL CELERY_RESULT_BACKEND = config.CELERY_RESULT_BACKEND diff --git a/surfsense_backend/app/config/__init__.py b/surfsense_backend/app/config/__init__.py index 73062775c..45ef065df 100644 --- a/surfsense_backend/app/config/__init__.py +++ b/surfsense_backend/app/config/__init__.py @@ -1179,6 +1179,19 @@ class Config: os.getenv("CRAWL_HEADED_XVFB_ENABLED", "FALSE").upper() == "TRUE" ) + # PostHog server-side product analytics (opt-in, mirrors the OTel pattern: + # no key set => the analytics wrapper is a silent no-op). Use the SAME + # project key as the frontend's NEXT_PUBLIC_POSTHOG_KEY so server events + # merge onto the persons the web app already identifies by user id. + POSTHOG_API_KEY = os.getenv("POSTHOG_API_KEY") + POSTHOG_HOST = os.getenv("POSTHOG_HOST", "https://us.i.posthog.com") + # When true (default), the LLM-analytics LangChain handler suppresses + # prompt/completion bodies ($ai_input / $ai_output_choices) and captures + # only metrics — chat content includes users' private documents. + POSTHOG_AI_PRIVACY_MODE = ( + os.getenv("POSTHOG_AI_PRIVACY_MODE", "TRUE").upper() == "TRUE" + ) + # Litellm TTS Configuration TTS_SERVICE = os.getenv("TTS_SERVICE") TTS_SERVICE_API_BASE = os.getenv("TTS_SERVICE_API_BASE") diff --git a/surfsense_backend/app/observability/__init__.py b/surfsense_backend/app/observability/__init__.py index a675b1dae..4ec9f1b03 100644 --- a/surfsense_backend/app/observability/__init__.py +++ b/surfsense_backend/app/observability/__init__.py @@ -6,4 +6,4 @@ wrapper is a no-op when OTEL is not configured, so importing it from performance-critical paths is safe. """ -__all__ = ["bootstrap", "metrics", "otel"] +__all__ = ["analytics", "bootstrap", "metrics", "otel"] diff --git a/surfsense_backend/app/observability/analytics.py b/surfsense_backend/app/observability/analytics.py new file mode 100644 index 000000000..312335d24 --- /dev/null +++ b/surfsense_backend/app/observability/analytics.py @@ -0,0 +1,202 @@ +"""Server-side PostHog product analytics for SurfSense. + +Opt-in, mirroring the OpenTelemetry bootstrap contract: when +``POSTHOG_API_KEY`` is unset every function here is a silent no-op, so it is +safe to call from hot paths (including async request handlers) and from +self-hosted installs that never configure telemetry. + +Design notes: +- The underlying ``posthog`` client enqueues events onto a background + consumer thread, so ``capture()`` is a non-blocking queue append; the only + network I/O happens off-thread. ``shutdown()`` flushes and joins that thread + and MUST run before a process exits or queued events are lost. +- The client is created lazily on first use, never at import time. This keeps + it fork-safe under Celery's prefork pool: a client (and its consumer thread) + created in the parent would not survive ``fork()``, so each worker process + builds its own on first capture. +- ``distinct_id`` is always ``str(user.id)`` so server events merge onto the + same PostHog persons the web frontend identifies (see + ``surfsense_web/components/providers/PostHogIdentify.tsx``). +- Every event passes ``disable_geoip=True``; without it PostHog would resolve + the *server's* IP and overwrite each person's real (client-derived) location. +""" + +from __future__ import annotations + +import logging +import threading +from typing import TYPE_CHECKING, Any + +from app.config import config + +if TYPE_CHECKING: + from app.auth.context import AuthContext + +logger = logging.getLogger(__name__) + +_client: Any | None = None +_init_attempted = False +_lock = threading.Lock() + +# Stamped on every backend event so client-observed (frontend) and +# server-truth events are always distinguishable in PostHog. +_SOURCE = "backend" + + +def _get_client() -> Any | None: + """Return the process-local PostHog client, or ``None`` when disabled. + + Lazy + fork-safe: built on first use inside whichever process (web worker + or Celery worker) calls it, never at import time. + """ + global _client, _init_attempted + + if _init_attempted: + return _client + + with _lock: + if _init_attempted: + return _client + _init_attempted = True + + api_key = config.POSTHOG_API_KEY + if not api_key: + # ponytail: opt-in like OTel — no key means telemetry is off, not + # a misconfiguration. Stay silent so self-hosters see no noise. + return None + + try: + from posthog import Posthog + + _client = Posthog( + project_api_key=api_key, + host=config.POSTHOG_HOST, + ) + except Exception: + logger.warning("PostHog analytics init failed; disabling", exc_info=True) + _client = None + + return _client + + +def is_enabled() -> bool: + """True when a PostHog client is configured and available.""" + return _get_client() is not None + + +def get_client() -> Any | None: + """Raw PostHog client for integrations that need it (e.g. the LLM handler).""" + return _get_client() + + +def _client_label(auth: AuthContext) -> str: + """Best-effort ``client`` property derived from the auth principal. + + ``session`` can't be split into web vs desktop from auth alone, so callers + that know better may override ``client`` in ``properties``. + """ + if auth.method == "system": + return auth.source or "system" + if auth.method == "pat": + return "pat" + return "web" + + +def capture( + event: str, + *, + distinct_id: str, + properties: dict[str, Any] | None = None, + groups: dict[str, str] | None = None, +) -> None: + """Capture a product event. No-op (and never raises) when disabled. + + Wrapped in try/except like the frontend ``safeCapture`` — analytics must + never break a request. ``posthog`` v6 signature is ``capture(event, + distinct_id=..., properties=...)`` (event first, distinct_id a kwarg). + """ + client = _get_client() + if client is None: + return + + try: + props = {"source": _SOURCE, **(properties or {})} + client.capture( + event, + distinct_id=distinct_id, + properties=props, + groups=groups, + disable_geoip=True, + ) + except Exception: + logger.debug("PostHog capture failed for %s", event, exc_info=True) + + +def capture_for( + auth: AuthContext, + event: str, + properties: dict[str, Any] | None = None, + groups: dict[str, str] | None = None, +) -> None: + """Capture an event attributed to an ``AuthContext`` principal. + + Derives ``distinct_id`` from the user id and stamps ``auth_method`` and a + best-effort ``client`` so events are attributable to their surface + (web/desktop/pat/gateway/automation). + """ + if _get_client() is None: + return + + props = { + "auth_method": auth.method, + "client": _client_label(auth), + **(properties or {}), + } + capture( + event, + distinct_id=str(auth.user.id), + properties=props, + groups=groups, + ) + + +def group_identify( + group_type: str, + group_key: str, + properties: dict[str, Any] | None = None, +) -> None: + """Upsert group properties (e.g. per-workspace metadata). No-op when disabled.""" + client = _get_client() + if client is None: + return + + try: + client.group_identify( + group_type=group_type, + group_key=group_key, + properties=properties or {}, + ) + except Exception: + logger.debug("PostHog group_identify failed for %s", group_type, exc_info=True) + + +def shutdown() -> None: + """Flush queued events and stop the consumer thread. Safe to call always.""" + global _client + client = _client + if client is None: + return + try: + client.shutdown() + except Exception: + logger.debug("PostHog shutdown failed", exc_info=True) + + +__all__ = [ + "capture", + "capture_for", + "get_client", + "group_identify", + "is_enabled", + "shutdown", +] diff --git a/surfsense_backend/app/podcasts/tasks/render.py b/surfsense_backend/app/podcasts/tasks/render.py index 7759691a9..0cd6fbe5a 100644 --- a/surfsense_backend/app/podcasts/tasks/render.py +++ b/surfsense_backend/app/podcasts/tasks/render.py @@ -11,7 +11,10 @@ import logging import tempfile from pathlib import Path +from sqlalchemy import select + from app.celery_app import celery_app +from app.observability import analytics as ph_analytics from app.podcasts.persistence import PodcastRepository from app.podcasts.rendering import PodcastRenderer from app.podcasts.service import ( @@ -76,6 +79,30 @@ async def _render_audio(podcast_id: int) -> dict: podcast, storage_backend=backend_name, storage_key=key ) await session.commit() + + # Credit-consuming deliverable; the frontend never confirms the + # render finished. Owner (workspace.user_id) resolved lazily so + # disabled installs pay nothing for the extra query. + if ph_analytics.is_enabled(): + # Local import: app.db <-> app.podcasts.persistence have a + # module-init cycle; deferring keeps this task importable. + from app.db import Workspace + + owner_id = await session.scalar( + select(Workspace.user_id).where( + Workspace.id == podcast.workspace_id + ) + ) + if owner_id: + ph_analytics.capture( + "podcast_generated", + distinct_id=str(owner_id), + properties={ + "workspace_id": podcast.workspace_id, + "podcast_id": podcast_id, + }, + groups={"workspace": str(podcast.workspace_id)}, + ) except InvalidTransitionError: # A user back-out won the race (e.g. the regeneration was # reverted): drop the stale render and leave the row alone. diff --git a/surfsense_backend/app/routes/anonymous_chat_routes.py b/surfsense_backend/app/routes/anonymous_chat_routes.py index 3ff2fd38a..f53217826 100644 --- a/surfsense_backend/app/routes/anonymous_chat_routes.py +++ b/surfsense_backend/app/routes/anonymous_chat_routes.py @@ -2,6 +2,7 @@ from __future__ import annotations +import contextlib import logging import secrets import uuid @@ -13,6 +14,10 @@ from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from app.config import config +from app.observability import analytics as ph_analytics +from app.tasks.chat.streaming.flows.shared.analytics import ( + build_llm_callback_handler, +) from app.etl_pipeline.file_classifier import ( DIRECT_CONVERT_EXTENSIONS, PLAINTEXT_EXTENSIONS, @@ -354,6 +359,7 @@ async def stream_anonymous_chat( accumulator = start_turn() streaming_service = VercelStreamingService() + anon_outcome = "success" try: async with shielded_async_session(): @@ -394,6 +400,21 @@ async def stream_anonymous_chat( "recursion_limit": 40, } + # PostHog LLM analytics for the free tier — model spend per + # model is the highest-value cost insight. distinct_id is the + # anon session id (not joined to any registered person). + _anon_llm_handler = build_llm_callback_handler( + distinct_id=session_id, + trace_id=anon_thread_id, + properties={ + "client": "anonymous", + "model_slug": body.model_slug, + "$ai_session_id": session_id, + }, + ) + if _anon_llm_handler is not None: + langgraph_config["callbacks"] = [_anon_llm_handler] + yield streaming_service.format_message_start() yield streaming_service.format_start_step() @@ -465,6 +486,7 @@ async def stream_anonymous_chat( except Exception as e: logger.exception("Anonymous chat stream error") + anon_outcome = "error" await TokenQuotaService.anon_release(session_key, ip_key, request_id) _, error_code, _, _, user_message, extra = classify_stream_exception( e, @@ -479,6 +501,26 @@ async def stream_anonymous_chat( finally: await TokenQuotaService.anon_release_stream_slot(client_ip) + # Server-truth free-tier volume/model/outcome/token-burn. distinct_id + # is the anon session id — deliberately NOT joined to the frontend's + # PostHog anonymous id (the point is server truth, not funnel merge). + with contextlib.suppress(Exception): + if ph_analytics.is_enabled(): + ph_analytics.capture( + "anon_chat_turn_completed", + distinct_id=session_id, + properties={ + "client": "anonymous", + "outcome": anon_outcome, + "model_slug": body.model_slug, + "model_name": model_cfg.get("model_name"), + "total_tokens": accumulator.grand_total, + "prompt_tokens": accumulator.total_prompt_tokens, + "completion_tokens": accumulator.total_completion_tokens, + "cost_micros": accumulator.total_cost_micros, + }, + ) + return StreamingResponse( _generate(), media_type="text/event-stream", diff --git a/surfsense_backend/app/routes/image_generation_routes.py b/surfsense_backend/app/routes/image_generation_routes.py index 5a6c71127..e4aab3f80 100644 --- a/surfsense_backend/app/routes/image_generation_routes.py +++ b/surfsense_backend/app/routes/image_generation_routes.py @@ -17,6 +17,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload from app.auth.context import AuthContext +from app.observability import analytics as ph_analytics from app.config import config from app.db import ( ImageGeneration, @@ -327,6 +328,20 @@ async def create_image_generation( await session.commit() await session.refresh(db_image_gen) + + # Credit-consuming deliverable; the frontend never confirms + # completion. ``error_message`` set => provider call failed. + ph_analytics.capture_for( + auth, + "image_generated", + { + "workspace_id": data.workspace_id, + "model": data.model, + "n": data.n, + "status": "error" if db_image_gen.error_message else "success", + }, + groups={"workspace": str(data.workspace_id)}, + ) return db_image_gen except HTTPException: diff --git a/surfsense_backend/app/routes/incentive_tasks_routes.py b/surfsense_backend/app/routes/incentive_tasks_routes.py index 2635df42f..a4c553097 100644 --- a/surfsense_backend/app/routes/incentive_tasks_routes.py +++ b/surfsense_backend/app/routes/incentive_tasks_routes.py @@ -15,6 +15,7 @@ from app.db import ( UserIncentiveTask, get_async_session, ) +from app.observability import analytics as ph_analytics from app.schemas.incentive_tasks import ( CompleteTaskResponse, IncentiveTaskInfo, @@ -125,6 +126,20 @@ async def complete_task( await session.commit() await session.refresh(user) + # Authoritative reward grant (migrated from earn-credits-content.tsx). + # Placed after the already-completed early-return so retries never + # double-count. + ph_analytics.capture_for( + auth, + "incentive_task_completed", + { + "task_type": task_type.value + if hasattr(task_type, "value") + else str(task_type), + "credit_micros_rewarded": credit_micros_reward, + }, + ) + return CompleteTaskResponse( success=True, message=f"Task completed! You earned ${credit_micros_reward / 1_000_000:.2f} of credit.", diff --git a/surfsense_backend/app/routes/new_chat_routes.py b/surfsense_backend/app/routes/new_chat_routes.py index 89d1cef95..71a70e9b5 100644 --- a/surfsense_backend/app/routes/new_chat_routes.py +++ b/surfsense_backend/app/routes/new_chat_routes.py @@ -51,6 +51,7 @@ from app.db import ( get_async_session, shielded_async_session, ) +from app.observability import analytics as ph_analytics from app.schemas.new_chat import ( AgentToolInfo, CancelActiveTurnResponse, @@ -823,6 +824,15 @@ async def create_thread( session.add(db_thread) await session.commit() await session.refresh(db_thread) + + # Authoritative thread creation (migrated from stream-engine/engine.ts). + # get_auth_context => covers PAT/MCP callers the frontend never sees. + ph_analytics.capture_for( + auth, + "chat_created", + {"workspace_id": db_thread.workspace_id, "chat_id": db_thread.id}, + groups={"workspace": str(db_thread.workspace_id)}, + ) return db_thread except HTTPException: diff --git a/surfsense_backend/app/routes/personal_access_tokens_routes.py b/surfsense_backend/app/routes/personal_access_tokens_routes.py index a7849a2fc..3d63703b4 100644 --- a/surfsense_backend/app/routes/personal_access_tokens_routes.py +++ b/surfsense_backend/app/routes/personal_access_tokens_routes.py @@ -8,6 +8,7 @@ from sqlalchemy.future import select from app.auth.context import AuthContext from app.config import config from app.db import PersonalAccessToken, get_async_session +from app.observability import analytics as ph_analytics from app.schemas.pat import PATCreate, PATCreated, PATRead from app.users import require_session_context from app.utils.pat import generate_pat, hash_pat, token_prefix @@ -57,6 +58,13 @@ async def create_personal_access_token( await session.commit() await session.refresh(pat) + # Leading indicator of MCP / programmatic-API adoption. + ph_analytics.capture_for( + auth, + "pat_created", + {"pat_id": pat.id, "has_expiry": pat.expires_at is not None}, + ) + return PATCreated( id=pat.id, label=pat.label, @@ -102,3 +110,5 @@ async def delete_personal_access_token( ) ) await session.commit() + + ph_analytics.capture_for(auth, "pat_revoked", {"pat_id": pat_id}) diff --git a/surfsense_backend/app/routes/public_chat_routes.py b/surfsense_backend/app/routes/public_chat_routes.py index 70f012911..f90050055 100644 --- a/surfsense_backend/app/routes/public_chat_routes.py +++ b/surfsense_backend/app/routes/public_chat_routes.py @@ -13,6 +13,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.auth.context import AuthContext from app.db import get_async_session +from app.observability import analytics as ph_analytics from app.schemas.new_chat import ( CloneResponse, PublicChatResponse, @@ -56,7 +57,20 @@ async def clone_public_chat( Creates thread and copies messages. Requires authentication. """ - return await clone_from_snapshot(session, share_token, user) + result = await clone_from_snapshot(session, share_token, user) + + # Share-link conversion — only observable server-side. + ph_analytics.capture_for( + auth, + "public_chat_cloned", + { + "workspace_id": result.workspace_id, + "chat_id": result.thread_id, + }, + groups={"workspace": str(result.workspace_id)}, + ) + + return result @router.get("/{share_token}/podcasts/{podcast_id}") diff --git a/surfsense_backend/app/routes/rbac_routes.py b/surfsense_backend/app/routes/rbac_routes.py index 2d8f58d33..8ab0e17c6 100644 --- a/surfsense_backend/app/routes/rbac_routes.py +++ b/surfsense_backend/app/routes/rbac_routes.py @@ -28,6 +28,7 @@ from app.db import ( WorkspaceRole, get_async_session, ) +from app.observability import analytics as ph_analytics from app.schemas import ( InviteAcceptRequest, InviteAcceptResponse, @@ -782,6 +783,19 @@ async def create_invite( ) db_invite = result.scalars().first() + # Authoritative invite creation (migrated from team-content.tsx). + ph_analytics.capture_for( + auth, + "workspace_invite_sent", + { + "workspace_id": workspace_id, + "role_name": db_invite.role.name if db_invite.role else None, + "has_expiry": db_invite.expires_at is not None, + "has_max_uses": db_invite.max_uses is not None, + }, + groups={"workspace": str(workspace_id)}, + ) + return db_invite except HTTPException: @@ -1091,6 +1105,16 @@ async def accept_invite( role_name = invite.role.name if invite.role else "Default" workspace_name = invite.workspace.name if invite.workspace else "" + # Authoritative join (migrated from app/invite/[invite_code]/page.tsx, + # which fired both events). workspace_name dropped — user content. + for _evt in ("workspace_invite_accepted", "workspace_user_added"): + ph_analytics.capture_for( + auth, + _evt, + {"workspace_id": invite.workspace_id, "role_name": role_name}, + groups={"workspace": str(invite.workspace_id)}, + ) + return InviteAcceptResponse( message="Successfully joined the workspace", workspace_id=invite.workspace_id, diff --git a/surfsense_backend/app/routes/search_source_connectors_routes.py b/surfsense_backend/app/routes/search_source_connectors_routes.py index 4cfa78af2..ba0899db3 100644 --- a/surfsense_backend/app/routes/search_source_connectors_routes.py +++ b/surfsense_backend/app/routes/search_source_connectors_routes.py @@ -29,6 +29,7 @@ import redis from dateutil.parser import isoparse from fastapi import APIRouter, Body, Depends, HTTPException, Query from pydantic import BaseModel, Field, ValidationError +from sqlalchemy import event as sa_event from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.future import select @@ -44,6 +45,7 @@ from app.db import ( get_async_session, ) from app.notifications.service import NotificationService +from app.observability import analytics as ph_analytics from app.observability import metrics as ot_metrics, otel as ot from app.schemas import ( GoogleDriveIndexRequest, @@ -135,6 +137,45 @@ def _get_heartbeat_key(notification_id: int) -> str: router = APIRouter() +def _connector_type_value(connector_type) -> str: + """Low-cardinality connector_type string for analytics (enum -> its value).""" + return getattr(connector_type, "value", None) or str(connector_type) + + +def _emit_connector_connected(_mapper, _connection, target) -> None: + """SQLAlchemy after_insert hook: one guard for every connector-creation path. + + Fires ``connector_connected`` for the form route here AND all 15 OAuth + callback routes (each builds ``SearchSourceConnector(...)`` inline with no + shared helper), so we don't instrument 16 call sites. ``distinct_id`` comes + off the row itself (``user_id``) — the same person the web app identifies. + + ponytail: fires during flush, so a post-flush rollback would over-count + (rare; connector creation commits immediately after add). No-op without a + PostHog key. Upgrade path: move to an after_commit collector if over-count + ever shows up in the data. + """ + if not ph_analytics.is_enabled(): + return + user_id = getattr(target, "user_id", None) + if not user_id: + return + workspace_id = getattr(target, "workspace_id", None) + ph_analytics.capture( + "connector_connected", + distinct_id=str(user_id), + properties={ + "workspace_id": workspace_id, + "connector_id": target.id, + "connector_type": _connector_type_value(target.connector_type), + }, + groups={"workspace": str(workspace_id)} if workspace_id is not None else None, + ) + + +sa_event.listen(SearchSourceConnector, "after_insert", _emit_connector_connected) + + # Use Pydantic's BaseModel here class GitHubPATRequest(BaseModel): github_pat: str = Field(..., description="GitHub Personal Access Token") @@ -244,6 +285,10 @@ async def create_search_source_connector( await session.commit() await session.refresh(db_connector) + # ``connector_connected`` is emitted by the after_insert listener below, + # so it fires once for this form route AND all 15 OAuth callback routes + # without instrumenting each — see _emit_connector_connected. + # Create periodic schedule if periodic indexing is enabled if ( db_connector.periodic_indexing_enabled @@ -679,10 +724,23 @@ async def delete_search_source_connector( # Delete the connector record workspace_id = db_connector.workspace_id + deleted_connector_type = db_connector.connector_type is_mcp = db_connector.connector_type == SearchSourceConnectorType.MCP_CONNECTOR await session.delete(db_connector) await session.commit() + # Authoritative deletion (migrated from use-connector-dialog.ts). + ph_analytics.capture_for( + auth, + "connector_deleted", + { + "workspace_id": workspace_id, + "connector_id": connector_id, + "connector_type": _connector_type_value(deleted_connector_type), + }, + groups={"workspace": str(workspace_id)}, + ) + if is_mcp: from app.agents.chat.multi_agent_chat.shared.tools.mcp.tool import ( invalidate_mcp_tools_cache, diff --git a/surfsense_backend/app/routes/stripe_routes.py b/surfsense_backend/app/routes/stripe_routes.py index 459b23d60..4da24c2af 100644 --- a/surfsense_backend/app/routes/stripe_routes.py +++ b/surfsense_backend/app/routes/stripe_routes.py @@ -27,6 +27,7 @@ from app.db import ( User, get_async_session, ) +from app.observability import analytics as ph_analytics from app.schemas.stripe import ( AutoReloadSettingsResponse, CreateAutoReloadSetupSessionRequest, @@ -47,6 +48,26 @@ logger = logging.getLogger(__name__) router = APIRouter(prefix="/stripe", tags=["stripe"]) +def _capture_credits_purchased(user: User, purchase: CreditPurchase) -> None: + """Emit ``credits_purchased`` — revenue events only exist server-side. + + Call only from the idempotent grant paths (which early-return on already + COMPLETED / non-PENDING rows) so Stripe retries never double-count. No-op + when PostHog is unconfigured. + """ + ph_analytics.capture( + "credits_purchased", + distinct_id=str(user.id), + properties={ + "credit_micros_granted": purchase.credit_micros_granted, + "quantity": purchase.quantity, + "amount_total": purchase.amount_total, + "currency": purchase.currency, + "source": purchase.source, + }, + ) + + def get_stripe_client() -> StripeClient: """Return a configured Stripe client or raise if Stripe is disabled.""" if not config.STRIPE_SECRET_KEY: @@ -309,6 +330,7 @@ async def _fulfill_completed_credit_purchase( ) await db_session.commit() + _capture_credits_purchased(user, purchase) return StripeWebhookResponse() @@ -448,6 +470,8 @@ async def _reconcile_auto_reload_payment_intent( purchase.status = CreditPurchaseStatus.FAILED await db_session.commit() + if succeeded: + _capture_credits_purchased(user, purchase) return StripeWebhookResponse() diff --git a/surfsense_backend/app/routes/workspaces_routes.py b/surfsense_backend/app/routes/workspaces_routes.py index 53d4344f7..26e08a225 100644 --- a/surfsense_backend/app/routes/workspaces_routes.py +++ b/surfsense_backend/app/routes/workspaces_routes.py @@ -15,6 +15,7 @@ from app.db import ( get_async_session, get_default_roles_config, ) +from app.observability import analytics as ph_analytics from app.routes.model_connections_routes import compute_llm_setup_status from app.schemas import ( WorkspaceApiAccessUpdate, @@ -112,6 +113,16 @@ async def create_workspace( await session.commit() await session.refresh(db_workspace) + # Authoritative creation event (migrated from the frontend + # CreateWorkspaceDialog). Workspace name is intentionally NOT sent — + # it's user content with no aggregation value. + ph_analytics.capture_for( + auth, + "workspace_created", + {"workspace_id": db_workspace.id}, + groups={"workspace": str(db_workspace.id)}, + ) + response = WorkspaceRead.model_validate(db_workspace) response.llm_setup = await compute_llm_setup_status( session, auth, db_workspace.id diff --git a/surfsense_backend/app/tasks/celery_tasks/connector_tasks.py b/surfsense_backend/app/tasks/celery_tasks/connector_tasks.py index ec078e031..fb08215b5 100644 --- a/surfsense_backend/app/tasks/celery_tasks/connector_tasks.py +++ b/surfsense_backend/app/tasks/celery_tasks/connector_tasks.py @@ -1,5 +1,6 @@ """Celery tasks for connector indexing.""" +import contextlib import logging import time import traceback @@ -8,6 +9,7 @@ from collections.abc import Awaitable, Callable from celery import current_task from app.celery_app import celery_app +from app.observability import analytics as ph_analytics from app.observability import metrics as ot_metrics, otel as ot from app.tasks.celery_tasks import ( get_celery_session_maker, @@ -17,8 +19,18 @@ from app.tasks.celery_tasks import ( logger = logging.getLogger(__name__) -def run_async_celery_task[T](coro_factory: Callable[[], Awaitable[T]]) -> T: - """Run connector sync work and record aggregate connector metrics.""" +def run_async_celery_task[T]( + coro_factory: Callable[[], Awaitable[T]], + *, + workspace_id: int | None = None, + user_id: str | None = None, +) -> T: + """Run connector sync work and record aggregate connector metrics. + + When ``workspace_id``/``user_id`` are provided, also emits the PostHog + ``connector_indexing_completed``/``_failed`` outcome (the frontend only + knows indexing *started*, never whether it worked). No-op without a key. + """ task_name = getattr(current_task, "name", None) or "unknown" t0 = time.perf_counter() status = "failed" @@ -45,6 +57,29 @@ def run_async_celery_task[T](coro_factory: Callable[[], Awaitable[T]]) -> T: status=status, error_category=error_category, ) + if user_id and ph_analytics.is_enabled(): + event = ( + "connector_indexing_completed" + if status == "success" + else "connector_indexing_failed" + ) + with contextlib.suppress(Exception): + ph_analytics.capture( + event, + distinct_id=str(user_id), + properties={ + "workspace_id": workspace_id, + # ``connector_type`` is the Celery task name, e.g. + # index_notion_pages / index_github_repos. + "connector_type": task_name, + "status": status, + "error_category": error_category, + "duration_ms": int(elapsed_s * 1000), + }, + groups={"workspace": str(workspace_id)} + if workspace_id is not None + else None, + ) def _handle_greenlet_error(e: Exception, task_name: str, connector_id: int) -> None: @@ -91,7 +126,9 @@ def index_notion_pages_task( return run_async_celery_task( lambda: _index_notion_pages( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) except Exception as e: _handle_greenlet_error(e, "index_notion_pages", connector_id) @@ -129,7 +166,9 @@ def index_github_repos_task( return run_async_celery_task( lambda: _index_github_repos( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -164,7 +203,9 @@ def index_confluence_pages_task( return run_async_celery_task( lambda: _index_confluence_pages( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -200,7 +241,9 @@ def index_google_calendar_events_task( return run_async_celery_task( lambda: _index_google_calendar_events( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) except Exception as e: _handle_greenlet_error(e, "index_google_calendar_events", connector_id) @@ -238,7 +281,9 @@ def index_google_gmail_messages_task( return run_async_celery_task( lambda: _index_google_gmail_messages( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -275,7 +320,9 @@ def index_google_drive_files_task( workspace_id, user_id, items_dict, - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -315,7 +362,9 @@ def index_onedrive_files_task( workspace_id, user_id, items_dict, - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -355,7 +404,9 @@ def index_dropbox_files_task( workspace_id, user_id, items_dict, - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -393,7 +444,9 @@ def index_elasticsearch_documents_task( return run_async_celery_task( lambda: _index_elasticsearch_documents( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -428,7 +481,9 @@ def index_bookstack_pages_task( return run_async_celery_task( lambda: _index_bookstack_pages( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) @@ -463,7 +518,9 @@ def index_composio_connector_task( return run_async_celery_task( lambda: _index_composio_connector( connector_id, workspace_id, user_id, start_date, end_date - ) + ), + workspace_id=workspace_id, + user_id=user_id, ) diff --git a/surfsense_backend/app/tasks/celery_tasks/document_tasks.py b/surfsense_backend/app/tasks/celery_tasks/document_tasks.py index e2ca5345e..8e20275b7 100644 --- a/surfsense_backend/app/tasks/celery_tasks/document_tasks.py +++ b/surfsense_backend/app/tasks/celery_tasks/document_tasks.py @@ -4,11 +4,13 @@ import asyncio import contextlib import logging import os +import time from uuid import UUID from app.celery_app import celery_app from app.config import config from app.notifications.service import NotificationService +from app.observability import analytics as ph_analytics from app.observability import metrics as ot_metrics from app.services.task_logging_service import TaskLoggingService from app.tasks.celery_tasks import get_celery_session_maker, run_async_celery_task @@ -22,6 +24,44 @@ from app.tasks.document_processors import ( logger = logging.getLogger(__name__) + +def _capture_doc_processing( + status: str, + *, + user_id: str | None, + workspace_id: int, + doc_type: str, + file_size: int | None = None, + duration_ms: int | None = None, +) -> None: + """Emit ``document_processing_completed``/``_failed`` from a Celery task. + + The frontend only knows the upload POST succeeded, never whether ingestion + actually worked — this is the authoritative outcome. No-op when PostHog is + unconfigured. ``distinct_id`` is the owning user's id so it joins the same + person the web app identifies. + """ + if not ph_analytics.is_enabled() or not user_id: + return + event = ( + "document_processing_completed" + if status == "success" + else "document_processing_failed" + ) + ph_analytics.capture( + event, + distinct_id=str(user_id), + properties={ + "workspace_id": workspace_id, + "doc_type": doc_type, + "file_size": file_size, + "duration_ms": duration_ms, + "status": status, + }, + groups={"workspace": str(workspace_id)}, + ) + + # ===== Redis heartbeat for document processing tasks ===== # Same mechanism as connector indexing heartbeats (search_source_connectors_routes.py). # A background coroutine refreshes a Redis key every 60s with a 2-min TTL. @@ -268,11 +308,30 @@ def process_extension_document_task( workspace_id: ID of the workspace user_id: ID of the user """ - return run_async_celery_task( - lambda: _process_extension_document( - individual_document_dict, workspace_id, user_id + _t0 = time.perf_counter() + try: + result = run_async_celery_task( + lambda: _process_extension_document( + individual_document_dict, workspace_id, user_id + ) ) + except Exception: + _capture_doc_processing( + "failed", + user_id=user_id, + workspace_id=workspace_id, + doc_type="extension", + duration_ms=int((time.perf_counter() - _t0) * 1000), + ) + raise + _capture_doc_processing( + "success", + user_id=user_id, + workspace_id=workspace_id, + doc_type="extension", + duration_ms=int((time.perf_counter() - _t0) * 1000), ) + return result async def _process_extension_document( @@ -430,12 +489,14 @@ def process_file_upload_task( ) return + file_size: int | None = None try: file_size = os.path.getsize(file_path) logger.info(f"[process_file_upload] File size: {file_size} bytes") except Exception as e: logger.warning(f"[process_file_upload] Could not get file size: {e}") + _t0 = time.perf_counter() try: run_async_celery_task( lambda: _process_file_upload(file_path, filename, workspace_id, user_id) @@ -448,7 +509,23 @@ def process_file_upload_task( f"[process_file_upload] Task failed for {filename}: {e}\n" f"Traceback:\n{traceback.format_exc()}" ) + _capture_doc_processing( + "failed", + user_id=user_id, + workspace_id=workspace_id, + doc_type="file_upload", + file_size=file_size, + duration_ms=int((time.perf_counter() - _t0) * 1000), + ) raise + _capture_doc_processing( + "success", + user_id=user_id, + workspace_id=workspace_id, + doc_type="file_upload", + file_size=file_size, + duration_ms=int((time.perf_counter() - _t0) * 1000), + ) async def _process_file_upload( @@ -682,6 +759,7 @@ def process_file_upload_with_document_task( ) return + _t0 = time.perf_counter() try: run_async_celery_task( lambda: _process_file_with_document( @@ -702,7 +780,21 @@ def process_file_upload_with_document_task( f"[process_file_upload_with_document] Task failed for {filename}: {e}\n" f"Traceback:\n{traceback.format_exc()}" ) + _capture_doc_processing( + "failed", + user_id=user_id, + workspace_id=workspace_id, + doc_type="file_upload_2phase", + duration_ms=int((time.perf_counter() - _t0) * 1000), + ) raise + _capture_doc_processing( + "success", + user_id=user_id, + workspace_id=workspace_id, + doc_type="file_upload_2phase", + duration_ms=int((time.perf_counter() - _t0) * 1000), + ) async def _mark_document_failed(document_id: int, reason: str): diff --git a/surfsense_backend/app/tasks/celery_tasks/video_presentation_tasks.py b/surfsense_backend/app/tasks/celery_tasks/video_presentation_tasks.py index 4f934cc27..c91d3031c 100644 --- a/surfsense_backend/app/tasks/celery_tasks/video_presentation_tasks.py +++ b/surfsense_backend/app/tasks/celery_tasks/video_presentation_tasks.py @@ -12,6 +12,7 @@ from app.agents.video_presentation.state import State as VideoPresentationState from app.celery_app import celery_app from app.config import config as app_config from app.db import VideoPresentation, VideoPresentationStatus +from app.observability import analytics as ph_analytics from app.services.billable_calls import ( BillingSettlementError, QuotaInsufficientError, @@ -239,6 +240,20 @@ async def _generate_video_presentation( logger.info(f"Successfully generated video presentation: {video_pres.id}") + # Credit-consuming deliverable — the frontend never confirms + # completion. Attributed to the workspace owner resolved above. + if owner_user_id: + ph_analytics.capture( + "video_presentation_generated", + distinct_id=str(owner_user_id), + properties={ + "workspace_id": workspace_id, + "video_presentation_id": video_pres.id, + "slide_count": len(serializable_slides), + }, + groups={"workspace": str(workspace_id)}, + ) + return { "status": "ready", "video_presentation_id": video_pres.id, diff --git a/surfsense_backend/app/tasks/chat/streaming/flows/new_chat/orchestrator.py b/surfsense_backend/app/tasks/chat/streaming/flows/new_chat/orchestrator.py index 4451b4b37..a6544b810 100644 --- a/surfsense_backend/app/tasks/chat/streaming/flows/new_chat/orchestrator.py +++ b/surfsense_backend/app/tasks/chat/streaming/flows/new_chat/orchestrator.py @@ -97,6 +97,10 @@ from app.tasks.chat.streaming.flows.shared.rate_limit_recovery import ( log_rate_limit_recovered, reroute_to_next_auto_pin, ) +from app.tasks.chat.streaming.flows.shared.analytics import ( + build_llm_callback_handler, + capture_chat_turn_completed, +) from app.tasks.chat.streaming.flows.shared.span import ( close_chat_request_span, open_chat_request_span, @@ -487,6 +491,23 @@ async def stream_new_chat( "recursion_limit": 10_000, } + # PostHog LLM analytics: attach a callback so the full agent trace + # tree (LLM calls, tools, subagents) is captured per turn. No-op when + # PostHog is unconfigured. + _llm_handler = build_llm_callback_handler( + distinct_id=user_id, + trace_id=stream_result.turn_id, + properties={ + "workspace_id": workspace_id, + "chat_id": chat_id, + "$ai_session_id": str(chat_id), + "flow": flow, + }, + groups={"workspace": str(workspace_id)}, + ) + if _llm_handler is not None: + config["callbacks"] = [_llm_handler] + # --- Block 4: First SSE frames --- for sse in iter_initial_frames( @@ -830,6 +851,27 @@ async def stream_new_chat( log_prefix="stream_new_chat", ) + # Authoritative server-side product event. Inside the shield so it + # survives client-disconnect cancellation (abandoned tabs still + # produce a turn event), and while ``stream_result`` is still live + # (it's dropped to None below for GC). + capture_chat_turn_completed( + flow=flow, + outcome=chat_outcome, + error_category=chat_error_category, + workspace_id=workspace_id, + chat_id=chat_id, + user_id=user_id, + auth_context=auth_context, + agent_mode=chat_agent_mode, + client_platform=fs_platform, + filesystem_mode=fs_mode, + turn_id=getattr(stream_result, "turn_id", None), + request_id=request_id, + duration_ms=int((time.perf_counter() - _t_total) * 1000), + accumulator=accumulator, + ) + # Persist any sandbox-produced files to local storage so they remain # downloadable after the Daytona sandbox auto-deletes. if stream_result and stream_result.sandbox_files: diff --git a/surfsense_backend/app/tasks/chat/streaming/flows/resume_chat/orchestrator.py b/surfsense_backend/app/tasks/chat/streaming/flows/resume_chat/orchestrator.py index 8b7956f4c..8dd749aee 100644 --- a/surfsense_backend/app/tasks/chat/streaming/flows/resume_chat/orchestrator.py +++ b/surfsense_backend/app/tasks/chat/streaming/flows/resume_chat/orchestrator.py @@ -76,6 +76,10 @@ from app.tasks.chat.streaming.flows.shared.rate_limit_recovery import ( log_rate_limit_recovered, reroute_to_next_auto_pin, ) +from app.tasks.chat.streaming.flows.shared.analytics import ( + build_llm_callback_handler, + capture_chat_turn_completed, +) from app.tasks.chat.streaming.flows.shared.span import ( close_chat_request_span, open_chat_request_span, @@ -387,6 +391,22 @@ async def stream_resume_chat( "recursion_limit": 10_000, } + # PostHog LLM analytics — same trace as the original turn's + # conversation via ``$ai_session_id``. No-op when unconfigured. + _llm_handler = build_llm_callback_handler( + distinct_id=user_id, + trace_id=stream_result.turn_id, + properties={ + "workspace_id": workspace_id, + "chat_id": chat_id, + "$ai_session_id": str(chat_id), + "flow": "resume", + }, + groups={"workspace": str(workspace_id)}, + ) + if _llm_handler is not None: + config["callbacks"] = [_llm_handler] + # --- First SSE frames --- for sse in iter_initial_frames( @@ -599,6 +619,25 @@ async def stream_resume_chat( log_prefix="stream_resume", ) + # Authoritative server-side product event (see new_chat + # orchestrator for rationale). ``flow="resume"``. + capture_chat_turn_completed( + flow="resume", + outcome=chat_outcome, + error_category=chat_error_category, + workspace_id=workspace_id, + chat_id=chat_id, + user_id=user_id, + auth_context=auth_context, + agent_mode=chat_agent_mode, + client_platform=fs_platform, + filesystem_mode=fs_mode, + turn_id=getattr(stream_result, "turn_id", None), + request_id=request_id, + duration_ms=int((time.perf_counter() - _t_total) * 1000), + accumulator=accumulator, + ) + # Release the lock from the original interrupted turn or any # re-interrupt/bailout. Skip on ``BusyError`` (lock not held here). if not busy_error_raised: diff --git a/surfsense_backend/app/tasks/chat/streaming/flows/shared/analytics.py b/surfsense_backend/app/tasks/chat/streaming/flows/shared/analytics.py new file mode 100644 index 000000000..2d311c26b --- /dev/null +++ b/surfsense_backend/app/tasks/chat/streaming/flows/shared/analytics.py @@ -0,0 +1,117 @@ +"""PostHog chat-turn analytics for streaming flows. + +Emits a single authoritative ``chat_turn_completed`` product event per turn, +shared by the new-chat and resume orchestrators so every chat source (web, +desktop, PAT scripts, gateway, automations) is tracked identically — +including sources the frontend can never observe. No-op when PostHog is +unconfigured. +""" + +from __future__ import annotations + +import logging +from typing import TYPE_CHECKING, Any + +from app.config import config +from app.observability import analytics + +if TYPE_CHECKING: + from app.auth.context import AuthContext + from app.services.token_tracking_service import TurnTokenAccumulator + +logger = logging.getLogger(__name__) + + +def build_llm_callback_handler( + *, + distinct_id: str | None, + trace_id: str | None, + properties: dict[str, Any] | None = None, + groups: dict[str, str] | None = None, +) -> Any | None: + """Build a PostHog LangChain ``CallbackHandler`` for a chat turn. + + Attaching this to the LangGraph ``config["callbacks"]`` captures the full + agent trace tree ($ai_trace / $ai_span / $ai_generation) — every LLM call, + tool, subagent, and retriever — joined to the ``chat_turn_completed`` event + via ``trace_id`` (the turn id) and grouped per conversation via + ``$ai_session_id`` (in ``properties``). + + Returns ``None`` when PostHog is disabled or the package/handler is + unavailable, so callers can simply skip attaching callbacks. ``privacy_mode`` + (default on) suppresses prompt/completion bodies — chat content includes + users' private documents. + """ + client_obj = analytics.get_client() + if client_obj is None or not distinct_id: + return None + + try: + from posthog.ai.langchain import CallbackHandler + + return CallbackHandler( + client=client_obj, + distinct_id=distinct_id, + trace_id=trace_id, + properties=properties or {}, + privacy_mode=config.POSTHOG_AI_PRIVACY_MODE, + groups=groups, + ) + except Exception: + logger.debug("PostHog LLM callback handler unavailable", exc_info=True) + return None + + +def capture_chat_turn_completed( + *, + flow: str, + outcome: str, + error_category: str | None, + workspace_id: int, + chat_id: int, + user_id: str | None, + auth_context: AuthContext | None, + agent_mode: str, + client_platform: str, + filesystem_mode: str, + turn_id: str | None, + request_id: str | None, + duration_ms: int, + accumulator: TurnTokenAccumulator, +) -> None: + """Capture ``chat_turn_completed``. Best-effort; never raises.""" + if not analytics.is_enabled() or not user_id: + return + + props: dict[str, Any] = { + "flow": flow, + "outcome": outcome, + "error_category": error_category, + "workspace_id": workspace_id, + "chat_id": chat_id, + "agent_mode": agent_mode, + "client_platform": client_platform, + "filesystem_mode": filesystem_mode, + "turn_id": turn_id, + "request_id": request_id, + "duration_ms": duration_ms, + # Cost is micro-USD (integer), matching TurnTokenAccumulator; do not + # convert to float dollars. + "total_tokens": accumulator.grand_total, + "prompt_tokens": accumulator.total_prompt_tokens, + "completion_tokens": accumulator.total_completion_tokens, + "cost_micros": accumulator.total_cost_micros, + } + groups = {"workspace": str(workspace_id)} + + if auth_context is not None: + analytics.capture_for( + auth_context, "chat_turn_completed", props, groups=groups + ) + else: + analytics.capture( + "chat_turn_completed", + distinct_id=user_id, + properties=props, + groups=groups, + ) diff --git a/surfsense_backend/app/users.py b/surfsense_backend/app/users.py index dbec93876..633b08443 100644 --- a/surfsense_backend/app/users.py +++ b/surfsense_backend/app/users.py @@ -20,6 +20,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.auth.context import AuthContext from app.auth.session_cookies import access_expires_at, write_session from app.config import config +from app.observability import analytics as ph_analytics from app.db import ( Prompt, User, @@ -144,6 +145,9 @@ class UserManager(UUIDIDMixin, BaseUserManager[User, uuid.UUID]): except Exception as e: logger.warning(f"Failed to update last_login for user {user.id}: {e}") + # Authoritative login event (vs. the frontend's optimistic capture). + ph_analytics.capture("auth_login_success", distinct_id=str(user.id)) + async def on_after_register(self, user: User, request: Request | None = None): """ Called after a user registers. Creates a default workspace for the user @@ -206,6 +210,17 @@ class UserManager(UUIDIDMixin, BaseUserManager[User, uuid.UUID]): logger.info( f"Created default workspace (ID: {default_workspace.id}) for user {user.id}" ) + + # Authoritative registration + auto-created default workspace. + ph_analytics.capture( + "auth_registration_success", distinct_id=str(user.id) + ) + ph_analytics.capture( + "workspace_created", + distinct_id=str(user.id), + properties={"client": "auto_register"}, + groups={"workspace": str(default_workspace.id)}, + ) except Exception as e: logger.error(f"Failed to create default workspace for user {user.id}: {e}") @@ -338,6 +353,13 @@ async def get_auth_context( FastAPI-Users still handles JWT mechanics; PATs are resolved here so RBAC receives the full SurfSense principal instead of a bare User. """ + def _stash(ctx: AuthContext) -> AuthContext: + # Expose the resolved principal on request.state so downstream + # middleware (e.g. PostHog pat_api_request attribution) can read it + # without re-resolving auth. + request.state.auth_context = ctx + return ctx + auth_header = request.headers.get("Authorization") if auth_header: scheme, _, credential = auth_header.partition(" ") @@ -348,7 +370,7 @@ async def get_auth_context( pat = await resolve_pat(session, token) if pat and pat.user and pat.user.is_active: maybe_touch_last_used(pat) - return AuthContext.pat_auth(pat.user, pat) + return _stash(AuthContext.pat_auth(pat.user, pat)) if is_bearer and _token_meets_epoch(token): try: @@ -358,7 +380,7 @@ async def get_auth_context( user = None if user and user.is_active: - return AuthContext.session(user) + return _stash(AuthContext.session(user)) cookie_token = request.cookies.get(config.SESSION_COOKIE_NAME) if cookie_token and _token_meets_epoch(cookie_token): @@ -369,7 +391,7 @@ async def get_auth_context( user = None if user and user.is_active: - return AuthContext.session(user) + return _stash(AuthContext.session(user)) raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, diff --git a/surfsense_backend/pyproject.toml b/surfsense_backend/pyproject.toml index bd346b306..955178320 100644 --- a/surfsense_backend/pyproject.toml +++ b/surfsense_backend/pyproject.toml @@ -86,6 +86,7 @@ dependencies = [ "python-telegram-bot>=22.7", "croniter>=2.0.0", "scrapling[fetchers]>=0.4.11", + "posthog>=6.0.0", ] [project.optional-dependencies] diff --git a/surfsense_backend/tests/unit/observability/test_analytics.py b/surfsense_backend/tests/unit/observability/test_analytics.py new file mode 100644 index 000000000..4cacb23c0 --- /dev/null +++ b/surfsense_backend/tests/unit/observability/test_analytics.py @@ -0,0 +1,149 @@ +"""Unit tests for the PostHog analytics wrapper. + +Covers the two properties the rest of the codebase relies on: +1. Opt-in no-op: with no client configured, every entry point is silent and + never raises (analytics must never break a request). +2. Correct stamping when a client IS present: ``source="backend"`` and + ``disable_geoip=True`` on every capture, plus ``auth_method`` / ``client`` + derived from the ``AuthContext`` principal by ``capture_for``. +""" + +from __future__ import annotations + +from types import SimpleNamespace + +import pytest + +from app.observability import analytics + +pytestmark = pytest.mark.unit + + +class _FakeClient: + """Records capture()/group_identify() calls instead of hitting the network.""" + + def __init__(self) -> None: + self.captures: list[dict] = [] + self.groups: list[dict] = [] + + def capture(self, event, *, distinct_id, properties, groups, disable_geoip): + self.captures.append( + { + "event": event, + "distinct_id": distinct_id, + "properties": properties, + "groups": groups, + "disable_geoip": disable_geoip, + } + ) + + def group_identify(self, *, group_type, group_key, properties): + self.groups.append( + {"group_type": group_type, "group_key": group_key, "properties": properties} + ) + + +@pytest.fixture(autouse=True) +def _reset_module_state(): + """Isolate the module-level lazy-client singleton between tests.""" + orig_client = analytics._client + orig_attempted = analytics._init_attempted + yield + analytics._client = orig_client + analytics._init_attempted = orig_attempted + + +def _use_fake_client() -> _FakeClient: + """Inject a fake client, bypassing lazy init (no posthog import, no key).""" + fake = _FakeClient() + analytics._client = fake + analytics._init_attempted = True + return fake + + +def _disable_client() -> None: + analytics._client = None + analytics._init_attempted = True + + +# ---- No-op behaviour when disabled ------------------------------------------------- + + +def test_capture_is_noop_without_client(): + _disable_client() + assert analytics.is_enabled() is False + # Must not raise even though there is no client. + analytics.capture("some_event", distinct_id="u1", properties={"a": 1}) + + +def test_capture_for_is_noop_without_client(): + _disable_client() + auth = SimpleNamespace(method="session", source=None, user=SimpleNamespace(id="u1")) + analytics.capture_for(auth, "some_event", {"a": 1}) # no raise + + +def test_shutdown_is_noop_without_client(): + _disable_client() + analytics.shutdown() # no raise + + +# ---- Stamping when a client is present --------------------------------------------- + + +def test_capture_stamps_source_and_disables_geoip(): + fake = _use_fake_client() + analytics.capture( + "chat_turn_completed", + distinct_id="user-123", + properties={"workspace_id": 7}, + groups={"workspace": "7"}, + ) + + assert len(fake.captures) == 1 + call = fake.captures[0] + assert call["event"] == "chat_turn_completed" + assert call["distinct_id"] == "user-123" + # source is always stamped so backend vs frontend events are separable. + assert call["properties"]["source"] == "backend" + assert call["properties"]["workspace_id"] == 7 + assert call["groups"] == {"workspace": "7"} + # Without this the server IP would overwrite each person's real location. + assert call["disable_geoip"] is True + + +def test_capture_never_raises_when_client_errors(): + class _Boom(_FakeClient): + def capture(self, *a, **k): + raise RuntimeError("network down") + + boom = _Boom() + analytics._client = boom + analytics._init_attempted = True + # Swallowed like the frontend safeCapture — analytics never breaks a request. + analytics.capture("evt", distinct_id="u1") + + +@pytest.mark.parametrize( + ("method", "source", "expected_client"), + [ + ("session", None, "web"), + ("pat", None, "pat"), + ("system", "gateway", "gateway"), + ("system", None, "system"), + ], +) +def test_capture_for_stamps_auth_method_and_client(method, source, expected_client): + fake = _use_fake_client() + auth = SimpleNamespace( + method=method, source=source, user=SimpleNamespace(id="user-abc") + ) + + analytics.capture_for(auth, "workspace_created", {"workspace_id": 1}) + + assert len(fake.captures) == 1 + props = fake.captures[0]["properties"] + assert fake.captures[0]["distinct_id"] == "user-abc" + assert props["auth_method"] == method + assert props["client"] == expected_client + assert props["source"] == "backend" # capture_for delegates to capture + assert props["workspace_id"] == 1 diff --git a/surfsense_backend/uv.lock b/surfsense_backend/uv.lock index 3a27770b6..5f8ec906a 100644 --- a/surfsense_backend/uv.lock +++ b/surfsense_backend/uv.lock @@ -48,6 +48,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'linux' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -94,6 +97,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'linux' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -140,6 +146,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -186,6 +195,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -232,6 +244,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -278,6 +293,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -324,6 +342,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -370,6 +391,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -416,6 +440,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -462,6 +489,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra == 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -523,6 +553,10 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'linux' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -569,6 +603,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'linux' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -615,6 +652,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'linux' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -661,6 +701,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -707,6 +750,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -753,6 +799,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -799,6 +848,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -845,6 +897,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -891,6 +946,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -937,6 +995,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform != 'linux' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -983,6 +1044,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra == 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1044,6 +1108,10 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'linux' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1090,6 +1158,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'linux' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1136,6 +1207,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'linux' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1182,6 +1256,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1228,6 +1305,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'emscripten' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1274,6 +1354,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1320,6 +1403,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1366,6 +1452,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'emscripten' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1412,6 +1501,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform != 'emscripten' and sys_platform != 'linux' and sys_platform != 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1458,6 +1550,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform != 'linux' and sys_platform != 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1504,6 +1599,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'win32' and extra == 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1565,6 +1663,10 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1611,6 +1713,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1702,6 +1807,12 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1748,6 +1859,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1794,6 +1908,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform == 'emscripten' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1885,6 +2002,12 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version == '3.13.*' and sys_platform != 'emscripten' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -1976,6 +2099,12 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform != 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", "python_version < '0'", "python_version < '0'", @@ -2022,6 +2151,9 @@ resolution-markers = [ "python_version < '0'", "python_version < '0'", "python_version < '0'", + "python_version < '0'", + "python_version < '0'", + "python_version < '0'", "python_full_version < '3.13' and sys_platform == 'win32' and extra != 'extra-16-surf-new-backend-cpu' and extra != 'extra-16-surf-new-backend-cu126' and extra != 'extra-16-surf-new-backend-cu128'", ] conflicts = [[ @@ -2557,6 +2689,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/77/f5/21d2de20e8b8b0408f0681956ca2c69f1320a3848ac50e6e7f39c6159675/babel-2.18.0-py3-none-any.whl", hash = "sha256:e2b422b277c2b9a9630c1d7903c2a00d0830c409c59ac8cae9081c92f1aeba35", size = 10196845 }, ] +[[package]] +name = "backoff" +version = "2.2.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/47/d7/5bbeb12c44d7c4f2fb5b56abce497eb5ed9f34d85701de869acedd602619/backoff-2.2.1.tar.gz", hash = "sha256:03f829f5bb1923180821643f8753b0502c3b682293992485b0eef2807afa5cba", size = 17001 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/df/73/b6e24bd22e6720ca8ee9a85a0c4a2971af8497d8f3193fa05390cbd46e09/backoff-2.2.1-py3-none-any.whl", hash = "sha256:63579f9a0628e06278f7e47b7d7d5b6ce20dc65c5e96a6f3ca99a6adca0396e8", size = 15148 }, +] + [[package]] name = "banks" version = "2.4.1" @@ -8650,6 +8791,21 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/4b/a6/38c8e2f318bf67d338f4d629e93b0b4b9af331f455f0390ea8ce4a099b26/portalocker-3.2.0-py3-none-any.whl", hash = "sha256:3cdc5f565312224bc570c49337bd21428bba0ef363bbcf58b9ef4a9f11779968", size = 22424 }, ] +[[package]] +name = "posthog" +version = "7.28.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "backoff" }, + { name = "distro" }, + { name = "requests" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/36/f0/3af875ac3fd5863ed4874c9618d85eab3f7fd24b395827d99da0ef90aca6/posthog-7.28.0.tar.gz", hash = "sha256:9e048dee58f27373db622c0744be30c1a2b7f1df31049956d6341c9646fb3833", size = 356360 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/1c/8b/d41d98e64bd6ce650d45690cffc718274b1eda65e3bfd79575a1cf95b45c/posthog-7.28.0-py3-none-any.whl", hash = "sha256:4cff10062807bbd8ae6ec2804a71584831a75b390ce7ba3e736bf4cc0fb25d28", size = 425147 }, +] + [[package]] name = "preshed" version = "3.0.13" @@ -10926,6 +11082,7 @@ dependencies = [ { name = "opentelemetry-sdk" }, { name = "opentelemetry-semantic-conventions" }, { name = "pgvector" }, + { name = "posthog" }, { name = "psycopg", extra = ["binary", "pool"] }, { name = "pyarrow" }, { name = "pyjwt" }, @@ -11041,6 +11198,7 @@ requires-dist = [ { name = "opentelemetry-sdk", specifier = ">=1.40.0" }, { name = "opentelemetry-semantic-conventions", specifier = ">=0.61b0" }, { name = "pgvector", specifier = ">=0.3.6" }, + { name = "posthog", specifier = ">=6.0.0" }, { name = "psycopg", extras = ["binary", "pool"], specifier = ">=3.3.2" }, { name = "pyarrow", specifier = ">=15.0.0,<19.0.0" }, { name = "pyjwt", specifier = ">=2.12.0" }, diff --git a/surfsense_mcp/mcp_server/core/client.py b/surfsense_mcp/mcp_server/core/client.py index d65b47b7e..2bdf405dc 100644 --- a/surfsense_mcp/mcp_server/core/client.py +++ b/surfsense_mcp/mcp_server/core/client.py @@ -37,7 +37,11 @@ class SurfSenseClient: self._fallback_api_key = fallback_api_key self._http = httpx.AsyncClient( base_url=api_base, - headers={"Accept": "application/json"}, + # ``X-SurfSense-Client`` lets the backend distinguish PAT traffic + # originating from this MCP server vs. raw PAT scripts, so + # "documents added via MCP" / "searches via MCP" are queryable. + # Server-to-server, so no CORS implications. + headers={"Accept": "application/json", "X-SurfSense-Client": "mcp"}, timeout=timeout, ) diff --git a/surfsense_web/app/(home)/login/LocalLoginForm.tsx b/surfsense_web/app/(home)/login/LocalLoginForm.tsx index dd415e10f..e60463b9d 100644 --- a/surfsense_web/app/(home)/login/LocalLoginForm.tsx +++ b/surfsense_web/app/(home)/login/LocalLoginForm.tsx @@ -13,7 +13,7 @@ import { Spinner } from "@/components/ui/spinner"; import { getAuthErrorDetails, isNetworkError } from "@/lib/auth-errors"; import { getPostLoginRedirectPath } from "@/lib/auth-utils"; import { ValidationError } from "@/lib/error"; -import { trackLoginAttempt, trackLoginFailure, trackLoginSuccess } from "@/lib/posthog/events"; +import { trackLoginAttempt, trackLoginFailure } from "@/lib/posthog/events"; export function LocalLoginForm() { const t = useTranslations("auth"); @@ -45,8 +45,8 @@ export function LocalLoginForm() { grant_type: "password", }); - // Track successful login - trackLoginSuccess("local"); + // auth_login_success is now emitted server-side + // (UserManager.on_after_login) — authoritative vs. optimistic. // Small delay to show success message setTimeout(() => { diff --git a/surfsense_web/app/(home)/register/page.tsx b/surfsense_web/app/(home)/register/page.tsx index 571103e79..5c91c7064 100644 --- a/surfsense_web/app/(home)/register/page.tsx +++ b/surfsense_web/app/(home)/register/page.tsx @@ -15,11 +15,7 @@ import { Spinner } from "@/components/ui/spinner"; import { useSession } from "@/hooks/use-session"; import { getAuthErrorDetails, isNetworkError, shouldRetry } from "@/lib/auth-errors"; import { AppError, ValidationError } from "@/lib/error"; -import { - trackRegistrationAttempt, - trackRegistrationFailure, - trackRegistrationSuccess, -} from "@/lib/posthog/events"; +import { trackRegistrationAttempt, trackRegistrationFailure } from "@/lib/posthog/events"; import { AmbientBackground } from "../login/AmbientBackground"; export default function RegisterPage() { @@ -81,8 +77,8 @@ export default function RegisterPage() { is_verified: false, }); - // Track successful registration - trackRegistrationSuccess(); + // auth_registration_success is now emitted server-side + // (UserManager.on_after_register) — authoritative vs. optimistic. // Success toast toast.success(t("register_success"), { diff --git a/surfsense_web/app/dashboard/[workspace_id]/team/team-content.tsx b/surfsense_web/app/dashboard/[workspace_id]/team/team-content.tsx index 571512305..dc45462e8 100644 --- a/surfsense_web/app/dashboard/[workspace_id]/team/team-content.tsx +++ b/surfsense_web/app/dashboard/[workspace_id]/team/team-content.tsx @@ -96,7 +96,6 @@ import type { Role } from "@/contracts/types/roles.types"; import { invitesApiService } from "@/lib/apis/invites-api.service"; import { rolesApiService } from "@/lib/apis/roles-api.service"; import { formatRelativeDate } from "@/lib/format-date"; -import { trackWorkspaceInviteSent, trackWorkspaceUsersViewed } from "@/lib/posthog/events"; import { cacheKeys } from "@/lib/query-client/cache-keys"; import { cn } from "@/lib/utils"; @@ -226,12 +225,7 @@ export function TeamContent({ workspaceId }: TeamContentProps) { const canPrev = pageIndex > 0; const canNext = displayEnd < totalItems; - useEffect(() => { - if (members.length > 0 && !membersLoading) { - const ownerCount = members.filter((m) => m.is_owner).length; - trackWorkspaceUsersViewed(workspaceId, members.length, ownerCount); - } - }, [members, membersLoading, workspaceId]); + // workspace_users_viewed removed — redundant with $pageview. if (accessLoading || membersLoading) { return ( @@ -342,11 +336,7 @@ export function TeamContent({ workspaceId }: TeamContentProps) { Invite members ) : ( - + )} {invitesLoading ? (