Compare commits

...

4 Commits

Author SHA1 Message Date
2569718930@qq.com 9c5a08dc1e Reduce false signup trial risk alerts 2026-05-31 20:21:57 +08:00
2569718930@qq.com ada5f274d3 Debounce terminal SSE subscription reconnects 2026-05-31 20:15:39 +08:00
2569718930@qq.com 78ea0326a5 Return partial city detail batches 2026-05-31 19:51:28 +08:00
2569718930@qq.com d3f444dbf6 Bind PolyWeather services to loopback 2026-05-31 19:46:00 +08:00
12 changed files with 360 additions and 23 deletions
+2 -2
View File
@@ -79,7 +79,7 @@ services:
timeout: 5s timeout: 5s
image: ghcr.io/yangyuan-zhen/polyweather-frontend:${IMAGE_TAG:-latest} image: ghcr.io/yangyuan-zhen/polyweather-frontend:${IMAGE_TAG:-latest}
ports: ports:
- 3001:3000 - "127.0.0.1:3001:3000"
restart: unless-stopped restart: unless-stopped
polyweather_web: polyweather_web:
command: python web/app.py command: python web/app.py
@@ -104,7 +104,7 @@ services:
timeout: 5s timeout: 5s
image: ghcr.io/yangyuan-zhen/polyweather-backend:${IMAGE_TAG:-latest} image: ghcr.io/yangyuan-zhen/polyweather-backend:${IMAGE_TAG:-latest}
ports: ports:
- 8000:8000 - "127.0.0.1:8000:8000"
restart: unless-stopped restart: unless-stopped
user: ${UID:-1000}:${GID:-1000} user: ${UID:-1000}:${GID:-1000}
volumes: volumes:
@@ -42,8 +42,13 @@ export async function GET(req: NextRequest) {
try { try {
return await proxyBackendJsonGet(req, { return await proxyBackendJsonGet(req, {
cacheControl: cachePolicy.responseCacheControl, cacheControl: cachePolicy.responseCacheControl,
fetchCache: cacheControlForData: (data) =>
cachePolicy.fetchMode === "no-store" ? "no-store" : undefined, data &&
typeof data === "object" &&
(data as { partial?: unknown }).partial === true
? "no-store, max-age=0"
: cachePolicy.responseCacheControl,
fetchCache: "no-store",
publicMessage: "Failed to fetch city detail batch", publicMessage: "Failed to fetch city detail batch",
revalidateSeconds: cachePolicy.revalidateSeconds, revalidateSeconds: cachePolicy.revalidateSeconds,
signal: controller.signal, signal: controller.signal,
@@ -96,6 +96,18 @@ export async function runTests() {
chartLogicSource.includes("primeCityDetailCache"), chartLogicSource.includes("primeCityDetailCache"),
"visible terminal chart detail fetches should be coalesced into one batch request and prime the shared chart cache", "visible terminal chart detail fetches should be coalesced into one batch request and prime the shared chart cache",
); );
assert(
chartLogicSource.includes("partial?: boolean") &&
chartLogicSource.includes("missing?: string[]"),
"frontend city detail batch payload should understand partial responses and missing city markers",
);
const flushCityDetailBatchBlock = chartLogicSource.match(/async function flushCityDetailBatch[\s\S]*?\r?\n}\r?\n\r?\nfunction fetchCityDetailBatchWithTimeout/)?.[0] || "";
assert(
flushCityDetailBatchBlock.includes("partialMissingCities") &&
flushCityDetailBatchBlock.includes("resolveBatchWaiters(waiters, null)") &&
flushCityDetailBatchBlock.includes("payload?.partial === true"),
"partial detail-batch misses should resolve without immediately issuing single-city fallback requests",
);
const fetchHourlyBlock = chartLogicSource.match(/async function fetchHourlyForecastForCity[\s\S]*?\r?\n}\r?\n\r?\nfunction fetchCityDetailWithTimeout/)?.[0] || ""; const fetchHourlyBlock = chartLogicSource.match(/async function fetchHourlyForecastForCity[\s\S]*?\r?\n}\r?\n\r?\nfunction fetchCityDetailWithTimeout/)?.[0] || "";
assert( assert(
fetchHourlyBlock.includes("queueCityDetailBatch(city, resParam)") && fetchHourlyBlock.includes("queueCityDetailBatch(city, resParam)") &&
@@ -99,6 +99,18 @@ export function runTests() {
assert(hook.includes("useLatestPatch"), "frontend patch hook must export useLatestPatch(city)"); assert(hook.includes("useLatestPatch"), "frontend patch hook must export useLatestPatch(city)");
assert(hook.includes("revision"), "frontend patch hook must track revisions and skip stale patches"); assert(hook.includes("revision"), "frontend patch hook must track revisions and skip stale patches");
assert(hook.includes("setTimeout"), "frontend patch hook must implement explicit reconnect backoff"); assert(hook.includes("setTimeout"), "frontend patch hook must implement explicit reconnect backoff");
assert(
hook.includes("SSE_SUBSCRIPTION_RECONNECT_DELAY_MS") &&
hook.includes("scheduleSubscriptionReconnect") &&
hook.includes("clearSubscriptionReconnectTimer"),
"frontend patch hook must debounce visible-city subscription changes before reopening SSE",
);
const subscriptionBlock = hook.match(/function registerCitySubscription[\s\S]*?\n}\r?\n\r?\nfunction normalizeLegacyPatch/)?.[0] || "";
assert(
subscriptionBlock.includes("scheduleSubscriptionReconnect()") &&
!subscriptionBlock.includes("ensureSsePatchConnection();"),
"city subscription mount/unmount should schedule one coalesced SSE reconnect instead of reconnecting per chart",
);
const bffEventsRoute = readFrontendFile("app", "api", "events", "route.ts"); const bffEventsRoute = readFrontendFile("app", "api", "events", "route.ts");
assert(bffEventsRoute.includes("searchParams"), "Next.js SSE proxy must forward query parameters to FastAPI"); assert(bffEventsRoute.includes("searchParams"), "Next.js SSE proxy must forward query parameters to FastAPI");
@@ -966,8 +966,11 @@ type HourlyForecastFetchOptions = {
}; };
type CityDetailBatchPayload = { type CityDetailBatchPayload = {
cities?: string[];
details?: Record<string, CityDetail | null | undefined>; details?: Record<string, CityDetail | null | undefined>;
errors?: Record<string, string>; errors?: Record<string, string>;
missing?: string[];
partial?: boolean;
}; };
type CityDetailBatchWaiter = { type CityDetailBatchWaiter = {
@@ -1112,6 +1115,10 @@ async function flushCityDetailBatch(resolution: string) {
try { try {
const payload = await fetchCityDetailBatchWithTimeout(cities, resolution); const payload = await fetchCityDetailBatchWithTimeout(cities, resolution);
const details = payload?.details || {}; const details = payload?.details || {};
const partialMissingCities =
payload?.partial === true
? new Set((payload.missing || []).map((city) => normalizeCityKey(city)))
: new Set<string>();
await Promise.all( await Promise.all(
cities.map(async (city) => { cities.map(async (city) => {
const waiters = queue.waiters.get(city); const waiters = queue.waiters.get(city);
@@ -1121,6 +1128,10 @@ async function flushCityDetailBatch(resolution: string) {
resolveBatchWaiters(waiters, data); resolveBatchWaiters(waiters, data);
return; return;
} }
if (partialMissingCities.has(normalizeCityKey(city))) {
resolveBatchWaiters(waiters, null);
return;
}
try { try {
resolveBatchWaiters( resolveBatchWaiters(
waiters, waiters,
@@ -41,6 +41,21 @@ export function runTests() {
const detailBatchProxy = readFrontend("app", "api", "cities", "detail-batch", "route.ts"); const detailBatchProxy = readFrontend("app", "api", "cities", "detail-batch", "route.ts");
assert.match(detailBatchProxy, /createProxyTimer\(req,\s*"city_detail_batch"\)/); assert.match(detailBatchProxy, /createProxyTimer\(req,\s*"city_detail_batch"\)/);
assert.match(detailBatchProxy, /timing:\s*timer/); assert.match(detailBatchProxy, /timing:\s*timer/);
assert.match(
detailBatchProxy,
/fetchCache:\s*"no-store"/,
"city detail batch proxy should avoid caching partial backend fetches in the Next data cache",
);
assert.match(
detailBatchProxy,
/cacheControlForData/,
"city detail batch proxy should be able to suppress response caching for partial payloads",
);
assert.match(
apiProxySource,
/cacheControlForData\?:/,
"generic backend JSON proxy should allow response cache policy to depend on parsed data",
);
const scanTerminalProxy = readFrontend("app", "api", "scan", "terminal", "route.ts"); const scanTerminalProxy = readFrontend("app", "api", "scan", "terminal", "route.ts");
assert.match(scanTerminalProxy, /createProxyTimer\(req,\s*"scan_terminal"\)/); assert.match(scanTerminalProxy, /createProxyTimer\(req,\s*"scan_terminal"\)/);
+25 -2
View File
@@ -4,6 +4,7 @@ import { useEffect, useSyncExternalStore } from "react";
import { resolveBackendApiUrl } from "@/lib/backend-api"; import { resolveBackendApiUrl } from "@/lib/backend-api";
const V1_EVENT_TYPE = "city_observation_patch.v1"; const V1_EVENT_TYPE = "city_observation_patch.v1";
const SSE_SUBSCRIPTION_RECONNECT_DELAY_MS = 150;
export type CityPatch = { export type CityPatch = {
type?: string; type?: string;
@@ -38,6 +39,7 @@ const subscribedCities = new Map<string, number>();
let eventSource: EventSource | null = null; let eventSource: EventSource | null = null;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null; let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let subscriptionReconnectTimer: ReturnType<typeof setTimeout> | null = null;
let reconnectAttempt = 0; let reconnectAttempt = 0;
let patchVersion = 0; let patchVersion = 0;
let resyncVersion = 0; let resyncVersion = 0;
@@ -74,6 +76,12 @@ function clearReconnectTimer() {
reconnectTimer = null; reconnectTimer = null;
} }
function clearSubscriptionReconnectTimer() {
if (!subscriptionReconnectTimer) return;
clearTimeout(subscriptionReconnectTimer);
subscriptionReconnectTimer = null;
}
function buildSseUrl(baseUrl: string) { function buildSseUrl(baseUrl: string) {
const params = new URLSearchParams(); const params = new URLSearchParams();
const cities = subscribedCityList(); const cities = subscribedCityList();
@@ -113,6 +121,7 @@ function closeEventSource() {
function reconnectNow() { function reconnectNow() {
if (typeof window === "undefined") return; if (typeof window === "undefined") return;
clearReconnectTimer(); clearReconnectTimer();
clearSubscriptionReconnectTimer();
closeEventSource(); closeEventSource();
connectSsePatches(); connectSsePatches();
} }
@@ -159,6 +168,7 @@ function connectSsePatches() {
function ensureSsePatchConnection() { function ensureSsePatchConnection() {
if (subscribedCities.size === 0) { if (subscribedCities.size === 0) {
clearSubscriptionReconnectTimer();
closeEventSource(); closeEventSource();
clearReconnectTimer(); clearReconnectTimer();
return; return;
@@ -167,6 +177,19 @@ function ensureSsePatchConnection() {
reconnectNow(); reconnectNow();
} }
function scheduleSubscriptionReconnect() {
if (typeof window === "undefined") return;
if (subscribedCities.size === 0) {
ensureSsePatchConnection();
return;
}
clearSubscriptionReconnectTimer();
subscriptionReconnectTimer = setTimeout(() => {
subscriptionReconnectTimer = null;
ensureSsePatchConnection();
}, SSE_SUBSCRIPTION_RECONNECT_DELAY_MS);
}
function registerCitySubscription(city: string) { function registerCitySubscription(city: string) {
const cityKey = normalizeCityKey(city); const cityKey = normalizeCityKey(city);
if (!cityKey) return () => {}; if (!cityKey) return () => {};
@@ -174,7 +197,7 @@ function registerCitySubscription(city: string) {
const previousCount = subscribedCities.get(cityKey) ?? 0; const previousCount = subscribedCities.get(cityKey) ?? 0;
subscribedCities.set(cityKey, previousCount + 1); subscribedCities.set(cityKey, previousCount + 1);
if (previousCount === 0) { if (previousCount === 0) {
ensureSsePatchConnection(); scheduleSubscriptionReconnect();
} }
return () => { return () => {
@@ -184,7 +207,7 @@ function registerCitySubscription(city: string) {
} else { } else {
subscribedCities.set(cityKey, nextCount); subscribedCities.set(cityKey, nextCount);
} }
ensureSsePatchConnection(); scheduleSubscriptionReconnect();
}; };
} }
+7 -4
View File
@@ -85,6 +85,7 @@ export async function proxyBackendJsonGet(
req: NextRequest, req: NextRequest,
options: { options: {
cacheControl?: string; cacheControl?: string;
cacheControlForData?: (data: unknown) => string | undefined;
conditionalResponse?: boolean; conditionalResponse?: boolean;
detailLimit?: number; detailLimit?: number;
error?: string; error?: string;
@@ -148,12 +149,14 @@ export async function proxyBackendJsonGet(
const data = await (timing const data = await (timing
? timing.measure("backend_read", () => res.json()) ? timing.measure("backend_read", () => res.json())
: res.json()); : res.json());
const responseCacheControl =
options.cacheControlForData?.(data) ?? options.cacheControl;
const response = const response =
options.cacheControl && options.conditionalResponse !== false responseCacheControl && options.conditionalResponse !== false
? buildCachedJsonResponse(req, data, options.cacheControl) ? buildCachedJsonResponse(req, data, responseCacheControl)
: NextResponse.json(data, { : NextResponse.json(data, {
headers: options.cacheControl headers: responseCacheControl
? { "Cache-Control": options.cacheControl } ? { "Cache-Control": responseCacheControl }
: undefined, : undefined,
}); });
const withCookies = applyAuthResponseCookies(response, auth.response); const withCookies = applyAuthResponseCookies(response, auth.response);
+9
View File
@@ -80,6 +80,15 @@ def test_deploy_script_retries_startup_smoke_checks():
assert 'smoke_check "frontend" "https://www.polyweather.top/" 15 3 5' in script assert 'smoke_check "frontend" "https://www.polyweather.top/" 15 3 5' in script
def test_docker_compose_keeps_polyweather_ports_on_loopback():
compose = (ROOT / "docker-compose.yml").read_text(encoding="utf-8")
assert "127.0.0.1:3001:3000" in compose
assert "127.0.0.1:8000:8000" in compose
assert "\n - 3001:3000" not in compose
assert "\n - 8000:8000" not in compose
def test_city_detail_builds_deb_hourly_consensus_before_peak_window(): def test_city_detail_builds_deb_hourly_consensus_before_peak_window():
source = (ROOT / "web" / "analysis_service.py").read_text(encoding="utf-8") source = (ROOT / "web" / "analysis_service.py").read_text(encoding="utf-8")
+149
View File
@@ -315,6 +315,15 @@ def test_ops_billing_risk_surfaces_trial_payment_referral_and_points(monkeypatch
DBManager, DBManager,
"list_app_analytics_events", "list_app_analytics_events",
lambda self, limit=20000, since_iso=None: [ lambda self, limit=20000, since_iso=None: [
{
"id": 10,
"event_type": "login_start",
"user_id": None,
"client_id": "c1",
"session_id": "session-gap",
"created_at": recent,
"payload": {"mode": "signup"},
},
{ {
"id": 11, "id": 11,
"event_type": "signup_success", "event_type": "signup_success",
@@ -419,6 +428,108 @@ def test_ops_billing_risk_does_not_flag_signup_when_backend_trial_exists(monkeyp
assert not any(issue["category"] == "signup_trial" for issue in payload["issues"]) assert not any(issue["category"] == "signup_trial" for issue in payload["issues"])
def test_ops_billing_risk_ignores_account_visit_without_signup_intent(monkeypatch):
from src.database.db_manager import DBManager
recent = datetime.now(timezone.utc).isoformat()
monkeypatch.setattr(ops_api.legacy_routes, "_require_ops_admin", lambda request: {"email": "ops@example.com"})
monkeypatch.setattr(ops_api, "_supabase_rest_rows", lambda table, params, *, timeout=10: [])
monkeypatch.setattr(
DBManager,
"list_app_analytics_events",
lambda self, limit=20000, since_iso=None: [
{
"id": 61,
"event_type": "signup_success",
"user_id": "returning-no-trial-user",
"client_id": "client-return",
"session_id": "session-return",
"created_at": recent,
"payload": {
"entry": "account_center",
"user_id": "returning-no-trial-user",
},
}
],
)
monkeypatch.setattr(
DBManager,
"list_payment_audit_events",
lambda self, limit=50, event_type=None: [],
)
payload = ops_api.get_ops_billing_risk(None, days=30, limit=20)
assert payload["summary"]["trial_gaps"] == 0
assert not any(issue["category"] == "signup_trial" for issue in payload["issues"])
def test_ops_billing_risk_treats_expired_signup_trial_subscription_as_backend_evidence(monkeypatch):
from src.database.db_manager import DBManager
now = datetime.now(timezone.utc)
recent = now.isoformat()
expired = (now - timedelta(days=2)).isoformat()
def fake_supabase_rows(table, params, *, timeout=10):
if table == "subscriptions" and params.get("or") == "(source.eq.signup_trial,plan_code.eq.signup_trial_3d)":
return [
{
"id": 81,
"user_id": "expired-trial-user",
"plan_code": "signup_trial_3d",
"source": "signup_trial",
"status": "expired",
"starts_at": (now - timedelta(days=5)).isoformat(),
"expires_at": expired,
"created_at": (now - timedelta(days=5)).isoformat(),
"updated_at": expired,
}
]
return []
monkeypatch.setattr(ops_api.legacy_routes, "_require_ops_admin", lambda request: {"email": "ops@example.com"})
monkeypatch.setattr(ops_api, "_supabase_rest_rows", fake_supabase_rows)
monkeypatch.setattr(
DBManager,
"list_app_analytics_events",
lambda self, limit=20000, since_iso=None: [
{
"id": 80,
"event_type": "login_start",
"user_id": None,
"client_id": "client-expired",
"session_id": "session-expired",
"created_at": recent,
"payload": {"mode": "signup"},
},
{
"id": 82,
"event_type": "signup_success",
"user_id": "expired-trial-user",
"client_id": "client-expired",
"session_id": "session-expired",
"created_at": recent,
"payload": {
"entry": "account_center",
"user_id": "expired-trial-user",
},
},
],
)
monkeypatch.setattr(
DBManager,
"list_payment_audit_events",
lambda self, limit=50, event_type=None: [],
)
payload = ops_api.get_ops_billing_risk(None, days=30, limit=20)
assert payload["summary"]["trial_gaps"] == 0
assert not any(issue["category"] == "signup_trial" for issue in payload["issues"])
def test_ops_payment_incidents_expose_top_level_reason_and_filters_resolved(monkeypatch): def test_ops_payment_incidents_expose_top_level_reason_and_filters_resolved(monkeypatch):
from src.database.db_manager import DBManager from src.database.db_manager import DBManager
@@ -618,6 +729,44 @@ def test_city_detail_batch_endpoint_limits_backend_concurrency(monkeypatch):
assert max_active <= 2 assert max_active <= 2
def test_city_detail_batch_returns_completed_details_when_one_city_is_slow(monkeypatch):
import asyncio
completed = []
async def build_batch_item(city, **kwargs):
if city == "slow":
await asyncio.sleep(0.08)
completed.append(city)
return city, {
"city": city,
"hourly": {"times": ["2026-05-30T00:00:00Z"], "temps": [20.0]},
"resolution": kwargs.get("resolution"),
}
monkeypatch.setenv("POLYWEATHER_CITY_DETAIL_BATCH_PARTIAL_TIMEOUT_MS", "20")
monkeypatch.setattr(city_api.legacy_routes, "_assert_entitlement", lambda request: None)
monkeypatch.setattr(city_api.legacy_routes, "_normalize_city_or_404", lambda name: name.strip().lower())
monkeypatch.setattr(city_api, "_build_city_detail_batch_item_async", build_batch_item)
payload = asyncio.run(
city_api.get_city_detail_batch_payload(
object(),
cities="fast,slow,other",
resolution="10m",
limit=3,
)
)
assert payload["cities"] == ["fast", "slow", "other"]
assert sorted(payload["details"]) == ["fast", "other"]
assert payload["details"]["fast"]["resolution"] == "10m"
assert payload["partial"] is True
assert payload["missing"] == ["slow"]
assert payload["errors"] == {}
assert "slow" not in completed
def test_concurrent_city_detail_requests_share_same_full_cache_refresh(monkeypatch): def test_concurrent_city_detail_requests_share_same_full_cache_refresh(monkeypatch):
import asyncio import asyncio
+49 -10
View File
@@ -500,6 +500,19 @@ def _city_detail_batch_concurrency() -> int:
return max(1, min(6, value)) return max(1, min(6, value))
def _city_detail_batch_partial_timeout_seconds() -> Optional[float]:
try:
timeout_ms = int(
os.getenv("POLYWEATHER_CITY_DETAIL_BATCH_PARTIAL_TIMEOUT_MS", "8500")
or "8500"
)
except ValueError:
timeout_ms = 8500
if timeout_ms <= 0:
return None
return max(0.001, min(60.0, timeout_ms / 1000.0))
async def get_city_detail_batch_payload( async def get_city_detail_batch_payload(
request: Request, request: Request,
*, *,
@@ -528,7 +541,13 @@ async def get_city_detail_batch_payload(
), ),
) )
if not city_names: if not city_names:
return {"cities": [], "details": {}, "errors": {}} return {
"cities": [],
"details": {},
"errors": {},
"missing": [],
"partial": False,
}
semaphore = asyncio.Semaphore(_city_detail_batch_concurrency()) semaphore = asyncio.Semaphore(_city_detail_batch_concurrency())
@@ -543,27 +562,47 @@ async def get_city_detail_batch_payload(
timing_recorder=timer, timing_recorder=timer,
) )
tasks = [ task_by_city = {
_build_with_limit(city) city: asyncio.create_task(_build_with_limit(city))
for city in city_names for city in city_names
] }
results = await timer.measure_async( task_city_lookup = {task: city for city, task in task_by_city.items()}
done, pending = await timer.measure_async(
"build_details", "build_details",
lambda: asyncio.gather(*tasks, return_exceptions=True), lambda: asyncio.wait(
task_by_city.values(),
timeout=_city_detail_batch_partial_timeout_seconds(),
),
) )
details: Dict[str, Any] = {} details: Dict[str, Any] = {}
errors: Dict[str, str] = {} errors: Dict[str, str] = {}
for city, result in zip(city_names, results): missing: List[str] = []
if isinstance(result, Exception): for task in done:
errors[city] = str(result) city = task_city_lookup[task]
try:
result_city, payload = task.result()
except Exception as exc:
errors[city] = str(exc)
continue continue
result_city, payload = result
details[result_city] = payload details[result_city] = payload
for task in pending:
city = task_city_lookup[task]
missing.append(city)
task.cancel()
missing_set = set(missing)
missing = [city for city in city_names if city in missing_set]
partial = bool(missing or errors)
if partial:
outcome = "partial"
return { return {
"cities": city_names, "cities": city_names,
"details": details, "details": details,
"errors": errors, "errors": errors,
"missing": missing,
"partial": partial,
} }
except HTTPException as exc: except HTTPException as exc:
outcome = f"http_{exc.status_code}" outcome = f"http_{exc.status_code}"
+62 -3
View File
@@ -545,18 +545,41 @@ def get_ops_billing_risk(
"limit": str(max(safe_limit * 10, 500)), "limit": str(max(safe_limit * 10, 500)),
}, },
) )
subscription_rows = collect( trial_subscription_rows = collect(
"subscriptions", "subscriptions",
{ {
"select": ( "select": (
"id,user_id,plan_code,source,status,starts_at,expires_at," "id,user_id,plan_code,source,status,starts_at,expires_at,"
"created_at,updated_at" "created_at,updated_at"
), ),
"or": "(source.eq.signup_trial,plan_code.eq.signup_trial_3d,status.eq.active)", "or": "(source.eq.signup_trial,plan_code.eq.signup_trial_3d)",
"order": "created_at.desc", "order": "created_at.desc",
"limit": str(max(safe_limit * 20, 1000)), "limit": str(max(safe_limit * 20, 1000)),
}, },
) )
active_subscription_rows = collect(
"subscriptions",
{
"select": (
"id,user_id,plan_code,source,status,starts_at,expires_at,"
"created_at,updated_at"
),
"status": "eq.active",
"order": "created_at.desc",
"limit": str(max(safe_limit * 20, 1000)),
},
)
subscription_rows: List[Dict[str, Any]] = []
seen_subscription_keys: set[str] = set()
for row in [*trial_subscription_rows, *active_subscription_rows]:
key = str(row.get("id") or "").strip() or (
f"{row.get('user_id')}:{row.get('plan_code')}:{row.get('source')}:"
f"{row.get('starts_at')}:{row.get('expires_at')}"
)
if key in seen_subscription_keys:
continue
seen_subscription_keys.add(key)
subscription_rows.append(row)
entitlement_trial_events = collect( entitlement_trial_events = collect(
"entitlement_events", "entitlement_events",
{ {
@@ -751,6 +774,40 @@ def get_ops_billing_risk(
def normalize_user_key(value: Any) -> str: def normalize_user_key(value: Any) -> str:
return str(value or "").strip().lower() return str(value or "").strip().lower()
def analytics_payload(row: Dict[str, Any]) -> Dict[str, Any]:
payload = row.get("payload")
return payload if isinstance(payload, dict) else {}
def analytics_correlation_keys(row: Dict[str, Any]) -> set[str]:
payload = analytics_payload(row)
keys: set[str] = set()
user_id = normalize_user_key(row.get("user_id") or payload.get("user_id"))
client_id = str(row.get("client_id") or "").strip()
session_id = str(row.get("session_id") or "").strip()
if user_id:
keys.add(f"user:{user_id}")
if client_id:
keys.add(f"client:{client_id}")
if session_id:
keys.add(f"session:{session_id}")
return keys
signup_intent_keys: set[str] = set()
for row in events:
if str(row.get("event_type") or "").strip().lower() != "login_start":
continue
payload = analytics_payload(row)
mode = str(payload.get("mode") or payload.get("auth_mode") or "").strip().lower()
if mode == "signup":
signup_intent_keys.update(analytics_correlation_keys(row))
def has_signup_intent(row: Dict[str, Any]) -> bool:
payload = analytics_payload(row)
mode = str(payload.get("mode") or payload.get("auth_mode") or "").strip().lower()
if mode == "signup" or payload.get("signup_intent") is True:
return True
return bool(analytics_correlation_keys(row).intersection(signup_intent_keys))
trial_actor_keys = { trial_actor_keys = {
_app_analytics_actor_key(row) _app_analytics_actor_key(row)
for row in events for row in events
@@ -817,10 +874,12 @@ def get_ops_billing_risk(
) )
for row in signup_rows[:300]: for row in signup_rows[:300]:
if not has_signup_intent(row):
continue
actor_key = _app_analytics_actor_key(row) actor_key = _app_analytics_actor_key(row)
if actor_key in trial_actor_keys: if actor_key in trial_actor_keys:
continue continue
payload = row.get("payload") if isinstance(row.get("payload"), dict) else {} payload = analytics_payload(row)
signup_user_id = normalize_user_key(row.get("user_id") or payload.get("user_id")) signup_user_id = normalize_user_key(row.get("user_id") or payload.get("user_id"))
if not signup_user_id: if not signup_user_id:
continue continue