Skip to content

Commit 33897bb

Browse files
committed
feat(webapp,clickhouse): decouple log search indexing from event inserts
Project closed source windows with durable watermarks and leases. Keep v2 reads and backfill disabled by default until sufficient history exists.
1 parent a8be3ae commit 33897bb

23 files changed

Lines changed: 1876 additions & 49 deletions

.server-changes/improve-global-log-search.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,4 +3,4 @@ area: webapp
33
type: improvement
44
---
55

6-
Global log search now supports a bounded search index and clearer time-range expansion while keeping existing search history available during rollout.
6+
Global log search now supports faster bounded substring matching and clearer time-range expansion. Existing search remains the default while the new index builds sufficient history.

apps/webapp/app/entry.server.tsx

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import { PassThrough } from "stream";
1010
import { initMollifierDrainerWorker } from "~/v3/mollifierDrainerWorker.server";
1111
import { initMollifierStaleSweepWorker } from "~/v3/mollifierStaleSweepWorker.server";
1212
import { initBillingLimitWorker } from "~/v3/billingLimitWorker.server";
13+
import { initLogsSearchProjectorWorker } from "~/v3/logsSearchProjectorWorker.server";
1314
import { initQueueMetricsConsumer, initQueueMetricsEmitter } from "~/v3/queueMetrics.server";
1415
import { bootstrap } from "./bootstrap";
1516
import { LocaleContextProvider } from "./components/primitives/LocaleProvider";
@@ -265,6 +266,7 @@ export const handleError = wrapHandleErrorWithSentry((error, { request }) => {
265266
initMollifierDrainerWorker();
266267
initMollifierStaleSweepWorker();
267268
initBillingLimitWorker();
269+
initLogsSearchProjectorWorker();
268270
initQueueMetricsEmitter();
269271
initQueueMetricsConsumer();
270272

apps/webapp/app/env.server.ts

Lines changed: 37 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1925,10 +1925,45 @@ const EnvironmentSchema = z
19251925
.nonnegative()
19261926
.optional(),
19271927

1928-
// v2 is populated forward-only. Keep reads on v1 until v2 has enough history or has been
1929-
// backfilled, then opt in explicitly per deployment.
1928+
// Keep reads on v1 until the scheduled v2 projector has enough history.
19301929
LOGS_SEARCH_TABLE_VERSION: z.enum(["v1", "v2"]).default("v1"),
19311930

1931+
// Scheduled logs-search projection. Disabled by default. The writer URL must reach both the
1932+
// task_events_v2 source and task_events_search_v2 destination tables.
1933+
LOGS_SEARCH_PROJECTOR_ENABLED: BoolEnv.default(false),
1934+
LOGS_SEARCH_PROJECTOR_CLICKHOUSE_URL: z
1935+
.string()
1936+
.optional()
1937+
.transform((v) => v ?? process.env.EVENTS_CLICKHOUSE_URL ?? process.env.CLICKHOUSE_URL),
1938+
LOGS_SEARCH_PROJECTOR_SAFETY_DELAY_SECONDS: z.coerce
1939+
.number()
1940+
.int()
1941+
.min(60)
1942+
.max(3600)
1943+
.default(120),
1944+
LOGS_SEARCH_PROJECTOR_MAX_WINDOWS_PER_TICK: z.coerce.number().int().min(1).max(20).default(5),
1945+
LOGS_SEARCH_PROJECTOR_MAX_EXECUTION_TIME_SECONDS: z.coerce
1946+
.number()
1947+
.int()
1948+
.min(1)
1949+
.max(300)
1950+
.default(120),
1951+
LOGS_SEARCH_PROJECTOR_MAX_ROWS_TO_READ: z.coerce.number().int().positive().default(10_000_000),
1952+
LOGS_SEARCH_PROJECTOR_MAX_MEMORY_USAGE: z.coerce
1953+
.number()
1954+
.int()
1955+
.positive()
1956+
.default(1_500_000_000),
1957+
LOGS_SEARCH_PROJECTOR_MAX_THREADS: z.coerce.number().int().min(1).max(8).default(2),
1958+
LOGS_SEARCH_PROJECTOR_BACKFILL_ENABLED: BoolEnv.default(false),
1959+
LOGS_SEARCH_PROJECTOR_MAX_BACKFILL_RANGE_DAYS: z.coerce
1960+
.number()
1961+
.int()
1962+
.min(1)
1963+
.max(90)
1964+
.default(7),
1965+
LOGS_SEARCH_PROJECTOR_MAX_BACKFILL_AGE_DAYS: z.coerce.number().int().min(1).max(90).default(90),
1966+
19321967
// Logs list pagination tuning.
19331968
LOGS_LIST_DEFAULT_PAGE_SIZE: z.coerce.number().int().positive().default(50),
19341969
LOGS_LIST_MAX_PAGE_SIZE: z.coerce.number().int().positive().default(100),
Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
import { type ActionFunctionArgs, type LoaderFunctionArgs, json } from "@remix-run/server-runtime";
2+
import { z } from "zod";
3+
import {
4+
LogsSearchProjectorConflictError,
5+
LogsSearchProjectorValidationError,
6+
} from "~/services/logsSearchProjector.server";
7+
import { getLogsSearchProjector } from "~/services/logsSearchProjectorInstance.server";
8+
import { logger } from "~/services/logger.server";
9+
import { requireAdminApiRequest } from "~/services/personalAccessToken.server";
10+
11+
const Body = z.discriminatedUnion("action", [
12+
z.object({ action: z.literal("pause") }),
13+
z.object({ action: z.literal("resume") }),
14+
z.object({ action: z.literal("cancelBackfill") }),
15+
z.object({
16+
action: z.literal("startBackfill"),
17+
from: z
18+
.string()
19+
.datetime()
20+
.transform((value) => new Date(value)),
21+
to: z
22+
.string()
23+
.datetime()
24+
.transform((value) => new Date(value)),
25+
}),
26+
]);
27+
28+
export async function loader({ request }: LoaderFunctionArgs) {
29+
await requireAdminApiRequest(request);
30+
return json(await getLogsSearchProjector().status());
31+
}
32+
33+
export async function action({ request }: ActionFunctionArgs) {
34+
const user = await requireAdminApiRequest(request);
35+
36+
try {
37+
const body = Body.parse(await request.json());
38+
const logsSearchProjector = getLogsSearchProjector();
39+
logger.info("Updating logs search projector", { userId: user.id, action: body.action });
40+
41+
switch (body.action) {
42+
case "pause":
43+
return json(await logsSearchProjector.pause());
44+
case "resume":
45+
return json(await logsSearchProjector.resume());
46+
case "cancelBackfill":
47+
return json(await logsSearchProjector.cancelBackfill());
48+
case "startBackfill":
49+
return json(await logsSearchProjector.startBackfill(body));
50+
}
51+
} catch (error) {
52+
if (error instanceof LogsSearchProjectorConflictError) {
53+
return json({ error: error.message }, { status: 409 });
54+
}
55+
if (
56+
error instanceof LogsSearchProjectorValidationError ||
57+
error instanceof z.ZodError ||
58+
error instanceof SyntaxError
59+
) {
60+
return json(
61+
{ error: error instanceof Error ? error.message : String(error) },
62+
{ status: 400 }
63+
);
64+
}
65+
throw error;
66+
}
67+
}

apps/webapp/app/services/clickhouse/clickhouseFactory.server.ts

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,33 @@ const defaultLogsClickhouseClient = singleton(
3737
initializeLogsClickhouseClient
3838
);
3939

40+
const logsSearchProjectorClickhouseClient = singleton(
41+
"logsSearchProjectorClickhouseClient",
42+
initializeLogsSearchProjectorClickhouseClient
43+
);
44+
45+
function initializeLogsSearchProjectorClickhouseClient() {
46+
if (!env.LOGS_SEARCH_PROJECTOR_CLICKHOUSE_URL) {
47+
throw new Error("LOGS_SEARCH_PROJECTOR_CLICKHOUSE_URL is not set");
48+
}
49+
50+
const url = new URL(env.LOGS_SEARCH_PROJECTOR_CLICKHOUSE_URL);
51+
url.searchParams.delete("secure");
52+
53+
return new ClickHouse({
54+
url: url.toString(),
55+
name: "logs-search-projector",
56+
keepAlive: {
57+
enabled: env.CLICKHOUSE_KEEP_ALIVE_ENABLED === "1",
58+
idleSocketTtl: env.CLICKHOUSE_KEEP_ALIVE_IDLE_SOCKET_TTL_MS,
59+
},
60+
logLevel: env.CLICKHOUSE_LOG_LEVEL,
61+
compression: { request: true },
62+
maxOpenConnections: Math.min(env.CLICKHOUSE_MAX_OPEN_CONNECTIONS, 2),
63+
requestTimeoutMs: (env.LOGS_SEARCH_PROJECTOR_MAX_EXECUTION_TIME_SECONDS + 30) * 1000,
64+
});
65+
}
66+
4067
function getLogsListClickhouseSettings() {
4168
return {
4269
max_memory_usage: env.CLICKHOUSE_LOGS_LIST_MAX_MEMORY_USAGE.toString(),
@@ -653,6 +680,10 @@ export function getDefaultLogsClickhouseClient(): ClickHouse {
653680
return defaultLogsClickhouseClient;
654681
}
655682

683+
export function getLogsSearchProjectorClickhouseClient(): ClickHouse {
684+
return logsSearchProjectorClickhouseClient;
685+
}
686+
656687
/** Queue-metrics client for callers with no organization in scope (the ingestion consumer). */
657688
export function getQueueMetricsClickhouseClient(): ClickHouse {
658689
return defaultQueueMetricsClickhouseClient;

0 commit comments

Comments
 (0)