Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 372d4366a8 |
+2
-2
@@ -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:
|
||||||
- "127.0.0.1:3001:3000"
|
- 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:
|
||||||
- "127.0.0.1:8000:8000"
|
- 8000:8000
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
user: ${UID:-1000}:${GID:-1000}
|
user: ${UID:-1000}:${GID:-1000}
|
||||||
volumes:
|
volumes:
|
||||||
|
|||||||
@@ -7,46 +7,25 @@ import {
|
|||||||
buildProxyExceptionResponse,
|
buildProxyExceptionResponse,
|
||||||
buildUpstreamErrorResponse,
|
buildUpstreamErrorResponse,
|
||||||
} from "@/lib/api-proxy";
|
} from "@/lib/api-proxy";
|
||||||
import {
|
|
||||||
createProxyTimer,
|
|
||||||
finishProxyTimedResponse,
|
|
||||||
} from "@/lib/proxy-timing";
|
|
||||||
|
|
||||||
const API_BASE = process.env.POLYWEATHER_API_BASE_URL;
|
const API_BASE = process.env.POLYWEATHER_API_BASE_URL;
|
||||||
const ANALYTICS_ENABLED =
|
const ANALYTICS_ENABLED =
|
||||||
process.env.NEXT_PUBLIC_POLYWEATHER_APP_ANALYTICS !== "false";
|
process.env.NEXT_PUBLIC_POLYWEATHER_APP_ANALYTICS !== "false";
|
||||||
const ANALYTICS_PROXY_TIMEOUT_MS = Math.max(
|
|
||||||
250,
|
|
||||||
Number(process.env.POLYWEATHER_ANALYTICS_PROXY_TIMEOUT_MS || "1500") || 1500,
|
|
||||||
);
|
|
||||||
|
|
||||||
export async function POST(req: NextRequest) {
|
export async function POST(req: NextRequest) {
|
||||||
const timer = createProxyTimer(req, "analytics_events");
|
|
||||||
if (!ANALYTICS_ENABLED) {
|
if (!ANALYTICS_ENABLED) {
|
||||||
return finishProxyTimedResponse(
|
return new NextResponse(null, { status: 204 });
|
||||||
new NextResponse(null, { status: 204 }),
|
|
||||||
timer,
|
|
||||||
"disabled",
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!API_BASE) {
|
if (!API_BASE) {
|
||||||
return finishProxyTimedResponse(
|
return NextResponse.json(
|
||||||
NextResponse.json(
|
{ error: "POLYWEATHER_API_BASE_URL is not configured" },
|
||||||
{ error: "POLYWEATHER_API_BASE_URL is not configured" },
|
{ status: 500 },
|
||||||
{ status: 500 },
|
|
||||||
),
|
|
||||||
timer,
|
|
||||||
"missing_api_base",
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
const controller = new AbortController();
|
|
||||||
const timeoutId = setTimeout(() => controller.abort(), ANALYTICS_PROXY_TIMEOUT_MS);
|
|
||||||
let auth: Awaited<ReturnType<typeof buildBackendRequestHeaders>> | null = null;
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const body = await timer.measure("request_read", () => req.json());
|
const body = await req.json();
|
||||||
const payload =
|
const payload =
|
||||||
body && typeof body.payload === "object" && body.payload != null
|
body && typeof body.payload === "object" && body.payload != null
|
||||||
? body.payload
|
? body.payload
|
||||||
@@ -63,65 +42,30 @@ export async function POST(req: NextRequest) {
|
|||||||
referer_header: req.headers.get("referer") || "",
|
referer_header: req.headers.get("referer") || "",
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
auth = await timer.measure("auth_headers", () =>
|
const auth = await buildBackendRequestHeaders(req, {
|
||||||
buildBackendRequestHeaders(req, {
|
includeSupabaseIdentity: false,
|
||||||
includeSupabaseIdentity: false,
|
});
|
||||||
}),
|
|
||||||
);
|
|
||||||
const headers = new Headers(auth.headers);
|
const headers = new Headers(auth.headers);
|
||||||
headers.set("Content-Type", "application/json");
|
headers.set("Content-Type", "application/json");
|
||||||
const res = await timer.measure("backend_fetch", () =>
|
const res = await fetch(`${API_BASE}/api/analytics/events`, {
|
||||||
fetch(`${API_BASE}/api/analytics/events`, {
|
method: "POST",
|
||||||
method: "POST",
|
headers,
|
||||||
headers,
|
body: JSON.stringify(enrichedBody),
|
||||||
body: JSON.stringify(enrichedBody),
|
cache: "no-store",
|
||||||
cache: "no-store",
|
});
|
||||||
signal: controller.signal,
|
|
||||||
}),
|
|
||||||
);
|
|
||||||
const backendServerTiming = res.headers.get("server-timing") || "";
|
|
||||||
if (!res.ok) {
|
if (!res.ok) {
|
||||||
const raw = await timer.measure("backend_read", () => res.text());
|
const raw = await res.text();
|
||||||
const response = buildUpstreamErrorResponse(res.status, raw, {
|
const response = buildUpstreamErrorResponse(res.status, raw, {
|
||||||
detailLimit: 260,
|
detailLimit: 260,
|
||||||
});
|
});
|
||||||
return finishProxyTimedResponse(
|
return applyAuthResponseCookies(response, auth.response);
|
||||||
applyAuthResponseCookies(response, auth.response),
|
|
||||||
timer,
|
|
||||||
`upstream_${res.status}`,
|
|
||||||
{ backendServerTiming },
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
const data = await timer.measure("backend_read", () => res.json());
|
const data = await res.json();
|
||||||
const response = NextResponse.json(data);
|
const response = NextResponse.json(data);
|
||||||
return finishProxyTimedResponse(
|
return applyAuthResponseCookies(response, auth.response);
|
||||||
applyAuthResponseCookies(response, auth.response),
|
|
||||||
timer,
|
|
||||||
"ok",
|
|
||||||
{ backendServerTiming },
|
|
||||||
);
|
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
const timedOut = controller.signal.aborted;
|
return buildProxyExceptionResponse(error, {
|
||||||
const response = timedOut
|
publicMessage: "Failed to track analytics event",
|
||||||
? NextResponse.json(
|
});
|
||||||
{
|
|
||||||
ok: false,
|
|
||||||
accepted: true,
|
|
||||||
dropped: true,
|
|
||||||
reason: "timeout",
|
|
||||||
},
|
|
||||||
{ status: 202 },
|
|
||||||
)
|
|
||||||
: buildProxyExceptionResponse(error, {
|
|
||||||
publicMessage: "Failed to track analytics event",
|
|
||||||
});
|
|
||||||
const withCookies = auth ? applyAuthResponseCookies(response, auth.response) : response;
|
|
||||||
return finishProxyTimedResponse(
|
|
||||||
withCookies,
|
|
||||||
timer,
|
|
||||||
timedOut ? "timeout_accepted" : "exception",
|
|
||||||
);
|
|
||||||
} finally {
|
|
||||||
clearTimeout(timeoutId);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -42,13 +42,8 @@ export async function GET(req: NextRequest) {
|
|||||||
try {
|
try {
|
||||||
return await proxyBackendJsonGet(req, {
|
return await proxyBackendJsonGet(req, {
|
||||||
cacheControl: cachePolicy.responseCacheControl,
|
cacheControl: cachePolicy.responseCacheControl,
|
||||||
cacheControlForData: (data) =>
|
fetchCache:
|
||||||
data &&
|
cachePolicy.fetchMode === "no-store" ? "no-store" : undefined,
|
||||||
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,
|
||||||
|
|||||||
@@ -55,10 +55,7 @@ import {
|
|||||||
mergeAccessStateWithAuthPayload,
|
mergeAccessStateWithAuthPayload,
|
||||||
type AuthProfilePayload,
|
type AuthProfilePayload,
|
||||||
} from "@/components/dashboard/scan-terminal/terminal-access-state";
|
} from "@/components/dashboard/scan-terminal/terminal-access-state";
|
||||||
import {
|
import { loadTerminalAuthProfile } from "@/components/dashboard/scan-terminal/terminal-auth-bootstrap";
|
||||||
createAuthProfileRequestCache,
|
|
||||||
loadTerminalAuthProfile,
|
|
||||||
} from "@/components/dashboard/scan-terminal/terminal-auth-bootstrap";
|
|
||||||
import {
|
import {
|
||||||
cityListItemsToScanRows,
|
cityListItemsToScanRows,
|
||||||
mergeScanRowsWithCityFallbackRows,
|
mergeScanRowsWithCityFallbackRows,
|
||||||
@@ -965,7 +962,7 @@ function ScanTerminalScreen() {
|
|||||||
createEmptyAccess(true),
|
createEmptyAccess(true),
|
||||||
);
|
);
|
||||||
|
|
||||||
const rawLoadAuthProfile = useCallback(
|
const loadAuthProfile = useCallback(
|
||||||
async (
|
async (
|
||||||
accessToken?: string | null,
|
accessToken?: string | null,
|
||||||
options?: { preferSnapshot?: boolean },
|
options?: { preferSnapshot?: boolean },
|
||||||
@@ -987,28 +984,6 @@ function ScanTerminalScreen() {
|
|||||||
},
|
},
|
||||||
[],
|
[],
|
||||||
);
|
);
|
||||||
const authProfileRequestCacheRef = useRef<{
|
|
||||||
load: typeof rawLoadAuthProfile;
|
|
||||||
cached: ReturnType<typeof createAuthProfileRequestCache>;
|
|
||||||
} | null>(null);
|
|
||||||
const loadAuthProfile = useCallback(
|
|
||||||
(
|
|
||||||
accessToken?: string | null,
|
|
||||||
options?: { preferSnapshot?: boolean },
|
|
||||||
): Promise<AuthProfilePayload> => {
|
|
||||||
const current = authProfileRequestCacheRef.current;
|
|
||||||
if (current?.load === rawLoadAuthProfile) {
|
|
||||||
return current.cached(accessToken, options);
|
|
||||||
}
|
|
||||||
const next = {
|
|
||||||
load: rawLoadAuthProfile,
|
|
||||||
cached: createAuthProfileRequestCache(rawLoadAuthProfile),
|
|
||||||
};
|
|
||||||
authProfileRequestCacheRef.current = next;
|
|
||||||
return next.cached(accessToken, options);
|
|
||||||
},
|
|
||||||
[rawLoadAuthProfile],
|
|
||||||
);
|
|
||||||
|
|
||||||
const refreshLiveAuthProfile = useCallback(async () => {
|
const refreshLiveAuthProfile = useCallback(async () => {
|
||||||
const supabaseEnabled = hasSupabasePublicEnv();
|
const supabaseEnabled = hasSupabasePublicEnv();
|
||||||
|
|||||||
@@ -96,18 +96,6 @@ 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)") &&
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
import {
|
import {
|
||||||
createAuthProfileRequestCache,
|
|
||||||
loadTerminalAuthProfile,
|
loadTerminalAuthProfile,
|
||||||
type TerminalAuthProfilePayload,
|
type TerminalAuthProfilePayload,
|
||||||
} from "@/components/dashboard/scan-terminal/terminal-auth-bootstrap";
|
} from "@/components/dashboard/scan-terminal/terminal-auth-bootstrap";
|
||||||
@@ -28,41 +27,6 @@ async function flushMicrotasks() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export async function runTests() {
|
export async function runTests() {
|
||||||
const pendingProfile = deferred<TerminalAuthProfilePayload>();
|
|
||||||
let underlyingProfileLoads = 0;
|
|
||||||
const cachedLoadAuthProfile = createAuthProfileRequestCache((
|
|
||||||
accessToken?: string | null,
|
|
||||||
options?: { preferSnapshot?: boolean },
|
|
||||||
) => {
|
|
||||||
underlyingProfileLoads += 1;
|
|
||||||
assert(
|
|
||||||
accessToken === "shared-token" && options?.preferSnapshot === true,
|
|
||||||
"auth profile request cache should pass through the original token and snapshot preference",
|
|
||||||
);
|
|
||||||
return pendingProfile.promise;
|
|
||||||
});
|
|
||||||
const firstCachedProfile = cachedLoadAuthProfile("shared-token", {
|
|
||||||
preferSnapshot: true,
|
|
||||||
});
|
|
||||||
const secondCachedProfile = cachedLoadAuthProfile("shared-token", {
|
|
||||||
preferSnapshot: true,
|
|
||||||
});
|
|
||||||
assert(
|
|
||||||
firstCachedProfile === secondCachedProfile && underlyingProfileLoads === 1,
|
|
||||||
"auth profile request cache should dedupe identical in-flight profile loads",
|
|
||||||
);
|
|
||||||
pendingProfile.resolve({
|
|
||||||
authenticated: true,
|
|
||||||
user_id: "shared-user",
|
|
||||||
subscription_active: true,
|
|
||||||
});
|
|
||||||
await firstCachedProfile;
|
|
||||||
await cachedLoadAuthProfile("shared-token", { preferSnapshot: true });
|
|
||||||
assert(
|
|
||||||
underlyingProfileLoads === 2,
|
|
||||||
"auth profile request cache should clear a key after the in-flight load settles",
|
|
||||||
);
|
|
||||||
|
|
||||||
const slowCookieProfile = deferred<TerminalAuthProfilePayload>();
|
const slowCookieProfile = deferred<TerminalAuthProfilePayload>();
|
||||||
const fastSession = deferred<{ data: { session: { access_token: string } } }>();
|
const fastSession = deferred<{ data: { session: { access_token: string } } }>();
|
||||||
const calls: string[] = [];
|
const calls: string[] = [];
|
||||||
|
|||||||
@@ -966,11 +966,8 @@ 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 = {
|
||||||
@@ -1115,10 +1112,6 @@ 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);
|
||||||
@@ -1128,10 +1121,6 @@ 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,
|
||||||
|
|||||||
@@ -57,32 +57,6 @@ function canResolveProfileImmediately(
|
|||||||
return payload?.authenticated === true && payload.subscription_active === true;
|
return payload?.authenticated === true && payload.subscription_active === true;
|
||||||
}
|
}
|
||||||
|
|
||||||
function authProfileRequestCacheKey(
|
|
||||||
accessToken?: string | null,
|
|
||||||
options?: { preferSnapshot?: boolean },
|
|
||||||
) {
|
|
||||||
const token = String(accessToken || "").trim();
|
|
||||||
const scope = token ? `bearer:${token}` : "cookie";
|
|
||||||
const mode = options?.preferSnapshot ? "snapshot" : "live";
|
|
||||||
return `${mode}:${scope}`;
|
|
||||||
}
|
|
||||||
|
|
||||||
export function createAuthProfileRequestCache(
|
|
||||||
loadAuthProfile: LoadTerminalAuthProfileOptions["loadAuthProfile"],
|
|
||||||
): LoadTerminalAuthProfileOptions["loadAuthProfile"] {
|
|
||||||
const pending = new Map<string, Promise<TerminalAuthProfilePayload>>();
|
|
||||||
return (accessToken, options) => {
|
|
||||||
const key = authProfileRequestCacheKey(accessToken, options);
|
|
||||||
const existing = pending.get(key);
|
|
||||||
if (existing) return existing;
|
|
||||||
const request = loadAuthProfile(accessToken, options).finally(() => {
|
|
||||||
pending.delete(key);
|
|
||||||
});
|
|
||||||
pending.set(key, request);
|
|
||||||
return request;
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
export async function loadTerminalAuthProfile({
|
export async function loadTerminalAuthProfile({
|
||||||
getSession,
|
getSession,
|
||||||
hasSupabasePublicEnv,
|
hasSupabasePublicEnv,
|
||||||
|
|||||||
@@ -41,21 +41,6 @@ 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"\)/);
|
||||||
@@ -72,26 +57,4 @@ export function runTests() {
|
|||||||
for (const stage of ["auth_headers", "ops_auth", "backend_fetch", "backend_read"]) {
|
for (const stage of ["auth_headers", "ops_auth", "backend_fetch", "backend_read"]) {
|
||||||
assert.match(onlineUsersProxy, new RegExp(stage));
|
assert.match(onlineUsersProxy, new RegExp(stage));
|
||||||
}
|
}
|
||||||
|
|
||||||
const analyticsProxy = readFrontend("app", "api", "analytics", "events", "route.ts");
|
|
||||||
assert.match(
|
|
||||||
analyticsProxy,
|
|
||||||
/ANALYTICS_PROXY_TIMEOUT_MS/,
|
|
||||||
"analytics event proxy should use a short dedicated timeout instead of waiting for long backend stalls",
|
|
||||||
);
|
|
||||||
assert.match(
|
|
||||||
analyticsProxy,
|
|
||||||
/createProxyTimer\(req,\s*"analytics_events"\)/,
|
|
||||||
"analytics event proxy should expose Server-Timing for HAR inspection",
|
|
||||||
);
|
|
||||||
assert.match(
|
|
||||||
analyticsProxy,
|
|
||||||
/AbortController/,
|
|
||||||
"analytics event proxy should abort slow upstream tracking requests",
|
|
||||||
);
|
|
||||||
assert.match(
|
|
||||||
analyticsProxy,
|
|
||||||
/status:\s*202/,
|
|
||||||
"analytics event proxy timeout should stay non-blocking for the fire-and-forget client event",
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -85,7 +85,6 @@ 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;
|
||||||
@@ -149,14 +148,12 @@ 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 =
|
||||||
responseCacheControl && options.conditionalResponse !== false
|
options.cacheControl && options.conditionalResponse !== false
|
||||||
? buildCachedJsonResponse(req, data, responseCacheControl)
|
? buildCachedJsonResponse(req, data, options.cacheControl)
|
||||||
: NextResponse.json(data, {
|
: NextResponse.json(data, {
|
||||||
headers: responseCacheControl
|
headers: options.cacheControl
|
||||||
? { "Cache-Control": responseCacheControl }
|
? { "Cache-Control": options.cacheControl }
|
||||||
: undefined,
|
: undefined,
|
||||||
});
|
});
|
||||||
const withCookies = applyAuthResponseCookies(response, auth.response);
|
const withCookies = applyAuthResponseCookies(response, auth.response);
|
||||||
|
|||||||
@@ -80,15 +80,6 @@ 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")
|
||||||
|
|
||||||
|
|||||||
@@ -618,44 +618,6 @@ 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
|
||||||
|
|
||||||
@@ -716,70 +678,6 @@ def test_concurrent_city_detail_requests_share_same_full_cache_refresh(monkeypat
|
|||||||
assert build_calls == 1
|
assert build_calls == 1
|
||||||
|
|
||||||
|
|
||||||
def test_stale_city_detail_uses_cached_full_payload_while_refreshing(monkeypatch):
|
|
||||||
import asyncio
|
|
||||||
|
|
||||||
refresh_calls = 0
|
|
||||||
build_inputs = []
|
|
||||||
|
|
||||||
class FakeCache:
|
|
||||||
def get_city_cache(self, kind, city):
|
|
||||||
assert kind == "full"
|
|
||||||
assert city == "paris"
|
|
||||||
return {
|
|
||||||
"payload": {
|
|
||||||
"city": "paris",
|
|
||||||
"hourly": {"times": ["2026-05-30T00:00:00Z"], "temps": [20.0]},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
async def fake_run_in_threadpool(fn, *args, **kwargs):
|
|
||||||
if fn is city_api.legacy_routes._refresh_city_full_cache:
|
|
||||||
await asyncio.sleep(0.01)
|
|
||||||
return fn(*args, **kwargs)
|
|
||||||
|
|
||||||
def refresh_full(city, force_refresh):
|
|
||||||
nonlocal refresh_calls
|
|
||||||
refresh_calls += 1
|
|
||||||
return {
|
|
||||||
"city": city,
|
|
||||||
"hourly": {"times": ["2026-05-30T00:00:00Z"], "temps": [21.0]},
|
|
||||||
}
|
|
||||||
|
|
||||||
def build_detail(data, market_slug, target_date, resolution):
|
|
||||||
build_inputs.append(data["hourly"]["temps"][0])
|
|
||||||
return {
|
|
||||||
"city": data["city"],
|
|
||||||
"live_temp": data["hourly"]["temps"][0],
|
|
||||||
"resolution": resolution,
|
|
||||||
}
|
|
||||||
|
|
||||||
city_api._CITY_FULL_REFRESH_INFLIGHT.clear()
|
|
||||||
city_api._CITY_DETAIL_PAYLOAD_CACHE.clear()
|
|
||||||
city_api._CITY_DETAIL_PAYLOAD_CACHE_TS.clear()
|
|
||||||
city_api._CITY_DETAIL_PAYLOAD_INFLIGHT.clear()
|
|
||||||
|
|
||||||
monkeypatch.setattr(city_api, "run_in_threadpool", fake_run_in_threadpool)
|
|
||||||
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.legacy_routes, "_CACHE_DB", FakeCache())
|
|
||||||
monkeypatch.setattr(city_api.legacy_routes, "_city_cache_is_fresh", lambda entry, ttl: False)
|
|
||||||
monkeypatch.setattr(city_api.legacy_routes, "_overlay_latest_wunderground_current", lambda city, payload: payload)
|
|
||||||
monkeypatch.setattr(city_api.legacy_routes, "_refresh_city_full_cache", refresh_full)
|
|
||||||
monkeypatch.setattr(city_api.legacy_routes, "_build_city_detail_payload", build_detail)
|
|
||||||
|
|
||||||
async def run_request():
|
|
||||||
payload = await city_api.get_city_detail_aggregate_payload(object(), "Paris", resolution="10m")
|
|
||||||
await asyncio.sleep(0.03)
|
|
||||||
return payload
|
|
||||||
|
|
||||||
result = asyncio.run(run_request())
|
|
||||||
|
|
||||||
assert result["live_temp"] == 20.0
|
|
||||||
assert build_inputs == [20.0]
|
|
||||||
assert refresh_calls == 1
|
|
||||||
|
|
||||||
|
|
||||||
def test_force_refresh_invalidates_short_city_detail_payload_cache(monkeypatch):
|
def test_force_refresh_invalidates_short_city_detail_payload_cache(monkeypatch):
|
||||||
import asyncio
|
import asyncio
|
||||||
|
|
||||||
|
|||||||
+14
-88
@@ -24,7 +24,6 @@ _RECENT_DEB_CACHE_TTL_SEC = max(
|
|||||||
int(os.getenv("POLYWEATHER_CITIES_DEB_RECENT_CACHE_TTL_SEC", "300") or "300"),
|
int(os.getenv("POLYWEATHER_CITIES_DEB_RECENT_CACHE_TTL_SEC", "300") or "300"),
|
||||||
)
|
)
|
||||||
_CITY_FULL_REFRESH_INFLIGHT: Dict[str, "asyncio.Task[Dict[str, Any]]"] = {}
|
_CITY_FULL_REFRESH_INFLIGHT: Dict[str, "asyncio.Task[Dict[str, Any]]"] = {}
|
||||||
_CITY_FULL_STALE_REFRESH_TASKS: Dict[str, "asyncio.Task[Dict[str, Any]]"] = {}
|
|
||||||
_CITY_FULL_REFRESH_LOCK = asyncio.Lock()
|
_CITY_FULL_REFRESH_LOCK = asyncio.Lock()
|
||||||
CityDetailPayloadCacheKey = Tuple[str, str, str, str, str, int]
|
CityDetailPayloadCacheKey = Tuple[str, str, str, str, str, int]
|
||||||
_CITY_DETAIL_PAYLOAD_CACHE: Dict[CityDetailPayloadCacheKey, Dict[str, Any]] = {}
|
_CITY_DETAIL_PAYLOAD_CACHE: Dict[CityDetailPayloadCacheKey, Dict[str, Any]] = {}
|
||||||
@@ -55,17 +54,9 @@ async def _refresh_city_full_cache_singleflight(city: str, force_refresh: bool)
|
|||||||
async with _CITY_FULL_REFRESH_LOCK:
|
async with _CITY_FULL_REFRESH_LOCK:
|
||||||
task = _CITY_FULL_REFRESH_INFLIGHT.get(key)
|
task = _CITY_FULL_REFRESH_INFLIGHT.get(key)
|
||||||
if task is None:
|
if task is None:
|
||||||
async def _run_refresh() -> Dict[str, Any]:
|
task = asyncio.create_task(
|
||||||
try:
|
run_in_threadpool(legacy_routes._refresh_city_full_cache, city, force_refresh),
|
||||||
return await run_in_threadpool(
|
)
|
||||||
legacy_routes._refresh_city_full_cache,
|
|
||||||
city,
|
|
||||||
force_refresh,
|
|
||||||
)
|
|
||||||
finally:
|
|
||||||
await _invalidate_city_detail_payload_cache(city)
|
|
||||||
|
|
||||||
task = asyncio.create_task(_run_refresh())
|
|
||||||
_CITY_FULL_REFRESH_INFLIGHT[key] = task
|
_CITY_FULL_REFRESH_INFLIGHT[key] = task
|
||||||
try:
|
try:
|
||||||
return await task
|
return await task
|
||||||
@@ -93,40 +84,14 @@ async def _refresh_city_full_data(city: str, force_refresh: bool) -> Dict[str, A
|
|||||||
return await _refresh_city_full_cache_singleflight(city, force_refresh)
|
return await _refresh_city_full_cache_singleflight(city, force_refresh)
|
||||||
|
|
||||||
|
|
||||||
def _start_city_full_stale_refresh(city: str) -> None:
|
|
||||||
normalized = str(city or "").strip().lower()
|
|
||||||
if not normalized:
|
|
||||||
return
|
|
||||||
existing = _CITY_FULL_STALE_REFRESH_TASKS.get(normalized)
|
|
||||||
if existing is not None and not existing.done():
|
|
||||||
return
|
|
||||||
|
|
||||||
task = asyncio.create_task(_refresh_city_full_data(city, False))
|
|
||||||
_CITY_FULL_STALE_REFRESH_TASKS[normalized] = task
|
|
||||||
|
|
||||||
def _cleanup(done: "asyncio.Task[Dict[str, Any]]") -> None:
|
|
||||||
if _CITY_FULL_STALE_REFRESH_TASKS.get(normalized) is done:
|
|
||||||
_CITY_FULL_STALE_REFRESH_TASKS.pop(normalized, None)
|
|
||||||
try:
|
|
||||||
done.result()
|
|
||||||
except Exception as exc: # pragma: no cover - defensive background guard
|
|
||||||
logger.warning("city full stale refresh failed city={}: {}", city, exc)
|
|
||||||
|
|
||||||
task.add_done_callback(_cleanup)
|
|
||||||
|
|
||||||
|
|
||||||
async def _get_city_full_data(city: str, *, force_refresh: bool) -> Dict[str, Any]:
|
async def _get_city_full_data(city: str, *, force_refresh: bool) -> Dict[str, Any]:
|
||||||
if force_refresh:
|
if force_refresh:
|
||||||
return await _refresh_city_full_data(city, True)
|
return await _refresh_city_full_data(city, True)
|
||||||
cached_entry = await run_in_threadpool(legacy_routes._CACHE_DB.get_city_cache, "full", city)
|
cached_entry = await run_in_threadpool(legacy_routes._CACHE_DB.get_city_cache, "full", city)
|
||||||
if cached_entry:
|
if cached_entry:
|
||||||
payload = cached_entry.get("payload") or {}
|
|
||||||
if not legacy_routes._city_cache_is_fresh(cached_entry, legacy_routes.CITY_FULL_CACHE_TTL_SEC):
|
if not legacy_routes._city_cache_is_fresh(cached_entry, legacy_routes.CITY_FULL_CACHE_TTL_SEC):
|
||||||
if payload:
|
|
||||||
_start_city_full_stale_refresh(city)
|
|
||||||
return await _overlay_cached_wunderground(city, payload)
|
|
||||||
return await _refresh_city_full_data(city, False)
|
return await _refresh_city_full_data(city, False)
|
||||||
return await _overlay_cached_wunderground(city, payload)
|
return await _overlay_cached_wunderground(city, cached_entry.get("payload") or {})
|
||||||
return await _refresh_city_full_data(city, False)
|
return await _refresh_city_full_data(city, False)
|
||||||
|
|
||||||
|
|
||||||
@@ -500,19 +465,6 @@ 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,
|
||||||
*,
|
*,
|
||||||
@@ -541,13 +493,7 @@ async def get_city_detail_batch_payload(
|
|||||||
),
|
),
|
||||||
)
|
)
|
||||||
if not city_names:
|
if not city_names:
|
||||||
return {
|
return {"cities": [], "details": {}, "errors": {}}
|
||||||
"cities": [],
|
|
||||||
"details": {},
|
|
||||||
"errors": {},
|
|
||||||
"missing": [],
|
|
||||||
"partial": False,
|
|
||||||
}
|
|
||||||
|
|
||||||
semaphore = asyncio.Semaphore(_city_detail_batch_concurrency())
|
semaphore = asyncio.Semaphore(_city_detail_batch_concurrency())
|
||||||
|
|
||||||
@@ -562,47 +508,27 @@ async def get_city_detail_batch_payload(
|
|||||||
timing_recorder=timer,
|
timing_recorder=timer,
|
||||||
)
|
)
|
||||||
|
|
||||||
task_by_city = {
|
tasks = [
|
||||||
city: asyncio.create_task(_build_with_limit(city))
|
_build_with_limit(city)
|
||||||
for city in city_names
|
for city in city_names
|
||||||
}
|
]
|
||||||
task_city_lookup = {task: city for city, task in task_by_city.items()}
|
results = await timer.measure_async(
|
||||||
done, pending = await timer.measure_async(
|
|
||||||
"build_details",
|
"build_details",
|
||||||
lambda: asyncio.wait(
|
lambda: asyncio.gather(*tasks, return_exceptions=True),
|
||||||
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] = {}
|
||||||
missing: List[str] = []
|
for city, result in zip(city_names, results):
|
||||||
for task in done:
|
if isinstance(result, Exception):
|
||||||
city = task_city_lookup[task]
|
errors[city] = str(result)
|
||||||
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}"
|
||||||
|
|||||||
Reference in New Issue
Block a user