Skip to content

Commit 3428b34

Browse files
committed
fix(access-logs): harden ClickHouse pipeline, parse handshake errors, i18n
- Fix MV regex escaping, replaceAll for dates, explicit UTC end-to-end - Batch dedup via insert token + dependent MV dedup; drop dead offset column - Versioned MV (v3) with auto-drop of legacy views on ensureSchema - Parse connection-error lines (from IP rejected <msg>) as rejected with source IP + outbound_tag=handshake-error; exclude empty dest from top lists - Cache resolved CH config; remove leftover DuckDB/Parquet code and locale keys - Add common.attention / common.restore locale keys (all languages)
1 parent 913b0b3 commit 3428b34

6 files changed

Lines changed: 62 additions & 15 deletions

File tree

‎scripts/test-access-logs-clickhouse.js‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,21 @@ async function offlineTests() {
7474
assert.ok(mvDdl.includes('replaceAll(ts_str'), 'date slashes replaced with replaceAll');
7575
// Unparsed lines must not land in 1970 (instantly TTL-dropped).
7676
assert.ok(mvDdl.includes("now('UTC')"), 'zero timestamps fall back to now()');
77+
// Connection-level error lines are parsed via the fallback regex and tagged.
78+
assert.ok(mvDdl.includes('handshake-error'), 'handshake-error fallback tag present');
79+
assert.ok(mvDdl.includes('ne > 0'), 'error-line group participates in parse_ok');
80+
81+
// The fallback regex matches an error line but NOT as an access record.
82+
const errLine = '2026/07/10 17:41:28.205208 from 95.24.24.226:9048 rejected proxy/vless/encoding: invalid request user id: 55837f55-c7ee-4533-b2a6-0ace8a266802';
83+
assert.strictEqual(parseLikeMv(errLine).parse_ok, 0, 'error line is not a normal access record');
84+
const errRe = new RegExp(clickhouse.CH_ERR_RE);
85+
const em = errRe.exec(errLine);
86+
assert.ok(em, 'error line matches fallback regex');
87+
assert.strictEqual(em[2], '95.24.24.226:9048', 'fallback captures source');
88+
assert.strictEqual(em[3], 'rejected', 'fallback captures action');
89+
// A real access line must still be handled by the primary parser, not misrouted.
90+
assert.ok(errRe.exec('2023/11/22 17:01:32 1.2.3.4:1122 accepted tcp:example.com:443 [in -> direct]'),
91+
'fallback also matches normal lines (primary takes precedence in the MV)');
7792

7893
// Retention clamps to sane bounds (0/NaN falls back to the 30-day default).
7994
assert.ok(clickhouse.schemaStatements(-5)[1].includes('INTERVAL 1 DAY'), 'retention floor');
@@ -152,6 +167,7 @@ async function onlineTests() {
152167
const batchId = 'test-batch-' + Date.now();
153168
await clickhouse.insertRaw([
154169
{ node_id: 'n1', raw: '2023/11/22 17:01:32 1.2.3.4:1122 accepted tcp:example.com:443 [vless-in -> direct] email: 42' },
170+
{ node_id: 'n1', raw: '2023/11/22 17:01:33 from 9.9.9.9:5000 rejected proxy/vless/encoding: invalid request user id: abc' },
155171
], batchId);
156172

157173
// Give the MV a moment (insert is synchronous, but read is eventually there).
@@ -161,6 +177,14 @@ async function onlineTests() {
161177
assert.strictEqual(res.rows[0].network, 'tcp');
162178
assert.strictEqual(res.rows[0].dest_host, 'example.com');
163179

180+
// The connection-error line parsed via fallback: rejected, source set, tagged.
181+
const errRes = await clickhouse.query(
182+
"SELECT source_ip, action, outbound_tag, parse_ok FROM access_events WHERE source_ip = '9.9.9.9' LIMIT 1");
183+
assert.ok(errRes.ok && errRes.rows.length >= 1, 'error line stored');
184+
assert.strictEqual(errRes.rows[0].action, 'rejected', 'error line action');
185+
assert.strictEqual(errRes.rows[0].outbound_tag, 'handshake-error', 'error line tagged');
186+
assert.strictEqual(Number(errRes.rows[0].parse_ok), 1, 'error line counts as parsed');
187+
164188
await clickhouse.truncate();
165189
void orig;
166190
console.log(' online: end-to-end OK');

‎src/locales/en.json‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,9 @@
5050
"server": "Server",
5151
"ports": "Ports",
5252
"clear": "Clear",
53-
"reset": "Reset"
53+
"reset": "Reset",
54+
"attention": "Attention",
55+
"restore": "Restore"
5456
},
5557
"nav": {
5658
"dashboard": "Dashboard",

‎src/locales/ru.json‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,9 @@
5050
"server": "Сервер",
5151
"ports": "Порты",
5252
"clear": "Очистить",
53-
"reset": "Сбросить"
53+
"reset": "Сбросить",
54+
"attention": "Внимание",
55+
"restore": "Восстановить"
5456
},
5557
"nav": {
5658
"dashboard": "Главная",

‎src/locales/zh-CN.json‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,9 @@
5050
"server": "服务器",
5151
"ports": "端口",
5252
"clear": "清空",
53-
"reset": "重置"
53+
"reset": "重置",
54+
"attention": "注意",
55+
"restore": "恢复"
5456
},
5557
"nav": {
5658
"dashboard": "仪表盘",

‎src/services/accessLogs/clickhouseService.js‎

Lines changed: 27 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,17 @@ const CH_LINE_RE =
5151
'(?:\\s+email:\\s*(\\S+))?' +
5252
'\\s*$';
5353

54+
// Fallback for connection-level error lines that are NOT access records, e.g.:
55+
// 2026/07/10 17:41:28 from 95.24.24.226:9048 rejected proxy/vless/encoding: invalid request user id: <uuid>
56+
// These carry a timestamp, source and action but no "tcp:/udp:" destination
57+
// (an attacker/scanner hitting an inbound with a bad UUID). Captures:
58+
// 1 ts, 2 src, 3 action. The tail (error text) stays in `raw`.
59+
const CH_ERR_RE =
60+
'^(\\d{4}/\\d{2}/\\d{2} \\d{2}:\\d{2}:\\d{2}(?:\\.\\d+)?)\\s+' +
61+
'(?:from\\s+)?' +
62+
'(\\S+?)\\s+' +
63+
'(accepted|rejected|blocked)(?:\\s|$)';
64+
5465
// Escape a JS string for use inside a single-quoted ClickHouse SQL literal.
5566
// ClickHouse collapses unknown escapes ('\d' -> 'd'), so every backslash must
5667
// be doubled or the inlined regex silently loses all its character classes.
@@ -63,9 +74,9 @@ function sqlString(s) {
6374
// the stale one). The version is part of the MV name; ensureSchema drops any
6475
// older names listed here. Bump MV_VERSION whenever the MV definition changes
6576
// and append the previous name to LEGACY_MV_NAMES.
66-
const MV_VERSION = 2;
77+
const MV_VERSION = 3;
6778
const MV_NAME = `access_events_mv_v${MV_VERSION}`;
68-
const LEGACY_MV_NAMES = ['access_events_mv'];
79+
const LEGACY_MV_NAMES = ['access_events_mv', 'access_events_mv_v2'];
6980

7081
// ── Config ────────────────────────────────────────────────────────────────
7182

@@ -248,16 +259,21 @@ function schemaStatements(retentionDays) {
248259
SETTINGS non_replicated_deduplication_window = 1000`,
249260

250261
// Parse raw -> structured on insert. Everything derives from `raw`.
251-
// A line that does not match the regex (or carries a broken timestamp)
252-
// still lands with parse_ok = 0 and event_time = now(), so no data is
253-
// lost and it stays searchable by raw text within the retention window.
262+
// Two shapes are recognised: (n) a normal access line and, as a fallback,
263+
// (ne) a connection-level error line ("from IP rejected <msg>", no
264+
// destination) which is tagged outbound_tag = 'handshake-error' so the
265+
// action counters stay accurate and the attacking source IP is visible.
266+
// A line matching neither still lands with parse_ok = 0 and
267+
// event_time = now(), so nothing is lost and it stays searchable by raw.
254268
`CREATE MATERIALIZED VIEW IF NOT EXISTS ${MV_NAME} TO access_events AS
255269
WITH
256270
extractGroups(raw, '${sqlString(CH_LINE_RE)}') AS g,
257271
length(g) AS n,
258-
if(n > 0, g[1], '') AS ts_str,
259-
if(n > 0, g[2], '') AS src,
260-
if(n > 0, g[3], '') AS act,
272+
extractGroups(raw, '${sqlString(CH_ERR_RE)}') AS ge,
273+
length(ge) AS ne,
274+
if(n > 0, g[1], if(ne > 0, ge[1], '')) AS ts_str,
275+
if(n > 0, g[2], if(ne > 0, ge[2], '')) AS src,
276+
if(n > 0, g[3], if(ne > 0, ge[3], '')) AS act,
261277
if(n > 0, g[4], '') AS net,
262278
if(n > 0, g[5], '') AS dst,
263279
if(n > 0, g[6], '') AS route,
@@ -283,10 +299,10 @@ function schemaStatements(retentionDays) {
283299
dst_port AS dest_port,
284300
net AS network,
285301
in_tag AS inbound_tag,
286-
out_tag AS outbound_tag,
302+
if(n > 0, out_tag, if(ne > 0, 'handshake-error', '')) AS outbound_tag,
287303
act AS action,
288304
raw,
289-
toUInt8(n > 0) AS parse_ok
305+
toUInt8(n > 0 OR ne > 0) AS parse_ok
290306
FROM access_ingest`,
291307
];
292308
}
@@ -392,6 +408,7 @@ async function truncate() {
392408

393409
module.exports = {
394410
CH_LINE_RE,
411+
CH_ERR_RE,
395412
readConfig,
396413
getClient,
397414
reset,

‎src/services/accessLogs/searchService.js‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -186,7 +186,7 @@ async function overview(filters = {}, opts = {}) {
186186

187187
const topDestSql = `
188188
SELECT ${DEST} AS dest, count() AS hits
189-
FROM access_events ${where}
189+
FROM access_events ${where} ${andWhere} ${DEST} != ''
190190
GROUP BY dest ORDER BY hits DESC LIMIT ${topN}`;
191191

192192
const topPortsSql = `
@@ -196,7 +196,7 @@ async function overview(filters = {}, opts = {}) {
196196

197197
const topBlockedSql = `
198198
SELECT ${DEST} AS dest, count() AS hits
199-
FROM access_events ${where} ${andWhere} action IN ('blocked','rejected')
199+
FROM access_events ${where} ${andWhere} action IN ('blocked','rejected') AND ${DEST} != ''
200200
GROUP BY dest ORDER BY hits DESC LIMIT ${topN}`;
201201

202202
const usersSql = `

0 commit comments

Comments
 (0)