Compare commits
36
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2b6ec59286 | ||
|
|
fdc7f01c26 | ||
|
|
f31b1f62d4 | ||
|
|
7be069d73f | ||
|
|
88373985ec | ||
|
|
7acc49e1fb | ||
|
|
81b3128bb0 | ||
|
|
1d0b3f422d | ||
|
|
0aef0b6c08 | ||
|
|
cd1d6b2d0d | ||
|
|
c888f23901 | ||
|
|
54d02098c8 | ||
|
|
5fc0e8b3cb | ||
|
|
045cdf1f75 | ||
|
|
deed7bdfb0 | ||
|
|
c4d9ade85e | ||
|
|
80daa9f045 | ||
|
|
ef7708bf7d | ||
|
|
750f3aa598 | ||
|
|
c57ee12da1 | ||
|
|
494e16b3b3 | ||
|
|
6bf3b40cc7 | ||
|
|
8743fcc0b5 | ||
|
|
e34dcd6bc6 | ||
|
|
f9fecfc144 | ||
|
|
a59f3132ee | ||
|
|
c7f53e4f7e | ||
|
|
8583bcdf17 | ||
|
|
b12eb0a038 | ||
|
|
b72423c64d | ||
|
|
6b34fc97ec | ||
|
|
8ce8978755 | ||
|
|
869cad88e4 | ||
|
|
f136cf3f6b | ||
|
|
cd3ee5b5d8 | ||
|
|
93288afa0b |
@@ -64,22 +64,11 @@ AI_LLM_API_KEY= # REQUIRED if AI_ANALYSIS_ENABLED=true. L
|
||||
AI_LLM_BASE_URL=http://100.121.180.82:20128/api/v1 # LLM API base URL (omniroute — OpenAI-compatible router on imrnes, Tailscale 100.121.180.82)
|
||||
AI_LLM_MODEL=text # LLM text model name (default: text)
|
||||
# AI_LLM_VISION_MODEL= # Vision model for image analysis (falls back to AI_LLM_MODEL)
|
||||
# AI_LLM_EMBEDDING_MODEL= # Embedding model for semantic moderation cache (optional; enables near-duplicate text reuse to save LLM calls)
|
||||
# AI_LLM_EMBEDDING_MIN_SIMILARITY=0.97 # Min cosine similarity to reuse a cached verdict (default: 0.97)
|
||||
QDRANT_URL=http://100.121.180.82:6333 # Qdrant vector store for embeddings (semantic cache); when set, vectors are stored/searched in Qdrant instead of Postgres
|
||||
# QDRANT_COLLECTION=gmw_text_moderation # Qdrant collection name (default: gmw_text_moderation)
|
||||
# QDRANT_API_KEY= # Qdrant API key (optional)
|
||||
AI_LLM_MAX_CONCURRENT=5 # Max concurrent LLM API calls (default: 5)
|
||||
AI_LLM_IMAGE_MAX_DIMENSION=1024 # Max image dimension in pixels before resize (default: 1024)
|
||||
AI_LLM_TEXT_BATCH_SIZE=20 # Max messages per text-only moderation batch (default: 20)
|
||||
AI_LLM_MEDIA_ANALYSIS_TIMEOUT_MS=60000 # Timeout in ms for media analysis calls (default: 60000)
|
||||
AI_LLM_TEXT_ANALYSIS_TIMEOUT_MS=30000 # Timeout in ms for text-only analysis calls (default: 30000)
|
||||
AI_LLM_JEV_ENABLED=false # Use TypeSafe Jev (System One) as PRIMARY text analyzer; LLM is fallback
|
||||
AI_LLM_JEV_API_KEY= # REQUIRED if AI_LLM_JEV_ENABLED=true. 9router/TypeSafe API key for /v1/systemone
|
||||
AI_LLM_JEV_BASE_URL=http://127.0.0.1:4014 # 9router base URL (default: local 9router; prod: https://9router.asepharyana.my.id)
|
||||
AI_LLM_JEV_MODEL=oc/jev-1.13-free # Jev model id (default: oc/jev-1.13-free)
|
||||
AI_LLM_JEV_TIMEOUT_MS=45000 # Timeout in ms for a Jev systemone batch call (default: 45000)
|
||||
AI_LLM_JEV_MIN_CONFIDENCE=0.9 # Min status-choice confidence to accept a Jev verdict (default: 0.9)
|
||||
|
||||
# === AI Analysis Tuning ===
|
||||
AI_ANALYSIS_DEBOUNCE_MS=500 # Debounce window for batching messages in ms (default: 500)
|
||||
|
||||
@@ -32,29 +32,26 @@ jobs:
|
||||
with:
|
||||
node-version: 22
|
||||
|
||||
- name: Install pnpm
|
||||
run: corepack enable && corepack prepare pnpm@11 --activate
|
||||
- name: Setup Bun
|
||||
uses: oven-sh/setup-bun@v2
|
||||
with:
|
||||
bun-version: 1.3.14
|
||||
|
||||
- name: Install deps (backend)
|
||||
working-directory: services/backend
|
||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
||||
|
||||
- name: Typecheck + test (backend)
|
||||
- name: Install deps + test (backend)
|
||||
working-directory: services/backend
|
||||
run: |
|
||||
bun install --frozen-lockfile
|
||||
./node_modules/.bin/tsc --noEmit
|
||||
# e2e.test.ts requires a live backend (API_BASE) — run unit tests only
|
||||
./node_modules/.bin/vitest run --exclude "src/e2e.test.ts"
|
||||
# src/e2e.test.ts requires a live backend (API_BASE) — unit tests
|
||||
# live in tests/ and are excluded by the bun test dir.
|
||||
bun test tests/
|
||||
|
||||
- name: Install deps (discord-gateway)
|
||||
working-directory: services/discord-gateway
|
||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
||||
|
||||
- name: Typecheck + test (discord-gateway)
|
||||
- name: Install deps + test (discord-gateway)
|
||||
working-directory: services/discord-gateway
|
||||
run: |
|
||||
bun install --frozen-lockfile
|
||||
./node_modules/.bin/tsc --noEmit
|
||||
./node_modules/.bin/vitest run
|
||||
bun test tests/
|
||||
|
||||
- name: Biome check (all services)
|
||||
run: |
|
||||
@@ -151,9 +148,17 @@ jobs:
|
||||
|
||||
attic_push_vps_hop() {
|
||||
echo "Fallback: VPS-hop attic push"
|
||||
# Recover the client binary BEFORE the fallback can use it: the
|
||||
# bootstrap cascade below resets ATTIC_BIN="" and never restores it
|
||||
# in the fallback branch, so `sudo $ATTIC_BIN push` used to run as
|
||||
# `sudo push` -> "sudo: 'push': command not found". On the VPS the
|
||||
# closure lives at the canonical ATTIC_DIR path.
|
||||
VPS_ATTIC="/nix/store/fygyy3yk4rqdknxkiwkqambpnhyax0k4-attic-0.1.0/bin/attic"
|
||||
ssh "$VPS_USER@$VPS_HOST" "test -x '$VPS_ATTIC'" \
|
||||
|| ssh "$VPS_USER@$VPS_HOST" "sudo /nix/var/nix/profiles/default/bin/nix-store --realise '$ATTIC_DIR'"
|
||||
# Copy closure to VPS (fast if attic already has it via substitute)
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo /nix/var/nix/profiles/default/bin/nix-store --realise '$STORE_PATH'" 2>/dev/null \
|
||||
|| nix copy --to "ssh://$VPS_USER@$VPS_HOST" "$STORE_PATH"
|
||||
|| nix copy --to "ssh://***@$VPS_HOST" "$STORE_PATH"
|
||||
# Push from VPS → Attic over Tailscale.
|
||||
# --ignore-upstream-cache-filter is REQUIRED: without it, attic skips
|
||||
# writing the narinfo to gmw when chunks exist in the upstream
|
||||
@@ -162,7 +167,7 @@ jobs:
|
||||
# sudo: attic must read root's config (~/.config/attic), which has
|
||||
# the imrnes-ts server → Tailscale. Non-root users' configs only
|
||||
# have the public `pub` server → "Server imrnes-ts does not exist".
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo $ATTIC_BIN push imrnes-ts:gmw '$STORE_PATH' --jobs 4 --ignore-upstream-cache-filter" \
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo $VPS_ATTIC push imrnes-ts:gmw '$STORE_PATH' --jobs 4 --ignore-upstream-cache-filter" \
|
||||
|| echo "attic push failed (non-fatal; ssh copy fallback below)"
|
||||
}
|
||||
|
||||
|
||||
@@ -23,3 +23,7 @@ result
|
||||
|
||||
# Playwright MCP artifacts
|
||||
.playwright-mcp/
|
||||
findings.md
|
||||
findings.md
|
||||
progress.md
|
||||
task_plan.md
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
{ "id": "gateway", "type": "backend", "label": "discord-gateway", "sublabel": "selfbot :4016", "pos": [260, 240], "size": [140, 60], "tag": "discord.js-selfbot-v13" },
|
||||
{ "id": "backend", "type": "backend", "label": "gmw-backend", "sublabel": "Express · oRPC · WS :4001", "pos": [640, 240], "size": [140, 60], "tag": "Drizzle ORM" },
|
||||
{ "id": "proxy", "type": "cloud", "label": "nginx proxy", "sublabel": "reverse proxy :4009", "pos": [960, 240], "size": [140, 60], "tag": "nginx" },
|
||||
{ "id": "ai", "type": "backend", "label": "AI Moderation", "sublabel": "LLM caller · embeddings", "pos": [260, 400], "size": [140, 60] },
|
||||
{ "id": "ai", "type": "backend", "label": "AI Moderation", "sublabel": "LLM caller · vision", "pos": [260, 400], "size": [140, 60] },
|
||||
{ "id": "frontend", "type": "frontend", "label": "gmw-frontend", "sublabel": "Next.js 16 SSR :4017", "pos": [960, 400], "size": [140, 60], "tag": "React 19 · Tailwind v4" },
|
||||
{ "id": "llm", "type": "cloud", "label": "9router LLM", "sublabel": "text + vision API", "pos": [260, 540], "size": [140, 60], "tag": "AI_LLM_BASE_URL" },
|
||||
{ "id": "browser", "type": "external", "label": "Dashboard Users", "sublabel": "browser · partysocket", "pos": [960, 540], "size": [140, 60] }
|
||||
|
||||
@@ -35,9 +35,11 @@
|
||||
|
||||
# ---- Shared build tools ----
|
||||
nodejs = pkgs.nodejs_22;
|
||||
pnpm = pkgs.pnpm.override { nodejs = nodejs; };
|
||||
# Bun for deps/install (replaces pnpm); keeps nodejs for the tsc +
|
||||
# fix-imports.mjs build path (Bun's own bundler is not used for dist).
|
||||
bun = pkgs.bun;
|
||||
|
||||
pnpmInstall = ''
|
||||
bunInstall = ''
|
||||
export HOME=$TMPDIR/home
|
||||
export npm_config_cache=$TMPDIR/npm-cache
|
||||
mkdir -p $npm_config_cache
|
||||
@@ -48,15 +50,9 @@
|
||||
export GIT_SSL_CAINFO=${pkgs.cacert}/etc/ssl/certs/ca-bundle.crt
|
||||
export NIX_SSL_CERT_FILE=${pkgs.cacert}/etc/ssl/certs/ca-bundle.crt
|
||||
|
||||
# pnpm uses node-gyp for native addons — provide build tools (kept for
|
||||
# the rare case a prebuilt is unavailable and it falls back to compile).
|
||||
export CPPFLAGS="-I${pkgs.lib.getDev pkgs.openssl}/include"
|
||||
export LDFLAGS="-L${pkgs.lib.getLib pkgs.openssl}/lib"
|
||||
|
||||
pnpm install --no-frozen-lockfile --ignore-scripts 2>&1
|
||||
|
||||
# Build native addons that need compilation
|
||||
pnpm rebuild 2>&1 || true
|
||||
# Build native addons (bun install runs postinstall scripts for
|
||||
# @discordjs/opus / sharp unless trustedDependencies restricts).
|
||||
bun install 2>&1
|
||||
'';
|
||||
|
||||
# Shrink the shipped node_modules to production deps only. The full
|
||||
@@ -74,28 +70,15 @@
|
||||
# Must run AFTER tsc (typescript is a devDep) and after native builds.
|
||||
pruneProd = ''
|
||||
echo "=== Pruning devDependencies (production-only node_modules) ==="
|
||||
pnpm list --prod --depth 999 --parseable 2>/dev/null \
|
||||
| grep -o '\.pnpm/[^/]*' | sort -u > $TMPDIR/prod-pnms.txt
|
||||
( cd node_modules/.pnpm \
|
||||
&& for d in */; do \
|
||||
d="''${d%/}"; \
|
||||
[ "$d" = "node_modules" ] && continue; \
|
||||
grep -qF ".pnpm/$d" $TMPDIR/prod-pnms.txt || rm -rf "$d"; \
|
||||
done ) || true
|
||||
# Drop runtime-dead packages that still land in the prod graph:
|
||||
# - `@types/*` (pure TypeScript declarations) get pulled in as
|
||||
# REAL dependencies by type-aware deps (discord-api-types ->
|
||||
# @types/node, pg-protocol -> @types/pg, ...) even though nothing
|
||||
# ever `require`s them at runtime. Safe to strip.
|
||||
# - `opusscript` is only a pure-JS fallback Opus engine that
|
||||
# prism-media's loader uses IF `@discordjs/opus` (native, always
|
||||
# present/prebuilt) fails to load. Since the native engine loads,
|
||||
# opusscript is never executed — dead weight pulled in via
|
||||
# discord.js-selfbot-v13's dependency. Strip it too.
|
||||
( cd node_modules/.pnpm && rm -rf @types+* opusscript@* 2>/dev/null ) || true
|
||||
# Drop symlinks whose .pnpm target was pruned (top-level, scoped dirs,
|
||||
# hoist, .bin — any depth). Mirrors stdenv's noBrokenSymlinks check,
|
||||
# which would otherwise fail the fixupPhase.
|
||||
# bun install's layout: node_modules/<pkg> for prod deps; devDeps are
|
||||
# also present during build (needed for tsc). Keep only what the prod
|
||||
# graph needs: simplest robust approach is `bun install --production`
|
||||
# semantics — but bun keeps the same flat layout; since the Nix build
|
||||
# already ran `bun install` (full, scripts on), prune dev-only top
|
||||
# entries that were only pulled by devDeps (typescript, biome, vitest,
|
||||
# drizzle-kit, tsx, @types/*).
|
||||
find node_modules -maxdepth 2 -type d \( -name 'typescript' -o -name '@biomejs' -o -name 'vitest' -o -name 'drizzle-kit' -o -name 'tsx' -o -name 'esbuild' \) -prune -exec rm -rf {} + 2>/dev/null || true
|
||||
rm -rf node_modules/.bin/tsc node_modules/.bin/vitest node_modules/.bin/biome node_modules/.bin/drizzle-kit 2>/dev/null || true
|
||||
find node_modules -type l ! -exec test -e {} \; -delete 2>/dev/null || true
|
||||
du -sh node_modules
|
||||
'';
|
||||
@@ -107,11 +90,11 @@
|
||||
|
||||
src = ./services/backend;
|
||||
|
||||
nativeBuildInputs = [ nodejs pnpm pkgs.python3 pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
nativeBuildInputs = [ nodejs bun pkgs.python3 pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Compiling TypeScript ==="
|
||||
npx tsc 2>&1
|
||||
./node_modules/.bin/tsc 2>&1
|
||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||
node scripts/fix-imports.mjs
|
||||
echo "=== Build complete ==="
|
||||
@@ -149,7 +132,7 @@ WRAPPER
|
||||
# libvips download), so no cmake or rust toolchain is needed.
|
||||
# python3/gnumake/gcc stay as node-gyp fallback for @discordjs/opus.
|
||||
nativeBuildInputs = [
|
||||
nodejs pnpm
|
||||
nodejs bun
|
||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||
pkgs.pkg-config
|
||||
pkgs.openssl
|
||||
@@ -177,22 +160,11 @@ WRAPPER
|
||||
# neither needed nor wanted here. Skip it entirely.
|
||||
dontFixup = true;
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
echo "=== Building native voice deps ==="
|
||||
# pnpm rebuild aborts on the first failing package and runs scripts
|
||||
# from the wrong cwd — build each native dep explicitly with its own
|
||||
# install script. Each failure is tolerated (|| true); the packages
|
||||
# @discordjs/opus ships prebuilt binaries for Node 22 (ABI node-v127,
|
||||
# linux-x64-glibc-2.35) — node-pre-gyp downloads the prebuilt .node
|
||||
# instead of compiling C++ from source. With build_from_source unset
|
||||
# (above), `pnpm rebuild` runs the package's own install script which
|
||||
# fetches the matching prebuilt; it only falls back to a source build
|
||||
# if the download fails. This keeps voice working without a per-build
|
||||
# native compile.
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Rebuilding @discordjs/opus (prebuilt download) ==="
|
||||
pnpm rebuild @discordjs/opus 2>&1 || true
|
||||
bun pm rebuild @discordjs/opus 2>&1 || true
|
||||
echo "=== Compiling TypeScript ===="
|
||||
npx tsc 2>&1
|
||||
./node_modules/.bin/tsc 2>&1
|
||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||
node scripts/fix-imports.mjs
|
||||
echo "=== Build complete ==="
|
||||
@@ -228,13 +200,13 @@ WRAPPER
|
||||
|
||||
src = frontendSrc;
|
||||
|
||||
nativeBuildInputs = [ nodejs pnpm pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
nativeBuildInputs = [ nodejs bun pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Building Next.js SSR (standalone) ==="
|
||||
export NEXT_TELEMETRY_DISABLED=1
|
||||
export GMW_BACKEND_URL=http://127.0.0.1:4001
|
||||
npx next build 2>&1
|
||||
./node_modules/.bin/next build 2>&1
|
||||
'';
|
||||
|
||||
installPhase = ''
|
||||
@@ -312,13 +284,13 @@ WRAPPER
|
||||
|
||||
devShells.default = pkgs.mkShell {
|
||||
buildInputs = [
|
||||
nodejs pnpm
|
||||
nodejs bun
|
||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||
pkgs.rustc pkgs.cargo
|
||||
pkgs.ffmpeg-headless
|
||||
];
|
||||
shellHook = ''
|
||||
echo "GMW dev shell ready — node $(node --version), pnpm $(pnpm --version)"
|
||||
echo "GMW dev shell ready — node $(node --version), bun $(bun --version)"
|
||||
'';
|
||||
};
|
||||
});
|
||||
|
||||
@@ -60,7 +60,7 @@ src/
|
||||
| moderation | Moderation actions & metrics | `ai_moderations`, `moderation_actions` |
|
||||
| media | Media file management | `media_attachments` |
|
||||
| dashboard | Stats aggregation | Various (read-only) |
|
||||
| knowledge | Semantic search | Qdrant vector DB |
|
||||
| knowledge | Channel cultures & glossary browser | `channel_cultures`, `term_glossary_cache` |
|
||||
| chatbot | AI chatbot with tools | `chatbot_history` |
|
||||
| health | Health checks + metrics | Various |
|
||||
| analysis | Text analysis cache | `text_analysis_cache` |
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,2 @@
|
||||
[test]
|
||||
preload = ["./tests/setup-env.ts"]
|
||||
@@ -5,7 +5,7 @@
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"packageManager": "pnpm@11.20.0",
|
||||
"packageManager": "bun@1.3.14",
|
||||
"engines": {
|
||||
"node": ">=22.12.0",
|
||||
"pnpm": ">=9.0.0"
|
||||
@@ -16,14 +16,15 @@
|
||||
"format": "biome format --write .",
|
||||
"lint": "biome check --diagnostic-level=error .",
|
||||
"start": "node dist/index.js",
|
||||
"test": "vitest run",
|
||||
"typecheck": "tsc --noEmit"
|
||||
"test": "bun test tests/",
|
||||
"typecheck": "tsc --noEmit",
|
||||
"test:e2e": "bun test src/e2e.test.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@orpc/server": "1.15.2",
|
||||
"@orpc/server": "1.15.3",
|
||||
"axios": "^1.20.0",
|
||||
"dotenv": "^18.0.1",
|
||||
"drizzle-orm": "^0.45.2",
|
||||
"drizzle-orm": "^0.45.3",
|
||||
"express": "^5.2.1",
|
||||
"helmet": "^8.1.0",
|
||||
"ioredis": "^6.0.0",
|
||||
@@ -41,6 +42,6 @@
|
||||
"@types/ws": "^8.18.1",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "^5.0.1"
|
||||
"@types/bun": "latest"
|
||||
}
|
||||
}
|
||||
}
|
||||
Generated
-2677
File diff suppressed because it is too large
Load Diff
@@ -1,10 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const searchQuerySchema = z.object({
|
||||
q: z.string().default(""),
|
||||
channelId: z.string().optional(),
|
||||
guildId: z.string().optional(),
|
||||
limit: z.coerce.number().int().positive().max(100).default(20),
|
||||
});
|
||||
|
||||
export type SearchQuery = z.infer<typeof searchQuerySchema>;
|
||||
@@ -1,5 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const healthCheckSchema = z.object({
|
||||
verbose: z.coerce.boolean().optional().default(false),
|
||||
});
|
||||
@@ -1,80 +0,0 @@
|
||||
/**
|
||||
* moderationMetrics.ts
|
||||
*
|
||||
* Prometheus metrics for AI moderation pipeline.
|
||||
* Defined in backend (where prom-client is installed + /api/metrics endpoint).
|
||||
*/
|
||||
import { Counter, Histogram } from "prom-client";
|
||||
|
||||
// ── LLM Call Metrics ──
|
||||
export const llmCallsTotal = new Counter({
|
||||
name: "moderation_llm_calls_total",
|
||||
help: "Total LLM moderation calls",
|
||||
labelNames: ["path", "model"] as const,
|
||||
});
|
||||
|
||||
export const llmCallDuration = new Histogram({
|
||||
name: "moderation_llm_call_duration_ms",
|
||||
help: "LLM moderation call duration (ms)",
|
||||
labelNames: ["path", "status"] as const,
|
||||
buckets: [500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 120000],
|
||||
});
|
||||
|
||||
export const llmTokensTotal = new Counter({
|
||||
name: "moderation_llm_tokens_total",
|
||||
help: "Total tokens consumed by LLM moderation",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
// ── Cache Metrics ──
|
||||
export const moderationCacheHits = new Counter({
|
||||
name: "moderation_cache_hits_total",
|
||||
help: "Moderation cache hits",
|
||||
labelNames: ["layer"] as const,
|
||||
});
|
||||
|
||||
export const moderationCacheMisses = new Counter({
|
||||
name: "moderation_cache_misses_total",
|
||||
help: "Moderation cache misses",
|
||||
labelNames: ["layer"] as const,
|
||||
});
|
||||
|
||||
// ── Media Analysis Metrics ──
|
||||
export const mediaAnalysesTotal = new Counter({
|
||||
name: "moderation_media_analyses_total",
|
||||
help: "Media analyses performed",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
export const mediaDownloadDuration = new Histogram({
|
||||
name: "moderation_media_download_duration_ms",
|
||||
help: "Media download duration (ms)",
|
||||
labelNames: ["source"] as const,
|
||||
buckets: [100, 500, 1000, 2000, 5000, 10000, 30000],
|
||||
});
|
||||
|
||||
// ── Batch & Error Metrics ──
|
||||
export const moderationBatchSize = new Histogram({
|
||||
name: "moderation_batch_size",
|
||||
help: "Messages per batch",
|
||||
labelNames: ["path"] as const,
|
||||
buckets: [1, 5, 10, 20, 50, 100],
|
||||
});
|
||||
|
||||
export const moderationErrors = new Counter({
|
||||
name: "moderation_errors_total",
|
||||
help: "Moderation errors",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
export const webSearchCalls = new Counter({
|
||||
name: "moderation_websearch_calls_total",
|
||||
help: "Wikipedia web-search calls",
|
||||
labelNames: ["status"] as const,
|
||||
});
|
||||
|
||||
export const autoDeleteActions = new Counter({
|
||||
name: "moderation_auto_delete_total",
|
||||
help: "Auto-delete actions",
|
||||
labelNames: ["action"] as const,
|
||||
});
|
||||
@@ -1,65 +0,0 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
|
||||
const logger = createChildLogger("messages-embed");
|
||||
|
||||
/** Max chars for a search query fed to the embedding model. */
|
||||
const MAX_QUERY_CHARS = 300;
|
||||
|
||||
/**
|
||||
* Normalize a user search query before embedding so it lands in the same
|
||||
* vector space as the archived content (which is normalized the same way on
|
||||
* write). Mirrors the gateway's normalizer: strip control/zero-width chars,
|
||||
* lowercase, collapse whitespace, cap length. Readable punctuation is kept —
|
||||
* a search query is already compact.
|
||||
*/
|
||||
export function normalizeEmbeddingQuery(raw: string): string {
|
||||
if (!raw) return "";
|
||||
return raw
|
||||
.replace(/[\p{Cc}\p{Cf}]/gu, " ")
|
||||
.toLowerCase()
|
||||
.replace(/\s+/g, " ")
|
||||
.trim()
|
||||
.slice(0, MAX_QUERY_CHARS);
|
||||
}
|
||||
|
||||
/**
|
||||
* Embed a search query with the configured OpenAI-compatible embedding model.
|
||||
* Uses raw fetch (the backend has no openai SDK dependency) and returns null
|
||||
* when embeddings are not configured (search unavailable).
|
||||
*
|
||||
* encoding_format: "float" is REQUIRED — Nvidia-backed models reject base64.
|
||||
*/
|
||||
export async function embedQuery(rawQuery: string): Promise<number[] | null> {
|
||||
if (!config.AI_LLM_API_KEY || !config.AI_LLM_EMBEDDING_MODEL) return null;
|
||||
const text = normalizeEmbeddingQuery(rawQuery);
|
||||
if (!text) return null;
|
||||
try {
|
||||
const res = await fetch(`${config.AI_LLM_BASE_URL}/embeddings`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Authorization: `Bearer ${config.AI_LLM_API_KEY}`,
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model: config.AI_LLM_EMBEDDING_MODEL,
|
||||
input: text,
|
||||
encoding_format: "float",
|
||||
}),
|
||||
});
|
||||
if (!res.ok) {
|
||||
logger.warn({ status: res.status }, "query embed HTTP error");
|
||||
return null;
|
||||
}
|
||||
const json = (await res.json()) as {
|
||||
data?: Array<{ embedding?: number[] }>;
|
||||
};
|
||||
return json.data?.[0]?.embedding ?? null;
|
||||
} catch (error) {
|
||||
logger.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"query embed failed",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -41,11 +41,3 @@ export const messageUpdateSchema = z.object({
|
||||
export type MessageQuery = z.infer<typeof messageQuerySchema>;
|
||||
export type MessageCreate = z.infer<typeof messageCreateSchema>;
|
||||
export type MessageUpdate = z.infer<typeof messageUpdateSchema>;
|
||||
|
||||
export const semanticSearchSchema = z.object({
|
||||
query: z.string().min(1).max(500),
|
||||
limit: z.coerce.number().int().positive().max(50).default(10),
|
||||
guildId: z.string().optional(),
|
||||
});
|
||||
|
||||
export type SemanticSearchQuery = z.infer<typeof semanticSearchSchema>;
|
||||
|
||||
@@ -1,10 +1,7 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { NotFoundError, ValidationError } from "@/shared/errors/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { embedQuery } from "./embed.js";
|
||||
import { type MessageRow, messagesRepository } from "./messages.repository.js";
|
||||
import type { MessageQuery, SemanticSearchQuery } from "./messages.schema.js";
|
||||
import { searchArchive } from "./qdrant.js";
|
||||
import type { MessageQuery } from "./messages.schema.js";
|
||||
|
||||
const logger = createChildLogger("messages.service");
|
||||
|
||||
@@ -102,32 +99,6 @@ export class MessagesService {
|
||||
return messagesRepository.getReviewMessages(channelId, limit);
|
||||
}
|
||||
|
||||
/**
|
||||
* Public, read-only semantic search over the persistent message archive.
|
||||
* Embeds the query, searches Qdrant, returns text + metadata. Best-effort:
|
||||
* if embeddings/Qdrant are unavailable, returns an empty result set.
|
||||
*/
|
||||
async semanticSearch(
|
||||
input: SemanticSearchQuery,
|
||||
): Promise<{ results: ReturnType<typeof mapSearchHit>[]; nextCursor: null }> {
|
||||
const vector = await embedQuery(input.query);
|
||||
if (!vector) {
|
||||
logger.debug(
|
||||
{ query: input.query },
|
||||
"semantic search skipped: no embedder",
|
||||
);
|
||||
return { results: [], nextCursor: null };
|
||||
}
|
||||
const hits = await searchArchive(
|
||||
vector,
|
||||
input.limit,
|
||||
config.AI_LLM_EMBEDDING_ARCHIVE_MIN_SIMILARITY,
|
||||
input.guildId,
|
||||
);
|
||||
const results = hits.map((h) => mapSearchHit(h));
|
||||
return { results, nextCursor: null };
|
||||
}
|
||||
|
||||
async getActivity(
|
||||
days = 30,
|
||||
): Promise<Awaited<ReturnType<typeof messagesRepository.getActivity>>> {
|
||||
@@ -157,36 +128,4 @@ export class MessagesService {
|
||||
}
|
||||
}
|
||||
|
||||
/** Shape returned to the frontend (text + rich metadata from the archive payload). */
|
||||
function mapSearchHit(hit: {
|
||||
score: number;
|
||||
payload: {
|
||||
text: string;
|
||||
content_hash?: string;
|
||||
analyzed_at: number;
|
||||
username?: string;
|
||||
channel_id?: string;
|
||||
guild_id?: string;
|
||||
thread_id?: string | null;
|
||||
channel_name?: string | null;
|
||||
thread_name?: string | null;
|
||||
created_at?: number;
|
||||
};
|
||||
}) {
|
||||
return {
|
||||
message_id: hit.payload.content_hash ?? null,
|
||||
content: hit.payload.text,
|
||||
score: hit.score,
|
||||
// Prefer the real message timestamp; fall back to embed time for old
|
||||
// points that predate rich metadata.
|
||||
created_at: hit.payload.created_at ?? hit.payload.analyzed_at,
|
||||
username: hit.payload.username ?? null,
|
||||
channel_id: hit.payload.channel_id ?? null,
|
||||
guild_id: hit.payload.guild_id ?? null,
|
||||
thread_id: hit.payload.thread_id ?? null,
|
||||
channel_name: hit.payload.channel_name ?? null,
|
||||
thread_name: hit.payload.thread_name ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
export const messagesService = new MessagesService();
|
||||
|
||||
@@ -1,114 +0,0 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
|
||||
const logger = createChildLogger("messages-qdrant");
|
||||
|
||||
export interface ArchiveHit {
|
||||
score: number;
|
||||
payload: {
|
||||
text: string;
|
||||
content_hash?: string;
|
||||
analyzed_at: number;
|
||||
expires_at: number;
|
||||
username?: string;
|
||||
channel_id?: string;
|
||||
guild_id?: string;
|
||||
thread_id?: string | null;
|
||||
channel_name?: string | null;
|
||||
thread_name?: string | null;
|
||||
created_at?: number;
|
||||
};
|
||||
}
|
||||
|
||||
function baseUrl(): string {
|
||||
return (config.QDRANT_URL ?? "http://100.121.180.82:6333").replace(
|
||||
/\/+$/,
|
||||
"",
|
||||
);
|
||||
}
|
||||
|
||||
function headers(): Record<string, string> {
|
||||
const h: Record<string, string> = { "Content-Type": "application/json" };
|
||||
if (config.QDRANT_API_KEY) h["api-key"] = config.QDRANT_API_KEY;
|
||||
return h;
|
||||
}
|
||||
|
||||
export const ARCHIVE_COLLECTION =
|
||||
config.QDRANT_ARCHIVE_COLLECTION ?? "gmw_message_archive";
|
||||
|
||||
async function request(
|
||||
method: string,
|
||||
path: string,
|
||||
body?: unknown,
|
||||
timeoutMs = 10_000,
|
||||
): Promise<unknown> {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), timeoutMs);
|
||||
try {
|
||||
const res = await fetch(`${baseUrl()}${path}`, {
|
||||
method,
|
||||
headers: headers(),
|
||||
body: body === undefined ? undefined : JSON.stringify(body),
|
||||
signal: controller.signal,
|
||||
});
|
||||
const text = await res.text();
|
||||
if (!res.ok) {
|
||||
throw new Error(
|
||||
`Qdrant ${method} ${path} -> ${res.status}: ${text.slice(0, 200)}`,
|
||||
);
|
||||
}
|
||||
return text ? JSON.parse(text) : null;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
/** Search the archive collection for the nearest vectors to `vector`. */
|
||||
export async function searchArchive(
|
||||
vector: number[],
|
||||
limit: number,
|
||||
scoreThreshold: number,
|
||||
guildId?: string,
|
||||
): Promise<ArchiveHit[]> {
|
||||
if (!config.QDRANT_URL) return [];
|
||||
try {
|
||||
const json = (await request(
|
||||
"POST",
|
||||
`/collections/${ARCHIVE_COLLECTION}/points/search`,
|
||||
{
|
||||
vector,
|
||||
limit,
|
||||
score_threshold: scoreThreshold,
|
||||
with_payload: true,
|
||||
// Optional scope: only return vectors from a specific guild's archive.
|
||||
// Old points (embedded before rich metadata) have no guild_id payload —
|
||||
// the `must` match simply excludes them, which is the correct behavior
|
||||
// for a guild-scoped search.
|
||||
...(guildId
|
||||
? {
|
||||
filter: {
|
||||
must: [{ key: "guild_id", match: { value: guildId } }],
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
)) as {
|
||||
result?: Array<{
|
||||
score?: number;
|
||||
payload?: ArchiveHit["payload"];
|
||||
}>;
|
||||
};
|
||||
return (json.result ?? [])
|
||||
.filter((h) => h.payload?.text)
|
||||
.map((h) => ({
|
||||
score: h.score ?? 0,
|
||||
payload: h.payload as ArchiveHit["payload"],
|
||||
}));
|
||||
} catch (error) {
|
||||
logger.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"archive search failed",
|
||||
);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
@@ -5,10 +5,7 @@ import { chatRequestSchema } from "../modules/chatbot/chatbot.schema";
|
||||
import { chatbotService } from "../modules/chatbot/chatbot.service";
|
||||
import { dashboardService } from "../modules/dashboard/dashboard.service";
|
||||
import { knowledgeService } from "../modules/knowledge/knowledge.service";
|
||||
import {
|
||||
messageQuerySchema,
|
||||
semanticSearchSchema,
|
||||
} from "../modules/messages/messages.schema";
|
||||
import { messageQuerySchema } from "../modules/messages/messages.schema";
|
||||
import { messagesService } from "../modules/messages/messages.service";
|
||||
import { moderationService } from "../modules/moderation/moderation.service";
|
||||
import { uiStateService } from "../modules/ui-state/ui-state.service";
|
||||
@@ -123,10 +120,6 @@ const messagesRouter = {
|
||||
);
|
||||
return { results: rows, limit: input.limit, cursor: null };
|
||||
}),
|
||||
// Public, read-only semantic search over the message archive.
|
||||
semanticSearch: os
|
||||
.input(semanticSearchSchema)
|
||||
.handler(({ input }) => messagesService.semanticSearch(input)),
|
||||
// Public, read-only activity heatmap data (per-hour volume by channel).
|
||||
activity: os
|
||||
.input(
|
||||
|
||||
@@ -1,25 +0,0 @@
|
||||
import type { CommandReply } from "./index.js";
|
||||
import { createChildLogger } from "./logger/index.js";
|
||||
|
||||
export { createChildLogger };
|
||||
|
||||
/**
|
||||
* Attempt a Redis command first; if it fails or times out, fall back.
|
||||
*
|
||||
* @param commandFn - Function that issues the publishCommand and returns the reply.
|
||||
* @param fallbackFn - Async fallback, typically reads from Redis status key.
|
||||
* @param commandLabel - Label used for logging (e.g. "voice:connect").
|
||||
*/
|
||||
export async function tryCommandThenFallback<T>(
|
||||
commandFn: () => Promise<CommandReply<T> | null>,
|
||||
fallbackFn: () => Promise<T>,
|
||||
commandLabel: string,
|
||||
): Promise<T> {
|
||||
const logger = createChildLogger(`command-helper:${commandLabel}`);
|
||||
const reply = await commandFn();
|
||||
if (reply?.success && reply.data !== undefined && reply.data !== null) {
|
||||
return reply.data;
|
||||
}
|
||||
logger.warn("discord-gateway unreachable, falling back");
|
||||
return fallbackFn();
|
||||
}
|
||||
@@ -96,21 +96,12 @@ export const configSchema = z
|
||||
.transform((v) => v === "true")
|
||||
.default(false),
|
||||
AI_LLM_API_KEY: z.string().optional(),
|
||||
AI_LLM_BASE_URL: z
|
||||
.string()
|
||||
.url()
|
||||
.default("http://100.121.180.82:20128/api/v1"),
|
||||
// 9router — OpenAI-compatible router on this host (127.0.0.1:4014).
|
||||
// Loopback on purpose: backend runs on the same machine as 9router, so no
|
||||
// TLS/proxy hop is needed.
|
||||
AI_LLM_BASE_URL: z.string().url().default("http://127.0.0.1:4014/v1"),
|
||||
AI_LLM_MODEL: z.string().default("text"),
|
||||
AI_LLM_VISION_MODEL: z.string().optional(),
|
||||
AI_LLM_EMBEDDING_MODEL: z.string().optional(),
|
||||
// Minimum cosine similarity for the public archive semantic search. Lower
|
||||
// = more (noisier) results; raise it to tighten precision. Tuned for a 1B
|
||||
// embedding model — re-tune if the model's dimensionality changes.
|
||||
AI_LLM_EMBEDDING_ARCHIVE_MIN_SIMILARITY: z.coerce
|
||||
.number()
|
||||
.min(0)
|
||||
.max(1)
|
||||
.default(0.6),
|
||||
AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(5),
|
||||
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
||||
.number()
|
||||
@@ -179,11 +170,6 @@ export const configSchema = z
|
||||
.default("https://api.openai.com/v1"),
|
||||
OPENAI_MODERATION_MODEL: z.string().default("omni-moderation-latest"),
|
||||
|
||||
// ── Qdrant (message archive for semantic search) ──────────────────
|
||||
QDRANT_URL: z.string().optional(),
|
||||
QDRANT_API_KEY: z.string().optional(),
|
||||
QDRANT_ARCHIVE_COLLECTION: z.string().default("gmw_message_archive"),
|
||||
|
||||
// ── Auto Delete ─────────────────────────────────────────────────────
|
||||
AUTO_DELETE_FLAGGED_ENABLED: z
|
||||
.string()
|
||||
|
||||
@@ -1,8 +0,0 @@
|
||||
export {
|
||||
broadcastBinary,
|
||||
broadcastEvent,
|
||||
clearBroadcastFunctions,
|
||||
setBroadcastFunctions,
|
||||
} from "./broadcast.js";
|
||||
export { startRedisBridge, stopRedisBridge } from "./redis-bridge.js";
|
||||
export { closeWebSocketServer, createWebSocketServer } from "./server.js";
|
||||
@@ -1,4 +1,4 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { tools } from "../src/modules/chatbot/chatbot.toolDefs.js";
|
||||
|
||||
const names = tools.map((t) => t.function.name);
|
||||
|
||||
@@ -1,6 +1,32 @@
|
||||
// ─── Shared Error Classes ────────────────────────────────────────────────────
|
||||
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
// bun:test compat facade — vitest's `vi` maps onto bun's `jest`/`mock`/`spyOn`.
|
||||
// bun:test 1.3.14 exports both `jest` (fn, useFakeTimers, spyOn) and `mock`
|
||||
// (module, restore). `vi.fn` -> `jest.fn`, `vi.useFakeTimers` -> `jest.useFakeTimers`,
|
||||
// `vi.waitFor` -> waitForCompat (poll until the assertion passes).
|
||||
|
||||
import { afterEach, describe, expect, it, jest } from "bun:test";
|
||||
|
||||
const useFakeTimers = () => jest.useFakeTimers();
|
||||
const useRealTimers = () => jest.useRealTimers();
|
||||
const advanceTimersByTime = (ms: number) => jest.advanceTimersByTime(ms);
|
||||
async function waitForCompat(fn: () => Promise<unknown>, timeoutMs = 2_000) {
|
||||
const start = Date.now();
|
||||
let lastErr: unknown;
|
||||
while (Date.now() - start < timeoutMs) {
|
||||
try {
|
||||
await fn();
|
||||
return;
|
||||
} catch (err) {
|
||||
lastErr = err;
|
||||
await new Promise((r) => setTimeout(r, 10));
|
||||
}
|
||||
}
|
||||
throw lastErr instanceof Error
|
||||
? lastErr
|
||||
: new Error("waitForCompat timed out");
|
||||
}
|
||||
|
||||
import {
|
||||
AppError,
|
||||
ConfigError,
|
||||
@@ -94,41 +120,41 @@ describe("AppError subclasses", () => {
|
||||
// ═══════════════════════════════════════════════════════════════════════════════
|
||||
describe("delay", () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
useRealTimers();
|
||||
});
|
||||
|
||||
it("resolves after the given time", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const promise = delay(500);
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
await expect(promise).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("rejects are not triggered on non-matching timer", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const promise = delay(1000);
|
||||
// Advance only part way — the timer should NOT fire yet
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
// The timer is still pending; the promise has not resolved yet
|
||||
// We advance the rest
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
await expect(promise).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
describe("retryWithBackoff", () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
useRealTimers();
|
||||
});
|
||||
|
||||
it("returns the result on first success without retrying", async () => {
|
||||
const fn = vi.fn().mockResolvedValue("ok");
|
||||
const fn = jest.fn().mockResolvedValue("ok");
|
||||
await expect(retryWithBackoff(fn)).resolves.toBe("ok");
|
||||
expect(fn).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("re-throws after exhausting all retries", async () => {
|
||||
const fn = vi.fn().mockRejectedValue(new Error("persistent"));
|
||||
const fn = jest.fn().mockRejectedValue(new Error("persistent"));
|
||||
await expect(
|
||||
retryWithBackoff(fn, { retries: 1, minTimeout: 1, maxTimeout: 5 }),
|
||||
).rejects.toThrow("persistent");
|
||||
@@ -139,7 +165,7 @@ describe("retryWithBackoff", () => {
|
||||
it("throws AbortError immediately when signal is already aborted", async () => {
|
||||
const ac = new AbortController();
|
||||
ac.abort();
|
||||
const fn = vi.fn().mockResolvedValue("ok");
|
||||
const fn = jest.fn().mockResolvedValue("ok");
|
||||
await expect(
|
||||
retryWithBackoff(fn, { retries: 3, signal: ac.signal }),
|
||||
).rejects.toThrow("Aborted");
|
||||
@@ -147,9 +173,9 @@ describe("retryWithBackoff", () => {
|
||||
});
|
||||
|
||||
it("respects abort signal during retry", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const ac = new AbortController();
|
||||
const fn = vi.fn().mockRejectedValue(new Error("fail"));
|
||||
const fn = jest.fn().mockRejectedValue(new Error("fail"));
|
||||
|
||||
const promise = retryWithBackoff(fn, {
|
||||
retries: 5,
|
||||
@@ -159,8 +185,8 @@ describe("retryWithBackoff", () => {
|
||||
|
||||
// Schedule abort after first failure + backoff starts
|
||||
setTimeout(() => ac.abort(), 150);
|
||||
vi.advanceTimersByTime(200);
|
||||
await vi.waitFor(async () => {
|
||||
advanceTimersByTime(200);
|
||||
await waitForCompat(async () => {
|
||||
await expect(promise).rejects.toThrow("Aborted");
|
||||
});
|
||||
});
|
||||
@@ -232,7 +258,7 @@ describe("asyncHandler", () => {
|
||||
const wrapped = asyncHandler(async () => {
|
||||
throw error;
|
||||
});
|
||||
const next = vi.fn();
|
||||
const next = jest.fn();
|
||||
|
||||
wrapped({} as any, {} as any, next);
|
||||
|
||||
@@ -246,7 +272,7 @@ describe("asyncHandler", () => {
|
||||
const wrapped = asyncHandler(async (_req: any, _res: any, _next: any) => {
|
||||
// no-op
|
||||
});
|
||||
const next = vi.fn();
|
||||
const next = jest.fn();
|
||||
|
||||
wrapped({} as any, {} as any, next);
|
||||
await Promise.resolve();
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
// bun test preload — nothing needed for backend unit tests today.
|
||||
@@ -1,4 +1,4 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { describe, expect, it } from "bun:test";
|
||||
|
||||
/**
|
||||
* Lock the contract that the WS `stream_messages` handler + frontend
|
||||
|
||||
@@ -1,16 +0,0 @@
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { defineConfig } from "vitest/config";
|
||||
|
||||
export default defineConfig({
|
||||
resolve: {
|
||||
alias: {
|
||||
"@": fileURLToPath(new URL("./src", import.meta.url)),
|
||||
},
|
||||
},
|
||||
test: {
|
||||
globals: true,
|
||||
environment: "node",
|
||||
include: ["src/**/*.test.ts", "tests/**/*.test.ts"],
|
||||
testTimeout: 15000,
|
||||
},
|
||||
});
|
||||
@@ -53,24 +53,27 @@ src/
|
||||
1. **LLM is the only judge.** Failed LLM → `status:"error"` + recovery retry.
|
||||
**Never** reintroduce regex/heuristic content classification.
|
||||
2. **Discord tokens sanitized** before reaching LLM (`discordTokens.ts`).
|
||||
3. **Semantic cache is batched** — one embed call + one Qdrant batch search.
|
||||
4. **Streaming is mandatory** against the omniroute base URL.
|
||||
3. **Streaming is mandatory** against the router base URL.
|
||||
|
||||
## AI moderation pipeline
|
||||
|
||||
```
|
||||
aiAnalyzer.ts → batchScheduler.ts → batchProcessor.ts → individualFallbackProcessor.ts
|
||||
↓ ↓ ↓ ↓
|
||||
moderationOrchestrator.ts → (hash cache → Qdrant → LLM)
|
||||
moderationOrchestrator.ts → (hash cache → LLM)
|
||||
↓ ↓ ↓
|
||||
textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||
embeddingClient.ts
|
||||
qdrantClient.ts
|
||||
```
|
||||
|
||||
- Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`, `startPendingAIAnalysisWorker`)
|
||||
- Concurrency: LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5)
|
||||
- Concurrency: **two per-lane LLM semaphores** (2026-09-24) — text
|
||||
(`AI_LLM_MAX_CONCURRENT`, default 8) and media/vision
|
||||
(`AI_LLM_MEDIA_MAX_CONCURRENT`, default 4); a media backlog can never
|
||||
consume text slots
|
||||
- Piscina: text pool (4 threads) + media pool (2 threads)
|
||||
- Locks are **per conversation per lane** (`conversationProcessing` maps key →
|
||||
lane → startedAt): the text lane of a conversation never waits on that
|
||||
conversation's media lane (this was the "image blocks the queue" bug)
|
||||
- **Each worker thread has its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`)
|
||||
|
||||
## Module: message-capture
|
||||
@@ -79,7 +82,6 @@ textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||
- `messageStore.ts` — DB operations
|
||||
- `messageMetadata.ts` — metadata extraction
|
||||
- `messagesDb.ts` / `messagesCrud.ts` — DB schema operations
|
||||
- `archiveEmbedder.ts` — Qdrant embedding (respect age-restricted guard)
|
||||
- `retentionDb.ts` / `reviewsDb.ts` / `attachmentsDb.ts` — auxiliary tables
|
||||
|
||||
## Redis channels (outbound to backend)
|
||||
|
||||
@@ -5,10 +5,9 @@ messages/attachments/reactions/threads/presence, runs LLM-based AI
|
||||
moderation, and publishes everything to Redis pub/sub for the backend to
|
||||
consume. The backend serves the HTTP/WS API to the frontend.
|
||||
|
||||
> NOTE: this doc is the source of truth for the module layout. The older
|
||||
> `MODULE_STRUCTURE.md` was stale (referenced `winston`, `mock-crc.ts`,
|
||||
> `indonesianTextNormalizer.ts`, and `aiAnalysisWorker.ts`/`llmModerationClient.ts`
|
||||
> which were renamed/merged). If they disagree, this file wins.
|
||||
> NOTE: this doc is the source of truth for the module layout. The old
|
||||
> `MODULE_STRUCTURE.md` was a stale duplicate and has been removed. `README.md`
|
||||
> only covers how to run the service.
|
||||
|
||||
## Top-level layout
|
||||
|
||||
@@ -16,72 +15,99 @@ consume. The backend serves the HTTP/WS API to the frontend.
|
||||
services/discord-gateway/
|
||||
├── src/
|
||||
│ ├── index.ts # Entry point → initializeDiscordGateway()
|
||||
│ ├── app/
|
||||
│ │ ├── bootstrap.ts # Wires client, DB, Redis, workers, schedulers
|
||||
│ │ ├── shutdown.ts # Graceful shutdown (SIGINT/SIGTERM + transient errors)
|
||||
│ ├── app/ # Process lifecycle
|
||||
│ │ ├── bootstrap.ts # Startup order: config → DB → services → metrics → login
|
||||
│ │ ├── lifecycle.ts # Everything wired on the Discord 'ready' hook
|
||||
│ │ ├── process-guards.ts # SIGINT/SIGTERM + uncaught-error policy
|
||||
│ │ ├── metrics-collector.ts # AI pipeline Prometheus gauges
|
||||
│ │ ├── shutdown.ts # Graceful shutdown sequence
|
||||
│ │ └── retention.ts # Expired-record cleanup scheduler
|
||||
│ ├── shared/
|
||||
│ ├── shared/ # Infrastructure — never imports from modules/
|
||||
│ │ ├── config/ # Zod-validated env (index.ts = schema+loader)
|
||||
│ │ ├── database/ # Drizzle ORM + pg Pool + migrations
|
||||
│ │ │ ├── init.ts drizzle.ts pool.ts migrate.ts migrateCli.ts
|
||||
│ │ │ └── schema/ # messages, cache, meta, analytics
|
||||
│ │ ├── logger/ # pino wrapper + createChildLogger()
|
||||
│ │ ├── errors/ # AppError / ConfigError ...
|
||||
│ │ ├── errors/ # AppError / ConfigError ... + errorMessage()
|
||||
│ │ │ # + isTransientStreamError()
|
||||
│ │ ├── utils/ # retry, pagination
|
||||
│ │ ├── discord/clientOptions.ts # discord.js-selfbot-v13 client options
|
||||
│ │ ├── uploader.ts # Shared attachment upload helper
|
||||
│ │ ├── redis-channels.ts # Redis channel-name constants
|
||||
│ │ ├── redis-channels.ts # Redis channel + command constants
|
||||
│ │ └── moderation-types.ts # Shared AI analysis domain types
|
||||
│ └── modules/
|
||||
│ └── modules/ # Feature modules, each with an index.ts facade
|
||||
│ ├── message-capture/ # Discord event listeners + DB store
|
||||
│ ├── ai-moderation/ # LLM moderation pipeline (see below)
|
||||
│ ├── attachment-upload/ # Download + (sharp) resize + upload
|
||||
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
||||
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
||||
│ ├── command-handler/ # Redis-subscribed backend→gateway commands
|
||||
│ ├── reaction-tracking/ thread-tracking/ user-presence/
|
||||
│ ├── channel-topic/ guild-member-events/
|
||||
│ └── gateway-metrics/ # Prometheus /metrics endpoint (port 4016)
|
||||
│ ├── channel-topic/ guild-member-events/ monitor/
|
||||
│ └── gateway-metrics/ # Prometheus /metrics endpoint (METRICS_PORT)
|
||||
```
|
||||
|
||||
Dependency direction is one-way: `index.ts` → `app/` → `modules/` → `shared/`.
|
||||
Code outside a module imports its `index.ts` facade, never an internal file;
|
||||
deep imports stay valid inside the module itself.
|
||||
|
||||
## AI moderation pipeline (`ai-moderation/`)
|
||||
|
||||
LLM-only judge — no regex/heuristic classification. One orchestrator call
|
||||
handles a whole batch (text + media split internally, parallel paths).
|
||||
handles a whole batch. **Independent text/media lanes** (2026-09-24): a
|
||||
conversation batch is split into a text lane (messages with no media) and a
|
||||
media lane (attachments/stickers/embeds) that are dispatched to separate
|
||||
pools, hold SEPARATE per-lane processing locks, and run under SEPARATE LLM
|
||||
concurrency semaphores. The text lane frees its lock and saves+broadcasts the
|
||||
moment text analysis finishes — it never waits on a slow vision/media batch
|
||||
of the same conversation, and vice versa.
|
||||
|
||||
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `getAnalysisQueueStatus`,
|
||||
`startPendingAIAnalysisWorker` (recovery worker + cache-prune).
|
||||
- `batchScheduler.ts` — per-conversation debounce → `processBatch`.
|
||||
- `batchProcessor.ts` — batch lock/circuit-breaker, fans failed targets to
|
||||
individual fallback.
|
||||
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `queueConversationAnalysis`,
|
||||
`getAnalysisQueueStatus`, `startPendingAIAnalysisWorker`. Short-circuits
|
||||
age-restricted and skip-list messages before any LLM work.
|
||||
- `recovery-worker.ts` — periodic sweep for stranded `pending` messages
|
||||
(re-scheduled per lane) and `error`/`analysis_incomplete` messages
|
||||
(individual fallback queue); prunes stale lane locks, per-conversation CB
|
||||
counters and individual in-flight markers.
|
||||
- `cache-prune.ts` — throttled (6h) expired-verdict sweep across Postgres,
|
||||
driven from the recovery interval.
|
||||
- `batchScheduler.ts` — per-conversation per-LANE debounce → `processBatch`
|
||||
(lane-aware). `splitMessagesByLane` / `laneOfMessage` live in
|
||||
`analysisLanes.ts` (pure, unit-testable).
|
||||
- `batchProcessor.ts` — per-lane batch lock/circuit-breaker, fans failed
|
||||
targets to individual fallback. `processBatch` releases ITS lane's lock the
|
||||
moment that lane's worker job finishes; the other lane owns its own lock.
|
||||
- `individualFallbackProcessor.ts` — one-message-at-a-time retry path, own CB.
|
||||
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation state,
|
||||
Piscina `workerPool`, `getConversationKey`.
|
||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` / `individual` jobs).
|
||||
Runs `runModerationAnalysis` off the main thread.
|
||||
- `moderationOrchestrator.ts` — exact-hash cache → batched semantic (Qdrant)
|
||||
cache → LLM. Text and media paths run in parallel.
|
||||
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation PER-LANE
|
||||
state (`conversationProcessing` holds a lane → startedAt map per key),
|
||||
Piscina `textWorkerPool`/`mediaWorkerPool`, `getConversationKey`.
|
||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` (lane) /
|
||||
`individual` jobs). Runs `runModerationAnalysis` off the main thread.
|
||||
- `moderationOrchestrator.ts` — exact-hash cache → LLM. Text and media paths
|
||||
run in parallel.
|
||||
- `textBatchProcessor.ts` / `mediaBatchProcessor.ts` — actual LLM calls
|
||||
(one call per sub-batch, not per message).
|
||||
(one call per sub-batch, not per message). `mediaBatchProcessor` routes its
|
||||
moderation LLM call through the MEDIA semaphore.
|
||||
- `llmClient.ts` — central OpenAI-compatible chat client (streaming, retries,
|
||||
thinking-disable injection). `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
||||
thinking-disable injection). TWO concurrency semaphores:
|
||||
`AI_LLM_MAX_CONCURRENT` (text lane, default 8) and
|
||||
`AI_LLM_MEDIA_MAX_CONCURRENT` (media lane, default 4) — a vision backlog
|
||||
can never consume text slots. `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
||||
share the same router/base URL (different model alias for vision).
|
||||
- `embeddingClient.ts` + `qdrantClient.ts` — semantic cache (one embed call +
|
||||
one batched Qdrant search for all uncached targets).
|
||||
- `textCacheStore.ts` / `channelCultureStore.ts` / `userProfileStore.ts` /
|
||||
`userProfileStore.ts` — caches learned user profile summaries (optional).
|
||||
|
||||
### Concurrency model
|
||||
|
||||
- Main thread owns the LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5) via
|
||||
`llmClient.withLlmConcurrency`.
|
||||
- Main thread owns TWO per-lane LLM semaphores (2026-09-24):
|
||||
`AI_LLM_MAX_CONCURRENT` (text, default 8) and `AI_LLM_MEDIA_MAX_CONCURRENT`
|
||||
(media, default 4) via `llmClient.withLlmConcurrency(fn, { lane })`.
|
||||
- Two Piscina pools run the heavy LLM work off the event loop: a text pool
|
||||
(`PISCINA_MAX_THREADS`, default 4) and a dedicated media pool
|
||||
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed to the media
|
||||
pool if ANY of its messages carries an attachment/sticker/embed — this
|
||||
keeps a slow image/vision batch from occupying every thread and blocking
|
||||
unrelated text-only batches behind it. **Each worker thread (in either
|
||||
pool) initializes its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`).
|
||||
See "Memory & connections" below.
|
||||
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed by lane to the
|
||||
matching pool — this keeps a slow image/vision batch from occupying every
|
||||
thread and blocking unrelated text-only batches behind it. **Each worker
|
||||
thread (in either pool) initializes its own pg Pool** (min 0, grows to
|
||||
`POSTGRES_POOL_MAX`). See "Memory & connections" below.
|
||||
|
||||
## Memory & DB connections
|
||||
|
||||
@@ -109,28 +135,53 @@ See `src/shared/redis-channels.ts` for the canonical names.
|
||||
|
||||
## Initialization flow
|
||||
|
||||
`bootstrap.ts` runs these steps in order (each is a named function):
|
||||
|
||||
1. Validate env (Zod). Refuse to start if `AI_ANALYSIS_ENABLED` but no key.
|
||||
2. `AUTO_MIGRATE_ON_STARTUP` → run pending Drizzle migrations.
|
||||
3. `initializeDatabase()` (pg Pool, min 0).
|
||||
4. Create discord.js-selfbot-v13 client; register listeners on `ready`.
|
||||
5. Start `gmw-discord-gateway` metrics server (port `METRICS_PORT`, default 4016).
|
||||
6. `client.login(token)`.
|
||||
→ `assertConfigIsUsable()`
|
||||
2. Build long-lived services: Discord client, `RedisEventPublisher` +
|
||||
`EventBroadcaster`, `CommandHandler`; install the shutdown handler.
|
||||
3. Connect infrastructure → `connectDatabase()`:
|
||||
`AUTO_MIGRATE_ON_STARTUP` runs pending Drizzle migrations, then
|
||||
`initializeDatabase()` (pg Pool, min 0).
|
||||
4. `registerClientDebugLogging()` — only client debug lines carrying signal.
|
||||
5. Install process guards (`registerProcessGuards`).
|
||||
6. Register pipeline gauges + start the metrics server (`METRICS_PORT`, code
|
||||
default 9090, set per deployment).
|
||||
7. `client.login(token)`.
|
||||
|
||||
On the Discord `ready` event, `lifecycle.ts` runs `startGatewayLifecycle()`:
|
||||
|
||||
1. Inject the event broadcaster into message-capture and moderation-actions
|
||||
(before any listener can fire).
|
||||
2. Register Discord listeners: message-capture, reaction, thread, presence,
|
||||
channel-topic, guild-member.
|
||||
3. Start background work: AI analysis worker + recovery worker, command
|
||||
handler, retention cleanup, weekly digest.
|
||||
|
||||
## Graceful shutdown
|
||||
|
||||
`SIGINT`/`SIGTERM` (and uncaught transient stream errors: EPIPE / ECONNRESET /
|
||||
ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END are treated as non-fatal):
|
||||
stop metrics → close event broadcaster (Redis) → close command handler →
|
||||
close DB → destroy client → exit.
|
||||
`process-guards.ts` owns the policy. `SIGINT`/`SIGTERM` and non-transient
|
||||
uncaught exceptions/rejections run `shutdown.ts`; transient stream errors
|
||||
(EPIPE / ECONNRESET / ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END, see
|
||||
`isTransientStreamError()`) are logged and IGNORED so the bot stays online.
|
||||
|
||||
Shutdown order: stop metrics → close event broadcaster (Redis) → close command
|
||||
handler → close DB → destroy client → exit.
|
||||
|
||||
## Observability
|
||||
|
||||
Prometheus scrapes `127.0.0.1:4016/metrics` (`bete_*` prefix). Collectors run
|
||||
Prometheus scrapes the metrics server at `127.0.0.1:$METRICS_PORT/metrics`
|
||||
(`bete_*` prefix; the code default is 9090 — deployments set it explicitly,
|
||||
this host uses 4018). Collectors run
|
||||
per-scrape and expose: process memory/uptime, and (when AI analysis is on) live
|
||||
pipeline gauges — `ai_analysis_queued_conversations`,
|
||||
`ai_analysis_active_batch_requests`, `ai_analysis_active_individual_requests`,
|
||||
`ai_analysis_individual_in_flight`, `ai_analysis_individual_circuit_breaker_active`,
|
||||
`ai_analysis_worker_threads`, `ai_analysis_worker_threads_active`.
|
||||
pipeline gauges registered by `app/metrics-collector.ts` —
|
||||
`ai_analysis_queued_conversations`, `ai_analysis_active_batch_requests`,
|
||||
`ai_analysis_active_text_requests`, `ai_analysis_active_media_requests`,
|
||||
`ai_analysis_active_individual_requests`, `ai_analysis_individual_in_flight`,
|
||||
`ai_analysis_individual_circuit_breaker_active`,
|
||||
`ai_analysis_worker_threads_{text,media}`,
|
||||
`ai_analysis_worker_threads_active_{text,media}`.
|
||||
|
||||
## Key invariants (do not break)
|
||||
|
||||
@@ -139,7 +190,5 @@ pipeline gauges — `ai_analysis_queued_conversations`,
|
||||
- **Discord tokens are sanitized** (`discordTokens.ts`: `<:emoji:id>` →
|
||||
`[emoji:name]`, `<@id>` → `@user`, etc.) before content reaches the LLM, so
|
||||
numeric snowflake IDs never trigger false positives.
|
||||
- **Semantic cache is batched** (one embed call + one Qdrant batch search),
|
||||
not N sequential round-trips. `ensureQdrantCollection` is memoized.
|
||||
- **Streaming is mandatory** against the omniroute base URL (non-stream waits for
|
||||
- **Streaming is mandatory** against the router base URL (non-stream waits for
|
||||
the full body and times out). `llmClient` aggregates SSE chunks.
|
||||
|
||||
@@ -1,73 +0,0 @@
|
||||
# Discord Gateway Service — Module Structure
|
||||
|
||||
> Kept as a compact module map. For the authoritative layout, design
|
||||
> decisions, and invariants, see `ARCHITECTURE.md`. This file was rewritten
|
||||
> on 2026-08-16 to fix stale references (`winston` → pino,
|
||||
> `mock-crc.ts`/`indonesianTextNormalizer.ts` removed,
|
||||
> `aiAnalysisWorker.ts` → `ai-analysis-worker.ts`,
|
||||
> `llmModerationClient.ts` → `llmClient.ts`).
|
||||
|
||||
## Top-level
|
||||
|
||||
```
|
||||
services/discord-gateway/
|
||||
├── src/
|
||||
│ ├── index.ts # Entry point
|
||||
│ ├── app/ # bootstrap, shutdown, retention
|
||||
│ ├── shared/ # config, database, logger, errors, utils, discord, uploader
|
||||
│ └── modules/
|
||||
│ ├── message-capture/ # Discord listeners + DB store + metadata
|
||||
│ ├── ai-moderation/ # LLM moderation pipeline (largest module)
|
||||
│ ├── attachment-upload/ # Download + sharp resize + upload
|
||||
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
||||
│ ├── command-handler/ # Backend→gateway Redis commands
|
||||
│ ├── reaction-tracking/ thread-tracking/ user-presence/
|
||||
│ ├── channel-topic/ guild-member-events/
|
||||
│ └── gateway-metrics/ # Prometheus /metrics (port 4016)
|
||||
├── tests/ # Vitest suites (129 tests)
|
||||
├── drizzle/ # Drizzle migration SQL + journal
|
||||
├── ARCHITECTURE.md README.md package.json tsconfig.json vitest.config.ts
|
||||
```
|
||||
|
||||
## Module responsibilities (summary)
|
||||
|
||||
### message-capture
|
||||
Captures `messageCreate`/`messageUpdate`/`messageDelete`, extracts metadata,
|
||||
stores to Postgres, publishes to Redis. Controller–Service–Repository split:
|
||||
`messageCapture.ts` (listener) → `messageStore.ts` (DB) + `messageMetadata.ts`
|
||||
(service).
|
||||
|
||||
### ai-moderation
|
||||
LLM-only moderation. Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`,
|
||||
`startPendingAIAnalysisWorker`, `getAnalysisQueueStatus`). Scheduling:
|
||||
`batchScheduler.ts` → `batchProcessor.ts` (batch lock + circuit breaker) →
|
||||
`individualFallbackProcessor.ts` (per-message retry). Heavy work runs in the
|
||||
Piscina pool via `ai-analysis-worker.ts` (jobs `batch` / `individual`).
|
||||
Orchestration/caching: `moderationOrchestrator.ts` (exact hash → batched
|
||||
semantic Qdrant → LLM), `textBatchProcessor.ts` / `mediaBatchProcessor.ts`
|
||||
(one LLM call per sub-batch), `llmClient.ts` (central streaming client),
|
||||
`embeddingClient.ts` + `qdrantClient.ts` (semantic cache), plus
|
||||
`channelCultureStore.ts` / `userProfileStore.ts`.
|
||||
|
||||
### attachment-upload
|
||||
`attachmentUploader.ts` (download → upload to storage) + `imageResizer.ts`
|
||||
(sharp resize). Emits `discord:attachment:*`.
|
||||
|
||||
### event-broadcaster
|
||||
`RedisEventPublisher` (ioredis publish) + `EventBroadcaster` (typed methods).
|
||||
Channel names in `src/shared/redis-channels.ts`.
|
||||
|
||||
### gateway-metrics
|
||||
`metrics.ts` Prometheus HTTP server on `METRICS_PORT` (4016). Collectors run
|
||||
per scrape; live pipeline gauges registered in `bootstrap.ts`.
|
||||
|
||||
## Shared infrastructure
|
||||
- **config** — Zod schema in `shared/config/index.ts` (single source of truth).
|
||||
- **database** — Drizzle ORM over `pg`; pool `min:0` (`shared/config`).
|
||||
- **logger** — `pino` wrapper, `createChildLogger()` for context loggers.
|
||||
- **errors** — `AppError` hierarchy (`ConfigError`, …).
|
||||
|
||||
## Notes
|
||||
- No HTTP server (other than the metrics endpoint). Pure event-driven.
|
||||
- `MODULE_STRUCTURE.md` is intentionally a sketch; `ARCHITECTURE.md` is the
|
||||
detailed reference. When they diverge, `ARCHITECTURE.md` wins.
|
||||
@@ -1,319 +1,63 @@
|
||||
# Discord Gateway Service - Extraction Complete
|
||||
# Discord Gateway
|
||||
|
||||
## Overview
|
||||
Event-driven selfbot service: captures Discord events, runs LLM moderation,
|
||||
publishes everything to Redis for the backend to consume.
|
||||
|
||||
Successfully extracted Discord Gateway service with **Modular MVC + Event-Driven Architecture** using Redis pub/sub for inter-service communication.
|
||||
> Architecture, invariants and the AI pipeline are documented in
|
||||
> **`ARCHITECTURE.md`** — that file is the source of truth. This README only
|
||||
> covers how to run it.
|
||||
|
||||
## Directory Structure
|
||||
## Commands
|
||||
|
||||
```
|
||||
services/discord-gateway/
|
||||
├── src/
|
||||
│ ├── app/
|
||||
│ │ ├── bootstrap.ts # Service initialization (Discord client, DB, Redis)
|
||||
│ │ └── shutdown.ts # Graceful shutdown handler
|
||||
│ ├── shared/ # Shared infrastructure layer
|
||||
│ │ ├── config/
|
||||
│ │ │ └── config.ts # Zod-validated environment config
|
||||
│ │ ├── database/
|
||||
│ │ │ ├── schema.ts # Drizzle ORM schema
|
||||
│ │ │ ├── drizzle.ts # PostgreSQL connection
|
||||
│ │ │ ├── migrate.ts # Migration runner
|
||||
│ │ ├── errors/
|
||||
│ │ │ └── errors.ts # Custom error classes
|
||||
│ │ ├── logger/
|
||||
│ │ │ ├── logger.ts # Winston logger wrapper
|
||||
│ │ │ └── serialization.ts # Log serialization
|
||||
│ │ ├── utils/
|
||||
│ │ │ └── retry.ts # Retry with exponential backoff
|
||||
│ │ └── discord/
|
||||
│ │ └── clientOptions.ts # Discord.js client config
|
||||
│ ├── modules/ # Feature modules (Modular MVC)
|
||||
│ │ ├── message-capture/ # Controller-Service-Repository
|
||||
│ │ │ ├── messageCapture.ts # Controller: Discord event listeners
|
||||
│ │ │ ├── messageStore.ts # Repository: DB operations
|
||||
│ │ │ ├── messageMetadata.ts # Service: Metadata extraction
|
||||
│ │ │ ├── types.ts # Domain types
|
||||
│ │ │ └── index.ts # Module exports
|
||||
│ │ ├── ai-moderation/ # Controller-Service-Repository
|
||||
│ │ │ ├── aiAnalyzer.ts # Controller: Analysis orchestration
|
||||
│ │ │ ├── llmModerationClient.ts # Service: LLM API client
|
||||
│ │ │ ├── aiAnalysisWorker.ts # Service: Worker pool
|
||||
│ │ │ ├── indonesianTextNormalizer.ts # Service: Text normalization
|
||||
│ │ │ ├── moderationPrompt.ts # Service: Prompt generation
|
||||
│ │ │ └── index.ts # Module exports
|
||||
│ │ ├── attachment-upload/ # Controller-Service-Repository
|
||||
│ │ │ ├── attachmentUploader.ts # Service: Upload orchestration
|
||||
│ │ │ ├── imageResizer.ts # Service: Image resizing
|
||||
│ │ │ └── index.ts # Module exports
|
||||
│ │ └── event-broadcaster/ # Event-driven layer
|
||||
│ │ ├── eventBroadcaster.ts # Service: Redis pub/sub publisher
|
||||
│ │ ├── eventTypes.ts # Domain: Event type definitions
|
||||
│ │ └── index.ts # Module exports
|
||||
│ ├── mock-crc.ts # CRC polyfill for discord.js
|
||||
│ └── index.ts # Service entry point
|
||||
├── ARCHITECTURE.md # Detailed architecture documentation
|
||||
├── package.json # Service dependencies
|
||||
└── tsconfig.json # TypeScript configuration (inherited)
|
||||
```bash
|
||||
pnpm install
|
||||
pnpm typecheck # tsc --noEmit
|
||||
pnpm lint # biome check --diagnostic-level=error .
|
||||
pnpm test # vitest run (138 tests)
|
||||
pnpm build # tsc — CI/prod builds run this inside nix, which also
|
||||
# runs scripts/fix-imports.mjs to rewrite @/ aliases and
|
||||
# extensionless imports for Node ESM
|
||||
pnpm dev # tsx watch src/index.ts
|
||||
pnpm start # node dist/index.js
|
||||
```
|
||||
|
||||
## Architecture Patterns
|
||||
Deployment is CI-only: `nix build .#discord-gateway` → Attic cache → systemd
|
||||
restart on the VPS. Do not build/hand-copy the artifact.
|
||||
|
||||
### 1. Modular MVC Structure
|
||||
Each feature module follows **Controller-Service-Repository** pattern:
|
||||
|
||||
**Message Capture Module**:
|
||||
- **Controller** (`messageCapture.ts`): Listens to Discord events (messageCreate, messageUpdate, messageDelete)
|
||||
- **Service** (`messageMetadata.ts`): Extracts and normalizes message metadata
|
||||
- **Repository** (`messageStore.ts`): Database CRUD operations
|
||||
|
||||
**AI Moderation Module**:
|
||||
- **Controller** (`aiAnalyzer.ts`): Orchestrates analysis workflow
|
||||
- **Service** (`llmModerationClient.ts`): LLM API integration
|
||||
- **Service** (`aiAnalysisWorker.ts`): Worker pool management
|
||||
- **Service** (`indonesianTextNormalizer.ts`): Text preprocessing
|
||||
|
||||
**Attachment Upload Module**:
|
||||
- **Service** (`attachmentUploader.ts`): Upload orchestration
|
||||
- **Service** (`imageResizer.ts`): Image processing
|
||||
|
||||
### 2. Event-Driven Architecture
|
||||
**Redis Pub/Sub** replaces WebSocket broadcaster:
|
||||
## Layout
|
||||
|
||||
```
|
||||
Discord Events → Discord Gateway Service → Redis Pub/Sub → Backend Service
|
||||
↓
|
||||
Event Channels:
|
||||
- discord:message:created
|
||||
- discord:message:updated
|
||||
- discord:message:deleted
|
||||
- discord:message:analyzed
|
||||
- discord:attachment:created
|
||||
- discord:attachment:uploaded
|
||||
- discord:analysis:queue_status
|
||||
src/
|
||||
├── index.ts # entry → initializeDiscordGateway()
|
||||
├── app/ # process lifecycle
|
||||
│ ├── bootstrap.ts # startup order: config → DB → services → metrics → login
|
||||
│ ├── lifecycle.ts # everything wired on the Discord 'ready' hook
|
||||
│ ├── process-guards.ts # SIGINT/SIGTERM + uncaught error policy
|
||||
│ ├── metrics-collector.ts # AI pipeline Prometheus gauges
|
||||
│ ├── shutdown.ts # graceful shutdown sequence
|
||||
│ └── retention.ts # expired-record cleanup scheduler
|
||||
├── shared/ # infrastructure — never imports from modules/
|
||||
│ ├── config/ database/ logger/ errors/ utils/
|
||||
│ ├── discord/clientOptions.ts
|
||||
│ ├── redis-channels.ts # canonical Redis channel + command constants
|
||||
│ └── moderation-types.ts # domain types shared across services
|
||||
└── modules/ # feature modules (each exposes an index.ts facade)
|
||||
├── ai-moderation/ # LLM moderation pipeline (largest module)
|
||||
├── message-capture/ # Discord listeners + message/attachment DB
|
||||
├── attachment-upload/ # download → resize → upload
|
||||
├── event-broadcaster/ # Redis pub/sub publisher
|
||||
├── command-handler/ # backend → gateway commands over Redis
|
||||
├── gateway-metrics/ # Prometheus /metrics (METRICS_PORT)
|
||||
├── monitor/ # weekly digest scheduler
|
||||
└── reaction-tracking/ thread-tracking/ user-presence/
|
||||
channel-topic/ guild-member-events/
|
||||
```
|
||||
|
||||
### 3. Shared Infrastructure Layer
|
||||
Centralized, reusable components:
|
||||
- **Config**: Zod-validated environment variables
|
||||
- **Logger**: Winston logger with context support
|
||||
- **Database**: Drizzle ORM with PostgreSQL
|
||||
- **Errors**: Custom error classes with codes and HTTP status codes
|
||||
- **Utils**: Retry logic with exponential backoff
|
||||
- **Discord**: Client configuration and options
|
||||
Dependency direction is one-way: `index.ts` → `app/` → `modules/` → `shared/`.
|
||||
Callers outside a module import its `index.ts` facade, never an internal file.
|
||||
|
||||
### 4. No HTTP Server
|
||||
- **Event-driven only**: No Express, WebSocket, or HTTP routes
|
||||
- **Redis pub/sub**: All inter-service communication via Redis
|
||||
- **Backend service**: Consumes events and serves HTTP API
|
||||
- **Frontend**: Continues to use Backend HTTP API
|
||||
## Testing
|
||||
|
||||
## Key Features
|
||||
|
||||
### Message Capture
|
||||
1. Discord emits `messageCreate`, `messageUpdate`, `messageDelete` events
|
||||
2. `messageCapture.ts` listener receives and validates event
|
||||
3. Extract metadata: user, channel, content, timestamp, attachments
|
||||
4. `messageStore.ts` inserts into PostgreSQL
|
||||
5. `eventBroadcaster.messageCreated()` publishes to Redis
|
||||
6. Backend service subscribes and processes
|
||||
|
||||
### AI Moderation
|
||||
1. `aiAnalyzer.ts` queues messages for analysis
|
||||
2. `llmModerationClient.ts` calls LLM API with context
|
||||
3. `indonesianTextNormalizer.ts` preprocesses text
|
||||
4. Results stored in database
|
||||
5. `eventBroadcaster.messageAnalyzed()` publishes results
|
||||
6. Backend service receives and updates UI
|
||||
|
||||
### Attachment Upload
|
||||
1. `messageCapture.ts` detects attachments
|
||||
2. `attachmentUploader.ts` downloads from Discord
|
||||
3. `imageResizer.ts` resizes images if needed
|
||||
4. Upload to external storage with retry logic
|
||||
5. `eventBroadcaster.attachmentUploaded()` publishes
|
||||
6. Backend service stores metadata
|
||||
|
||||
## Initialization Flow
|
||||
|
||||
```
|
||||
1. Load environment config (Zod validation)
|
||||
↓
|
||||
2. Initialize PostgreSQL connection
|
||||
↓
|
||||
3. Run pending database migrations
|
||||
↓
|
||||
4. Create Discord client with optimized cache
|
||||
↓
|
||||
5. Initialize Redis event broadcaster
|
||||
↓
|
||||
6. Register Discord event listeners
|
||||
- messageCapture (message events)
|
||||
- aiAnalyzer (analysis worker)
|
||||
↓
|
||||
7. Login to Discord
|
||||
↓
|
||||
8. Listen for graceful shutdown signals
|
||||
```
|
||||
|
||||
## Graceful Shutdown
|
||||
|
||||
On SIGINT/SIGTERM/uncaughtException/unhandledRejection:
|
||||
1. Close PostgreSQL connection
|
||||
2. Close Redis connection
|
||||
3. Destroy Discord client
|
||||
4. Exit process (code 0 for clean, 1 for error)
|
||||
|
||||
## Dependencies
|
||||
|
||||
**Core Discord**:
|
||||
- `discord.js-selfbot-v13` — Discord client (selfbot variant)
|
||||
|
||||
**Media Processing**:
|
||||
- `sharp` — Image resizing
|
||||
|
||||
**Data & Config**:
|
||||
- `drizzle-orm` — Type-safe ORM
|
||||
- `pg` — PostgreSQL driver
|
||||
- `zod` — Config validation
|
||||
- `ioredis` — Redis client
|
||||
|
||||
**Logging & Utilities**:
|
||||
- `winston` — Structured logging
|
||||
- `p-retry` — Retry with backoff
|
||||
- `p-limit` — Concurrency limiting
|
||||
- `piscina` — Worker pool
|
||||
|
||||
## No Breaking Changes
|
||||
|
||||
- Original `src/` remains untouched
|
||||
- Discord Gateway is a **new service** in `services/discord-gateway/`
|
||||
- Can run alongside existing monolith during transition
|
||||
- Backend service will consume Redis events
|
||||
- Frontend continues to use Backend HTTP API
|
||||
|
||||
## Next Steps
|
||||
|
||||
1. **Create Backend service** (`services/backend/`)
|
||||
- HTTP API endpoints
|
||||
- Redis event subscribers
|
||||
- Database models
|
||||
- WebSocket broadcaster
|
||||
|
||||
2. **Update Frontend** (`frontend/`)
|
||||
- Connect to Backend HTTP API
|
||||
- Subscribe to WebSocket events
|
||||
|
||||
3. **Nix & CI/CD**
|
||||
- flake.nix package for Discord Gateway
|
||||
- systemd services (gmw-backend, gmw-discord-gateway)
|
||||
- GitHub Actions for build/deploy (nix copy → systemctl restart)
|
||||
|
||||
4. **Documentation**
|
||||
- API documentation
|
||||
- Event schema documentation
|
||||
- Deployment guide
|
||||
|
||||
## Files Created
|
||||
|
||||
**Total: 43 files**
|
||||
|
||||
### Shared Infrastructure (9 files)
|
||||
- `src/shared/config/config.ts`
|
||||
- `src/shared/database/` (5 files)
|
||||
- `@bete/shared/errors` (shared package)
|
||||
- `src/shared/logger/logger.ts`
|
||||
- `src/shared/logger/serialization.ts`
|
||||
- `src/shared/utils/retry.ts`
|
||||
- `src/shared/discord/clientOptions.ts`
|
||||
|
||||
### Modules (28 files)
|
||||
- `src/modules/message-capture/` (5 files)
|
||||
- `src/modules/ai-moderation/` (6 files)
|
||||
- `src/modules/attachment-upload/` (3 files)
|
||||
- `src/modules/event-broadcaster/` (3 files)
|
||||
|
||||
### App & Entry (4 files)
|
||||
- `src/app/bootstrap.ts`
|
||||
- `src/app/shutdown.ts`
|
||||
- `src/index.ts`
|
||||
- `src/mock-crc.ts`
|
||||
|
||||
### Configuration (2 files)
|
||||
- `package.json`
|
||||
- `ARCHITECTURE.md`
|
||||
|
||||
## Verification Checklist
|
||||
|
||||
✅ Directory structure created
|
||||
✅ Shared infrastructure migrated
|
||||
✅ Message capture module migrated
|
||||
✅ AI moderation module migrated
|
||||
✅ Attachment upload module migrated
|
||||
✅ Event broadcaster module created (Redis pub/sub)
|
||||
✅ Bootstrap and entry point created
|
||||
✅ Package.json with dependencies
|
||||
✅ No HTTP server code (Express, WebSocket removed)
|
||||
✅ Event-driven architecture implemented
|
||||
✅ Graceful shutdown handler
|
||||
✅ Module index files for clean exports
|
||||
✅ Architecture documentation
|
||||
|
||||
## Event Flow Diagram
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────────────────────────────┐
|
||||
│ Discord Gateway Service │
|
||||
├─────────────────────────────────────────────────────────────────┤
|
||||
│ │
|
||||
│ ┌────────────────────────────┐ ┌────────────────────────────┐ │
|
||||
│ │ Message Capture │ │ AI Moderation │ │
|
||||
│ │ (Controller) │ │ (Controller) │ │
|
||||
│ └──────────────┬─────────────┘ └──────────────┬─────────────┘ │
|
||||
│ │ │ │
|
||||
│ ├───────────────────────────────┤ │
|
||||
│ │ │ │
|
||||
│ ▼ ▼ │
|
||||
│ ┌───────────────────────────────────────────────────────────┐ │
|
||||
│ │ Event Broadcaster (Redis Pub/Sub) │ │
|
||||
│ │ - discord:message:created │ │
|
||||
│ │ - discord:message:updated │ │
|
||||
│ │ - discord:message:deleted │ │
|
||||
│ │ - discord:message:analyzed │ │
|
||||
│ │ - discord:attachment:created │ │
|
||||
│ │ - discord:attachment:uploaded │ │
|
||||
│ │ - discord:analysis:queue_status │ │
|
||||
│ └───────────────────────────────────────────────────────────┘ │
|
||||
│ │ │
|
||||
└────────────────────────────────┼────────────────────────────────┘
|
||||
│
|
||||
│ Redis Pub/Sub
|
||||
│
|
||||
▼
|
||||
┌─────────────────────────────────────────────────────────────────┐
|
||||
│ Backend Service │
|
||||
│ (Subscribes to events, serves HTTP API, manages WebSocket) │
|
||||
└─────────────────────────────────────────────────────────────────┘
|
||||
│
|
||||
│ HTTP API
|
||||
│
|
||||
▼
|
||||
┌─────────────────────────────────────────────────────────────────┐
|
||||
│ Frontend Application │
|
||||
│ (React SPA, real-time updates via WebSocket) │
|
||||
└─────────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
## Summary
|
||||
|
||||
The Discord Gateway service has been successfully extracted with:
|
||||
- **Modular MVC architecture** for clean separation of concerns
|
||||
- **Event-driven design** using Redis pub/sub for inter-service communication
|
||||
- **Shared infrastructure layer** for reusable components
|
||||
- **No HTTP server** — pure event-driven service
|
||||
- **Graceful shutdown** handling
|
||||
- **Type-safe configuration** with Zod validation
|
||||
- **Structured logging** with Winston
|
||||
- **PostgreSQL integration** with Drizzle ORM
|
||||
|
||||
The service is ready for integration with the Backend service, which will consume Redis events and serve the HTTP API to the Frontend.
|
||||
Vitest, tests in `tests/`. Config supplies dummy env vars so the suite runs
|
||||
without live Postgres/Redis; external services are mocked. `llmE2e.test.ts`
|
||||
is skipped by default and needs real credentials (`pnpm test:e2e:live`).
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,2 @@
|
||||
[test]
|
||||
preload = ["./tests/setup-env.ts"]
|
||||
@@ -17,13 +17,10 @@
|
||||
"typecheck": "tsc --noEmit",
|
||||
"lint": "biome check --diagnostic-level=error .",
|
||||
"format": "biome format --write .",
|
||||
"test": "vitest run",
|
||||
"test:unit": "vitest run --exclude \"tests/llmE2e.test.ts\"",
|
||||
"test:e2e": "vitest run tests/llmE2e.test.ts",
|
||||
"test:e2e:live": "bash scripts/run-llm-e2e.sh"
|
||||
"test": "bun test tests/",
|
||||
"test:e2e": "bun test tests/llmE2e.test.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@typesafe-ai/sdk": "^0.6.0",
|
||||
"axios": "^1.20.0",
|
||||
"discord.js-selfbot-v13": "^3.7.1",
|
||||
"dotenv": "^18.0.0",
|
||||
@@ -44,12 +41,13 @@
|
||||
},
|
||||
"devDependencies": {
|
||||
"@biomejs/biome": "latest",
|
||||
"@types/node": "^26.4.0",
|
||||
"@types/node": "^26.6.2",
|
||||
"@types/pg": "^8.23.1",
|
||||
"@types/ws": "^8.18.1",
|
||||
"drizzle-kit": "^0.31.10",
|
||||
"tsx": "^4.23.13",
|
||||
"drizzle-kit": "^0.31.11",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "latest"
|
||||
}
|
||||
}
|
||||
"@types/bun": "latest"
|
||||
},
|
||||
"packageManager": "bun@1.3.14"
|
||||
}
|
||||
Generated
-4101
File diff suppressed because it is too large
Load Diff
@@ -1,19 +0,0 @@
|
||||
allowBuilds:
|
||||
"@discordjs/opus": true
|
||||
"@lng2004/node-datachannel": true
|
||||
esbuild: true
|
||||
node-av: true
|
||||
sharp: true
|
||||
zeromq: true
|
||||
# pnpm 11 requires build-script approvals here (the legacy `pnpm` field in
|
||||
# package.json is ignored). Native voice deps need their postinstall build.
|
||||
# NOTE: sharp sengaja TIDAK ada — binary-nya dari @img/sharp-linux-x64
|
||||
# (prebuilt), install script-nya cuma validasi dan gagal di Nix sandbox.
|
||||
# Kalau script sharp dijalankan pnpm rebuild abort sebelum opus/datachannel
|
||||
# kebangun. node-crc dihapus dari deps (tidak pernah di-import).
|
||||
onlyBuiltDependencies:
|
||||
- "@discordjs/opus"
|
||||
- "@lng2004/node-datachannel"
|
||||
- esbuild
|
||||
- node-av
|
||||
- zeromq
|
||||
@@ -1,82 +1,54 @@
|
||||
import { Client } from "discord.js-selfbot-v13";
|
||||
import { ConfigError, DatabaseError } from "@/shared/errors/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import {
|
||||
getAnalysisQueueStatus,
|
||||
startPendingAIAnalysisWorker,
|
||||
} from "../modules/ai-moderation/aiAnalyzer.js";
|
||||
import {
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "../modules/ai-moderation/circuitBreaker.js";
|
||||
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
|
||||
ConfigError,
|
||||
DatabaseError,
|
||||
errorMessage,
|
||||
} from "@/shared/errors/index.js";
|
||||
import { createChildLogger } from "@/shared/logger/index.js";
|
||||
import { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||
import {
|
||||
EventBroadcaster,
|
||||
RedisEventPublisher,
|
||||
} from "../modules/event-broadcaster/index.js";
|
||||
import {
|
||||
registerCollector,
|
||||
setGauge,
|
||||
startMetricsServer,
|
||||
stopMetricsServer,
|
||||
} from "../modules/gateway-metrics/index.js";
|
||||
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
|
||||
import {
|
||||
registerMessageCapture,
|
||||
setEventBroadcaster as setMessageCaptureEventBroadcaster,
|
||||
} from "../modules/message-capture/messageCapture.js";
|
||||
import { setModerationEventBroadcaster } from "../modules/message-capture/moderationActionsDb.js";
|
||||
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
|
||||
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
|
||||
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
|
||||
import { registerPresenceCapture } from "../modules/user-presence/index.js";
|
||||
import { config } from "../shared/config/config.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import {
|
||||
closeDatabase,
|
||||
initializeDatabase,
|
||||
} from "../shared/database/drizzle.js";
|
||||
import { runMigrations } from "../shared/database/migrate.js";
|
||||
import { createDiscordClientOptions } from "../shared/discord/clientOptions.js";
|
||||
import { startRetentionCleanup } from "./retention.js";
|
||||
import { startGatewayLifecycle } from "./lifecycle.js";
|
||||
import { registerPipelineMetrics } from "./metrics-collector.js";
|
||||
import { registerProcessGuards } from "./process-guards.js";
|
||||
import { createGracefulShutdown } from "./shutdown.js";
|
||||
|
||||
const logger = createChildLogger("discord-gateway");
|
||||
|
||||
// ─── Bootstrap ─────────────────────────────────────────────────────────────
|
||||
//
|
||||
// Startup order:
|
||||
// 1. validate config (fail fast on missing AI credentials)
|
||||
// 2. connect infrastructure (migrations → DB pool)
|
||||
// 3. build long-lived services (Discord client, Redis publisher, command
|
||||
// handler) + install shutdown/process guards
|
||||
// 4. start observability (pipeline gauges → metrics server)
|
||||
// 5. log in (ready-hook wires listeners via lifecycle.ts)
|
||||
|
||||
export async function initializeDiscordGateway() {
|
||||
/** Refuse to start when AI analysis is on but no LLM credentials exist. */
|
||||
function assertConfigIsUsable(): void {
|
||||
if (config.AI_ANALYSIS_ENABLED && !config.AI_LLM_API_KEY) {
|
||||
throw new ConfigError(
|
||||
"AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const token = config.DISCORD_TOKEN;
|
||||
logger.info(
|
||||
{ hasToken: token.length > 0, tokenLength: token.length },
|
||||
"Config loaded",
|
||||
);
|
||||
|
||||
logger.info("Creating Discord client");
|
||||
const client = new Client(createDiscordClientOptions());
|
||||
|
||||
// Initialize Redis event broadcaster
|
||||
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
|
||||
const eventBroadcaster = new EventBroadcaster(redisPublisher);
|
||||
|
||||
// Initialize Redis command handler for backend→gateway commands
|
||||
const commandHandler = new CommandHandler();
|
||||
|
||||
const gracefulShutdown = createGracefulShutdown({
|
||||
logger,
|
||||
closeDatabase,
|
||||
client,
|
||||
eventBroadcaster,
|
||||
commandHandler,
|
||||
stopMetricsServer,
|
||||
});
|
||||
|
||||
/** Run migrations (when enabled) then open the PostgreSQL pool. */
|
||||
async function connectDatabase(): Promise<void> {
|
||||
try {
|
||||
if (config.AUTO_MIGRATE_ON_STARTUP) {
|
||||
logger.info(
|
||||
@@ -90,175 +62,79 @@ export async function initializeDiscordGateway() {
|
||||
logger.info("PostgreSQL database initialized");
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
|
||||
{ err, errorMsg: errorMessage(err) },
|
||||
"Failed to initialize database",
|
||||
);
|
||||
throw new DatabaseError(
|
||||
`Database initialization failed: ${err instanceof Error ? err.message : String(err)}`,
|
||||
`Database initialization failed: ${errorMessage(err)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/** Log only client debug lines that carry signal (errors/streams, or VERBOSE). */
|
||||
function registerClientDebugLogging(client: Client): void {
|
||||
client.on("debug", (msg) => {
|
||||
if (
|
||||
msg.toLowerCase().includes("error") ||
|
||||
msg.toLowerCase().includes("stream")
|
||||
) {
|
||||
const lower = msg.toLowerCase();
|
||||
if (lower.includes("error") || lower.includes("stream")) {
|
||||
logger.info({ debugMsg: msg }, "Discord Client Debug");
|
||||
} else if (config.VERBOSE) {
|
||||
logger.debug({ debugMsg: msg }, "Discord Client Debug");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
client.on("ready", async () => {
|
||||
export async function initializeDiscordGateway() {
|
||||
assertConfigIsUsable();
|
||||
|
||||
const token = config.DISCORD_TOKEN;
|
||||
logger.info(
|
||||
{ hasToken: token.length > 0, tokenLength: token.length },
|
||||
"Config loaded",
|
||||
);
|
||||
|
||||
logger.info("Creating Discord client");
|
||||
const client = new Client(createDiscordClientOptions());
|
||||
|
||||
// Long-lived services: Redis event broadcaster (gateway → backend) and the
|
||||
// Redis command handler (backend → gateway).
|
||||
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
|
||||
const eventBroadcaster = new EventBroadcaster(redisPublisher);
|
||||
const commandHandler = new CommandHandler();
|
||||
|
||||
const gracefulShutdown = createGracefulShutdown({
|
||||
logger,
|
||||
closeDatabase,
|
||||
client,
|
||||
eventBroadcaster,
|
||||
commandHandler,
|
||||
stopMetricsServer,
|
||||
});
|
||||
|
||||
await connectDatabase();
|
||||
|
||||
registerClientDebugLogging(client);
|
||||
|
||||
client.on("ready", () => {
|
||||
logger.info({ user: client.user?.tag }, "Bot logged in");
|
||||
setMessageCaptureEventBroadcaster(eventBroadcaster);
|
||||
setModerationEventBroadcaster(eventBroadcaster);
|
||||
registerMessageCapture(client);
|
||||
startPendingAIAnalysisWorker(client, eventBroadcaster);
|
||||
|
||||
// Register new event captures
|
||||
registerReactionCapture(client, eventBroadcaster);
|
||||
registerThreadCapture(client, eventBroadcaster);
|
||||
registerPresenceCapture(client, eventBroadcaster);
|
||||
registerChannelTopicCapture(client, eventBroadcaster);
|
||||
registerGuildMemberEvents(client, eventBroadcaster);
|
||||
|
||||
// Start command handler after Discord is ready
|
||||
commandHandler.start(client);
|
||||
logger.info("Command handler started");
|
||||
|
||||
// Start retention cleanup scheduler
|
||||
startRetentionCleanup();
|
||||
// Start weekly moderation digest (public, automated)
|
||||
startDigestScheduler();
|
||||
startGatewayLifecycle({
|
||||
client,
|
||||
eventBroadcaster,
|
||||
commandHandler,
|
||||
logger,
|
||||
});
|
||||
});
|
||||
|
||||
client.on("error", (err) => {
|
||||
logger.error(
|
||||
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
|
||||
"Client error",
|
||||
);
|
||||
logger.error({ err, errorMsg: errorMessage(err) }, "Client error");
|
||||
});
|
||||
|
||||
process.on("SIGINT", () => {
|
||||
gracefulShutdown("SIGINT");
|
||||
});
|
||||
registerProcessGuards(logger, gracefulShutdown);
|
||||
|
||||
process.on("SIGTERM", () => {
|
||||
gracefulShutdown("SIGTERM");
|
||||
});
|
||||
|
||||
process.on("uncaughtException", (err) => {
|
||||
const code =
|
||||
typeof (err as NodeJS.ErrnoException).code === "string"
|
||||
? (err as NodeJS.ErrnoException).code
|
||||
: "";
|
||||
// Transient stream-teardown errors (voice stop/disconnect races, child
|
||||
// process stdin closed while we still write) are NOT fatal — crashing the
|
||||
// gateway on EPIPE takes the whole bot offline mid-music. Log + continue.
|
||||
if (
|
||||
code === "EPIPE" ||
|
||||
code === "ERR_STREAM_DESTROYED" ||
|
||||
code === "ERR_STREAM_WRITE_AFTER_END" ||
|
||||
code === "ECONNRESET"
|
||||
) {
|
||||
logger.warn(
|
||||
{ error: err },
|
||||
"Uncaught transient stream error — continuing",
|
||||
);
|
||||
return;
|
||||
}
|
||||
logger.error(
|
||||
{
|
||||
err,
|
||||
errorMsg: err instanceof Error ? err.message : String(err),
|
||||
stack: err?.stack,
|
||||
},
|
||||
"Uncaught exception",
|
||||
);
|
||||
gracefulShutdown("uncaughtException");
|
||||
});
|
||||
|
||||
process.on("unhandledRejection", (reason) => {
|
||||
const err =
|
||||
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
|
||||
const code = (err as NodeJS.ErrnoException).code ?? "";
|
||||
// Same transient-teardown policy as uncaughtException: a rejection that
|
||||
// fires while a stream is being torn down (EPIPE after ffmpeg stdin
|
||||
// closes, write-after-destroy, socket reset) must NOT take the whole
|
||||
// gateway offline. Log detail + continue. Everything else still shuts
|
||||
// down so real bugs surface.
|
||||
if (
|
||||
code === "EPIPE" ||
|
||||
code === "ERR_STREAM_DESTROYED" ||
|
||||
code === "ERR_STREAM_WRITE_AFTER_END" ||
|
||||
code === "ECONNRESET"
|
||||
) {
|
||||
logger.warn(
|
||||
{ error: err },
|
||||
"Unhandled rejection transient stream error — continuing",
|
||||
);
|
||||
return;
|
||||
}
|
||||
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
|
||||
gracefulShutdown("unhandledRejection");
|
||||
});
|
||||
|
||||
// ── Metrics: register live pipeline collectors before starting server ──
|
||||
// These refresh on every scrape so Prometheus sees real AI-analysis
|
||||
// queue depth, concurrency, and DB pool state instead of an empty stub.
|
||||
registerCollector(() => {
|
||||
if (!config.AI_ANALYSIS_ENABLED) return;
|
||||
try {
|
||||
const status = getAnalysisQueueStatus();
|
||||
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
|
||||
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
|
||||
setGauge(
|
||||
"ai_analysis_active_individual_requests",
|
||||
status.activeIndividualRequests,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_individual_in_flight",
|
||||
status.individualInFlightCount,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_individual_circuit_breaker_active",
|
||||
status.individualCircuitBreakerActive ? 1 : 0,
|
||||
);
|
||||
if (typeof status.lastError === "string") {
|
||||
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
|
||||
}
|
||||
type PoolState = { _poolState?: { size: number; active: number } };
|
||||
const textPool = textWorkerPool as unknown as PoolState;
|
||||
const mediaPool = mediaWorkerPool as unknown as PoolState;
|
||||
// Reported per queue (2026-08-31 text/media pool split) so the text
|
||||
// and media backlogs are distinguishable in dashboards/alerts instead
|
||||
// of one combined "worker threads" number.
|
||||
if (textPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_text",
|
||||
textPool._poolState.active,
|
||||
);
|
||||
}
|
||||
if (mediaPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_media",
|
||||
mediaPool._poolState.active,
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.warn({ error: String(err) }, "AI metrics collector failed");
|
||||
}
|
||||
});
|
||||
|
||||
// Start metrics server
|
||||
// Metrics: register live pipeline collectors before starting the server.
|
||||
registerPipelineMetrics(logger);
|
||||
startMetricsServer();
|
||||
|
||||
logger.info("Calling Discord client.login");
|
||||
|
||||
// Fix: use await + try/catch instead of .then().catch()
|
||||
try {
|
||||
await client.login(token);
|
||||
logger.info("Discord client logged in successfully");
|
||||
|
||||
@@ -0,0 +1,61 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import type { Logger } from "@/shared/logger/index.js";
|
||||
import { startPendingAIAnalysisWorker } from "../modules/ai-moderation/index.js";
|
||||
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
|
||||
import type { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
|
||||
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
|
||||
import {
|
||||
registerMessageCapture,
|
||||
setEventBroadcaster as setMessageCaptureEventBroadcaster,
|
||||
setModerationEventBroadcaster,
|
||||
} from "../modules/message-capture/index.js";
|
||||
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
|
||||
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
|
||||
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
|
||||
import { registerPresenceCapture } from "../modules/user-presence/index.js";
|
||||
import { startRetentionCleanup } from "./retention.js";
|
||||
|
||||
export interface GatewayLifecycleOptions {
|
||||
client: Client;
|
||||
eventBroadcaster: EventBroadcaster;
|
||||
commandHandler: CommandHandler;
|
||||
logger: Logger;
|
||||
}
|
||||
|
||||
/**
|
||||
* Wires everything that must start once Discord is connected.
|
||||
*
|
||||
* Ordering matters:
|
||||
* 1. Inject the event broadcaster into the modules that publish events —
|
||||
* they must be able to publish before their listeners are registered.
|
||||
* 2. Register the Discord event listeners (capture modules).
|
||||
* 3. Start the background workers/schedulers.
|
||||
*/
|
||||
export function startGatewayLifecycle({
|
||||
client,
|
||||
eventBroadcaster,
|
||||
commandHandler,
|
||||
logger,
|
||||
}: GatewayLifecycleOptions): void {
|
||||
// 1. Inject broadcaster first so no captured event is dropped.
|
||||
setMessageCaptureEventBroadcaster(eventBroadcaster);
|
||||
setModerationEventBroadcaster(eventBroadcaster);
|
||||
|
||||
// 2. Discord event listeners.
|
||||
registerMessageCapture(client);
|
||||
registerReactionCapture(client, eventBroadcaster);
|
||||
registerThreadCapture(client, eventBroadcaster);
|
||||
registerPresenceCapture(client, eventBroadcaster);
|
||||
registerChannelTopicCapture(client, eventBroadcaster);
|
||||
registerGuildMemberEvents(client, eventBroadcaster);
|
||||
|
||||
// 3. Background workers + schedulers.
|
||||
startPendingAIAnalysisWorker(client, eventBroadcaster);
|
||||
commandHandler.start(client);
|
||||
logger.info("Command handler started");
|
||||
|
||||
startRetentionCleanup();
|
||||
// Weekly moderation digest (public, automated)
|
||||
startDigestScheduler();
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
import type { Logger } from "@/shared/logger/index.js";
|
||||
import {
|
||||
getAnalysisQueueStatus,
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "../modules/ai-moderation/index.js";
|
||||
import {
|
||||
registerCollector,
|
||||
setGauge,
|
||||
} from "../modules/gateway-metrics/index.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
|
||||
/** Piscina exposes its live thread counters on `_poolState`. */
|
||||
type PoolState = { _poolState?: { size: number; active: number } };
|
||||
|
||||
/**
|
||||
* Registers the AI-pipeline Prometheus gauges.
|
||||
*
|
||||
* The collector refreshes on every scrape, so Prometheus sees real queue
|
||||
* depth / concurrency / worker-thread state instead of an empty stub.
|
||||
* Registered before the metrics server starts.
|
||||
*/
|
||||
export function registerPipelineMetrics(logger: Logger): void {
|
||||
registerCollector(() => {
|
||||
if (!config.AI_ANALYSIS_ENABLED) return;
|
||||
try {
|
||||
const status = getAnalysisQueueStatus();
|
||||
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
|
||||
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
|
||||
setGauge(
|
||||
"ai_analysis_active_text_requests",
|
||||
status.activeTextRequests ?? status.activeRequests,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_active_media_requests",
|
||||
status.activeMediaRequests ?? 0,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_active_individual_requests",
|
||||
status.activeIndividualRequests,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_individual_in_flight",
|
||||
status.individualInFlightCount,
|
||||
);
|
||||
setGauge(
|
||||
"ai_analysis_individual_circuit_breaker_active",
|
||||
status.individualCircuitBreakerActive ? 1 : 0,
|
||||
);
|
||||
if (typeof status.lastError === "string") {
|
||||
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
|
||||
}
|
||||
|
||||
// Reported per queue (2026-08-31 text/media pool split) so the text
|
||||
// and media backlogs are distinguishable in dashboards/alerts instead
|
||||
// of one combined "worker threads" number.
|
||||
const textPool = textWorkerPool as unknown as PoolState;
|
||||
const mediaPool = mediaWorkerPool as unknown as PoolState;
|
||||
if (textPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_text",
|
||||
textPool._poolState.active,
|
||||
);
|
||||
}
|
||||
if (mediaPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_media",
|
||||
mediaPool._poolState.active,
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.warn({ error: String(err) }, "AI metrics collector failed");
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
import { errorMessage, isTransientStreamError } from "@/shared/errors/index.js";
|
||||
import type { Logger } from "@/shared/logger/index.js";
|
||||
import type { GracefulShutdown } from "./shutdown.js";
|
||||
|
||||
/**
|
||||
* Process-level signal + error guards.
|
||||
*
|
||||
* Extracted from bootstrap so the "what keeps the gateway alive vs what
|
||||
* shuts it down" policy lives in exactly one place.
|
||||
*
|
||||
* Policy: transient stream-teardown failures (EPIPE / ERR_STREAM_DESTROYED /
|
||||
* ERR_STREAM_WRITE_AFTER_END / ECONNRESET) are logged and IGNORED — crashing
|
||||
* the gateway on them (voice stop/disconnect races, a child process stdin
|
||||
* closed while we still write) takes the whole bot offline mid-operation.
|
||||
* Anything else is a real bug: log with stack and shut down cleanly.
|
||||
*/
|
||||
export function registerProcessGuards(
|
||||
logger: Logger,
|
||||
gracefulShutdown: GracefulShutdown,
|
||||
): void {
|
||||
process.on("SIGINT", () => {
|
||||
gracefulShutdown("SIGINT");
|
||||
});
|
||||
|
||||
process.on("SIGTERM", () => {
|
||||
gracefulShutdown("SIGTERM");
|
||||
});
|
||||
|
||||
process.on("uncaughtException", (err) => {
|
||||
if (isTransientStreamError(err)) {
|
||||
logger.warn(
|
||||
{ error: err },
|
||||
"Uncaught transient stream error — continuing",
|
||||
);
|
||||
return;
|
||||
}
|
||||
logger.error(
|
||||
{
|
||||
err,
|
||||
errorMsg: errorMessage(err),
|
||||
stack: err?.stack,
|
||||
},
|
||||
"Uncaught exception",
|
||||
);
|
||||
gracefulShutdown("uncaughtException");
|
||||
});
|
||||
|
||||
process.on("unhandledRejection", (reason) => {
|
||||
const err =
|
||||
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
|
||||
if (isTransientStreamError(err)) {
|
||||
logger.warn(
|
||||
{ error: err },
|
||||
"Unhandled rejection transient stream error — continuing",
|
||||
);
|
||||
return;
|
||||
}
|
||||
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
|
||||
gracefulShutdown("unhandledRejection");
|
||||
});
|
||||
}
|
||||
@@ -1,10 +1,7 @@
|
||||
import { inArray, lt } from "drizzle-orm";
|
||||
import type {
|
||||
NodePgDatabase,
|
||||
NodePgQueryResultHKT,
|
||||
} from "drizzle-orm/node-postgres";
|
||||
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../shared/config/config.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import { getDatabase } from "../shared/database/drizzle.js";
|
||||
import type * as schema from "../shared/database/schema.js";
|
||||
import { attachmentsTable, messagesTable } from "../shared/database/schema.js";
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import type { createChildLogger } from "@/shared/logger/index";
|
||||
import {
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "../modules/ai-moderation/circuitBreaker.js";
|
||||
import type { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
|
||||
import type { stopMetricsServer } from "../modules/gateway-metrics/index.js";
|
||||
@@ -45,6 +49,30 @@ export function createGracefulShutdown(
|
||||
options.logger.info("Closing command handler...");
|
||||
await options.commandHandler.close();
|
||||
|
||||
// ½. Tear down AI-analysis worker pools BEFORE closing the DB.
|
||||
// Piscina worker threads survive process.exit() as orphans otherwise —
|
||||
// they keep holding DB connections/locks after the main process is gone.
|
||||
// (Two live gateways fighting over the same rows was the root cause of
|
||||
// messages stuck in ai_status='processing'.)
|
||||
options.logger.info("Destroying AI worker pools...");
|
||||
const destroyPool = (pool: { destroy: () => Promise<void> }) =>
|
||||
Promise.race([
|
||||
pool.destroy(),
|
||||
new Promise<void>((resolve) =>
|
||||
setTimeout(() => {
|
||||
options.logger.warn(
|
||||
"Timed out destroying worker pool; exiting anyway",
|
||||
);
|
||||
resolve();
|
||||
}, 5000),
|
||||
),
|
||||
]);
|
||||
await Promise.allSettled([
|
||||
destroyPool(textWorkerPool),
|
||||
destroyPool(mediaWorkerPool),
|
||||
]);
|
||||
options.logger.info("AI worker pools destroyed");
|
||||
|
||||
// 2. DB pool
|
||||
options.logger.info("Closing database...");
|
||||
await options.closeDatabase();
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
*/
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { initializeDatabase } from "../../shared/database/drizzle.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
@@ -86,7 +86,12 @@ export interface MessageBatch {
|
||||
|
||||
// Worker job types (Piscina entry point)
|
||||
type WorkerJob =
|
||||
| { type: "batch"; conversationKey: string; messages: MessageRecord[] }
|
||||
| {
|
||||
type: "batch";
|
||||
conversationKey: string;
|
||||
lane: "text" | "media";
|
||||
messages: MessageRecord[];
|
||||
}
|
||||
| { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean };
|
||||
|
||||
type BatchOkResponse = {
|
||||
@@ -263,6 +268,7 @@ function normalizeResult(
|
||||
async function processBatch(job: {
|
||||
type: "batch";
|
||||
conversationKey: string;
|
||||
lane: "text" | "media";
|
||||
messages: MessageRecord[];
|
||||
}): Promise<BatchOkResponse | BatchErrorResponse> {
|
||||
const { conversationKey, messages } = job;
|
||||
|
||||
@@ -1,34 +1,25 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { EventBroadcaster } from "../event-broadcaster/index.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { AnalysisQueueStatus } from "../message-capture/types.js";
|
||||
import {
|
||||
activeMediaRequests,
|
||||
activeRequests,
|
||||
activeTextRequests,
|
||||
buildAgeRestrictedSkipResult,
|
||||
buildSkipAnalysisUserResult,
|
||||
isAgeRestrictedMessage,
|
||||
isSkipAnalysisUser,
|
||||
skipAgeRestrictedMessages,
|
||||
skipAnalysisUserMessages,
|
||||
} from "./batchProcessor.js";
|
||||
import { scheduleConversationAnalysis } from "./batchScheduler.js";
|
||||
import { getConversationKey } from "./circuitBreaker.js";
|
||||
import {
|
||||
conversationConsecutiveErrors,
|
||||
conversationDebounceTimers,
|
||||
conversationErrorCooldown,
|
||||
conversationProcessing,
|
||||
isConversationProcessingLocked,
|
||||
} from "./conversationState.js";
|
||||
import { conversationDebounceTimers } from "./conversationState.js";
|
||||
import {
|
||||
activeIndividualRequests,
|
||||
enqueueIndividualFallbacks,
|
||||
individualCooldownUntil,
|
||||
individualInFlight,
|
||||
individualInFlightByConversation,
|
||||
individualInFlightLastTouched,
|
||||
} from "./individualFallbackProcessor.js";
|
||||
import {
|
||||
broadcastAnalysisCompleted,
|
||||
@@ -36,31 +27,20 @@ import {
|
||||
setModerationClient,
|
||||
setSharedEventBroadcaster,
|
||||
} from "./moderationState.js";
|
||||
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
|
||||
import { pruneExpiredTexts } from "./textCacheStore.js";
|
||||
import { startRecoveryWorker } from "./recovery-worker.js";
|
||||
|
||||
const logger = createChildLogger("ai-analyzer");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Cache hygiene (expired verdict sweep)
|
||||
// ---------------------------------------------------------------------------
|
||||
const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
|
||||
let lastCachePruneAt = 0;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Re-exports from sub-modules (preserving original public API)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export { pickBatchWithinBudget } from "./batchProcessor.js";
|
||||
export { getConversationKey } from "./circuitBreaker.js";
|
||||
export { onCircuitBreakerAlert } from "./conversationState.js";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Public API
|
||||
// Public API — queueing, status, worker startup
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Queues a message for analysis (debounced by conversation).
|
||||
*
|
||||
* Messages that never need an LLM call are short-circuited here and recorded
|
||||
* with their skip verdict: age-restricted messages and configured skip-list
|
||||
* users.
|
||||
*/
|
||||
export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
||||
if (!config.AI_ANALYSIS_ENABLED) return;
|
||||
@@ -73,13 +53,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
||||
}
|
||||
|
||||
if (isAgeRestrictedMessage(message)) {
|
||||
const updated = await messageStore.updateMessageAIAnalysis(
|
||||
message.id,
|
||||
buildAgeRestrictedSkipResult(),
|
||||
);
|
||||
if (updated) {
|
||||
broadcastAnalysisCompleted(updated);
|
||||
}
|
||||
await recordSkip(message.id, buildAgeRestrictedSkipResult());
|
||||
logger.debug(
|
||||
{ messageId },
|
||||
"Skipped AI analysis for age-restricted message",
|
||||
@@ -88,13 +62,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
||||
}
|
||||
|
||||
if (isSkipAnalysisUser(message)) {
|
||||
const updated = await messageStore.updateMessageAIAnalysis(
|
||||
message.id,
|
||||
buildSkipAnalysisUserResult(),
|
||||
);
|
||||
if (updated) {
|
||||
broadcastAnalysisCompleted(updated);
|
||||
}
|
||||
await recordSkip(message.id, buildSkipAnalysisUserResult());
|
||||
logger.debug(
|
||||
{ messageId, userId: message.user_id },
|
||||
"Skipped AI analysis for configured skip-list user",
|
||||
@@ -114,6 +82,17 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
/** Persist a skip verdict and broadcast it so the dashboard reflects it. */
|
||||
async function recordSkip(
|
||||
messageId: string,
|
||||
result: Parameters<typeof messageStore.updateMessageAIAnalysis>[1],
|
||||
): Promise<void> {
|
||||
const updated = await messageStore.updateMessageAIAnalysis(messageId, result);
|
||||
if (updated) {
|
||||
broadcastAnalysisCompleted(updated);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Queues a conversation for analysis (debounced).
|
||||
*/
|
||||
@@ -129,6 +108,8 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
|
||||
return {
|
||||
queuedConversations: conversationDebounceTimers.size,
|
||||
activeRequests,
|
||||
activeTextRequests,
|
||||
activeMediaRequests,
|
||||
activeIndividualRequests,
|
||||
individualInFlightCount: individualInFlight.size,
|
||||
individualCircuitBreakerActive: Date.now() < individualCooldownUntil,
|
||||
@@ -137,11 +118,12 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
|
||||
}
|
||||
|
||||
/**
|
||||
* Starts the periodic recovery worker.
|
||||
* Starts the background workers behind the analysis pipeline:
|
||||
* - the recovery worker (stranded pending / incomplete messages + cache prune)
|
||||
* - the optional culture and user-profile learners.
|
||||
*
|
||||
* Now also recovers messages stuck in `error/analysis_incomplete`
|
||||
* state (not just `pending`), and skips conversations that already have
|
||||
* individual fallback work in progress to avoid DB last-write-wins races.
|
||||
* Also injects the Discord client and event broadcaster into the pipeline
|
||||
* state so downstream modules can act and publish.
|
||||
*/
|
||||
export function startPendingAIAnalysisWorker(
|
||||
client?: Client,
|
||||
@@ -160,131 +142,5 @@ export function startPendingAIAnalysisWorker(
|
||||
.catch(console.error);
|
||||
}
|
||||
|
||||
setInterval(() => {
|
||||
// [D] Periodic cache hygiene: purge expired moderation verdicts from
|
||||
// Postgres and Qdrant. Expired entries are never reused (filters check
|
||||
// expires_at) but accumulate forever without this sweep.
|
||||
const now = Date.now();
|
||||
if (now - lastCachePruneAt >= CACHE_PRUNE_INTERVAL_MS) {
|
||||
lastCachePruneAt = now;
|
||||
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
|
||||
.then(([pgDeleted, qdDeleted]) => {
|
||||
if (pgDeleted > 0 || qdDeleted > 0) {
|
||||
logger.info(
|
||||
{ pgDeleted, qdDeleted },
|
||||
"Expired moderation cache pruned",
|
||||
);
|
||||
}
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
logger.warn({ error: String(err) }, "Moderation cache prune failed");
|
||||
});
|
||||
}
|
||||
|
||||
// Only revert stuck processing messages if there's active processing.
|
||||
// Avoids a DB query every recovery interval when the pipeline is idle.
|
||||
if (conversationProcessing.size > 0) {
|
||||
messageStore
|
||||
.revertStuckProcessingMessages(300000)
|
||||
.catch((err: unknown) => {
|
||||
logger.error(
|
||||
{ error: String(err) },
|
||||
"Failed to run stuck processing recovery",
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
Promise.all([
|
||||
messageStore.getPendingConversationKeys(500),
|
||||
messageStore.getConversationKeysWithIncompleteAnalysis(200),
|
||||
])
|
||||
.then(([pendingKeys, incompleteKeys]) => {
|
||||
const now = Date.now();
|
||||
|
||||
for (const [key, expiry] of conversationErrorCooldown) {
|
||||
if (now >= expiry) conversationErrorCooldown.delete(key);
|
||||
}
|
||||
for (const [key, startedAt] of conversationProcessing) {
|
||||
if (now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS) {
|
||||
conversationProcessing.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
const staleThreshold = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS * 2;
|
||||
for (const [key, lastTouched] of individualInFlightLastTouched) {
|
||||
if (now - lastTouched >= staleThreshold) {
|
||||
individualInFlightLastTouched.delete(key);
|
||||
individualInFlightByConversation.delete(key);
|
||||
logger.warn(
|
||||
{ key },
|
||||
"Pruned stale individualInFlightByConversation entry",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Also prune stale per-conversation CB error counts that have cooled
|
||||
// down so old conversations can be retried.
|
||||
for (const [key] of conversationConsecutiveErrors) {
|
||||
const cbExpire = conversationErrorCooldown.get(key) ?? 0;
|
||||
if (cbExpire && now >= cbExpire) {
|
||||
conversationConsecutiveErrors.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
const incompleteKeySet = new Set(incompleteKeys);
|
||||
|
||||
// --- Batch recovery for pending messages ---
|
||||
for (const key of pendingKeys) {
|
||||
if (conversationDebounceTimers.has(key)) continue;
|
||||
if (isConversationProcessingLocked(key)) continue;
|
||||
if (individualInFlightByConversation.has(key)) continue;
|
||||
if (incompleteKeySet.has(key)) continue;
|
||||
const cooldownUntil = conversationErrorCooldown.get(key);
|
||||
if (cooldownUntil && now < cooldownUntil) continue;
|
||||
scheduleConversationAnalysis(key);
|
||||
}
|
||||
|
||||
// --- Individual recovery for error/analysis_incomplete messages ---
|
||||
// Circuit breaker check: no point iterating if individual CB is active.
|
||||
if (now >= individualCooldownUntil) {
|
||||
const promises: Promise<void>[] = [];
|
||||
for (const key of incompleteKeys) {
|
||||
// Skip if individual work is already running for this conversation.
|
||||
if (individualInFlightByConversation.has(key)) continue;
|
||||
// Skip if batch processing is running.
|
||||
if (isConversationProcessingLocked(key)) continue;
|
||||
|
||||
promises.push(
|
||||
messageStore
|
||||
.getIncompleteMessagesByConversation(key, 500)
|
||||
.then(async (msgs) => {
|
||||
const processableMessages = await skipAnalysisUserMessages(
|
||||
await skipAgeRestrictedMessages(msgs),
|
||||
);
|
||||
return processableMessages;
|
||||
})
|
||||
.then((msgs) => {
|
||||
if (msgs.length > 0) {
|
||||
enqueueIndividualFallbacks(msgs);
|
||||
}
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
logger.error(
|
||||
{ key, error: String(err) },
|
||||
"Failed to fetch incomplete messages for recovery",
|
||||
);
|
||||
}),
|
||||
);
|
||||
}
|
||||
// Errors are handled per-key; return the combined promise for observability.
|
||||
return Promise.all(promises);
|
||||
}
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
logger.error(
|
||||
{ error: err instanceof Error ? err.message : String(err) },
|
||||
"Pending AI analysis recovery worker failed",
|
||||
);
|
||||
});
|
||||
}, config.AI_ANALYSIS_RECOVERY_INTERVAL_MS);
|
||||
startRecoveryWorker();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
/**
|
||||
* analysisLanes.ts
|
||||
*
|
||||
* Pure lane helpers for the AI-analysis queue. Kept free of any import chain
|
||||
* that pulls Piscina/worker/DB so they can be unit-tested in isolation (the
|
||||
* scheduler's `splitMessagesByLane` used to live in batchScheduler.ts, which
|
||||
* transitively imports the worker pool).
|
||||
*/
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
import type { AnalysisLane } from "./conversationState.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
|
||||
export type { AnalysisLane } from "./conversationState.js";
|
||||
|
||||
/** True when this message belongs to the media lane (has attachment/sticker/embed). */
|
||||
export function laneOfMessage(message: MessageRecord): AnalysisLane {
|
||||
return hasMediaContent(message) ? "media" : "text";
|
||||
}
|
||||
|
||||
/**
|
||||
* Splits an arbitrary message array into per-lane lists. Used when the
|
||||
* scheduler runs a conversation-wide pass (lane omitted): each lane gets its
|
||||
* own subset so text and media never share a worker job.
|
||||
*/
|
||||
export function splitMessagesByLane(messages: MessageRecord[]): {
|
||||
text: MessageRecord[];
|
||||
media: MessageRecord[];
|
||||
} {
|
||||
const text: MessageRecord[] = [];
|
||||
const media: MessageRecord[] = [];
|
||||
for (const m of messages) {
|
||||
(laneOfMessage(m) === "media" ? media : text).push(m);
|
||||
}
|
||||
return { text, media };
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type {
|
||||
AnalysisResult,
|
||||
MessageRecord,
|
||||
@@ -35,8 +35,17 @@ export function deriveSeverity(msg: MessageRecord): string {
|
||||
|
||||
/** Derive recommended action from legacy messages that lack structured AI fields. */
|
||||
export function deriveRecommendedAction(msg: MessageRecord): string {
|
||||
if (msg.ai_recommended_action) return msg.ai_recommended_action;
|
||||
const severity = deriveSeverity(msg);
|
||||
// Flagged at high/critical severity is ALWAYS delete — the stored
|
||||
// recommended_action from the LLM is conservative (often "review") and
|
||||
// must not override severity for severe violations.
|
||||
if (
|
||||
msg.ai_status === "flagged" &&
|
||||
(severity === "critical" || severity === "high")
|
||||
) {
|
||||
return "delete";
|
||||
}
|
||||
if (msg.ai_recommended_action) return msg.ai_recommended_action;
|
||||
if (
|
||||
msg.ai_status === "flagged" &&
|
||||
(severity === "critical" || severity === "high" || severity === "medium")
|
||||
@@ -177,19 +186,36 @@ export function isEligibleForAutoDelete(
|
||||
return false;
|
||||
}
|
||||
|
||||
// Recommended action check
|
||||
const recommendedAction =
|
||||
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
||||
// Recommended action check.
|
||||
// CRITICAL: the LLM's `recommended_action` is CONSERVATIVE — for a flagged
|
||||
// message at high/critical severity it frequently emits "review" (it sees a
|
||||
// screenshot/context ambiguity and hedges) even when the violation itself is
|
||||
// severe. Trusting that value lets serious violations (harassment, SARA,
|
||||
// threats) slip through undeleted. So: flagged + high/critical severity is
|
||||
// ALWAYS eligible regardless of the LLM's recommended action. The action
|
||||
// check only gates warn/flagged-medium (where a review is legitimate).
|
||||
if (
|
||||
recommendedAction !== "delete" &&
|
||||
recommendedAction !== "escalate" &&
|
||||
recommendedAction !== "warn"
|
||||
status === "flagged" &&
|
||||
(severity === "high" || severity === "critical")
|
||||
) {
|
||||
logger.debug(
|
||||
{ messageId: message.id, recommendedAction },
|
||||
"Message not eligible for auto-delete: recommended action is not delete/escalate/warn",
|
||||
{ messageId: message.id, status, severity },
|
||||
"Message eligible for auto-delete: flagged with high/critical severity",
|
||||
);
|
||||
return false;
|
||||
} else {
|
||||
const recommendedAction =
|
||||
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
||||
if (
|
||||
recommendedAction !== "delete" &&
|
||||
recommendedAction !== "escalate" &&
|
||||
recommendedAction !== "warn"
|
||||
) {
|
||||
logger.debug(
|
||||
{ messageId: message.id, recommendedAction },
|
||||
"Message not eligible for auto-delete: recommended action is not delete/escalate/warn",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Categories check
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Guild } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
interface ChannelWithSend {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import type { Client, PermissionString } from "discord.js-selfbot-v13";
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { parseRichMessageMetadata } from "../message-capture/messageMetadata.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
const logger = createChildLogger("auto-delete-notify");
|
||||
|
||||
@@ -43,3 +43,22 @@ export function pickBatchWithinBudget(
|
||||
|
||||
return batch;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the messages that were fetched/claimed but did NOT make it into the
|
||||
* trimmed batch (i.e. the tail past the token budget).
|
||||
*
|
||||
* The DB claim step flips every fetched pending row to `processing`; the batch
|
||||
* trim may then stop early on the token budget. Those tail rows would stay
|
||||
* stuck in `processing` forever unless the caller explicitly un-claims them —
|
||||
* this helper identifies exactly which rows that is, so the caller can write
|
||||
* them back to `pending` for the next wave.
|
||||
*/
|
||||
export function computeBudgetOverflowMessages(
|
||||
claimed: MessageRecord[],
|
||||
trimmed: MessageRecord[],
|
||||
): MessageRecord[] {
|
||||
if (trimmed.length === 0) return claimed;
|
||||
const trimmedIds = new Set(trimmed.map((m) => m.id));
|
||||
return claimed.filter((m) => !trimmedIds.has(m.id));
|
||||
}
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user