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_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_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_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_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_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_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_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_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 Tuning ===
|
||||||
AI_ANALYSIS_DEBOUNCE_MS=500 # Debounce window for batching messages in ms (default: 500)
|
AI_ANALYSIS_DEBOUNCE_MS=500 # Debounce window for batching messages in ms (default: 500)
|
||||||
|
|||||||
@@ -32,29 +32,26 @@ jobs:
|
|||||||
with:
|
with:
|
||||||
node-version: 22
|
node-version: 22
|
||||||
|
|
||||||
- name: Install pnpm
|
- name: Setup Bun
|
||||||
run: corepack enable && corepack prepare pnpm@11 --activate
|
uses: oven-sh/setup-bun@v2
|
||||||
|
with:
|
||||||
|
bun-version: 1.3.14
|
||||||
|
|
||||||
- name: Install deps (backend)
|
- name: Install deps + test (backend)
|
||||||
working-directory: services/backend
|
|
||||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
|
||||||
|
|
||||||
- name: Typecheck + test (backend)
|
|
||||||
working-directory: services/backend
|
working-directory: services/backend
|
||||||
run: |
|
run: |
|
||||||
|
bun install --frozen-lockfile
|
||||||
./node_modules/.bin/tsc --noEmit
|
./node_modules/.bin/tsc --noEmit
|
||||||
# e2e.test.ts requires a live backend (API_BASE) — run unit tests only
|
# src/e2e.test.ts requires a live backend (API_BASE) — unit tests
|
||||||
./node_modules/.bin/vitest run --exclude "src/e2e.test.ts"
|
# live in tests/ and are excluded by the bun test dir.
|
||||||
|
bun test tests/
|
||||||
|
|
||||||
- name: Install deps (discord-gateway)
|
- name: Install deps + test (discord-gateway)
|
||||||
working-directory: services/discord-gateway
|
|
||||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
|
||||||
|
|
||||||
- name: Typecheck + test (discord-gateway)
|
|
||||||
working-directory: services/discord-gateway
|
working-directory: services/discord-gateway
|
||||||
run: |
|
run: |
|
||||||
|
bun install --frozen-lockfile
|
||||||
./node_modules/.bin/tsc --noEmit
|
./node_modules/.bin/tsc --noEmit
|
||||||
./node_modules/.bin/vitest run
|
bun test tests/
|
||||||
|
|
||||||
- name: Biome check (all services)
|
- name: Biome check (all services)
|
||||||
run: |
|
run: |
|
||||||
@@ -151,9 +148,17 @@ jobs:
|
|||||||
|
|
||||||
attic_push_vps_hop() {
|
attic_push_vps_hop() {
|
||||||
echo "Fallback: VPS-hop attic push"
|
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)
|
# 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 \
|
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.
|
# Push from VPS → Attic over Tailscale.
|
||||||
# --ignore-upstream-cache-filter is REQUIRED: without it, attic skips
|
# --ignore-upstream-cache-filter is REQUIRED: without it, attic skips
|
||||||
# writing the narinfo to gmw when chunks exist in the upstream
|
# 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
|
# sudo: attic must read root's config (~/.config/attic), which has
|
||||||
# the imrnes-ts server → Tailscale. Non-root users' configs only
|
# the imrnes-ts server → Tailscale. Non-root users' configs only
|
||||||
# have the public `pub` server → "Server imrnes-ts does not exist".
|
# 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)"
|
|| echo "attic push failed (non-fatal; ssh copy fallback below)"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -23,3 +23,7 @@ result
|
|||||||
|
|
||||||
# Playwright MCP artifacts
|
# Playwright MCP artifacts
|
||||||
.playwright-mcp/
|
.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": "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": "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": "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": "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": "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] }
|
{ "id": "browser", "type": "external", "label": "Dashboard Users", "sublabel": "browser · partysocket", "pos": [960, 540], "size": [140, 60] }
|
||||||
|
|||||||
@@ -35,9 +35,11 @@
|
|||||||
|
|
||||||
# ---- Shared build tools ----
|
# ---- Shared build tools ----
|
||||||
nodejs = pkgs.nodejs_22;
|
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 HOME=$TMPDIR/home
|
||||||
export npm_config_cache=$TMPDIR/npm-cache
|
export npm_config_cache=$TMPDIR/npm-cache
|
||||||
mkdir -p $npm_config_cache
|
mkdir -p $npm_config_cache
|
||||||
@@ -48,15 +50,9 @@
|
|||||||
export GIT_SSL_CAINFO=${pkgs.cacert}/etc/ssl/certs/ca-bundle.crt
|
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
|
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
|
# Build native addons (bun install runs postinstall scripts for
|
||||||
# the rare case a prebuilt is unavailable and it falls back to compile).
|
# @discordjs/opus / sharp unless trustedDependencies restricts).
|
||||||
export CPPFLAGS="-I${pkgs.lib.getDev pkgs.openssl}/include"
|
bun install 2>&1
|
||||||
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
|
|
||||||
'';
|
'';
|
||||||
|
|
||||||
# Shrink the shipped node_modules to production deps only. The full
|
# 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.
|
# Must run AFTER tsc (typescript is a devDep) and after native builds.
|
||||||
pruneProd = ''
|
pruneProd = ''
|
||||||
echo "=== Pruning devDependencies (production-only node_modules) ==="
|
echo "=== Pruning devDependencies (production-only node_modules) ==="
|
||||||
pnpm list --prod --depth 999 --parseable 2>/dev/null \
|
# bun install's layout: node_modules/<pkg> for prod deps; devDeps are
|
||||||
| grep -o '\.pnpm/[^/]*' | sort -u > $TMPDIR/prod-pnms.txt
|
# also present during build (needed for tsc). Keep only what the prod
|
||||||
( cd node_modules/.pnpm \
|
# graph needs: simplest robust approach is `bun install --production`
|
||||||
&& for d in */; do \
|
# semantics — but bun keeps the same flat layout; since the Nix build
|
||||||
d="''${d%/}"; \
|
# already ran `bun install` (full, scripts on), prune dev-only top
|
||||||
[ "$d" = "node_modules" ] && continue; \
|
# entries that were only pulled by devDeps (typescript, biome, vitest,
|
||||||
grep -qF ".pnpm/$d" $TMPDIR/prod-pnms.txt || rm -rf "$d"; \
|
# drizzle-kit, tsx, @types/*).
|
||||||
done ) || true
|
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
|
||||||
# Drop runtime-dead packages that still land in the prod graph:
|
rm -rf node_modules/.bin/tsc node_modules/.bin/vitest node_modules/.bin/biome node_modules/.bin/drizzle-kit 2>/dev/null || true
|
||||||
# - `@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.
|
|
||||||
find node_modules -type l ! -exec test -e {} \; -delete 2>/dev/null || true
|
find node_modules -type l ! -exec test -e {} \; -delete 2>/dev/null || true
|
||||||
du -sh node_modules
|
du -sh node_modules
|
||||||
'';
|
'';
|
||||||
@@ -107,11 +90,11 @@
|
|||||||
|
|
||||||
src = ./services/backend;
|
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 ==="
|
echo "=== Compiling TypeScript ==="
|
||||||
npx tsc 2>&1
|
./node_modules/.bin/tsc 2>&1
|
||||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||||
node scripts/fix-imports.mjs
|
node scripts/fix-imports.mjs
|
||||||
echo "=== Build complete ==="
|
echo "=== Build complete ==="
|
||||||
@@ -149,7 +132,7 @@ WRAPPER
|
|||||||
# libvips download), so no cmake or rust toolchain is needed.
|
# libvips download), so no cmake or rust toolchain is needed.
|
||||||
# python3/gnumake/gcc stay as node-gyp fallback for @discordjs/opus.
|
# python3/gnumake/gcc stay as node-gyp fallback for @discordjs/opus.
|
||||||
nativeBuildInputs = [
|
nativeBuildInputs = [
|
||||||
nodejs pnpm
|
nodejs bun
|
||||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||||
pkgs.pkg-config
|
pkgs.pkg-config
|
||||||
pkgs.openssl
|
pkgs.openssl
|
||||||
@@ -177,22 +160,11 @@ WRAPPER
|
|||||||
# neither needed nor wanted here. Skip it entirely.
|
# neither needed nor wanted here. Skip it entirely.
|
||||||
dontFixup = true;
|
dontFixup = true;
|
||||||
|
|
||||||
buildPhase = pnpmInstall + ''
|
buildPhase = bunInstall + ''
|
||||||
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.
|
|
||||||
echo "=== Rebuilding @discordjs/opus (prebuilt download) ==="
|
echo "=== Rebuilding @discordjs/opus (prebuilt download) ==="
|
||||||
pnpm rebuild @discordjs/opus 2>&1 || true
|
bun pm rebuild @discordjs/opus 2>&1 || true
|
||||||
echo "=== Compiling TypeScript ===="
|
echo "=== Compiling TypeScript ===="
|
||||||
npx tsc 2>&1
|
./node_modules/.bin/tsc 2>&1
|
||||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||||
node scripts/fix-imports.mjs
|
node scripts/fix-imports.mjs
|
||||||
echo "=== Build complete ==="
|
echo "=== Build complete ==="
|
||||||
@@ -228,13 +200,13 @@ WRAPPER
|
|||||||
|
|
||||||
src = frontendSrc;
|
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) ==="
|
echo "=== Building Next.js SSR (standalone) ==="
|
||||||
export NEXT_TELEMETRY_DISABLED=1
|
export NEXT_TELEMETRY_DISABLED=1
|
||||||
export GMW_BACKEND_URL=http://127.0.0.1:4001
|
export GMW_BACKEND_URL=http://127.0.0.1:4001
|
||||||
npx next build 2>&1
|
./node_modules/.bin/next build 2>&1
|
||||||
'';
|
'';
|
||||||
|
|
||||||
installPhase = ''
|
installPhase = ''
|
||||||
@@ -312,13 +284,13 @@ WRAPPER
|
|||||||
|
|
||||||
devShells.default = pkgs.mkShell {
|
devShells.default = pkgs.mkShell {
|
||||||
buildInputs = [
|
buildInputs = [
|
||||||
nodejs pnpm
|
nodejs bun
|
||||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||||
pkgs.rustc pkgs.cargo
|
pkgs.rustc pkgs.cargo
|
||||||
pkgs.ffmpeg-headless
|
pkgs.ffmpeg-headless
|
||||||
];
|
];
|
||||||
shellHook = ''
|
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` |
|
| moderation | Moderation actions & metrics | `ai_moderations`, `moderation_actions` |
|
||||||
| media | Media file management | `media_attachments` |
|
| media | Media file management | `media_attachments` |
|
||||||
| dashboard | Stats aggregation | Various (read-only) |
|
| 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` |
|
| chatbot | AI chatbot with tools | `chatbot_history` |
|
||||||
| health | Health checks + metrics | Various |
|
| health | Health checks + metrics | Various |
|
||||||
| analysis | Text analysis cache | `text_analysis_cache` |
|
| 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,
|
"private": true,
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"main": "dist/index.js",
|
"main": "dist/index.js",
|
||||||
"packageManager": "pnpm@11.20.0",
|
"packageManager": "bun@1.3.14",
|
||||||
"engines": {
|
"engines": {
|
||||||
"node": ">=22.12.0",
|
"node": ">=22.12.0",
|
||||||
"pnpm": ">=9.0.0"
|
"pnpm": ">=9.0.0"
|
||||||
@@ -16,14 +16,15 @@
|
|||||||
"format": "biome format --write .",
|
"format": "biome format --write .",
|
||||||
"lint": "biome check --diagnostic-level=error .",
|
"lint": "biome check --diagnostic-level=error .",
|
||||||
"start": "node dist/index.js",
|
"start": "node dist/index.js",
|
||||||
"test": "vitest run",
|
"test": "bun test tests/",
|
||||||
"typecheck": "tsc --noEmit"
|
"typecheck": "tsc --noEmit",
|
||||||
|
"test:e2e": "bun test src/e2e.test.ts"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@orpc/server": "1.15.2",
|
"@orpc/server": "1.15.3",
|
||||||
"axios": "^1.20.0",
|
"axios": "^1.20.0",
|
||||||
"dotenv": "^18.0.1",
|
"dotenv": "^18.0.1",
|
||||||
"drizzle-orm": "^0.45.2",
|
"drizzle-orm": "^0.45.3",
|
||||||
"express": "^5.2.1",
|
"express": "^5.2.1",
|
||||||
"helmet": "^8.1.0",
|
"helmet": "^8.1.0",
|
||||||
"ioredis": "^6.0.0",
|
"ioredis": "^6.0.0",
|
||||||
@@ -41,6 +42,6 @@
|
|||||||
"@types/ws": "^8.18.1",
|
"@types/ws": "^8.18.1",
|
||||||
"tsx": "^4.23.15",
|
"tsx": "^4.23.15",
|
||||||
"typescript": "^7.0.2",
|
"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 MessageQuery = z.infer<typeof messageQuerySchema>;
|
||||||
export type MessageCreate = z.infer<typeof messageCreateSchema>;
|
export type MessageCreate = z.infer<typeof messageCreateSchema>;
|
||||||
export type MessageUpdate = z.infer<typeof messageUpdateSchema>;
|
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 { NotFoundError, ValidationError } from "@/shared/errors/index";
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
import { createChildLogger } from "@/shared/logger/index";
|
||||||
import { embedQuery } from "./embed.js";
|
|
||||||
import { type MessageRow, messagesRepository } from "./messages.repository.js";
|
import { type MessageRow, messagesRepository } from "./messages.repository.js";
|
||||||
import type { MessageQuery, SemanticSearchQuery } from "./messages.schema.js";
|
import type { MessageQuery } from "./messages.schema.js";
|
||||||
import { searchArchive } from "./qdrant.js";
|
|
||||||
|
|
||||||
const logger = createChildLogger("messages.service");
|
const logger = createChildLogger("messages.service");
|
||||||
|
|
||||||
@@ -102,32 +99,6 @@ export class MessagesService {
|
|||||||
return messagesRepository.getReviewMessages(channelId, limit);
|
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(
|
async getActivity(
|
||||||
days = 30,
|
days = 30,
|
||||||
): Promise<Awaited<ReturnType<typeof messagesRepository.getActivity>>> {
|
): 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();
|
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 { chatbotService } from "../modules/chatbot/chatbot.service";
|
||||||
import { dashboardService } from "../modules/dashboard/dashboard.service";
|
import { dashboardService } from "../modules/dashboard/dashboard.service";
|
||||||
import { knowledgeService } from "../modules/knowledge/knowledge.service";
|
import { knowledgeService } from "../modules/knowledge/knowledge.service";
|
||||||
import {
|
import { messageQuerySchema } from "../modules/messages/messages.schema";
|
||||||
messageQuerySchema,
|
|
||||||
semanticSearchSchema,
|
|
||||||
} from "../modules/messages/messages.schema";
|
|
||||||
import { messagesService } from "../modules/messages/messages.service";
|
import { messagesService } from "../modules/messages/messages.service";
|
||||||
import { moderationService } from "../modules/moderation/moderation.service";
|
import { moderationService } from "../modules/moderation/moderation.service";
|
||||||
import { uiStateService } from "../modules/ui-state/ui-state.service";
|
import { uiStateService } from "../modules/ui-state/ui-state.service";
|
||||||
@@ -123,10 +120,6 @@ const messagesRouter = {
|
|||||||
);
|
);
|
||||||
return { results: rows, limit: input.limit, cursor: null };
|
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).
|
// Public, read-only activity heatmap data (per-hour volume by channel).
|
||||||
activity: os
|
activity: os
|
||||||
.input(
|
.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")
|
.transform((v) => v === "true")
|
||||||
.default(false),
|
.default(false),
|
||||||
AI_LLM_API_KEY: z.string().optional(),
|
AI_LLM_API_KEY: z.string().optional(),
|
||||||
AI_LLM_BASE_URL: z
|
// 9router — OpenAI-compatible router on this host (127.0.0.1:4014).
|
||||||
.string()
|
// Loopback on purpose: backend runs on the same machine as 9router, so no
|
||||||
.url()
|
// TLS/proxy hop is needed.
|
||||||
.default("http://100.121.180.82:20128/api/v1"),
|
AI_LLM_BASE_URL: z.string().url().default("http://127.0.0.1:4014/v1"),
|
||||||
AI_LLM_MODEL: z.string().default("text"),
|
AI_LLM_MODEL: z.string().default("text"),
|
||||||
AI_LLM_VISION_MODEL: z.string().optional(),
|
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_MAX_CONCURRENT: z.coerce.number().int().positive().default(5),
|
||||||
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
||||||
.number()
|
.number()
|
||||||
@@ -179,11 +170,6 @@ export const configSchema = z
|
|||||||
.default("https://api.openai.com/v1"),
|
.default("https://api.openai.com/v1"),
|
||||||
OPENAI_MODERATION_MODEL: z.string().default("omni-moderation-latest"),
|
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 ─────────────────────────────────────────────────────
|
||||||
AUTO_DELETE_FLAGGED_ENABLED: z
|
AUTO_DELETE_FLAGGED_ENABLED: z
|
||||||
.string()
|
.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";
|
import { tools } from "../src/modules/chatbot/chatbot.toolDefs.js";
|
||||||
|
|
||||||
const names = tools.map((t) => t.function.name);
|
const names = tools.map((t) => t.function.name);
|
||||||
|
|||||||
@@ -1,6 +1,32 @@
|
|||||||
// ─── Shared Error Classes ────────────────────────────────────────────────────
|
// ─── 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 {
|
import {
|
||||||
AppError,
|
AppError,
|
||||||
ConfigError,
|
ConfigError,
|
||||||
@@ -94,41 +120,41 @@ describe("AppError subclasses", () => {
|
|||||||
// ═══════════════════════════════════════════════════════════════════════════════
|
// ═══════════════════════════════════════════════════════════════════════════════
|
||||||
describe("delay", () => {
|
describe("delay", () => {
|
||||||
afterEach(() => {
|
afterEach(() => {
|
||||||
vi.useRealTimers();
|
useRealTimers();
|
||||||
});
|
});
|
||||||
|
|
||||||
it("resolves after the given time", async () => {
|
it("resolves after the given time", async () => {
|
||||||
vi.useFakeTimers();
|
useFakeTimers();
|
||||||
const promise = delay(500);
|
const promise = delay(500);
|
||||||
vi.advanceTimersByTime(500);
|
advanceTimersByTime(500);
|
||||||
await expect(promise).resolves.toBeUndefined();
|
await expect(promise).resolves.toBeUndefined();
|
||||||
});
|
});
|
||||||
|
|
||||||
it("rejects are not triggered on non-matching timer", async () => {
|
it("rejects are not triggered on non-matching timer", async () => {
|
||||||
vi.useFakeTimers();
|
useFakeTimers();
|
||||||
const promise = delay(1000);
|
const promise = delay(1000);
|
||||||
// Advance only part way — the timer should NOT fire yet
|
// 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
|
// The timer is still pending; the promise has not resolved yet
|
||||||
// We advance the rest
|
// We advance the rest
|
||||||
vi.advanceTimersByTime(500);
|
advanceTimersByTime(500);
|
||||||
await expect(promise).resolves.toBeUndefined();
|
await expect(promise).resolves.toBeUndefined();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("retryWithBackoff", () => {
|
describe("retryWithBackoff", () => {
|
||||||
afterEach(() => {
|
afterEach(() => {
|
||||||
vi.useRealTimers();
|
useRealTimers();
|
||||||
});
|
});
|
||||||
|
|
||||||
it("returns the result on first success without retrying", async () => {
|
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");
|
await expect(retryWithBackoff(fn)).resolves.toBe("ok");
|
||||||
expect(fn).toHaveBeenCalledTimes(1);
|
expect(fn).toHaveBeenCalledTimes(1);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("re-throws after exhausting all retries", async () => {
|
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(
|
await expect(
|
||||||
retryWithBackoff(fn, { retries: 1, minTimeout: 1, maxTimeout: 5 }),
|
retryWithBackoff(fn, { retries: 1, minTimeout: 1, maxTimeout: 5 }),
|
||||||
).rejects.toThrow("persistent");
|
).rejects.toThrow("persistent");
|
||||||
@@ -139,7 +165,7 @@ describe("retryWithBackoff", () => {
|
|||||||
it("throws AbortError immediately when signal is already aborted", async () => {
|
it("throws AbortError immediately when signal is already aborted", async () => {
|
||||||
const ac = new AbortController();
|
const ac = new AbortController();
|
||||||
ac.abort();
|
ac.abort();
|
||||||
const fn = vi.fn().mockResolvedValue("ok");
|
const fn = jest.fn().mockResolvedValue("ok");
|
||||||
await expect(
|
await expect(
|
||||||
retryWithBackoff(fn, { retries: 3, signal: ac.signal }),
|
retryWithBackoff(fn, { retries: 3, signal: ac.signal }),
|
||||||
).rejects.toThrow("Aborted");
|
).rejects.toThrow("Aborted");
|
||||||
@@ -147,9 +173,9 @@ describe("retryWithBackoff", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("respects abort signal during retry", async () => {
|
it("respects abort signal during retry", async () => {
|
||||||
vi.useFakeTimers();
|
useFakeTimers();
|
||||||
const ac = new AbortController();
|
const ac = new AbortController();
|
||||||
const fn = vi.fn().mockRejectedValue(new Error("fail"));
|
const fn = jest.fn().mockRejectedValue(new Error("fail"));
|
||||||
|
|
||||||
const promise = retryWithBackoff(fn, {
|
const promise = retryWithBackoff(fn, {
|
||||||
retries: 5,
|
retries: 5,
|
||||||
@@ -159,8 +185,8 @@ describe("retryWithBackoff", () => {
|
|||||||
|
|
||||||
// Schedule abort after first failure + backoff starts
|
// Schedule abort after first failure + backoff starts
|
||||||
setTimeout(() => ac.abort(), 150);
|
setTimeout(() => ac.abort(), 150);
|
||||||
vi.advanceTimersByTime(200);
|
advanceTimersByTime(200);
|
||||||
await vi.waitFor(async () => {
|
await waitForCompat(async () => {
|
||||||
await expect(promise).rejects.toThrow("Aborted");
|
await expect(promise).rejects.toThrow("Aborted");
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
@@ -232,7 +258,7 @@ describe("asyncHandler", () => {
|
|||||||
const wrapped = asyncHandler(async () => {
|
const wrapped = asyncHandler(async () => {
|
||||||
throw error;
|
throw error;
|
||||||
});
|
});
|
||||||
const next = vi.fn();
|
const next = jest.fn();
|
||||||
|
|
||||||
wrapped({} as any, {} as any, next);
|
wrapped({} as any, {} as any, next);
|
||||||
|
|
||||||
@@ -246,7 +272,7 @@ describe("asyncHandler", () => {
|
|||||||
const wrapped = asyncHandler(async (_req: any, _res: any, _next: any) => {
|
const wrapped = asyncHandler(async (_req: any, _res: any, _next: any) => {
|
||||||
// no-op
|
// no-op
|
||||||
});
|
});
|
||||||
const next = vi.fn();
|
const next = jest.fn();
|
||||||
|
|
||||||
wrapped({} as any, {} as any, next);
|
wrapped({} as any, {} as any, next);
|
||||||
await Promise.resolve();
|
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
|
* 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.
|
1. **LLM is the only judge.** Failed LLM → `status:"error"` + recovery retry.
|
||||||
**Never** reintroduce regex/heuristic content classification.
|
**Never** reintroduce regex/heuristic content classification.
|
||||||
2. **Discord tokens sanitized** before reaching LLM (`discordTokens.ts`).
|
2. **Discord tokens sanitized** before reaching LLM (`discordTokens.ts`).
|
||||||
3. **Semantic cache is batched** — one embed call + one Qdrant batch search.
|
3. **Streaming is mandatory** against the router base URL.
|
||||||
4. **Streaming is mandatory** against the omniroute base URL.
|
|
||||||
|
|
||||||
## AI moderation pipeline
|
## AI moderation pipeline
|
||||||
|
|
||||||
```
|
```
|
||||||
aiAnalyzer.ts → batchScheduler.ts → batchProcessor.ts → individualFallbackProcessor.ts
|
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
|
textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||||
embeddingClient.ts
|
|
||||||
qdrantClient.ts
|
|
||||||
```
|
```
|
||||||
|
|
||||||
- Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`, `startPendingAIAnalysisWorker`)
|
- 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)
|
- 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`)
|
- **Each worker thread has its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`)
|
||||||
|
|
||||||
## Module: message-capture
|
## Module: message-capture
|
||||||
@@ -79,7 +82,6 @@ textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
|||||||
- `messageStore.ts` — DB operations
|
- `messageStore.ts` — DB operations
|
||||||
- `messageMetadata.ts` — metadata extraction
|
- `messageMetadata.ts` — metadata extraction
|
||||||
- `messagesDb.ts` / `messagesCrud.ts` — DB schema operations
|
- `messagesDb.ts` / `messagesCrud.ts` — DB schema operations
|
||||||
- `archiveEmbedder.ts` — Qdrant embedding (respect age-restricted guard)
|
|
||||||
- `retentionDb.ts` / `reviewsDb.ts` / `attachmentsDb.ts` — auxiliary tables
|
- `retentionDb.ts` / `reviewsDb.ts` / `attachmentsDb.ts` — auxiliary tables
|
||||||
|
|
||||||
## Redis channels (outbound to backend)
|
## 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
|
moderation, and publishes everything to Redis pub/sub for the backend to
|
||||||
consume. The backend serves the HTTP/WS API to the frontend.
|
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
|
> NOTE: this doc is the source of truth for the module layout. The old
|
||||||
> `MODULE_STRUCTURE.md` was stale (referenced `winston`, `mock-crc.ts`,
|
> `MODULE_STRUCTURE.md` was a stale duplicate and has been removed. `README.md`
|
||||||
> `indonesianTextNormalizer.ts`, and `aiAnalysisWorker.ts`/`llmModerationClient.ts`
|
> only covers how to run the service.
|
||||||
> which were renamed/merged). If they disagree, this file wins.
|
|
||||||
|
|
||||||
## Top-level layout
|
## Top-level layout
|
||||||
|
|
||||||
@@ -16,72 +15,99 @@ consume. The backend serves the HTTP/WS API to the frontend.
|
|||||||
services/discord-gateway/
|
services/discord-gateway/
|
||||||
├── src/
|
├── src/
|
||||||
│ ├── index.ts # Entry point → initializeDiscordGateway()
|
│ ├── index.ts # Entry point → initializeDiscordGateway()
|
||||||
│ ├── app/
|
│ ├── app/ # Process lifecycle
|
||||||
│ │ ├── bootstrap.ts # Wires client, DB, Redis, workers, schedulers
|
│ │ ├── bootstrap.ts # Startup order: config → DB → services → metrics → login
|
||||||
│ │ ├── shutdown.ts # Graceful shutdown (SIGINT/SIGTERM + transient errors)
|
│ │ ├── 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
|
│ │ └── retention.ts # Expired-record cleanup scheduler
|
||||||
│ ├── shared/
|
│ ├── shared/ # Infrastructure — never imports from modules/
|
||||||
│ │ ├── config/ # Zod-validated env (index.ts = schema+loader)
|
│ │ ├── config/ # Zod-validated env (index.ts = schema+loader)
|
||||||
│ │ ├── database/ # Drizzle ORM + pg Pool + migrations
|
│ │ ├── database/ # Drizzle ORM + pg Pool + migrations
|
||||||
│ │ │ ├── init.ts drizzle.ts pool.ts migrate.ts migrateCli.ts
|
│ │ │ ├── init.ts drizzle.ts pool.ts migrate.ts migrateCli.ts
|
||||||
│ │ │ └── schema/ # messages, cache, meta, analytics
|
│ │ │ └── schema/ # messages, cache, meta, analytics
|
||||||
│ │ ├── logger/ # pino wrapper + createChildLogger()
|
│ │ ├── logger/ # pino wrapper + createChildLogger()
|
||||||
│ │ ├── errors/ # AppError / ConfigError ...
|
│ │ ├── errors/ # AppError / ConfigError ... + errorMessage()
|
||||||
|
│ │ │ # + isTransientStreamError()
|
||||||
│ │ ├── utils/ # retry, pagination
|
│ │ ├── utils/ # retry, pagination
|
||||||
│ │ ├── discord/clientOptions.ts # discord.js-selfbot-v13 client options
|
│ │ ├── discord/clientOptions.ts # discord.js-selfbot-v13 client options
|
||||||
│ │ ├── uploader.ts # Shared attachment upload helper
|
│ │ ├── 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
|
│ │ └── 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
|
│ ├── message-capture/ # Discord event listeners + DB store
|
||||||
│ ├── ai-moderation/ # LLM moderation pipeline (see below)
|
│ ├── ai-moderation/ # LLM moderation pipeline (see below)
|
||||||
│ ├── attachment-upload/ # Download + (sharp) resize + upload
|
│ ├── attachment-upload/ # Download + (sharp) resize + upload
|
||||||
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
||||||
│ ├── command-handler/ # Redis-subscribed backend→gateway commands
|
│ ├── command-handler/ # Redis-subscribed backend→gateway commands
|
||||||
│ ├── reaction-tracking/ thread-tracking/ user-presence/
|
│ ├── reaction-tracking/ thread-tracking/ user-presence/
|
||||||
│ ├── channel-topic/ guild-member-events/
|
│ ├── channel-topic/ guild-member-events/ monitor/
|
||||||
│ └── gateway-metrics/ # Prometheus /metrics endpoint (port 4016)
|
│ └── 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/`)
|
## AI moderation pipeline (`ai-moderation/`)
|
||||||
|
|
||||||
LLM-only judge — no regex/heuristic classification. One orchestrator call
|
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`,
|
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `queueConversationAnalysis`,
|
||||||
`startPendingAIAnalysisWorker` (recovery worker + cache-prune).
|
`getAnalysisQueueStatus`, `startPendingAIAnalysisWorker`. Short-circuits
|
||||||
- `batchScheduler.ts` — per-conversation debounce → `processBatch`.
|
age-restricted and skip-list messages before any LLM work.
|
||||||
- `batchProcessor.ts` — batch lock/circuit-breaker, fans failed targets to
|
- `recovery-worker.ts` — periodic sweep for stranded `pending` messages
|
||||||
individual fallback.
|
(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.
|
- `individualFallbackProcessor.ts` — one-message-at-a-time retry path, own CB.
|
||||||
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation state,
|
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation PER-LANE
|
||||||
Piscina `workerPool`, `getConversationKey`.
|
state (`conversationProcessing` holds a lane → startedAt map per key),
|
||||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` / `individual` jobs).
|
Piscina `textWorkerPool`/`mediaWorkerPool`, `getConversationKey`.
|
||||||
Runs `runModerationAnalysis` off the main thread.
|
- `ai-analysis-worker.ts` — Piscina entry point (`batch` (lane) /
|
||||||
- `moderationOrchestrator.ts` — exact-hash cache → batched semantic (Qdrant)
|
`individual` jobs). Runs `runModerationAnalysis` off the main thread.
|
||||||
cache → LLM. Text and media paths run in parallel.
|
- `moderationOrchestrator.ts` — exact-hash cache → LLM. Text and media paths
|
||||||
|
run in parallel.
|
||||||
- `textBatchProcessor.ts` / `mediaBatchProcessor.ts` — actual LLM calls
|
- `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,
|
- `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).
|
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` /
|
- `textCacheStore.ts` / `channelCultureStore.ts` / `userProfileStore.ts` /
|
||||||
`userProfileStore.ts` — caches learned user profile summaries (optional).
|
`userProfileStore.ts` — caches learned user profile summaries (optional).
|
||||||
|
|
||||||
### Concurrency model
|
### Concurrency model
|
||||||
|
|
||||||
- Main thread owns the LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5) via
|
- Main thread owns TWO per-lane LLM semaphores (2026-09-24):
|
||||||
`llmClient.withLlmConcurrency`.
|
`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
|
- 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_MAX_THREADS`, default 4) and a dedicated media pool
|
||||||
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed to the media
|
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed by lane to the
|
||||||
pool if ANY of its messages carries an attachment/sticker/embed — this
|
matching pool — this keeps a slow image/vision batch from occupying every
|
||||||
keeps a slow image/vision batch from occupying every thread and blocking
|
thread and blocking unrelated text-only batches behind it. **Each worker
|
||||||
unrelated text-only batches behind it. **Each worker thread (in either
|
thread (in either pool) initializes its own pg Pool** (min 0, grows to
|
||||||
pool) initializes its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`).
|
`POSTGRES_POOL_MAX`). See "Memory & connections" below.
|
||||||
See "Memory & connections" below.
|
|
||||||
|
|
||||||
## Memory & DB connections
|
## Memory & DB connections
|
||||||
|
|
||||||
@@ -109,28 +135,53 @@ See `src/shared/redis-channels.ts` for the canonical names.
|
|||||||
|
|
||||||
## Initialization flow
|
## 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.
|
1. Validate env (Zod). Refuse to start if `AI_ANALYSIS_ENABLED` but no key.
|
||||||
2. `AUTO_MIGRATE_ON_STARTUP` → run pending Drizzle migrations.
|
→ `assertConfigIsUsable()`
|
||||||
3. `initializeDatabase()` (pg Pool, min 0).
|
2. Build long-lived services: Discord client, `RedisEventPublisher` +
|
||||||
4. Create discord.js-selfbot-v13 client; register listeners on `ready`.
|
`EventBroadcaster`, `CommandHandler`; install the shutdown handler.
|
||||||
5. Start `gmw-discord-gateway` metrics server (port `METRICS_PORT`, default 4016).
|
3. Connect infrastructure → `connectDatabase()`:
|
||||||
6. `client.login(token)`.
|
`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
|
## Graceful shutdown
|
||||||
|
|
||||||
`SIGINT`/`SIGTERM` (and uncaught transient stream errors: EPIPE / ECONNRESET /
|
`process-guards.ts` owns the policy. `SIGINT`/`SIGTERM` and non-transient
|
||||||
ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END are treated as non-fatal):
|
uncaught exceptions/rejections run `shutdown.ts`; transient stream errors
|
||||||
stop metrics → close event broadcaster (Redis) → close command handler →
|
(EPIPE / ECONNRESET / ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END, see
|
||||||
close DB → destroy client → exit.
|
`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
|
## 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
|
per-scrape and expose: process memory/uptime, and (when AI analysis is on) live
|
||||||
pipeline gauges — `ai_analysis_queued_conversations`,
|
pipeline gauges registered by `app/metrics-collector.ts` —
|
||||||
`ai_analysis_active_batch_requests`, `ai_analysis_active_individual_requests`,
|
`ai_analysis_queued_conversations`, `ai_analysis_active_batch_requests`,
|
||||||
`ai_analysis_individual_in_flight`, `ai_analysis_individual_circuit_breaker_active`,
|
`ai_analysis_active_text_requests`, `ai_analysis_active_media_requests`,
|
||||||
`ai_analysis_worker_threads`, `ai_analysis_worker_threads_active`.
|
`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)
|
## Key invariants (do not break)
|
||||||
|
|
||||||
@@ -139,7 +190,5 @@ pipeline gauges — `ai_analysis_queued_conversations`,
|
|||||||
- **Discord tokens are sanitized** (`discordTokens.ts`: `<:emoji:id>` →
|
- **Discord tokens are sanitized** (`discordTokens.ts`: `<:emoji:id>` →
|
||||||
`[emoji:name]`, `<@id>` → `@user`, etc.) before content reaches the LLM, so
|
`[emoji:name]`, `<@id>` → `@user`, etc.) before content reaches the LLM, so
|
||||||
numeric snowflake IDs never trigger false positives.
|
numeric snowflake IDs never trigger false positives.
|
||||||
- **Semantic cache is batched** (one embed call + one Qdrant batch search),
|
- **Streaming is mandatory** against the router base URL (non-stream waits for
|
||||||
not N sequential round-trips. `ensureQdrantCollection` is memoized.
|
|
||||||
- **Streaming is mandatory** against the omniroute base URL (non-stream waits for
|
|
||||||
the full body and times out). `llmClient` aggregates SSE chunks.
|
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
|
||||||
|
|
||||||
```
|
```bash
|
||||||
services/discord-gateway/
|
pnpm install
|
||||||
├── src/
|
pnpm typecheck # tsc --noEmit
|
||||||
│ ├── app/
|
pnpm lint # biome check --diagnostic-level=error .
|
||||||
│ │ ├── bootstrap.ts # Service initialization (Discord client, DB, Redis)
|
pnpm test # vitest run (138 tests)
|
||||||
│ │ └── shutdown.ts # Graceful shutdown handler
|
pnpm build # tsc — CI/prod builds run this inside nix, which also
|
||||||
│ ├── shared/ # Shared infrastructure layer
|
# runs scripts/fix-imports.mjs to rewrite @/ aliases and
|
||||||
│ │ ├── config/
|
# extensionless imports for Node ESM
|
||||||
│ │ │ └── config.ts # Zod-validated environment config
|
pnpm dev # tsx watch src/index.ts
|
||||||
│ │ ├── database/
|
pnpm start # node dist/index.js
|
||||||
│ │ │ ├── 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)
|
|
||||||
```
|
```
|
||||||
|
|
||||||
## 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
|
## Layout
|
||||||
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:
|
|
||||||
|
|
||||||
```
|
```
|
||||||
Discord Events → Discord Gateway Service → Redis Pub/Sub → Backend Service
|
src/
|
||||||
↓
|
├── index.ts # entry → initializeDiscordGateway()
|
||||||
Event Channels:
|
├── app/ # process lifecycle
|
||||||
- discord:message:created
|
│ ├── bootstrap.ts # startup order: config → DB → services → metrics → login
|
||||||
- discord:message:updated
|
│ ├── lifecycle.ts # everything wired on the Discord 'ready' hook
|
||||||
- discord:message:deleted
|
│ ├── process-guards.ts # SIGINT/SIGTERM + uncaught error policy
|
||||||
- discord:message:analyzed
|
│ ├── metrics-collector.ts # AI pipeline Prometheus gauges
|
||||||
- discord:attachment:created
|
│ ├── shutdown.ts # graceful shutdown sequence
|
||||||
- discord:attachment:uploaded
|
│ └── retention.ts # expired-record cleanup scheduler
|
||||||
- discord:analysis:queue_status
|
├── 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
|
Dependency direction is one-way: `index.ts` → `app/` → `modules/` → `shared/`.
|
||||||
Centralized, reusable components:
|
Callers outside a module import its `index.ts` facade, never an internal file.
|
||||||
- **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
|
|
||||||
|
|
||||||
### 4. No HTTP Server
|
## Testing
|
||||||
- **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
|
|
||||||
|
|
||||||
## Key Features
|
Vitest, tests in `tests/`. Config supplies dummy env vars so the suite runs
|
||||||
|
without live Postgres/Redis; external services are mocked. `llmE2e.test.ts`
|
||||||
### Message Capture
|
is skipped by default and needs real credentials (`pnpm test:e2e:live`).
|
||||||
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.
|
|
||||||
|
|||||||
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",
|
"typecheck": "tsc --noEmit",
|
||||||
"lint": "biome check --diagnostic-level=error .",
|
"lint": "biome check --diagnostic-level=error .",
|
||||||
"format": "biome format --write .",
|
"format": "biome format --write .",
|
||||||
"test": "vitest run",
|
"test": "bun test tests/",
|
||||||
"test:unit": "vitest run --exclude \"tests/llmE2e.test.ts\"",
|
"test:e2e": "bun test tests/llmE2e.test.ts"
|
||||||
"test:e2e": "vitest run tests/llmE2e.test.ts",
|
|
||||||
"test:e2e:live": "bash scripts/run-llm-e2e.sh"
|
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@typesafe-ai/sdk": "^0.6.0",
|
|
||||||
"axios": "^1.20.0",
|
"axios": "^1.20.0",
|
||||||
"discord.js-selfbot-v13": "^3.7.1",
|
"discord.js-selfbot-v13": "^3.7.1",
|
||||||
"dotenv": "^18.0.0",
|
"dotenv": "^18.0.0",
|
||||||
@@ -44,12 +41,13 @@
|
|||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@biomejs/biome": "latest",
|
"@biomejs/biome": "latest",
|
||||||
"@types/node": "^26.4.0",
|
"@types/node": "^26.6.2",
|
||||||
"@types/pg": "^8.23.1",
|
"@types/pg": "^8.23.1",
|
||||||
"@types/ws": "^8.18.1",
|
"@types/ws": "^8.18.1",
|
||||||
"drizzle-kit": "^0.31.10",
|
"drizzle-kit": "^0.31.11",
|
||||||
"tsx": "^4.23.13",
|
"tsx": "^4.23.15",
|
||||||
"typescript": "^7.0.2",
|
"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 { Client } from "discord.js-selfbot-v13";
|
||||||
import { ConfigError, DatabaseError } from "@/shared/errors/index";
|
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
|
||||||
import {
|
import {
|
||||||
getAnalysisQueueStatus,
|
ConfigError,
|
||||||
startPendingAIAnalysisWorker,
|
DatabaseError,
|
||||||
} from "../modules/ai-moderation/aiAnalyzer.js";
|
errorMessage,
|
||||||
import {
|
} from "@/shared/errors/index.js";
|
||||||
mediaWorkerPool,
|
import { createChildLogger } from "@/shared/logger/index.js";
|
||||||
textWorkerPool,
|
|
||||||
} from "../modules/ai-moderation/circuitBreaker.js";
|
|
||||||
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
|
|
||||||
import { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
import { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||||
import {
|
import {
|
||||||
EventBroadcaster,
|
EventBroadcaster,
|
||||||
RedisEventPublisher,
|
RedisEventPublisher,
|
||||||
} from "../modules/event-broadcaster/index.js";
|
} from "../modules/event-broadcaster/index.js";
|
||||||
import {
|
import {
|
||||||
registerCollector,
|
|
||||||
setGauge,
|
|
||||||
startMetricsServer,
|
startMetricsServer,
|
||||||
stopMetricsServer,
|
stopMetricsServer,
|
||||||
} from "../modules/gateway-metrics/index.js";
|
} from "../modules/gateway-metrics/index.js";
|
||||||
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
|
import { config } from "../shared/config/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 {
|
import {
|
||||||
closeDatabase,
|
closeDatabase,
|
||||||
initializeDatabase,
|
initializeDatabase,
|
||||||
} from "../shared/database/drizzle.js";
|
} from "../shared/database/drizzle.js";
|
||||||
import { runMigrations } from "../shared/database/migrate.js";
|
import { runMigrations } from "../shared/database/migrate.js";
|
||||||
import { createDiscordClientOptions } from "../shared/discord/clientOptions.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";
|
import { createGracefulShutdown } from "./shutdown.js";
|
||||||
|
|
||||||
const logger = createChildLogger("discord-gateway");
|
const logger = createChildLogger("discord-gateway");
|
||||||
|
|
||||||
// ─── Bootstrap ─────────────────────────────────────────────────────────────
|
// ─── 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) {
|
if (config.AI_ANALYSIS_ENABLED && !config.AI_LLM_API_KEY) {
|
||||||
throw new ConfigError(
|
throw new ConfigError(
|
||||||
"AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.",
|
"AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.",
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
const token = config.DISCORD_TOKEN;
|
/** Run migrations (when enabled) then open the PostgreSQL pool. */
|
||||||
logger.info(
|
async function connectDatabase(): Promise<void> {
|
||||||
{ 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,
|
|
||||||
});
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
if (config.AUTO_MIGRATE_ON_STARTUP) {
|
if (config.AUTO_MIGRATE_ON_STARTUP) {
|
||||||
logger.info(
|
logger.info(
|
||||||
@@ -90,175 +62,79 @@ export async function initializeDiscordGateway() {
|
|||||||
logger.info("PostgreSQL database initialized");
|
logger.info("PostgreSQL database initialized");
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
logger.error(
|
logger.error(
|
||||||
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
|
{ err, errorMsg: errorMessage(err) },
|
||||||
"Failed to initialize database",
|
"Failed to initialize database",
|
||||||
);
|
);
|
||||||
throw new DatabaseError(
|
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) => {
|
client.on("debug", (msg) => {
|
||||||
if (
|
const lower = msg.toLowerCase();
|
||||||
msg.toLowerCase().includes("error") ||
|
if (lower.includes("error") || lower.includes("stream")) {
|
||||||
msg.toLowerCase().includes("stream")
|
|
||||||
) {
|
|
||||||
logger.info({ debugMsg: msg }, "Discord Client Debug");
|
logger.info({ debugMsg: msg }, "Discord Client Debug");
|
||||||
} else if (config.VERBOSE) {
|
} else if (config.VERBOSE) {
|
||||||
logger.debug({ debugMsg: msg }, "Discord Client Debug");
|
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");
|
logger.info({ user: client.user?.tag }, "Bot logged in");
|
||||||
setMessageCaptureEventBroadcaster(eventBroadcaster);
|
startGatewayLifecycle({
|
||||||
setModerationEventBroadcaster(eventBroadcaster);
|
client,
|
||||||
registerMessageCapture(client);
|
eventBroadcaster,
|
||||||
startPendingAIAnalysisWorker(client, eventBroadcaster);
|
commandHandler,
|
||||||
|
logger,
|
||||||
// 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();
|
|
||||||
});
|
});
|
||||||
|
|
||||||
client.on("error", (err) => {
|
client.on("error", (err) => {
|
||||||
logger.error(
|
logger.error({ err, errorMsg: errorMessage(err) }, "Client error");
|
||||||
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
|
|
||||||
"Client error",
|
|
||||||
);
|
|
||||||
});
|
});
|
||||||
|
|
||||||
process.on("SIGINT", () => {
|
registerProcessGuards(logger, gracefulShutdown);
|
||||||
gracefulShutdown("SIGINT");
|
|
||||||
});
|
|
||||||
|
|
||||||
process.on("SIGTERM", () => {
|
// Metrics: register live pipeline collectors before starting the server.
|
||||||
gracefulShutdown("SIGTERM");
|
registerPipelineMetrics(logger);
|
||||||
});
|
|
||||||
|
|
||||||
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
|
|
||||||
startMetricsServer();
|
startMetricsServer();
|
||||||
|
|
||||||
logger.info("Calling Discord client.login");
|
logger.info("Calling Discord client.login");
|
||||||
|
|
||||||
// Fix: use await + try/catch instead of .then().catch()
|
|
||||||
try {
|
try {
|
||||||
await client.login(token);
|
await client.login(token);
|
||||||
logger.info("Discord client logged in successfully");
|
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 { inArray, lt } from "drizzle-orm";
|
||||||
import type {
|
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
|
||||||
NodePgDatabase,
|
|
||||||
NodePgQueryResultHKT,
|
|
||||||
} from "drizzle-orm/node-postgres";
|
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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 { getDatabase } from "../shared/database/drizzle.js";
|
||||||
import type * as schema from "../shared/database/schema.js";
|
import type * as schema from "../shared/database/schema.js";
|
||||||
import { attachmentsTable, messagesTable } 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 { Client } from "discord.js-selfbot-v13";
|
||||||
import type { createChildLogger } from "@/shared/logger/index";
|
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 { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||||
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
|
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
|
||||||
import type { stopMetricsServer } from "../modules/gateway-metrics/index.js";
|
import type { stopMetricsServer } from "../modules/gateway-metrics/index.js";
|
||||||
@@ -45,6 +49,30 @@ export function createGracefulShutdown(
|
|||||||
options.logger.info("Closing command handler...");
|
options.logger.info("Closing command handler...");
|
||||||
await options.commandHandler.close();
|
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
|
// 2. DB pool
|
||||||
options.logger.info("Closing database...");
|
options.logger.info("Closing database...");
|
||||||
await options.closeDatabase();
|
await options.closeDatabase();
|
||||||
|
|||||||
@@ -17,7 +17,7 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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 { initializeDatabase } from "../../shared/database/drizzle.js";
|
||||||
import { messageStore } from "../message-capture/messageStore.js";
|
import { messageStore } from "../message-capture/messageStore.js";
|
||||||
import type { MessageRecord } from "../message-capture/types.js";
|
import type { MessageRecord } from "../message-capture/types.js";
|
||||||
@@ -86,7 +86,12 @@ export interface MessageBatch {
|
|||||||
|
|
||||||
// Worker job types (Piscina entry point)
|
// Worker job types (Piscina entry point)
|
||||||
type WorkerJob =
|
type WorkerJob =
|
||||||
| { type: "batch"; conversationKey: string; messages: MessageRecord[] }
|
| {
|
||||||
|
type: "batch";
|
||||||
|
conversationKey: string;
|
||||||
|
lane: "text" | "media";
|
||||||
|
messages: MessageRecord[];
|
||||||
|
}
|
||||||
| { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean };
|
| { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean };
|
||||||
|
|
||||||
type BatchOkResponse = {
|
type BatchOkResponse = {
|
||||||
@@ -263,6 +268,7 @@ function normalizeResult(
|
|||||||
async function processBatch(job: {
|
async function processBatch(job: {
|
||||||
type: "batch";
|
type: "batch";
|
||||||
conversationKey: string;
|
conversationKey: string;
|
||||||
|
lane: "text" | "media";
|
||||||
messages: MessageRecord[];
|
messages: MessageRecord[];
|
||||||
}): Promise<BatchOkResponse | BatchErrorResponse> {
|
}): Promise<BatchOkResponse | BatchErrorResponse> {
|
||||||
const { conversationKey, messages } = job;
|
const { conversationKey, messages } = job;
|
||||||
|
|||||||
@@ -1,34 +1,25 @@
|
|||||||
import type { Client } from "discord.js-selfbot-v13";
|
import type { Client } from "discord.js-selfbot-v13";
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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 type { EventBroadcaster } from "../event-broadcaster/index.js";
|
||||||
import { messageStore } from "../message-capture/messageStore.js";
|
import { messageStore } from "../message-capture/messageStore.js";
|
||||||
import type { AnalysisQueueStatus } from "../message-capture/types.js";
|
import type { AnalysisQueueStatus } from "../message-capture/types.js";
|
||||||
import {
|
import {
|
||||||
|
activeMediaRequests,
|
||||||
activeRequests,
|
activeRequests,
|
||||||
|
activeTextRequests,
|
||||||
buildAgeRestrictedSkipResult,
|
buildAgeRestrictedSkipResult,
|
||||||
buildSkipAnalysisUserResult,
|
buildSkipAnalysisUserResult,
|
||||||
isAgeRestrictedMessage,
|
isAgeRestrictedMessage,
|
||||||
isSkipAnalysisUser,
|
isSkipAnalysisUser,
|
||||||
skipAgeRestrictedMessages,
|
|
||||||
skipAnalysisUserMessages,
|
|
||||||
} from "./batchProcessor.js";
|
} from "./batchProcessor.js";
|
||||||
import { scheduleConversationAnalysis } from "./batchScheduler.js";
|
import { scheduleConversationAnalysis } from "./batchScheduler.js";
|
||||||
import { getConversationKey } from "./circuitBreaker.js";
|
import { getConversationKey } from "./circuitBreaker.js";
|
||||||
import {
|
import { conversationDebounceTimers } from "./conversationState.js";
|
||||||
conversationConsecutiveErrors,
|
|
||||||
conversationDebounceTimers,
|
|
||||||
conversationErrorCooldown,
|
|
||||||
conversationProcessing,
|
|
||||||
isConversationProcessingLocked,
|
|
||||||
} from "./conversationState.js";
|
|
||||||
import {
|
import {
|
||||||
activeIndividualRequests,
|
activeIndividualRequests,
|
||||||
enqueueIndividualFallbacks,
|
|
||||||
individualCooldownUntil,
|
individualCooldownUntil,
|
||||||
individualInFlight,
|
individualInFlight,
|
||||||
individualInFlightByConversation,
|
|
||||||
individualInFlightLastTouched,
|
|
||||||
} from "./individualFallbackProcessor.js";
|
} from "./individualFallbackProcessor.js";
|
||||||
import {
|
import {
|
||||||
broadcastAnalysisCompleted,
|
broadcastAnalysisCompleted,
|
||||||
@@ -36,31 +27,20 @@ import {
|
|||||||
setModerationClient,
|
setModerationClient,
|
||||||
setSharedEventBroadcaster,
|
setSharedEventBroadcaster,
|
||||||
} from "./moderationState.js";
|
} from "./moderationState.js";
|
||||||
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
|
import { startRecoveryWorker } from "./recovery-worker.js";
|
||||||
import { pruneExpiredTexts } from "./textCacheStore.js";
|
|
||||||
|
|
||||||
const logger = createChildLogger("ai-analyzer");
|
const logger = createChildLogger("ai-analyzer");
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
// Cache hygiene (expired verdict sweep)
|
// Public API — queueing, status, worker startup
|
||||||
// ---------------------------------------------------------------------------
|
|
||||||
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
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Queues a message for analysis (debounced by conversation).
|
* 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> {
|
export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
||||||
if (!config.AI_ANALYSIS_ENABLED) return;
|
if (!config.AI_ANALYSIS_ENABLED) return;
|
||||||
@@ -73,13 +53,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (isAgeRestrictedMessage(message)) {
|
if (isAgeRestrictedMessage(message)) {
|
||||||
const updated = await messageStore.updateMessageAIAnalysis(
|
await recordSkip(message.id, buildAgeRestrictedSkipResult());
|
||||||
message.id,
|
|
||||||
buildAgeRestrictedSkipResult(),
|
|
||||||
);
|
|
||||||
if (updated) {
|
|
||||||
broadcastAnalysisCompleted(updated);
|
|
||||||
}
|
|
||||||
logger.debug(
|
logger.debug(
|
||||||
{ messageId },
|
{ messageId },
|
||||||
"Skipped AI analysis for age-restricted message",
|
"Skipped AI analysis for age-restricted message",
|
||||||
@@ -88,13 +62,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (isSkipAnalysisUser(message)) {
|
if (isSkipAnalysisUser(message)) {
|
||||||
const updated = await messageStore.updateMessageAIAnalysis(
|
await recordSkip(message.id, buildSkipAnalysisUserResult());
|
||||||
message.id,
|
|
||||||
buildSkipAnalysisUserResult(),
|
|
||||||
);
|
|
||||||
if (updated) {
|
|
||||||
broadcastAnalysisCompleted(updated);
|
|
||||||
}
|
|
||||||
logger.debug(
|
logger.debug(
|
||||||
{ messageId, userId: message.user_id },
|
{ messageId, userId: message.user_id },
|
||||||
"Skipped AI analysis for configured skip-list user",
|
"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).
|
* Queues a conversation for analysis (debounced).
|
||||||
*/
|
*/
|
||||||
@@ -129,6 +108,8 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
|
|||||||
return {
|
return {
|
||||||
queuedConversations: conversationDebounceTimers.size,
|
queuedConversations: conversationDebounceTimers.size,
|
||||||
activeRequests,
|
activeRequests,
|
||||||
|
activeTextRequests,
|
||||||
|
activeMediaRequests,
|
||||||
activeIndividualRequests,
|
activeIndividualRequests,
|
||||||
individualInFlightCount: individualInFlight.size,
|
individualInFlightCount: individualInFlight.size,
|
||||||
individualCircuitBreakerActive: Date.now() < individualCooldownUntil,
|
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`
|
* Also injects the Discord client and event broadcaster into the pipeline
|
||||||
* state (not just `pending`), and skips conversations that already have
|
* state so downstream modules can act and publish.
|
||||||
* individual fallback work in progress to avoid DB last-write-wins races.
|
|
||||||
*/
|
*/
|
||||||
export function startPendingAIAnalysisWorker(
|
export function startPendingAIAnalysisWorker(
|
||||||
client?: Client,
|
client?: Client,
|
||||||
@@ -160,131 +142,5 @@ export function startPendingAIAnalysisWorker(
|
|||||||
.catch(console.error);
|
.catch(console.error);
|
||||||
}
|
}
|
||||||
|
|
||||||
setInterval(() => {
|
startRecoveryWorker();
|
||||||
// [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);
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 { createChildLogger } from "@/shared/logger/index";
|
||||||
import { config } from "../../shared/config/config.js";
|
import { config } from "../../shared/config/index.js";
|
||||||
import type {
|
import type {
|
||||||
AnalysisResult,
|
AnalysisResult,
|
||||||
MessageRecord,
|
MessageRecord,
|
||||||
@@ -35,8 +35,17 @@ export function deriveSeverity(msg: MessageRecord): string {
|
|||||||
|
|
||||||
/** Derive recommended action from legacy messages that lack structured AI fields. */
|
/** Derive recommended action from legacy messages that lack structured AI fields. */
|
||||||
export function deriveRecommendedAction(msg: MessageRecord): string {
|
export function deriveRecommendedAction(msg: MessageRecord): string {
|
||||||
if (msg.ai_recommended_action) return msg.ai_recommended_action;
|
|
||||||
const severity = deriveSeverity(msg);
|
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 (
|
if (
|
||||||
msg.ai_status === "flagged" &&
|
msg.ai_status === "flagged" &&
|
||||||
(severity === "critical" || severity === "high" || severity === "medium")
|
(severity === "critical" || severity === "high" || severity === "medium")
|
||||||
@@ -177,7 +186,23 @@ export function isEligibleForAutoDelete(
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Recommended action check
|
// 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 (
|
||||||
|
status === "flagged" &&
|
||||||
|
(severity === "high" || severity === "critical")
|
||||||
|
) {
|
||||||
|
logger.debug(
|
||||||
|
{ messageId: message.id, status, severity },
|
||||||
|
"Message eligible for auto-delete: flagged with high/critical severity",
|
||||||
|
);
|
||||||
|
} else {
|
||||||
const recommendedAction =
|
const recommendedAction =
|
||||||
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
||||||
if (
|
if (
|
||||||
@@ -191,6 +216,7 @@ export function isEligibleForAutoDelete(
|
|||||||
);
|
);
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Categories check
|
// Categories check
|
||||||
const allowedCategories = parseStringList(
|
const allowedCategories = parseStringList(
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import type { Guild } from "discord.js-selfbot-v13";
|
import type { Guild } from "discord.js-selfbot-v13";
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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";
|
import type { MessageRecord } from "../message-capture/types.js";
|
||||||
|
|
||||||
interface ChannelWithSend {
|
interface ChannelWithSend {
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
import type { Client, PermissionString } from "discord.js-selfbot-v13";
|
import type { Client, PermissionString } from "discord.js-selfbot-v13";
|
||||||
import { LRUCache } from "lru-cache";
|
import { LRUCache } from "lru-cache";
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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 { parseRichMessageMetadata } from "../message-capture/messageMetadata.js";
|
||||||
import { messageStore } from "../message-capture/messageStore.js";
|
import { messageStore } from "../message-capture/messageStore.js";
|
||||||
import type { MessageRecord } from "../message-capture/types.js";
|
import type { MessageRecord } from "../message-capture/types.js";
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import type { Client } from "discord.js-selfbot-v13";
|
import type { Client } from "discord.js-selfbot-v13";
|
||||||
import { createChildLogger } from "@/shared/logger/index";
|
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";
|
import type { MessageRecord } from "../message-capture/types.js";
|
||||||
|
|
||||||
const logger = createChildLogger("auto-delete-notify");
|
const logger = createChildLogger("auto-delete-notify");
|
||||||
|
|||||||
@@ -43,3 +43,22 @@ export function pickBatchWithinBudget(
|
|||||||
|
|
||||||
return batch;
|
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