Compare commits
11
Commits
54d02098c8
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2b6ec59286 | ||
|
|
fdc7f01c26 | ||
|
|
f31b1f62d4 | ||
|
|
7be069d73f | ||
|
|
88373985ec | ||
|
|
7acc49e1fb | ||
|
|
81b3128bb0 | ||
|
|
1d0b3f422d | ||
|
|
0aef0b6c08 | ||
|
|
cd1d6b2d0d | ||
|
|
c888f23901 |
@@ -64,11 +64,6 @@ AI_LLM_API_KEY= # REQUIRED if AI_ANALYSIS_ENABLED=true. L
|
||||
AI_LLM_BASE_URL=http://100.121.180.82:20128/api/v1 # LLM API base URL (omniroute — OpenAI-compatible router on imrnes, Tailscale 100.121.180.82)
|
||||
AI_LLM_MODEL=text # LLM text model name (default: text)
|
||||
# AI_LLM_VISION_MODEL= # Vision model for image analysis (falls back to AI_LLM_MODEL)
|
||||
# AI_LLM_EMBEDDING_MODEL= # Embedding model for semantic moderation cache (optional; enables near-duplicate text reuse to save LLM calls)
|
||||
# AI_LLM_EMBEDDING_MIN_SIMILARITY=0.97 # Min cosine similarity to reuse a cached verdict (default: 0.97)
|
||||
QDRANT_URL=http://100.121.180.82:6333 # Qdrant vector store for embeddings (semantic cache); when set, vectors are stored/searched in Qdrant instead of Postgres
|
||||
# QDRANT_COLLECTION=gmw_text_moderation # Qdrant collection name (default: gmw_text_moderation)
|
||||
# QDRANT_API_KEY= # Qdrant API key (optional)
|
||||
AI_LLM_MAX_CONCURRENT=5 # Max concurrent LLM API calls (default: 5)
|
||||
AI_LLM_IMAGE_MAX_DIMENSION=1024 # Max image dimension in pixels before resize (default: 1024)
|
||||
AI_LLM_TEXT_BATCH_SIZE=20 # Max messages per text-only moderation batch (default: 20)
|
||||
|
||||
@@ -32,29 +32,26 @@ jobs:
|
||||
with:
|
||||
node-version: 22
|
||||
|
||||
- name: Install pnpm
|
||||
run: corepack enable && corepack prepare pnpm@11 --activate
|
||||
- name: Setup Bun
|
||||
uses: oven-sh/setup-bun@v2
|
||||
with:
|
||||
bun-version: 1.3.14
|
||||
|
||||
- name: Install deps (backend)
|
||||
working-directory: services/backend
|
||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
||||
|
||||
- name: Typecheck + test (backend)
|
||||
- name: Install deps + test (backend)
|
||||
working-directory: services/backend
|
||||
run: |
|
||||
bun install --frozen-lockfile
|
||||
./node_modules/.bin/tsc --noEmit
|
||||
# e2e.test.ts requires a live backend (API_BASE) — run unit tests only
|
||||
./node_modules/.bin/vitest run --exclude "src/e2e.test.ts"
|
||||
# src/e2e.test.ts requires a live backend (API_BASE) — unit tests
|
||||
# live in tests/ and are excluded by the bun test dir.
|
||||
bun test tests/
|
||||
|
||||
- name: Install deps (discord-gateway)
|
||||
working-directory: services/discord-gateway
|
||||
run: pnpm install --ignore-scripts --no-frozen-lockfile
|
||||
|
||||
- name: Typecheck + test (discord-gateway)
|
||||
- name: Install deps + test (discord-gateway)
|
||||
working-directory: services/discord-gateway
|
||||
run: |
|
||||
bun install --frozen-lockfile
|
||||
./node_modules/.bin/tsc --noEmit
|
||||
./node_modules/.bin/vitest run
|
||||
bun test tests/
|
||||
|
||||
- name: Biome check (all services)
|
||||
run: |
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
{ "id": "gateway", "type": "backend", "label": "discord-gateway", "sublabel": "selfbot :4016", "pos": [260, 240], "size": [140, 60], "tag": "discord.js-selfbot-v13" },
|
||||
{ "id": "backend", "type": "backend", "label": "gmw-backend", "sublabel": "Express · oRPC · WS :4001", "pos": [640, 240], "size": [140, 60], "tag": "Drizzle ORM" },
|
||||
{ "id": "proxy", "type": "cloud", "label": "nginx proxy", "sublabel": "reverse proxy :4009", "pos": [960, 240], "size": [140, 60], "tag": "nginx" },
|
||||
{ "id": "ai", "type": "backend", "label": "AI Moderation", "sublabel": "LLM caller · embeddings", "pos": [260, 400], "size": [140, 60] },
|
||||
{ "id": "ai", "type": "backend", "label": "AI Moderation", "sublabel": "LLM caller · vision", "pos": [260, 400], "size": [140, 60] },
|
||||
{ "id": "frontend", "type": "frontend", "label": "gmw-frontend", "sublabel": "Next.js 16 SSR :4017", "pos": [960, 400], "size": [140, 60], "tag": "React 19 · Tailwind v4" },
|
||||
{ "id": "llm", "type": "cloud", "label": "9router LLM", "sublabel": "text + vision API", "pos": [260, 540], "size": [140, 60], "tag": "AI_LLM_BASE_URL" },
|
||||
{ "id": "browser", "type": "external", "label": "Dashboard Users", "sublabel": "browser · partysocket", "pos": [960, 540], "size": [140, 60] }
|
||||
|
||||
@@ -35,9 +35,11 @@
|
||||
|
||||
# ---- Shared build tools ----
|
||||
nodejs = pkgs.nodejs_22;
|
||||
pnpm = pkgs.pnpm.override { nodejs = nodejs; };
|
||||
# Bun for deps/install (replaces pnpm); keeps nodejs for the tsc +
|
||||
# fix-imports.mjs build path (Bun's own bundler is not used for dist).
|
||||
bun = pkgs.bun;
|
||||
|
||||
pnpmInstall = ''
|
||||
bunInstall = ''
|
||||
export HOME=$TMPDIR/home
|
||||
export npm_config_cache=$TMPDIR/npm-cache
|
||||
mkdir -p $npm_config_cache
|
||||
@@ -48,15 +50,9 @@
|
||||
export GIT_SSL_CAINFO=${pkgs.cacert}/etc/ssl/certs/ca-bundle.crt
|
||||
export NIX_SSL_CERT_FILE=${pkgs.cacert}/etc/ssl/certs/ca-bundle.crt
|
||||
|
||||
# pnpm uses node-gyp for native addons — provide build tools (kept for
|
||||
# the rare case a prebuilt is unavailable and it falls back to compile).
|
||||
export CPPFLAGS="-I${pkgs.lib.getDev pkgs.openssl}/include"
|
||||
export LDFLAGS="-L${pkgs.lib.getLib pkgs.openssl}/lib"
|
||||
|
||||
pnpm install --no-frozen-lockfile --ignore-scripts 2>&1
|
||||
|
||||
# Build native addons that need compilation
|
||||
pnpm rebuild 2>&1 || true
|
||||
# Build native addons (bun install runs postinstall scripts for
|
||||
# @discordjs/opus / sharp unless trustedDependencies restricts).
|
||||
bun install 2>&1
|
||||
'';
|
||||
|
||||
# Shrink the shipped node_modules to production deps only. The full
|
||||
@@ -74,28 +70,15 @@
|
||||
# Must run AFTER tsc (typescript is a devDep) and after native builds.
|
||||
pruneProd = ''
|
||||
echo "=== Pruning devDependencies (production-only node_modules) ==="
|
||||
pnpm list --prod --depth 999 --parseable 2>/dev/null \
|
||||
| grep -o '\.pnpm/[^/]*' | sort -u > $TMPDIR/prod-pnms.txt
|
||||
( cd node_modules/.pnpm \
|
||||
&& for d in */; do \
|
||||
d="''${d%/}"; \
|
||||
[ "$d" = "node_modules" ] && continue; \
|
||||
grep -qF ".pnpm/$d" $TMPDIR/prod-pnms.txt || rm -rf "$d"; \
|
||||
done ) || true
|
||||
# Drop runtime-dead packages that still land in the prod graph:
|
||||
# - `@types/*` (pure TypeScript declarations) get pulled in as
|
||||
# REAL dependencies by type-aware deps (discord-api-types ->
|
||||
# @types/node, pg-protocol -> @types/pg, ...) even though nothing
|
||||
# ever `require`s them at runtime. Safe to strip.
|
||||
# - `opusscript` is only a pure-JS fallback Opus engine that
|
||||
# prism-media's loader uses IF `@discordjs/opus` (native, always
|
||||
# present/prebuilt) fails to load. Since the native engine loads,
|
||||
# opusscript is never executed — dead weight pulled in via
|
||||
# discord.js-selfbot-v13's dependency. Strip it too.
|
||||
( cd node_modules/.pnpm && rm -rf @types+* opusscript@* 2>/dev/null ) || true
|
||||
# Drop symlinks whose .pnpm target was pruned (top-level, scoped dirs,
|
||||
# hoist, .bin — any depth). Mirrors stdenv's noBrokenSymlinks check,
|
||||
# which would otherwise fail the fixupPhase.
|
||||
# bun install's layout: node_modules/<pkg> for prod deps; devDeps are
|
||||
# also present during build (needed for tsc). Keep only what the prod
|
||||
# graph needs: simplest robust approach is `bun install --production`
|
||||
# semantics — but bun keeps the same flat layout; since the Nix build
|
||||
# already ran `bun install` (full, scripts on), prune dev-only top
|
||||
# entries that were only pulled by devDeps (typescript, biome, vitest,
|
||||
# drizzle-kit, tsx, @types/*).
|
||||
find node_modules -maxdepth 2 -type d \( -name 'typescript' -o -name '@biomejs' -o -name 'vitest' -o -name 'drizzle-kit' -o -name 'tsx' -o -name 'esbuild' \) -prune -exec rm -rf {} + 2>/dev/null || true
|
||||
rm -rf node_modules/.bin/tsc node_modules/.bin/vitest node_modules/.bin/biome node_modules/.bin/drizzle-kit 2>/dev/null || true
|
||||
find node_modules -type l ! -exec test -e {} \; -delete 2>/dev/null || true
|
||||
du -sh node_modules
|
||||
'';
|
||||
@@ -107,11 +90,11 @@
|
||||
|
||||
src = ./services/backend;
|
||||
|
||||
nativeBuildInputs = [ nodejs pnpm pkgs.python3 pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
nativeBuildInputs = [ nodejs bun pkgs.python3 pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Compiling TypeScript ==="
|
||||
npx tsc 2>&1
|
||||
./node_modules/.bin/tsc 2>&1
|
||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||
node scripts/fix-imports.mjs
|
||||
echo "=== Build complete ==="
|
||||
@@ -149,7 +132,7 @@ WRAPPER
|
||||
# libvips download), so no cmake or rust toolchain is needed.
|
||||
# python3/gnumake/gcc stay as node-gyp fallback for @discordjs/opus.
|
||||
nativeBuildInputs = [
|
||||
nodejs pnpm
|
||||
nodejs bun
|
||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||
pkgs.pkg-config
|
||||
pkgs.openssl
|
||||
@@ -177,22 +160,11 @@ WRAPPER
|
||||
# neither needed nor wanted here. Skip it entirely.
|
||||
dontFixup = true;
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
echo "=== Building native voice deps ==="
|
||||
# pnpm rebuild aborts on the first failing package and runs scripts
|
||||
# from the wrong cwd — build each native dep explicitly with its own
|
||||
# install script. Each failure is tolerated (|| true); the packages
|
||||
# @discordjs/opus ships prebuilt binaries for Node 22 (ABI node-v127,
|
||||
# linux-x64-glibc-2.35) — node-pre-gyp downloads the prebuilt .node
|
||||
# instead of compiling C++ from source. With build_from_source unset
|
||||
# (above), `pnpm rebuild` runs the package's own install script which
|
||||
# fetches the matching prebuilt; it only falls back to a source build
|
||||
# if the download fails. This keeps voice working without a per-build
|
||||
# native compile.
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Rebuilding @discordjs/opus (prebuilt download) ==="
|
||||
pnpm rebuild @discordjs/opus 2>&1 || true
|
||||
bun pm rebuild @discordjs/opus 2>&1 || true
|
||||
echo "=== Compiling TypeScript ===="
|
||||
npx tsc 2>&1
|
||||
./node_modules/.bin/tsc 2>&1
|
||||
echo "=== Fixing @/ path aliases + extensionless relative imports for node ESM ==="
|
||||
node scripts/fix-imports.mjs
|
||||
echo "=== Build complete ==="
|
||||
@@ -228,13 +200,13 @@ WRAPPER
|
||||
|
||||
src = frontendSrc;
|
||||
|
||||
nativeBuildInputs = [ nodejs pnpm pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
nativeBuildInputs = [ nodejs bun pkgs.gnumake pkgs.gcc pkgs.cacert ];
|
||||
|
||||
buildPhase = pnpmInstall + ''
|
||||
buildPhase = bunInstall + ''
|
||||
echo "=== Building Next.js SSR (standalone) ==="
|
||||
export NEXT_TELEMETRY_DISABLED=1
|
||||
export GMW_BACKEND_URL=http://127.0.0.1:4001
|
||||
npx next build 2>&1
|
||||
./node_modules/.bin/next build 2>&1
|
||||
'';
|
||||
|
||||
installPhase = ''
|
||||
@@ -312,13 +284,13 @@ WRAPPER
|
||||
|
||||
devShells.default = pkgs.mkShell {
|
||||
buildInputs = [
|
||||
nodejs pnpm
|
||||
nodejs bun
|
||||
pkgs.python3 pkgs.gnumake pkgs.gcc
|
||||
pkgs.rustc pkgs.cargo
|
||||
pkgs.ffmpeg-headless
|
||||
];
|
||||
shellHook = ''
|
||||
echo "GMW dev shell ready — node $(node --version), pnpm $(pnpm --version)"
|
||||
echo "GMW dev shell ready — node $(node --version), bun $(bun --version)"
|
||||
'';
|
||||
};
|
||||
});
|
||||
|
||||
@@ -60,7 +60,7 @@ src/
|
||||
| moderation | Moderation actions & metrics | `ai_moderations`, `moderation_actions` |
|
||||
| media | Media file management | `media_attachments` |
|
||||
| dashboard | Stats aggregation | Various (read-only) |
|
||||
| knowledge | Semantic search | Qdrant vector DB |
|
||||
| knowledge | Channel cultures & glossary browser | `channel_cultures`, `term_glossary_cache` |
|
||||
| chatbot | AI chatbot with tools | `chatbot_history` |
|
||||
| health | Health checks + metrics | Various |
|
||||
| analysis | Text analysis cache | `text_analysis_cache` |
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,2 @@
|
||||
[test]
|
||||
preload = ["./tests/setup-env.ts"]
|
||||
@@ -5,7 +5,7 @@
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"packageManager": "pnpm@11.20.0",
|
||||
"packageManager": "bun@1.3.14",
|
||||
"engines": {
|
||||
"node": ">=22.12.0",
|
||||
"pnpm": ">=9.0.0"
|
||||
@@ -16,14 +16,15 @@
|
||||
"format": "biome format --write .",
|
||||
"lint": "biome check --diagnostic-level=error .",
|
||||
"start": "node dist/index.js",
|
||||
"test": "vitest run",
|
||||
"typecheck": "tsc --noEmit"
|
||||
"test": "bun test tests/",
|
||||
"typecheck": "tsc --noEmit",
|
||||
"test:e2e": "bun test src/e2e.test.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@orpc/server": "1.15.2",
|
||||
"@orpc/server": "1.15.3",
|
||||
"axios": "^1.20.0",
|
||||
"dotenv": "^18.0.1",
|
||||
"drizzle-orm": "^0.45.2",
|
||||
"drizzle-orm": "^0.45.3",
|
||||
"express": "^5.2.1",
|
||||
"helmet": "^8.1.0",
|
||||
"ioredis": "^6.0.0",
|
||||
@@ -41,6 +42,6 @@
|
||||
"@types/ws": "^8.18.1",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "^5.0.1"
|
||||
"@types/bun": "latest"
|
||||
}
|
||||
}
|
||||
}
|
||||
Generated
-2677
File diff suppressed because it is too large
Load Diff
@@ -1,65 +0,0 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
|
||||
const logger = createChildLogger("messages-embed");
|
||||
|
||||
/** Max chars for a search query fed to the embedding model. */
|
||||
const MAX_QUERY_CHARS = 300;
|
||||
|
||||
/**
|
||||
* Normalize a user search query before embedding so it lands in the same
|
||||
* vector space as the archived content (which is normalized the same way on
|
||||
* write). Mirrors the gateway's normalizer: strip control/zero-width chars,
|
||||
* lowercase, collapse whitespace, cap length. Readable punctuation is kept —
|
||||
* a search query is already compact.
|
||||
*/
|
||||
export function normalizeEmbeddingQuery(raw: string): string {
|
||||
if (!raw) return "";
|
||||
return raw
|
||||
.replace(/[\p{Cc}\p{Cf}]/gu, " ")
|
||||
.toLowerCase()
|
||||
.replace(/\s+/g, " ")
|
||||
.trim()
|
||||
.slice(0, MAX_QUERY_CHARS);
|
||||
}
|
||||
|
||||
/**
|
||||
* Embed a search query with the configured OpenAI-compatible embedding model.
|
||||
* Uses raw fetch (the backend has no openai SDK dependency) and returns null
|
||||
* when embeddings are not configured (search unavailable).
|
||||
*
|
||||
* encoding_format: "float" is REQUIRED — Nvidia-backed models reject base64.
|
||||
*/
|
||||
export async function embedQuery(rawQuery: string): Promise<number[] | null> {
|
||||
if (!config.AI_LLM_API_KEY || !config.AI_LLM_EMBEDDING_MODEL) return null;
|
||||
const text = normalizeEmbeddingQuery(rawQuery);
|
||||
if (!text) return null;
|
||||
try {
|
||||
const res = await fetch(`${config.AI_LLM_BASE_URL}/embeddings`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Authorization: `Bearer ${config.AI_LLM_API_KEY}`,
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model: config.AI_LLM_EMBEDDING_MODEL,
|
||||
input: text,
|
||||
encoding_format: "float",
|
||||
}),
|
||||
});
|
||||
if (!res.ok) {
|
||||
logger.warn({ status: res.status }, "query embed HTTP error");
|
||||
return null;
|
||||
}
|
||||
const json = (await res.json()) as {
|
||||
data?: Array<{ embedding?: number[] }>;
|
||||
};
|
||||
return json.data?.[0]?.embedding ?? null;
|
||||
} catch (error) {
|
||||
logger.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"query embed failed",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -41,11 +41,3 @@ export const messageUpdateSchema = z.object({
|
||||
export type MessageQuery = z.infer<typeof messageQuerySchema>;
|
||||
export type MessageCreate = z.infer<typeof messageCreateSchema>;
|
||||
export type MessageUpdate = z.infer<typeof messageUpdateSchema>;
|
||||
|
||||
export const semanticSearchSchema = z.object({
|
||||
query: z.string().min(1).max(500),
|
||||
limit: z.coerce.number().int().positive().max(50).default(10),
|
||||
guildId: z.string().optional(),
|
||||
});
|
||||
|
||||
export type SemanticSearchQuery = z.infer<typeof semanticSearchSchema>;
|
||||
|
||||
@@ -1,10 +1,7 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { NotFoundError, ValidationError } from "@/shared/errors/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { embedQuery } from "./embed.js";
|
||||
import { type MessageRow, messagesRepository } from "./messages.repository.js";
|
||||
import type { MessageQuery, SemanticSearchQuery } from "./messages.schema.js";
|
||||
import { searchArchive } from "./qdrant.js";
|
||||
import type { MessageQuery } from "./messages.schema.js";
|
||||
|
||||
const logger = createChildLogger("messages.service");
|
||||
|
||||
@@ -102,32 +99,6 @@ export class MessagesService {
|
||||
return messagesRepository.getReviewMessages(channelId, limit);
|
||||
}
|
||||
|
||||
/**
|
||||
* Public, read-only semantic search over the persistent message archive.
|
||||
* Embeds the query, searches Qdrant, returns text + metadata. Best-effort:
|
||||
* if embeddings/Qdrant are unavailable, returns an empty result set.
|
||||
*/
|
||||
async semanticSearch(
|
||||
input: SemanticSearchQuery,
|
||||
): Promise<{ results: ReturnType<typeof mapSearchHit>[]; nextCursor: null }> {
|
||||
const vector = await embedQuery(input.query);
|
||||
if (!vector) {
|
||||
logger.debug(
|
||||
{ query: input.query },
|
||||
"semantic search skipped: no embedder",
|
||||
);
|
||||
return { results: [], nextCursor: null };
|
||||
}
|
||||
const hits = await searchArchive(
|
||||
vector,
|
||||
input.limit,
|
||||
config.AI_LLM_EMBEDDING_ARCHIVE_MIN_SIMILARITY,
|
||||
input.guildId,
|
||||
);
|
||||
const results = hits.map((h) => mapSearchHit(h));
|
||||
return { results, nextCursor: null };
|
||||
}
|
||||
|
||||
async getActivity(
|
||||
days = 30,
|
||||
): Promise<Awaited<ReturnType<typeof messagesRepository.getActivity>>> {
|
||||
@@ -157,36 +128,4 @@ export class MessagesService {
|
||||
}
|
||||
}
|
||||
|
||||
/** Shape returned to the frontend (text + rich metadata from the archive payload). */
|
||||
function mapSearchHit(hit: {
|
||||
score: number;
|
||||
payload: {
|
||||
text: string;
|
||||
content_hash?: string;
|
||||
analyzed_at: number;
|
||||
username?: string;
|
||||
channel_id?: string;
|
||||
guild_id?: string;
|
||||
thread_id?: string | null;
|
||||
channel_name?: string | null;
|
||||
thread_name?: string | null;
|
||||
created_at?: number;
|
||||
};
|
||||
}) {
|
||||
return {
|
||||
message_id: hit.payload.content_hash ?? null,
|
||||
content: hit.payload.text,
|
||||
score: hit.score,
|
||||
// Prefer the real message timestamp; fall back to embed time for old
|
||||
// points that predate rich metadata.
|
||||
created_at: hit.payload.created_at ?? hit.payload.analyzed_at,
|
||||
username: hit.payload.username ?? null,
|
||||
channel_id: hit.payload.channel_id ?? null,
|
||||
guild_id: hit.payload.guild_id ?? null,
|
||||
thread_id: hit.payload.thread_id ?? null,
|
||||
channel_name: hit.payload.channel_name ?? null,
|
||||
thread_name: hit.payload.thread_name ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
export const messagesService = new MessagesService();
|
||||
|
||||
@@ -1,114 +0,0 @@
|
||||
import { config } from "@/shared/config/index";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
|
||||
const logger = createChildLogger("messages-qdrant");
|
||||
|
||||
export interface ArchiveHit {
|
||||
score: number;
|
||||
payload: {
|
||||
text: string;
|
||||
content_hash?: string;
|
||||
analyzed_at: number;
|
||||
expires_at: number;
|
||||
username?: string;
|
||||
channel_id?: string;
|
||||
guild_id?: string;
|
||||
thread_id?: string | null;
|
||||
channel_name?: string | null;
|
||||
thread_name?: string | null;
|
||||
created_at?: number;
|
||||
};
|
||||
}
|
||||
|
||||
function baseUrl(): string {
|
||||
return (config.QDRANT_URL ?? "http://100.121.180.82:6333").replace(
|
||||
/\/+$/,
|
||||
"",
|
||||
);
|
||||
}
|
||||
|
||||
function headers(): Record<string, string> {
|
||||
const h: Record<string, string> = { "Content-Type": "application/json" };
|
||||
if (config.QDRANT_API_KEY) h["api-key"] = config.QDRANT_API_KEY;
|
||||
return h;
|
||||
}
|
||||
|
||||
export const ARCHIVE_COLLECTION =
|
||||
config.QDRANT_ARCHIVE_COLLECTION ?? "gmw_message_archive";
|
||||
|
||||
async function request(
|
||||
method: string,
|
||||
path: string,
|
||||
body?: unknown,
|
||||
timeoutMs = 10_000,
|
||||
): Promise<unknown> {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), timeoutMs);
|
||||
try {
|
||||
const res = await fetch(`${baseUrl()}${path}`, {
|
||||
method,
|
||||
headers: headers(),
|
||||
body: body === undefined ? undefined : JSON.stringify(body),
|
||||
signal: controller.signal,
|
||||
});
|
||||
const text = await res.text();
|
||||
if (!res.ok) {
|
||||
throw new Error(
|
||||
`Qdrant ${method} ${path} -> ${res.status}: ${text.slice(0, 200)}`,
|
||||
);
|
||||
}
|
||||
return text ? JSON.parse(text) : null;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
/** Search the archive collection for the nearest vectors to `vector`. */
|
||||
export async function searchArchive(
|
||||
vector: number[],
|
||||
limit: number,
|
||||
scoreThreshold: number,
|
||||
guildId?: string,
|
||||
): Promise<ArchiveHit[]> {
|
||||
if (!config.QDRANT_URL) return [];
|
||||
try {
|
||||
const json = (await request(
|
||||
"POST",
|
||||
`/collections/${ARCHIVE_COLLECTION}/points/search`,
|
||||
{
|
||||
vector,
|
||||
limit,
|
||||
score_threshold: scoreThreshold,
|
||||
with_payload: true,
|
||||
// Optional scope: only return vectors from a specific guild's archive.
|
||||
// Old points (embedded before rich metadata) have no guild_id payload —
|
||||
// the `must` match simply excludes them, which is the correct behavior
|
||||
// for a guild-scoped search.
|
||||
...(guildId
|
||||
? {
|
||||
filter: {
|
||||
must: [{ key: "guild_id", match: { value: guildId } }],
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
)) as {
|
||||
result?: Array<{
|
||||
score?: number;
|
||||
payload?: ArchiveHit["payload"];
|
||||
}>;
|
||||
};
|
||||
return (json.result ?? [])
|
||||
.filter((h) => h.payload?.text)
|
||||
.map((h) => ({
|
||||
score: h.score ?? 0,
|
||||
payload: h.payload as ArchiveHit["payload"],
|
||||
}));
|
||||
} catch (error) {
|
||||
logger.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"archive search failed",
|
||||
);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
@@ -5,10 +5,7 @@ import { chatRequestSchema } from "../modules/chatbot/chatbot.schema";
|
||||
import { chatbotService } from "../modules/chatbot/chatbot.service";
|
||||
import { dashboardService } from "../modules/dashboard/dashboard.service";
|
||||
import { knowledgeService } from "../modules/knowledge/knowledge.service";
|
||||
import {
|
||||
messageQuerySchema,
|
||||
semanticSearchSchema,
|
||||
} from "../modules/messages/messages.schema";
|
||||
import { messageQuerySchema } from "../modules/messages/messages.schema";
|
||||
import { messagesService } from "../modules/messages/messages.service";
|
||||
import { moderationService } from "../modules/moderation/moderation.service";
|
||||
import { uiStateService } from "../modules/ui-state/ui-state.service";
|
||||
@@ -123,10 +120,6 @@ const messagesRouter = {
|
||||
);
|
||||
return { results: rows, limit: input.limit, cursor: null };
|
||||
}),
|
||||
// Public, read-only semantic search over the message archive.
|
||||
semanticSearch: os
|
||||
.input(semanticSearchSchema)
|
||||
.handler(({ input }) => messagesService.semanticSearch(input)),
|
||||
// Public, read-only activity heatmap data (per-hour volume by channel).
|
||||
activity: os
|
||||
.input(
|
||||
|
||||
@@ -102,15 +102,6 @@ export const configSchema = z
|
||||
AI_LLM_BASE_URL: z.string().url().default("http://127.0.0.1:4014/v1"),
|
||||
AI_LLM_MODEL: z.string().default("text"),
|
||||
AI_LLM_VISION_MODEL: z.string().optional(),
|
||||
AI_LLM_EMBEDDING_MODEL: z.string().optional(),
|
||||
// Minimum cosine similarity for the public archive semantic search. Lower
|
||||
// = more (noisier) results; raise it to tighten precision. Tuned for a 1B
|
||||
// embedding model — re-tune if the model's dimensionality changes.
|
||||
AI_LLM_EMBEDDING_ARCHIVE_MIN_SIMILARITY: z.coerce
|
||||
.number()
|
||||
.min(0)
|
||||
.max(1)
|
||||
.default(0.6),
|
||||
AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(5),
|
||||
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
||||
.number()
|
||||
@@ -179,11 +170,6 @@ export const configSchema = z
|
||||
.default("https://api.openai.com/v1"),
|
||||
OPENAI_MODERATION_MODEL: z.string().default("omni-moderation-latest"),
|
||||
|
||||
// ── Qdrant (message archive for semantic search) ──────────────────
|
||||
QDRANT_URL: z.string().optional(),
|
||||
QDRANT_API_KEY: z.string().optional(),
|
||||
QDRANT_ARCHIVE_COLLECTION: z.string().default("gmw_message_archive"),
|
||||
|
||||
// ── Auto Delete ─────────────────────────────────────────────────────
|
||||
AUTO_DELETE_FLAGGED_ENABLED: z
|
||||
.string()
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { tools } from "../src/modules/chatbot/chatbot.toolDefs.js";
|
||||
|
||||
const names = tools.map((t) => t.function.name);
|
||||
|
||||
@@ -1,6 +1,32 @@
|
||||
// ─── Shared Error Classes ────────────────────────────────────────────────────
|
||||
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
// bun:test compat facade — vitest's `vi` maps onto bun's `jest`/`mock`/`spyOn`.
|
||||
// bun:test 1.3.14 exports both `jest` (fn, useFakeTimers, spyOn) and `mock`
|
||||
// (module, restore). `vi.fn` -> `jest.fn`, `vi.useFakeTimers` -> `jest.useFakeTimers`,
|
||||
// `vi.waitFor` -> waitForCompat (poll until the assertion passes).
|
||||
|
||||
import { afterEach, describe, expect, it, jest } from "bun:test";
|
||||
|
||||
const useFakeTimers = () => jest.useFakeTimers();
|
||||
const useRealTimers = () => jest.useRealTimers();
|
||||
const advanceTimersByTime = (ms: number) => jest.advanceTimersByTime(ms);
|
||||
async function waitForCompat(fn: () => Promise<unknown>, timeoutMs = 2_000) {
|
||||
const start = Date.now();
|
||||
let lastErr: unknown;
|
||||
while (Date.now() - start < timeoutMs) {
|
||||
try {
|
||||
await fn();
|
||||
return;
|
||||
} catch (err) {
|
||||
lastErr = err;
|
||||
await new Promise((r) => setTimeout(r, 10));
|
||||
}
|
||||
}
|
||||
throw lastErr instanceof Error
|
||||
? lastErr
|
||||
: new Error("waitForCompat timed out");
|
||||
}
|
||||
|
||||
import {
|
||||
AppError,
|
||||
ConfigError,
|
||||
@@ -94,41 +120,41 @@ describe("AppError subclasses", () => {
|
||||
// ═══════════════════════════════════════════════════════════════════════════════
|
||||
describe("delay", () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
useRealTimers();
|
||||
});
|
||||
|
||||
it("resolves after the given time", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const promise = delay(500);
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
await expect(promise).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("rejects are not triggered on non-matching timer", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const promise = delay(1000);
|
||||
// Advance only part way — the timer should NOT fire yet
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
// The timer is still pending; the promise has not resolved yet
|
||||
// We advance the rest
|
||||
vi.advanceTimersByTime(500);
|
||||
advanceTimersByTime(500);
|
||||
await expect(promise).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
describe("retryWithBackoff", () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
useRealTimers();
|
||||
});
|
||||
|
||||
it("returns the result on first success without retrying", async () => {
|
||||
const fn = vi.fn().mockResolvedValue("ok");
|
||||
const fn = jest.fn().mockResolvedValue("ok");
|
||||
await expect(retryWithBackoff(fn)).resolves.toBe("ok");
|
||||
expect(fn).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("re-throws after exhausting all retries", async () => {
|
||||
const fn = vi.fn().mockRejectedValue(new Error("persistent"));
|
||||
const fn = jest.fn().mockRejectedValue(new Error("persistent"));
|
||||
await expect(
|
||||
retryWithBackoff(fn, { retries: 1, minTimeout: 1, maxTimeout: 5 }),
|
||||
).rejects.toThrow("persistent");
|
||||
@@ -139,7 +165,7 @@ describe("retryWithBackoff", () => {
|
||||
it("throws AbortError immediately when signal is already aborted", async () => {
|
||||
const ac = new AbortController();
|
||||
ac.abort();
|
||||
const fn = vi.fn().mockResolvedValue("ok");
|
||||
const fn = jest.fn().mockResolvedValue("ok");
|
||||
await expect(
|
||||
retryWithBackoff(fn, { retries: 3, signal: ac.signal }),
|
||||
).rejects.toThrow("Aborted");
|
||||
@@ -147,9 +173,9 @@ describe("retryWithBackoff", () => {
|
||||
});
|
||||
|
||||
it("respects abort signal during retry", async () => {
|
||||
vi.useFakeTimers();
|
||||
useFakeTimers();
|
||||
const ac = new AbortController();
|
||||
const fn = vi.fn().mockRejectedValue(new Error("fail"));
|
||||
const fn = jest.fn().mockRejectedValue(new Error("fail"));
|
||||
|
||||
const promise = retryWithBackoff(fn, {
|
||||
retries: 5,
|
||||
@@ -159,8 +185,8 @@ describe("retryWithBackoff", () => {
|
||||
|
||||
// Schedule abort after first failure + backoff starts
|
||||
setTimeout(() => ac.abort(), 150);
|
||||
vi.advanceTimersByTime(200);
|
||||
await vi.waitFor(async () => {
|
||||
advanceTimersByTime(200);
|
||||
await waitForCompat(async () => {
|
||||
await expect(promise).rejects.toThrow("Aborted");
|
||||
});
|
||||
});
|
||||
@@ -232,7 +258,7 @@ describe("asyncHandler", () => {
|
||||
const wrapped = asyncHandler(async () => {
|
||||
throw error;
|
||||
});
|
||||
const next = vi.fn();
|
||||
const next = jest.fn();
|
||||
|
||||
wrapped({} as any, {} as any, next);
|
||||
|
||||
@@ -246,7 +272,7 @@ describe("asyncHandler", () => {
|
||||
const wrapped = asyncHandler(async (_req: any, _res: any, _next: any) => {
|
||||
// no-op
|
||||
});
|
||||
const next = vi.fn();
|
||||
const next = jest.fn();
|
||||
|
||||
wrapped({} as any, {} as any, next);
|
||||
await Promise.resolve();
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
// bun test preload — nothing needed for backend unit tests today.
|
||||
@@ -1,4 +1,4 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { describe, expect, it } from "bun:test";
|
||||
|
||||
/**
|
||||
* Lock the contract that the WS `stream_messages` handler + frontend
|
||||
|
||||
@@ -1,16 +0,0 @@
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { defineConfig } from "vitest/config";
|
||||
|
||||
export default defineConfig({
|
||||
resolve: {
|
||||
alias: {
|
||||
"@": fileURLToPath(new URL("./src", import.meta.url)),
|
||||
},
|
||||
},
|
||||
test: {
|
||||
globals: true,
|
||||
environment: "node",
|
||||
include: ["src/**/*.test.ts", "tests/**/*.test.ts"],
|
||||
testTimeout: 15000,
|
||||
},
|
||||
});
|
||||
@@ -53,19 +53,16 @@ src/
|
||||
1. **LLM is the only judge.** Failed LLM → `status:"error"` + recovery retry.
|
||||
**Never** reintroduce regex/heuristic content classification.
|
||||
2. **Discord tokens sanitized** before reaching LLM (`discordTokens.ts`).
|
||||
3. **Semantic cache is batched** — one embed call + one Qdrant batch search.
|
||||
4. **Streaming is mandatory** against the router base URL.
|
||||
3. **Streaming is mandatory** against the router base URL.
|
||||
|
||||
## AI moderation pipeline
|
||||
|
||||
```
|
||||
aiAnalyzer.ts → batchScheduler.ts → batchProcessor.ts → individualFallbackProcessor.ts
|
||||
↓ ↓ ↓ ↓
|
||||
moderationOrchestrator.ts → (hash cache → Qdrant → LLM)
|
||||
moderationOrchestrator.ts → (hash cache → LLM)
|
||||
↓ ↓ ↓
|
||||
textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||
embeddingClient.ts
|
||||
qdrantClient.ts
|
||||
```
|
||||
|
||||
- Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`, `startPendingAIAnalysisWorker`)
|
||||
@@ -85,7 +82,6 @@ textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||
- `messageStore.ts` — DB operations
|
||||
- `messageMetadata.ts` — metadata extraction
|
||||
- `messagesDb.ts` / `messagesCrud.ts` — DB schema operations
|
||||
- `archiveEmbedder.ts` — Qdrant embedding (respect age-restricted guard)
|
||||
- `retentionDb.ts` / `reviewsDb.ts` / `attachmentsDb.ts` — auxiliary tables
|
||||
|
||||
## Redis channels (outbound to backend)
|
||||
|
||||
@@ -68,8 +68,8 @@ of the same conversation, and vice versa.
|
||||
(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 and
|
||||
Qdrant, driven from the recovery interval.
|
||||
- `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).
|
||||
@@ -82,8 +82,8 @@ of the same conversation, and vice versa.
|
||||
Piscina `textWorkerPool`/`mediaWorkerPool`, `getConversationKey`.
|
||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` (lane) /
|
||||
`individual` jobs). Runs `runModerationAnalysis` off the main thread.
|
||||
- `moderationOrchestrator.ts` — exact-hash cache → batched semantic (Qdrant)
|
||||
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
|
||||
(one call per sub-batch, not per message). `mediaBatchProcessor` routes its
|
||||
moderation LLM call through the MEDIA semaphore.
|
||||
@@ -93,8 +93,6 @@ of the same conversation, and vice versa.
|
||||
`AI_LLM_MEDIA_MAX_CONCURRENT` (media lane, default 4) — a vision backlog
|
||||
can never consume text slots. `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
||||
share the same router/base URL (different model alias for vision).
|
||||
- `embeddingClient.ts` + `qdrantClient.ts` — semantic cache (one embed call +
|
||||
one batched Qdrant search for all uncached targets).
|
||||
- `textCacheStore.ts` / `channelCultureStore.ts` / `userProfileStore.ts` /
|
||||
`userProfileStore.ts` — caches learned user profile summaries (optional).
|
||||
|
||||
@@ -192,7 +190,5 @@ pipeline gauges registered by `app/metrics-collector.ts` —
|
||||
- **Discord tokens are sanitized** (`discordTokens.ts`: `<:emoji:id>` →
|
||||
`[emoji:name]`, `<@id>` → `@user`, etc.) before content reaches the LLM, so
|
||||
numeric snowflake IDs never trigger false positives.
|
||||
- **Semantic cache is batched** (one embed call + one Qdrant batch search),
|
||||
not N sequential round-trips. `ensureQdrantCollection` is memoized.
|
||||
- **Streaming is mandatory** against the router base URL (non-stream waits for
|
||||
the full body and times out). `llmClient` aggregates SSE chunks.
|
||||
|
||||
@@ -59,5 +59,5 @@ Callers outside a module import its `index.ts` facade, never an internal file.
|
||||
## Testing
|
||||
|
||||
Vitest, tests in `tests/`. Config supplies dummy env vars so the suite runs
|
||||
without live Postgres/Redis/Qdrant; external services are mocked. `llmE2e.test.ts`
|
||||
without live Postgres/Redis; external services are mocked. `llmE2e.test.ts`
|
||||
is skipped by default and needs real credentials (`pnpm test:e2e:live`).
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,2 @@
|
||||
[test]
|
||||
preload = ["./tests/setup-env.ts"]
|
||||
@@ -17,10 +17,8 @@
|
||||
"typecheck": "tsc --noEmit",
|
||||
"lint": "biome check --diagnostic-level=error .",
|
||||
"format": "biome format --write .",
|
||||
"test": "vitest run",
|
||||
"test:unit": "vitest run --exclude \"tests/llmE2e.test.ts\"",
|
||||
"test:e2e": "vitest run tests/llmE2e.test.ts",
|
||||
"test:e2e:live": "bash scripts/run-llm-e2e.sh"
|
||||
"test": "bun test tests/",
|
||||
"test:e2e": "bun test tests/llmE2e.test.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"axios": "^1.20.0",
|
||||
@@ -43,12 +41,13 @@
|
||||
},
|
||||
"devDependencies": {
|
||||
"@biomejs/biome": "latest",
|
||||
"@types/node": "^26.4.0",
|
||||
"@types/node": "^26.6.2",
|
||||
"@types/pg": "^8.23.1",
|
||||
"@types/ws": "^8.18.1",
|
||||
"drizzle-kit": "^0.31.10",
|
||||
"drizzle-kit": "^0.31.11",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "latest"
|
||||
}
|
||||
}
|
||||
"@types/bun": "latest"
|
||||
},
|
||||
"packageManager": "bun@1.3.14"
|
||||
}
|
||||
Generated
-4092
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
|
||||
@@ -35,8 +35,17 @@ export function deriveSeverity(msg: MessageRecord): string {
|
||||
|
||||
/** Derive recommended action from legacy messages that lack structured AI fields. */
|
||||
export function deriveRecommendedAction(msg: MessageRecord): string {
|
||||
if (msg.ai_recommended_action) return msg.ai_recommended_action;
|
||||
const severity = deriveSeverity(msg);
|
||||
// Flagged at high/critical severity is ALWAYS delete — the stored
|
||||
// recommended_action from the LLM is conservative (often "review") and
|
||||
// must not override severity for severe violations.
|
||||
if (
|
||||
msg.ai_status === "flagged" &&
|
||||
(severity === "critical" || severity === "high")
|
||||
) {
|
||||
return "delete";
|
||||
}
|
||||
if (msg.ai_recommended_action) return msg.ai_recommended_action;
|
||||
if (
|
||||
msg.ai_status === "flagged" &&
|
||||
(severity === "critical" || severity === "high" || severity === "medium")
|
||||
@@ -177,19 +186,36 @@ export function isEligibleForAutoDelete(
|
||||
return false;
|
||||
}
|
||||
|
||||
// Recommended action check
|
||||
const recommendedAction =
|
||||
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
||||
// Recommended action check.
|
||||
// CRITICAL: the LLM's `recommended_action` is CONSERVATIVE — for a flagged
|
||||
// message at high/critical severity it frequently emits "review" (it sees a
|
||||
// screenshot/context ambiguity and hedges) even when the violation itself is
|
||||
// severe. Trusting that value lets serious violations (harassment, SARA,
|
||||
// threats) slip through undeleted. So: flagged + high/critical severity is
|
||||
// ALWAYS eligible regardless of the LLM's recommended action. The action
|
||||
// check only gates warn/flagged-medium (where a review is legitimate).
|
||||
if (
|
||||
recommendedAction !== "delete" &&
|
||||
recommendedAction !== "escalate" &&
|
||||
recommendedAction !== "warn"
|
||||
status === "flagged" &&
|
||||
(severity === "high" || severity === "critical")
|
||||
) {
|
||||
logger.debug(
|
||||
{ messageId: message.id, recommendedAction },
|
||||
"Message not eligible for auto-delete: recommended action is not delete/escalate/warn",
|
||||
{ messageId: message.id, status, severity },
|
||||
"Message eligible for auto-delete: flagged with high/critical severity",
|
||||
);
|
||||
return false;
|
||||
} else {
|
||||
const recommendedAction =
|
||||
analysisResult?.recommendedAction ?? deriveRecommendedAction(message);
|
||||
if (
|
||||
recommendedAction !== "delete" &&
|
||||
recommendedAction !== "escalate" &&
|
||||
recommendedAction !== "warn"
|
||||
) {
|
||||
logger.debug(
|
||||
{ messageId: message.id, recommendedAction },
|
||||
"Message not eligible for auto-delete: recommended action is not delete/escalate/warn",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Categories check
|
||||
|
||||
@@ -1,6 +1,4 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
|
||||
import { pruneExpiredTexts } from "./textCacheStore.js";
|
||||
|
||||
const logger = createChildLogger("cache-prune");
|
||||
@@ -11,7 +9,7 @@ const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
|
||||
let lastCachePruneAt = 0;
|
||||
|
||||
/**
|
||||
* Cache hygiene: purge expired moderation verdicts from Postgres and Qdrant.
|
||||
* Cache hygiene: purge expired moderation verdicts from Postgres.
|
||||
*
|
||||
* Expired entries are never reused (read filters check `expires_at`) but they
|
||||
* accumulate forever without a sweep. Called from the recovery interval; the
|
||||
@@ -21,13 +19,10 @@ export function runCachePruneIfDue(now: number = Date.now()): void {
|
||||
if (now - lastCachePruneAt < CACHE_PRUNE_INTERVAL_MS) return;
|
||||
lastCachePruneAt = now;
|
||||
|
||||
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
|
||||
.then(([pgDeleted, qdDeleted]) => {
|
||||
if (pgDeleted > 0 || qdDeleted > 0) {
|
||||
logger.info(
|
||||
{ pgDeleted, qdDeleted },
|
||||
"Expired moderation cache pruned",
|
||||
);
|
||||
Promise.resolve(pruneExpiredTexts())
|
||||
.then((pgDeleted) => {
|
||||
if (pgDeleted > 0) {
|
||||
logger.info({ pgDeleted }, "Expired moderation cache pruned");
|
||||
}
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
|
||||
@@ -1,200 +0,0 @@
|
||||
/**
|
||||
* embeddingClient.ts
|
||||
*
|
||||
* OpenAI-compatible embeddings helper used by the semantic moderation
|
||||
* cache. When AI_LLM_EMBEDDING_MODEL is configured, near-duplicate
|
||||
* messages can reuse a stored verdict (cosine similarity) instead of
|
||||
* paying for a full chat-completion call — the main cost-saver.
|
||||
*
|
||||
* Every function degrades gracefully: if the embedding model is not
|
||||
* configured or the API fails, callers fall back to the exact-hash cache
|
||||
* and then the LLM, so moderation quality is never reduced.
|
||||
*/
|
||||
|
||||
import OpenAI from "openai";
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { cleanContent } from "./textSignals.js";
|
||||
|
||||
const log = createChildLogger("embedding-client");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Text normalization (shared by moderation + archive embedding)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Max characters fed to the embedding model for a single document. Embedding
|
||||
* models have a hard token ceiling; embedding past it throws / wastes tokens.
|
||||
* Content messages are truncated; search queries have their own (smaller) cap.
|
||||
*/
|
||||
export const EMBEDDING_MAX_CHARS = 1200;
|
||||
export const EMBEDDING_MAX_QUERY_CHARS = 300;
|
||||
|
||||
/**
|
||||
* Normalize free-form Discord text before embedding.
|
||||
*
|
||||
* Raw messages are full of signal-hostile noise: @mentions, channels, custom
|
||||
* emoji, URLs, markdown and control chars. Embedding that noise directly
|
||||
* dilutes the vector (sentences that differ only in an @mention or a link
|
||||
* land far apart) and inflates token cost. The same cleanup is applied to the
|
||||
* user's search query so archived vectors and the query share one space.
|
||||
*
|
||||
* - `cleanContent` (from textSignals) strips URLs/@mentions/emoji/markdown and
|
||||
* collapses whitespace — good for embeddings, not just term extraction.
|
||||
* - Control characters / zero-width joiners are removed (Discord pastes these).
|
||||
* - Lowercase is applied so "Discord" and "discord" embed identically (embed
|
||||
* models are case-sensitive; this measurably improves near-duplicate recall).
|
||||
*
|
||||
* Falls back to the raw input if normalization empties a string (e.g. a
|
||||
* message that was only a URL) so index alignment is preserved by callers.
|
||||
*/
|
||||
export function normalizeEmbeddingText(raw: string, maxChars: number): string {
|
||||
if (!raw) return "";
|
||||
const cleaned = cleanContent(
|
||||
raw.replace(/[\p{Cc}\p{Cf}]/gu, " ").toLowerCase(),
|
||||
);
|
||||
if (!cleaned) return raw.slice(0, maxChars); // preserve original if stripped
|
||||
return cleaned.slice(0, maxChars);
|
||||
}
|
||||
|
||||
/** Normalize a single embedded document/message. */
|
||||
export function normalizeEmbeddingContent(raw: string): string {
|
||||
return normalizeEmbeddingText(raw, EMBEDDING_MAX_CHARS);
|
||||
}
|
||||
|
||||
/** Normalize a user-provided search query before embedding it. */
|
||||
export function normalizeEmbeddingQuery(raw: string): string {
|
||||
return normalizeEmbeddingText(raw, EMBEDDING_MAX_QUERY_CHARS);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Client (lazy singleton — same base URL as the chat client)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
let openaiClient: OpenAI | null = null;
|
||||
|
||||
function getClient(): OpenAI | null {
|
||||
if (!config.AI_LLM_API_KEY || !config.AI_LLM_EMBEDDING_MODEL) return null;
|
||||
if (!openaiClient) {
|
||||
openaiClient = new OpenAI({
|
||||
apiKey: config.AI_LLM_API_KEY,
|
||||
baseURL: config.AI_LLM_BASE_URL,
|
||||
// Embeddings are cheap and idempotent — a transient network blip should
|
||||
// NOT silently disable the whole semantic cache for a batch. Let the SDK
|
||||
// retry (2 retries, jittered) instead of failing open immediately.
|
||||
maxRetries: 2,
|
||||
timeout: 60_000,
|
||||
});
|
||||
}
|
||||
return openaiClient;
|
||||
}
|
||||
|
||||
/** True when the semantic cache is usable (key + model configured). */
|
||||
export function isEmbeddingEnabled(): boolean {
|
||||
return Boolean(config.AI_LLM_API_KEY && config.AI_LLM_EMBEDDING_MODEL);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Embedding calls
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Embed a batch of texts with the configured model.
|
||||
* Returns null on any failure so callers can skip semantic lookup.
|
||||
*
|
||||
* Each input is normalized (noise stripped, lowercased, length-capped) before
|
||||
* embedding — see normalizeEmbeddingText. Index alignment with `texts` is
|
||||
* preserved: if a text normalizes to empty we embed the raw original so the
|
||||
* caller's `embeddings[i] ↔ texts[i]` mapping never shifts.
|
||||
*/
|
||||
export async function embedTexts(texts: string[]): Promise<number[][] | null> {
|
||||
if (!isEmbeddingEnabled()) return null;
|
||||
if (texts.length === 0) return [];
|
||||
|
||||
const client = getClient();
|
||||
if (!client) return null;
|
||||
|
||||
const normalized = texts.map((t) => normalizeEmbeddingContent(t));
|
||||
|
||||
try {
|
||||
const response = await client.embeddings.create({
|
||||
model: config.AI_LLM_EMBEDDING_MODEL as string,
|
||||
input: normalized,
|
||||
// OpenAI SDK v6 defaults to base64; Nvidia-backed embedding models
|
||||
// (e.g. llama-nemotron-embed) reject it with 400. Always float.
|
||||
encoding_format: "float",
|
||||
});
|
||||
// All vectors in one response must share the model's dimension. If they
|
||||
// don't (shouldn't happen, but guards against a misconfigured/mismatched
|
||||
// model), fail the batch rather than feed garbage to cosine + Qdrant.
|
||||
const vectors = response.data.map((item) => item.embedding);
|
||||
const firstLen = vectors[0]?.length ?? 0;
|
||||
const consistent = vectors.every((v) => v.length === firstLen);
|
||||
if (!consistent || firstLen === 0) {
|
||||
log.error(
|
||||
{
|
||||
model: config.AI_LLM_EMBEDDING_MODEL,
|
||||
dims: vectors.map((v) => v.length),
|
||||
},
|
||||
"Embedding response had inconsistent/empty dimensions — treating as failure",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
return vectors;
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Embedding request failed — semantic cache disabled for this call",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Embed a single text; returns null on failure. */
|
||||
export async function embedText(text: string): Promise<number[] | null> {
|
||||
const vectors = await embedTexts([text]);
|
||||
return vectors?.[0] ?? null;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Similarity
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Cosine similarity between two equal-length vectors. */
|
||||
export function cosineSimilarity(a: number[], b: number[]): number {
|
||||
if (a.length === 0 || a.length !== b.length) return 0;
|
||||
let dot = 0;
|
||||
let normA = 0;
|
||||
let normB = 0;
|
||||
for (let i = 0; i < a.length; i++) {
|
||||
dot += a[i] * b[i];
|
||||
normA += a[i] * a[i];
|
||||
normB += b[i] * b[i];
|
||||
}
|
||||
if (normA === 0 || normB === 0) return 0;
|
||||
return dot / (Math.sqrt(normA) * Math.sqrt(normB));
|
||||
}
|
||||
|
||||
/**
|
||||
* Pick the best match above `minSimilarity`, or null.
|
||||
* Returns { index, similarity } relative to the candidates array.
|
||||
*/
|
||||
export function findBestEmbeddingMatch(
|
||||
vector: number[],
|
||||
candidates: number[][],
|
||||
minSimilarity: number,
|
||||
): { index: number; similarity: number } | null {
|
||||
let bestIndex = -1;
|
||||
let bestSimilarity = minSimilarity;
|
||||
for (let i = 0; i < candidates.length; i++) {
|
||||
const sim = cosineSimilarity(vector, candidates[i]);
|
||||
if (sim > bestSimilarity) {
|
||||
bestSimilarity = sim;
|
||||
bestIndex = i;
|
||||
}
|
||||
}
|
||||
return bestIndex >= 0
|
||||
? { index: bestIndex, similarity: bestSimilarity }
|
||||
: null;
|
||||
}
|
||||
@@ -14,10 +14,12 @@ export {
|
||||
downloadAndExtractFrame,
|
||||
sniffImageMimeType,
|
||||
} from "./mediaDownloader.js";
|
||||
export type {
|
||||
MessageImagePart,
|
||||
PreparedMediaMessage,
|
||||
} from "./visionAnalyzer.js";
|
||||
export {
|
||||
analyzeSingleMediaImage,
|
||||
hasMediaContent,
|
||||
MessageImagePart,
|
||||
PreparedMediaMessage,
|
||||
prepareMediaMessage,
|
||||
} from "./visionAnalyzer.js";
|
||||
|
||||
@@ -16,25 +16,19 @@ import type {
|
||||
MessageRecord,
|
||||
} from "../message-capture/types.js";
|
||||
import { initCacheStore } from "./cacheStore.js";
|
||||
import { embedTexts, isEmbeddingEnabled } from "./embeddingClient.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
import { runMediaBatch } from "./mediaBatchProcessor.js";
|
||||
import { isQdrantConfigured, searchQdrantBatch } from "./qdrantClient.js";
|
||||
import { logCacheEvent } from "./responseLogger.js";
|
||||
import { runTextOnlyBatch } from "./textBatchProcessor.js";
|
||||
import {
|
||||
bumpTextModerationHitCounts,
|
||||
ERROR_ARTIFACT_FLAGS,
|
||||
findSimilarTextModeration,
|
||||
getCachedTextModerations,
|
||||
isGloballyReusableCleanVerdict,
|
||||
isSemanticBandAccepted,
|
||||
makeModerationContextKey,
|
||||
makeTextModerationCacheKey,
|
||||
parseQdrantVerdict,
|
||||
type StoredModerationVerdict,
|
||||
setCachedTextModeration,
|
||||
upsertBareKeyToQdrant,
|
||||
} from "./textCacheStore.js";
|
||||
|
||||
const log = createChildLogger("moderationOrchestrator");
|
||||
@@ -74,12 +68,9 @@ export interface ModerationOutput {
|
||||
* Runs LLM-based moderation analysis on messages.
|
||||
* Splits text-only vs media, runs both paths in parallel, applies caching.
|
||||
*
|
||||
* Cache strategy (two-phase, batched):
|
||||
* 1. Exact-hash lookups (no API) — key is content + conversation context
|
||||
* (channel/thread) because LLM verdicts depend on context.
|
||||
* 2. Semantic near-duplicate lookup — ONE embeddings call for all uncached
|
||||
* text targets, then ONE Qdrant batch search (index-aligned), instead of
|
||||
* N sequential embed→search round-trips.
|
||||
* Cache strategy: exact-hash lookups only (no API) — key is content +
|
||||
* conversation context (channel/thread) because LLM verdicts depend on
|
||||
* context. Every miss goes to the LLM.
|
||||
*/
|
||||
export async function runModerationAnalysis(
|
||||
input: ModerationInput,
|
||||
@@ -100,9 +91,6 @@ export async function runModerationAnalysis(
|
||||
const uncachedTargets: MessageRecord[] = [];
|
||||
// cacheKey → representative result for identical-content dedupe
|
||||
const hitByKey = new Map<string, AnalysisResult>();
|
||||
// Embedding per exact cache key — computed once during lookup, reused
|
||||
// when the fresh LLM verdict is written back to the semantic cache.
|
||||
const embeddingsByKey = new Map<string, number[]>();
|
||||
|
||||
interface ExactCandidate {
|
||||
target: MessageRecord;
|
||||
@@ -251,142 +239,6 @@ export async function runModerationAnalysis(
|
||||
}
|
||||
incrementCounterBy("moderation_cache_misses", uncachedTargets.length);
|
||||
|
||||
// ── Phase 2: semantic cache — batched (one embed call + one Qdrant
|
||||
// batch search for ALL uncached text targets) ─────────────────────────
|
||||
if (isEmbeddingEnabled()) {
|
||||
const semanticCandidates = uncachedTargets
|
||||
.map((t) => ({
|
||||
target: t,
|
||||
cacheKey: makeTextModerationCacheKey(
|
||||
t.edited_content ?? t.content,
|
||||
makeModerationContextKey(t),
|
||||
),
|
||||
}))
|
||||
.filter(({ target }) => {
|
||||
const raw = (target.edited_content ?? target.content).trim();
|
||||
if (raw.length < 5) return false;
|
||||
if (hasMediaContent(target, attachments)) return false;
|
||||
return !hitByKey.has(
|
||||
makeTextModerationCacheKey(raw, makeModerationContextKey(target)),
|
||||
);
|
||||
});
|
||||
|
||||
if (semanticCandidates.length > 0) {
|
||||
const texts = semanticCandidates.map(
|
||||
({ target }) => target.edited_content ?? target.content,
|
||||
);
|
||||
const embeddings = await embedTexts(texts);
|
||||
if (embeddings && embeddings.length === texts.length) {
|
||||
// index-aligned with semanticCandidates
|
||||
for (let i = 0; i < semanticCandidates.length; i++) {
|
||||
const { cacheKey } = semanticCandidates[i];
|
||||
embeddingsByKey.set(cacheKey, embeddings[i]);
|
||||
}
|
||||
|
||||
if (isQdrantConfigured()) {
|
||||
// ONE batch search at the LOOSER threshold; per-hit re-classification
|
||||
// enforces the strict band for actionable verdicts.
|
||||
const batchHits = await searchQdrantBatch(
|
||||
embeddings,
|
||||
config.AI_LLM_EMBEDDING_MAX_CANDIDATES,
|
||||
config.AI_LLM_EMBEDDING_MIN_SIMILARITY_CLEAN,
|
||||
);
|
||||
for (let i = 0; i < semanticCandidates.length; i++) {
|
||||
const { target, cacheKey } = semanticCandidates[i];
|
||||
const hits = batchHits[i] ?? [];
|
||||
if (hits.length === 0) continue;
|
||||
const verdict = parseQdrantVerdict(hits[0].payload, hits[0].score);
|
||||
if (!verdict) continue;
|
||||
if (!isSemanticBandAccepted(verdict, verdict.similarity)) continue;
|
||||
log.debug(
|
||||
{
|
||||
messageId: target.id,
|
||||
similarity: Number(verdict.similarity.toFixed(4)),
|
||||
status: verdict.status,
|
||||
},
|
||||
"Semantic moderation cache hit — reusing stored verdict",
|
||||
);
|
||||
const hit: AnalysisResult = {
|
||||
messageId: target.id,
|
||||
status: verdict.status,
|
||||
flags: verdict.flags,
|
||||
score: verdict.score,
|
||||
analysis: verdict.analysis,
|
||||
categories: verdict.categories,
|
||||
severity: verdict.severity as AnalysisResult["severity"],
|
||||
confidence: verdict.confidence,
|
||||
recommendedAction:
|
||||
verdict.recommendedAction as AnalysisResult["recommendedAction"],
|
||||
policyVersion: "semantic-cache-2026-07",
|
||||
evidence: [],
|
||||
};
|
||||
cacheHits.push(hit);
|
||||
hitByKey.set(cacheKey, hit);
|
||||
servedCacheKeys.add(cacheKey); // bump hit_count for metrics
|
||||
logCacheEvent("hit", cacheKey, "text");
|
||||
incrementCounterBy("moderation_cache_hits", 1, {
|
||||
type: "semantic-qdrant",
|
||||
});
|
||||
}
|
||||
} else {
|
||||
// Legacy Postgres fallback path (no Qdrant): per-candidate scan.
|
||||
for (let i = 0; i < semanticCandidates.length; i++) {
|
||||
const { target, cacheKey } = semanticCandidates[i];
|
||||
const semantic = await findSimilarTextModeration(
|
||||
embeddings[i],
|
||||
config.AI_LLM_EMBEDDING_MIN_SIMILARITY_CLEAN,
|
||||
config.AI_LLM_EMBEDDING_MAX_CANDIDATES,
|
||||
);
|
||||
if (!semantic) continue;
|
||||
if (!isSemanticBandAccepted(semantic, semantic.similarity))
|
||||
continue;
|
||||
log.debug(
|
||||
{
|
||||
messageId: target.id,
|
||||
similarity: Number(semantic.similarity.toFixed(4)),
|
||||
status: semantic.status,
|
||||
},
|
||||
"Semantic moderation cache hit (PG fallback) — reusing stored verdict",
|
||||
);
|
||||
const hit: AnalysisResult = {
|
||||
messageId: target.id,
|
||||
status: semantic.status,
|
||||
flags: semantic.flags,
|
||||
score: semantic.score,
|
||||
analysis: semantic.analysis,
|
||||
categories: semantic.categories,
|
||||
severity: semantic.severity as AnalysisResult["severity"],
|
||||
confidence: semantic.confidence,
|
||||
recommendedAction:
|
||||
semantic.recommendedAction as AnalysisResult["recommendedAction"],
|
||||
policyVersion: "semantic-cache-2026-07",
|
||||
evidence: [],
|
||||
};
|
||||
cacheHits.push(hit);
|
||||
hitByKey.set(cacheKey, hit);
|
||||
servedCacheKeys.add(cacheKey); // bump hit_count for metrics
|
||||
logCacheEvent("hit", cacheKey, "text");
|
||||
incrementCounterBy("moderation_cache_hits", 1, {
|
||||
type: "semantic-pg",
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Drop semantic hits from the LLM work queue.
|
||||
for (let i = uncachedTargets.length - 1; i >= 0; i--) {
|
||||
const t = uncachedTargets[i];
|
||||
const key = makeTextModerationCacheKey(
|
||||
t.edited_content ?? t.content,
|
||||
makeModerationContextKey(t),
|
||||
);
|
||||
if (hitByKey.has(key)) {
|
||||
uncachedTargets.splice(i, 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (cacheHits.length > 0) {
|
||||
// Metrics: one bulk UPDATE for every exact-cache key actually served.
|
||||
bumpTextModerationHitCounts(Array.from(servedCacheKeys));
|
||||
@@ -466,11 +318,7 @@ export async function runModerationAnalysis(
|
||||
recommendedAction: result.recommendedAction ?? "none",
|
||||
status: result.status,
|
||||
};
|
||||
setCachedTextModeration(
|
||||
cacheKey,
|
||||
stored,
|
||||
embeddingsByKey.get(cacheKey),
|
||||
).catch((err: unknown) => {
|
||||
setCachedTextModeration(cacheKey, stored).catch((err: unknown) => {
|
||||
log.warn({ cacheKey, error: String(err) }, "Cache write failed");
|
||||
});
|
||||
|
||||
@@ -479,13 +327,6 @@ export async function runModerationAnalysis(
|
||||
// under the context-free bare key so repeats in OTHER channels hit the
|
||||
// exact cache instead of paying a new LLM call. Same guard as the read
|
||||
// path — only non-actionable clean verdicts may cross channels.
|
||||
//
|
||||
// 2026-08-25 cache-hit fix: the bare key is ALSO upserted to Qdrant
|
||||
// (via upsertBareKeyToQdrant) with the SAME embedding already computed
|
||||
// at lookup time. Previously the bare key was only PG-written with
|
||||
// embedding=null — bare clean verdicts were DB-only and invisible to
|
||||
// searchQdrantBatch, capping the semantic hit-rate below the exact-cache
|
||||
// hit-rate for cross-channel repeats.
|
||||
const bareKey = makeTextModerationCacheKey(rawContent);
|
||||
if (
|
||||
bareKey !== cacheKey &&
|
||||
@@ -505,11 +346,7 @@ export async function runModerationAnalysis(
|
||||
)
|
||||
) {
|
||||
globalBareKeysWritten.set(bareKey, true);
|
||||
setCachedTextModeration(bareKey, stored, null).catch(() => {});
|
||||
const bareEmbedding = embeddingsByKey.get(cacheKey);
|
||||
if (bareEmbedding && bareEmbedding.length > 0) {
|
||||
upsertBareKeyToQdrant(bareKey, stored, bareEmbedding).catch(() => {});
|
||||
}
|
||||
setCachedTextModeration(bareKey, stored).catch(() => {});
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,589 +0,0 @@
|
||||
/**
|
||||
* qdrantClient.ts
|
||||
*
|
||||
* Minimal Qdrant REST client (zero dependencies, fetch-based) used by the
|
||||
* semantic moderation cache. Embedding vectors + verdict payloads live in
|
||||
* Qdrant instead of the Postgres `embedding` column (legacy, kept for
|
||||
* backward-compatible fallback reads).
|
||||
*
|
||||
* All functions degrade gracefully: failures return null / empty results so
|
||||
* callers fall back to the LLM — moderation quality is never reduced.
|
||||
*/
|
||||
|
||||
import { createHash } from "node:crypto";
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
|
||||
const log = createChildLogger("qdrant");
|
||||
|
||||
// ensureQdrantCollection performs a network round-trip (GET, possibly
|
||||
// DELETE+PUT). Running it on every upsert adds 1-3 HTTP calls per
|
||||
// moderation verdict, which under Qdrant load pushes the upsert past the
|
||||
// request timeout and aborts it ("This operation was aborted"). Memoise the
|
||||
// result so the collection is only verified once per process lifetime.
|
||||
let ensureCollectionPromise: Promise<boolean> | null = null;
|
||||
|
||||
/** Reset the memoised ensure result (used by tests / config reload). */
|
||||
export function resetQdrantCollectionCache(): void {
|
||||
ensureCollectionPromise = null;
|
||||
}
|
||||
|
||||
export interface QdrantVerdictPayload {
|
||||
text: string;
|
||||
flags: string; // JSON string of the full moderation result
|
||||
analyzed_at: number;
|
||||
expires_at: number;
|
||||
/** Bare content hash (16 hex chars) — enables content-based invalidation
|
||||
* regardless of the (context-scoped) point id. */
|
||||
content_hash?: string;
|
||||
// Archive metadata (archiveEmbedder writes these; moderation cache doesn't).
|
||||
// Optional so the same payload shape serves both the moderation cache
|
||||
// collection and the public message archive collection.
|
||||
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 collectionName(): string {
|
||||
return config.QDRANT_COLLECTION ?? "gmw_text_moderation";
|
||||
}
|
||||
|
||||
/** Persistent archive collection for semantic message search (no TTL). */
|
||||
export const ARCHIVE_COLLECTION =
|
||||
config.QDRANT_ARCHIVE_COLLECTION ?? "gmw_message_archive";
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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();
|
||||
let json: unknown = null;
|
||||
try {
|
||||
json = text ? JSON.parse(text) : null;
|
||||
} catch {
|
||||
json = null;
|
||||
}
|
||||
if (!res.ok) {
|
||||
throw new Error(
|
||||
`Qdrant ${method} ${path} -> ${res.status}: ${text.slice(0, 200)}`,
|
||||
);
|
||||
}
|
||||
return json;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
/** True for transient errors worth retrying (408 timeout, ECONNRESET, aborts). */
|
||||
function isTransientQdrantError(error: unknown): boolean {
|
||||
if (error instanceof Error && error.name === "AbortError") return true;
|
||||
const msg = error instanceof Error ? error.message : String(error);
|
||||
return (
|
||||
msg.includes("-> 408") ||
|
||||
msg.includes("aborted") ||
|
||||
msg.includes("ECONNRESET") ||
|
||||
msg.includes("ETIMEDOUT") ||
|
||||
msg.includes("fetch failed")
|
||||
);
|
||||
}
|
||||
|
||||
/** Retry with exponential backoff around `request`, abort-aware. */
|
||||
async function requestWithRetry(
|
||||
method: string,
|
||||
path: string,
|
||||
body?: unknown,
|
||||
timeoutMs = 10_000,
|
||||
attempts = 3,
|
||||
): Promise<unknown> {
|
||||
let lastError: unknown;
|
||||
for (let attempt = 0; attempt < attempts; attempt++) {
|
||||
if (attempt > 0) {
|
||||
// Exponential backoff: 500ms → 1s → 2s (jittered ±20%).
|
||||
const base = 500 * 2 ** (attempt - 1);
|
||||
const delayMs = base + Math.floor(Math.random() * base * 0.2);
|
||||
await new Promise((resolve) => setTimeout(resolve, delayMs));
|
||||
}
|
||||
try {
|
||||
return await request(method, path, body, timeoutMs);
|
||||
} catch (error) {
|
||||
lastError = error;
|
||||
if (!isTransientQdrantError(error)) throw error;
|
||||
log.debug(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
attempt: attempt + 1,
|
||||
method,
|
||||
path,
|
||||
},
|
||||
"Qdrant transient error — retrying with backoff",
|
||||
);
|
||||
}
|
||||
}
|
||||
throw lastError;
|
||||
}
|
||||
|
||||
/** Deterministic uint64 point id from the exact-hash cache key. */
|
||||
export function qdrantPointId(cacheKey: string): number {
|
||||
const digest = createHash("sha256").update(cacheKey).digest();
|
||||
// First 8 bytes as BigInt, then clamp into Qdrant's uint64 space.
|
||||
const big = digest.readBigUInt64BE(0);
|
||||
return Number(big & 0x7fffffffffffffffn);
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensure the collection exists with the right vector size. If the size
|
||||
* changed (embedding model swapped), recreate — stale vectors are useless
|
||||
* anyway and cosine scores would be meaningless across dimensions.
|
||||
*/
|
||||
export async function ensureQdrantCollection(
|
||||
vectorSize: number,
|
||||
): Promise<boolean> {
|
||||
if (ensureCollectionPromise) return ensureCollectionPromise;
|
||||
ensureCollectionPromise = (async () => {
|
||||
try {
|
||||
// 404 = collection doesn't exist yet → create it.
|
||||
let existing: {
|
||||
result?: { config?: { params?: { vectors?: { size?: number } } } };
|
||||
} | null = null;
|
||||
try {
|
||||
existing = (await request(
|
||||
"GET",
|
||||
`/collections/${collectionName()}`,
|
||||
)) as {
|
||||
result?: { config?: { params?: { vectors?: { size?: number } } } };
|
||||
};
|
||||
} catch (error) {
|
||||
if (!(error instanceof Error) || !error.message.includes("-> 404")) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
const size = existing?.result?.config?.params?.vectors?.size;
|
||||
if (size === vectorSize) return true;
|
||||
|
||||
if (size !== undefined && size !== vectorSize) {
|
||||
log.warn(
|
||||
{ collection: collectionName(), oldSize: size, newSize: vectorSize },
|
||||
"Qdrant collection vector size changed — recreating collection",
|
||||
);
|
||||
await request("DELETE", `/collections/${collectionName()}`);
|
||||
}
|
||||
|
||||
await request("PUT", `/collections/${collectionName()}`, {
|
||||
vectors: { size: vectorSize, distance: "Cosine" },
|
||||
});
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.error(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
collection: collectionName(),
|
||||
},
|
||||
"Failed to ensure Qdrant collection",
|
||||
);
|
||||
// The memoised promise is permanently sticky: once it rejects (e.g. a
|
||||
// recreate aborted mid-way — DELETE done, PUT failed), every later call
|
||||
// returns the same rejected promise and the collection is never
|
||||
// re-created until process restart. Reset so the next call retries.
|
||||
ensureCollectionPromise = null;
|
||||
return false;
|
||||
}
|
||||
})();
|
||||
return ensureCollectionPromise;
|
||||
}
|
||||
|
||||
/** Upsert one embedding + verdict payload point. Returns false on failure. */
|
||||
export async function upsertQdrantPoint(
|
||||
cacheKey: string,
|
||||
vector: number[],
|
||||
payload: QdrantVerdictPayload,
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
if (!(await ensureQdrantCollection(vector.length))) return false;
|
||||
await requestWithRetry(
|
||||
"PUT",
|
||||
`/collections/${collectionName()}/points`,
|
||||
{
|
||||
points: [{ id: qdrantPointId(cacheKey), vector, payload }],
|
||||
wait: true,
|
||||
},
|
||||
30_000,
|
||||
3,
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant upsert failed — semantic entry skipped",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
export interface QdrantSearchHit {
|
||||
cacheKey: string;
|
||||
score: number;
|
||||
payload: QdrantVerdictPayload;
|
||||
}
|
||||
|
||||
/**
|
||||
* Search the nearest stored vector. Returns hits sorted by score desc,
|
||||
* filtered to unexpired payloads. Empty array on failure.
|
||||
*/
|
||||
export async function searchQdrant(
|
||||
vector: number[],
|
||||
limit: number,
|
||||
scoreThreshold: number,
|
||||
): Promise<QdrantSearchHit[]> {
|
||||
try {
|
||||
const json = (await request(
|
||||
"POST",
|
||||
`/collections/${collectionName()}/points/search`,
|
||||
{
|
||||
vector,
|
||||
limit,
|
||||
score_threshold: scoreThreshold,
|
||||
with_payload: true,
|
||||
filter: {
|
||||
must: [
|
||||
{
|
||||
key: "expires_at",
|
||||
range: { gte: Date.now() },
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
)) as {
|
||||
result?: Array<{
|
||||
id?: number;
|
||||
score?: number;
|
||||
payload?: QdrantVerdictPayload;
|
||||
}>;
|
||||
};
|
||||
|
||||
return (json.result ?? [])
|
||||
.filter((hit) => hit.payload?.flags)
|
||||
.map((hit) => ({
|
||||
cacheKey: `qdrant:${hit.id ?? "?"}`,
|
||||
score: hit.score ?? 0,
|
||||
payload: hit.payload as QdrantVerdictPayload,
|
||||
}));
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant search failed — semantic cache skipped",
|
||||
);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Batch search: one HTTP round-trip for N vectors (Qdrant
|
||||
* `/points/search/batch`). Result is index-aligned with `vectors` — each
|
||||
* entry is the top hits for that vector (or [] on per-vector failure).
|
||||
* Used by the orchestrator to avoid N sequential embed→search round-trips.
|
||||
*/
|
||||
export async function searchQdrantBatch(
|
||||
vectors: number[][],
|
||||
limit: number,
|
||||
scoreThreshold: number,
|
||||
): Promise<QdrantSearchHit[][]> {
|
||||
if (vectors.length === 0) return [];
|
||||
try {
|
||||
const json = (await request(
|
||||
"POST",
|
||||
`/collections/${collectionName()}/points/search/batch`,
|
||||
{
|
||||
searches: vectors.map((vector) => ({
|
||||
vector,
|
||||
limit,
|
||||
score_threshold: scoreThreshold,
|
||||
with_payload: true,
|
||||
filter: {
|
||||
must: [
|
||||
{
|
||||
key: "expires_at",
|
||||
range: { gte: Date.now() },
|
||||
},
|
||||
],
|
||||
},
|
||||
})),
|
||||
},
|
||||
)) as {
|
||||
result?: Array<{
|
||||
result?: Array<{
|
||||
id?: number;
|
||||
score?: number;
|
||||
payload?: QdrantVerdictPayload;
|
||||
}>;
|
||||
}>;
|
||||
};
|
||||
|
||||
return (json.result ?? []).map((entry) =>
|
||||
(entry.result ?? [])
|
||||
.filter((hit) => hit.payload?.flags)
|
||||
.map((hit) => ({
|
||||
cacheKey: `qdrant:${hit.id ?? "?"}`,
|
||||
score: hit.score ?? 0,
|
||||
payload: hit.payload as QdrantVerdictPayload,
|
||||
})),
|
||||
);
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant batch search failed — semantic cache skipped",
|
||||
);
|
||||
return vectors.map(() => []);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete expired verdict points from the collection. Best-effort: 404
|
||||
* (collection missing) and failures are swallowed — the periodic pruner
|
||||
* just retries next sweep.
|
||||
*/
|
||||
export async function deleteExpiredQdrantPoints(): Promise<number> {
|
||||
try {
|
||||
const json = (await request(
|
||||
"POST",
|
||||
`/collections/${collectionName()}/points/delete`,
|
||||
{
|
||||
filter: {
|
||||
must: [
|
||||
{
|
||||
key: "expires_at",
|
||||
range: { lt: Date.now() },
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
)) as { result?: { deleted?: number } | null };
|
||||
|
||||
return json.result?.deleted ?? 0;
|
||||
} catch (error) {
|
||||
if (error instanceof Error && error.message.includes("-> 404")) {
|
||||
log.debug({}, "Qdrant collection absent — nothing to prune");
|
||||
} else {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant expired-point prune failed",
|
||||
);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete the verdict point for an exact cache key (used by cache
|
||||
* invalidation when a moderator corrects a verdict).
|
||||
*/
|
||||
export async function deleteQdrantPoint(cacheKey: string): Promise<boolean> {
|
||||
try {
|
||||
await request("POST", `/collections/${collectionName()}/points/delete`, {
|
||||
points: [qdrantPointId(cacheKey)],
|
||||
});
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant point delete failed",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete all verdict points whose payload carries a given bare content hash.
|
||||
* Used by cache invalidation for corrected verdicts — matches context-scoped
|
||||
* points that share the same content regardless of their point ids.
|
||||
*/
|
||||
export async function deleteQdrantPointsByContentHash(
|
||||
bareHash: string,
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
await request("POST", `/collections/${collectionName()}/points/delete`, {
|
||||
filter: {
|
||||
must: [
|
||||
{
|
||||
key: "content_hash",
|
||||
match: { value: bareHash },
|
||||
},
|
||||
],
|
||||
},
|
||||
});
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Qdrant content-hash point delete failed",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** True when Qdrant is configured (non-empty URL). */
|
||||
export function isQdrantConfigured(): boolean {
|
||||
return Boolean(config.QDRANT_URL);
|
||||
}
|
||||
|
||||
// ─── Archive variants (collection-aware, for persistent message search) ───
|
||||
// These mirror the cache functions but take an explicit collection name so the
|
||||
// semantic-search archive (gmw_message_archive) can live alongside the
|
||||
// TTL-bounded automod cache without disturbing it.
|
||||
|
||||
/** Ensure an arbitrary collection exists with the right vector size. */
|
||||
export async function ensureQdrantCollectionV2(
|
||||
name: string,
|
||||
vectorSize: number,
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
let existing: {
|
||||
result?: { config?: { params?: { vectors?: { size?: number } } } };
|
||||
} | null = null;
|
||||
try {
|
||||
existing = (await request("GET", `/collections/${name}`)) as {
|
||||
result?: { config?: { params?: { vectors?: { size?: number } } } };
|
||||
} | null;
|
||||
} catch (error) {
|
||||
if (!(error instanceof Error) || !error.message.includes("-> 404")) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
const size = existing?.result?.config?.params?.vectors?.size;
|
||||
if (size === vectorSize) return true;
|
||||
|
||||
if (size !== undefined && size !== vectorSize) {
|
||||
log.warn(
|
||||
{ collection: name, oldSize: size, newSize: vectorSize },
|
||||
"Qdrant archive collection vector size changed — recreating collection",
|
||||
);
|
||||
await request("DELETE", `/collections/${name}`);
|
||||
}
|
||||
|
||||
await request("PUT", `/collections/${name}`, {
|
||||
vectors: { size: vectorSize, distance: "Cosine" },
|
||||
});
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.error(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
collection: name,
|
||||
},
|
||||
"Failed to ensure Qdrant archive collection",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** Upsert one embedding + payload point into a named collection. */
|
||||
export async function upsertQdrantPointV2(
|
||||
name: string,
|
||||
pointId: number,
|
||||
vector: number[],
|
||||
payload: QdrantVerdictPayload,
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
if (!(await ensureQdrantCollectionV2(name, vector.length))) return false;
|
||||
await requestWithRetry(
|
||||
"PUT",
|
||||
`/collections/${name}/points`,
|
||||
{
|
||||
points: [{ id: pointId, vector, payload }],
|
||||
wait: true,
|
||||
},
|
||||
30_000,
|
||||
3,
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
collection: name,
|
||||
} as Record<string, unknown>,
|
||||
"Qdrant archive upsert failed — entry skipped",
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
export interface QdrantArchiveHit {
|
||||
pointId: number;
|
||||
score: number;
|
||||
payload: QdrantVerdictPayload;
|
||||
}
|
||||
|
||||
/** Search a named collection for the nearest stored vector. */
|
||||
export async function searchQdrantV2(
|
||||
name: string,
|
||||
vector: number[],
|
||||
limit: number,
|
||||
scoreThreshold: number,
|
||||
): Promise<QdrantArchiveHit[]> {
|
||||
try {
|
||||
const json = (await request("POST", `/collections/${name}/points/search`, {
|
||||
vector,
|
||||
limit,
|
||||
score_threshold: scoreThreshold,
|
||||
with_payload: true,
|
||||
})) as {
|
||||
result?: Array<{
|
||||
id?: number;
|
||||
score?: number;
|
||||
payload?: QdrantVerdictPayload;
|
||||
}>;
|
||||
};
|
||||
|
||||
return (json.result ?? [])
|
||||
.filter((hit) => hit.payload?.text)
|
||||
.map((hit) => ({
|
||||
pointId: hit.id ?? 0,
|
||||
score: hit.score ?? 0,
|
||||
payload: hit.payload as QdrantVerdictPayload,
|
||||
}));
|
||||
} catch (error) {
|
||||
log.warn(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
collection: name,
|
||||
} as Record<string, unknown>,
|
||||
"Qdrant archive search failed — semantic search skipped",
|
||||
);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
@@ -2,15 +2,6 @@ import { createHash } from "node:crypto";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { executeAll, executeGet } from "../../shared/database/drizzle.js";
|
||||
import { findBestEmbeddingMatch } from "./embeddingClient.js";
|
||||
import {
|
||||
deleteQdrantPoint,
|
||||
deleteQdrantPointsByContentHash,
|
||||
isQdrantConfigured,
|
||||
type QdrantVerdictPayload,
|
||||
searchQdrant,
|
||||
upsertQdrantPoint,
|
||||
} from "./qdrantClient.js";
|
||||
|
||||
const logger = createChildLogger("text-cache-store");
|
||||
|
||||
@@ -263,8 +254,8 @@ export function makeModerationContextKey(message: {
|
||||
|
||||
/**
|
||||
* Invalidate cached moderation verdicts for a piece of content: removes
|
||||
* matching Postgres rows AND Qdrant points. Called when a moderator
|
||||
* corrects a verdict so a stale/wrong cached decision cannot resurface.
|
||||
* matching Postgres rows. Called when a moderator corrects a verdict so a
|
||||
* stale/wrong cached decision cannot resurface.
|
||||
*
|
||||
* Handles both key formats:
|
||||
* - legacy `text_mod:<hash>` (content-only, pre-context keys)
|
||||
@@ -277,20 +268,10 @@ export async function invalidateTextModerationCache(
|
||||
.update(content)
|
||||
.digest("hex")
|
||||
.slice(0, 16);
|
||||
const legacyKey = `text_mod:${bareHash}`;
|
||||
|
||||
const queries: Promise<unknown>[] = [
|
||||
executeAll(`DELETE FROM text_analysis_cache WHERE text LIKE $1`, [
|
||||
`text_mod:%${bareHash}`,
|
||||
]).catch(() => {}),
|
||||
];
|
||||
if (isQdrantConfigured()) {
|
||||
queries.push(
|
||||
deleteQdrantPoint(legacyKey).catch(() => {}),
|
||||
deleteQdrantPointsByContentHash(bareHash).catch(() => {}),
|
||||
);
|
||||
}
|
||||
await Promise.all(queries).catch(() => {});
|
||||
await executeAll(`DELETE FROM text_analysis_cache WHERE text LIKE $1`, [
|
||||
`text_mod:%${bareHash}`,
|
||||
]).catch(() => {});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -488,31 +469,6 @@ export function bumpTextModerationHitCounts(cacheKeys: string[]): void {
|
||||
).catch(() => {});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Semantic two-band acceptance
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* True when a semantic-cache hit may be reused given its verdict class.
|
||||
* Two bands (2026-08-24): non-actionable verdicts (clean / flagless /
|
||||
* action=none) are accepted from the LOOSER clean band; actionable verdicts
|
||||
* (warn/flagged or any flags/action) keep the strict historical gate.
|
||||
* Between the bands → reject → the message falls through to the LLM
|
||||
* (fail-open toward accuracy).
|
||||
*/
|
||||
export function isSemanticBandAccepted(
|
||||
verdict: StoredModerationVerdict,
|
||||
similarity: number,
|
||||
): boolean {
|
||||
const isNonActionable =
|
||||
verdict.status === "clean" &&
|
||||
verdict.flags.length === 0 &&
|
||||
(verdict.recommendedAction ?? "none") === "none";
|
||||
return isNonActionable
|
||||
? similarity >= config.AI_LLM_EMBEDDING_MIN_SIMILARITY_CLEAN
|
||||
: similarity >= config.AI_LLM_EMBEDDING_MIN_SIMILARITY;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Global exact-cache reuse guard (context-free fallback)
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -542,181 +498,9 @@ export function isGloballyReusableCleanVerdict(
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a Qdrant verdict payload into the result shape shared by the
|
||||
* semantic cache lookups. Returns null on malformed payloads (callers then
|
||||
* fall through to the LLM).
|
||||
*/
|
||||
export function parseQdrantVerdict(
|
||||
payload: QdrantVerdictPayload,
|
||||
similarity: number,
|
||||
):
|
||||
| (StoredModerationVerdict & {
|
||||
text: string;
|
||||
similarity: number;
|
||||
})
|
||||
| null {
|
||||
const parsed = parseStoredVerdictRow({ flags: payload.flags });
|
||||
if (!parsed) return null;
|
||||
|
||||
return {
|
||||
...parsed,
|
||||
text: payload.text,
|
||||
similarity,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Semantic moderation cache lookup.
|
||||
*
|
||||
* Primary: Qdrant vector search (when QDRANT_URL configured) — nearest
|
||||
* unexpired verdict above `minSimilarity`. Fallback: Postgres embedding
|
||||
* column (legacy rows written before Qdrant was wired in).
|
||||
* Returns null on no match or any failure — callers then proceed to the LLM.
|
||||
*/
|
||||
export async function findSimilarTextModeration(
|
||||
embedding: number[],
|
||||
minSimilarity: number,
|
||||
limit: number,
|
||||
): Promise<
|
||||
(StoredModerationVerdict & { text: string; similarity: number }) | null
|
||||
> {
|
||||
// Qdrant path (primary)
|
||||
if (isQdrantConfigured()) {
|
||||
const hits = await searchQdrant(embedding, limit, minSimilarity);
|
||||
if (hits.length > 0) {
|
||||
const hit = hits[0];
|
||||
return parseQdrantVerdict(hit.payload, hit.score);
|
||||
}
|
||||
// No Qdrant hit — fall through to Postgres legacy rows.
|
||||
}
|
||||
|
||||
try {
|
||||
const rows = await executeAll(
|
||||
`SELECT text, flags, embedding
|
||||
FROM text_analysis_cache
|
||||
WHERE source = 'user_moderation'
|
||||
AND embedding IS NOT NULL
|
||||
AND expires_at > $1
|
||||
ORDER BY analyzed_at DESC
|
||||
LIMIT $2`,
|
||||
[Date.now(), limit],
|
||||
);
|
||||
if (!rows || rows.length === 0) return null;
|
||||
|
||||
const candidates = rows.flatMap((row) => {
|
||||
let embeddingArr: number[] = [];
|
||||
let parsed: Record<string, unknown>;
|
||||
try {
|
||||
embeddingArr = JSON.parse(row.embedding) as number[];
|
||||
parsed = JSON.parse(row.flags) as Record<string, unknown>;
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
// Skip processing locks / malformed entries — never reuse an
|
||||
// in-flight or non-verdict row.
|
||||
const storedStatus = parsed.status as string | undefined;
|
||||
if (storedStatus === "processing" || storedStatus === undefined) {
|
||||
return [];
|
||||
}
|
||||
if (!Array.isArray(parsed.flags)) return [];
|
||||
return [
|
||||
{
|
||||
text: row.text,
|
||||
embedding: embeddingArr,
|
||||
parsed,
|
||||
},
|
||||
];
|
||||
});
|
||||
|
||||
const match = findBestEmbeddingMatch(
|
||||
embedding,
|
||||
candidates.map((c) => c.embedding),
|
||||
minSimilarity,
|
||||
);
|
||||
if (!match) return null;
|
||||
|
||||
const hit = candidates[match.index];
|
||||
const parsed = hit.parsed;
|
||||
const flags = (parsed.flags as string[]) ?? [];
|
||||
const status = normalizeStoredStatus(
|
||||
parsed.status as string | undefined,
|
||||
flags,
|
||||
);
|
||||
return {
|
||||
text: hit.text,
|
||||
similarity: match.similarity,
|
||||
status,
|
||||
flags,
|
||||
score: (parsed.score as number) ?? 0,
|
||||
analysis: (parsed.analysis as string) ?? "",
|
||||
categories: (parsed.categories as string[]) ?? [],
|
||||
severity: (parsed.severity as string) ?? "none",
|
||||
confidence: (parsed.confidence as number) ?? 0,
|
||||
recommendedAction: (parsed.recommendedAction as string) ?? "none",
|
||||
};
|
||||
} catch (error) {
|
||||
logger.error(
|
||||
{ error: error instanceof Error ? error.message : String(error) },
|
||||
"Failed semantic text moderation lookup",
|
||||
);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Upsert a bare (context-free) clean verdict to the Qdrant vector store,
|
||||
* making global-reuse clean verdicts discoverable by semantic search.
|
||||
*
|
||||
* Why: the main `setCachedTextModeration` writes bare-key rows to Postgres
|
||||
* with embedding=null (deliberate — no duplicate PG embedding column), but
|
||||
* a bare clean verdict that never reaches Qdrant is invisible to
|
||||
* searchQdrantBatch. So two messages with identical clean content in
|
||||
* DIFFERENT channels never match semantically — the semantic hit-rate is
|
||||
* capped below the exact-cache hit-rate. This helper shares the embedding
|
||||
* already computed at lookup time so the bare point is semantically
|
||||
* findable.
|
||||
*
|
||||
* Guard: only non-actionable clean verdicts qualify (same guard as the
|
||||
* read path and as the orchestrator's bare-key write-back). No-op when
|
||||
* Qdrant is disabled or no embedding is available.
|
||||
*/
|
||||
export async function upsertBareKeyToQdrant(
|
||||
bareKey: string,
|
||||
result: {
|
||||
status: string;
|
||||
flags: string[];
|
||||
score: number;
|
||||
analysis: string;
|
||||
categories: string[];
|
||||
severity: string;
|
||||
confidence: number;
|
||||
recommendedAction: string;
|
||||
},
|
||||
embedding: number[] | null | undefined,
|
||||
): Promise<void> {
|
||||
if (!isQdrantConfigured() || !embedding || embedding.length === 0) return;
|
||||
if (!isGloballyReusableCleanVerdict(result, undefined)) return;
|
||||
const now = Date.now();
|
||||
const USER_MOD_CACHE_TTL_MS = 24 * 60 * 60 * 1000;
|
||||
await upsertQdrantPoint(bareKey, embedding, {
|
||||
text: bareKey,
|
||||
flags: JSON.stringify(result),
|
||||
analyzed_at: now,
|
||||
expires_at: now + USER_MOD_CACHE_TTL_MS,
|
||||
content_hash: bareKey.split(":").pop() ?? "",
|
||||
}).catch((err: unknown) => {
|
||||
logger.error(
|
||||
{ error: err instanceof Error ? err.message : String(err), bareKey },
|
||||
"Failed to upsert bare-key clean verdict to Qdrant",
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Store a moderation result for a (user, content) pair.
|
||||
* The `flags` field stores the full result object as JSON.
|
||||
* `embedding` (optional) is stored for semantic near-duplicate lookups.
|
||||
*/
|
||||
export async function setCachedTextModeration(
|
||||
cacheKey: string,
|
||||
@@ -730,41 +514,25 @@ export async function setCachedTextModeration(
|
||||
recommendedAction: string;
|
||||
status?: "clean" | "warn" | "flagged" | "processing";
|
||||
},
|
||||
embedding?: number[] | null,
|
||||
): Promise<void> {
|
||||
const now = Date.now();
|
||||
const USER_MOD_CACHE_TTL_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
try {
|
||||
// Qdrant is the primary vector store when configured: upsert the point
|
||||
// with the verdict payload; skip the Postgres embedding column entirely.
|
||||
if (isQdrantConfigured() && embedding && embedding.length > 0) {
|
||||
await upsertQdrantPoint(cacheKey, embedding, {
|
||||
text: cacheKey,
|
||||
flags: JSON.stringify(result),
|
||||
analyzed_at: now,
|
||||
expires_at: now + USER_MOD_CACHE_TTL_MS,
|
||||
content_hash: cacheKey.split(":").pop() ?? "",
|
||||
});
|
||||
}
|
||||
|
||||
await executeAll(
|
||||
`INSERT INTO text_analysis_cache (text, flags, source, analyzed_at, expires_at, hit_count, embedding)
|
||||
VALUES ($1, $2, $3, $4, $5, 0, $6)
|
||||
`INSERT INTO text_analysis_cache (text, flags, source, analyzed_at, expires_at, hit_count)
|
||||
VALUES ($1, $2, $3, $4, $5, 0)
|
||||
ON CONFLICT (text) DO UPDATE SET
|
||||
flags = EXCLUDED.flags,
|
||||
source = EXCLUDED.source,
|
||||
analyzed_at = EXCLUDED.analyzed_at,
|
||||
expires_at = EXCLUDED.expires_at,
|
||||
embedding = COALESCE(EXCLUDED.embedding, text_analysis_cache.embedding)`,
|
||||
expires_at = EXCLUDED.expires_at`,
|
||||
[
|
||||
cacheKey,
|
||||
JSON.stringify(result),
|
||||
"user_moderation",
|
||||
now,
|
||||
now + USER_MOD_CACHE_TTL_MS,
|
||||
// Postgres embedding stays as legacy fallback; Qdrant is primary.
|
||||
embedding && embedding.length > 0 ? JSON.stringify(embedding) : null,
|
||||
],
|
||||
);
|
||||
} catch (error) {
|
||||
@@ -833,11 +601,10 @@ export async function getRecentCorrectedModerations(
|
||||
/**
|
||||
* Store a corrected moderation entry for future few-shot injection.
|
||||
*
|
||||
* Also invalidates any cached verdicts for the corrected content (both
|
||||
* Postgres rows and Qdrant points) so the corrected decision propagates
|
||||
* immediately instead of being shadowed by a stale cache entry. Full
|
||||
* content is looked up by message_id when available — more precise than
|
||||
* the (possibly truncated) snippet.
|
||||
* Also invalidates any cached verdicts for the corrected content so the
|
||||
* corrected decision propagates immediately instead of being shadowed by a
|
||||
* stale cache entry. Full content is looked up by message_id when available
|
||||
* — more precise than the (possibly truncated) snippet.
|
||||
*/
|
||||
export async function insertCorrectedModeration(entry: {
|
||||
messageId: string;
|
||||
|
||||
@@ -75,7 +75,7 @@ export const KNOWN_SAFE_TERMS = new Set(
|
||||
(
|
||||
"discord youtube google facebook instagram twitter tiktok whatsapp telegram netflix spotify steam github gitlab bitbucket chatgpt openai anthropic claude deepseek gemini llama copilot cursor vscode vscodium jetbrains intellij pycharm webstorm sublime codeblocks" +
|
||||
" docker kubernetes k8s linux ubuntu debian arch fedora manjaro kali windows macos android ios chrome firefox safari edge opera brave" +
|
||||
" react nextjs next vue svelte angular node nodejs deno bun pnpm yarn npm javascript typescript python golang go rust java kotlin swift cplusplus cpp css html json xml yaml toml regex backend frontend database mysql postgres postgresql mongodb redis qdrant sqlite nosql graphql rest websocket webhook" +
|
||||
" react nextjs next vue svelte angular node nodejs deno bun pnpm yarn npm javascript typescript python golang go rust java kotlin swift cplusplus cpp css html json xml yaml toml regex backend frontend database mysql postgres postgresql mongodb redis sqlite nosql graphql rest websocket webhook" +
|
||||
" bug crash error debug fix issue pr merge commit push pull branch main master dev staging production server client app website web browser" +
|
||||
" stream streaming video audio voice call camera screen share screenshare gameplay gaming game play steam epic xbox playstation nintendo switch console" +
|
||||
" bot discordbot moderation moderator admin member user profile avatar channel server guild message chat dm reply forward embed sticker emoji role permission" +
|
||||
|
||||
@@ -1,122 +0,0 @@
|
||||
import {
|
||||
embedText,
|
||||
normalizeEmbeddingContent,
|
||||
} from "@/modules/ai-moderation/embeddingClient";
|
||||
import {
|
||||
ARCHIVE_COLLECTION,
|
||||
qdrantPointId,
|
||||
upsertQdrantPointV2,
|
||||
} from "@/modules/ai-moderation/qdrantClient";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
|
||||
const log = createChildLogger("archive-embedder");
|
||||
|
||||
export interface ArchiveMessage {
|
||||
id: string;
|
||||
content: string;
|
||||
username: string;
|
||||
channel_id: string;
|
||||
guild_id: string;
|
||||
thread_id: string | null;
|
||||
created_at: number;
|
||||
/** JSON string of RichMessageMetadata (parsed for channel/thread names). */
|
||||
metadata?: string | null;
|
||||
/** True when the message came from an age-restricted (NSFW) channel. NSFW
|
||||
* content is deliberately NOT embedded into the public archive so it can't
|
||||
* be found via public semantic search. Defaults to false. */
|
||||
isAgeRestricted?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract a human-readable channel label from the message's metadata JSON.
|
||||
* The gateway captures `metadata.channel.{channelName,threadName}` per message;
|
||||
* prefer the thread name (thread messages read better by their thread title),
|
||||
* then the channel name. Returns null when unavailable (old messages without
|
||||
* the metadata field, or malformed JSON).
|
||||
*/
|
||||
export function extractChannelLabel(metadata: string | null | undefined): {
|
||||
channel_name: string | null;
|
||||
thread_name: string | null;
|
||||
} {
|
||||
if (!metadata) return { channel_name: null, thread_name: null };
|
||||
try {
|
||||
const m = JSON.parse(metadata) as {
|
||||
channel?: {
|
||||
channelName?: string | null;
|
||||
threadName?: string | null;
|
||||
};
|
||||
};
|
||||
return {
|
||||
channel_name: m?.channel?.channelName ?? null,
|
||||
thread_name: m?.channel?.threadName ?? null,
|
||||
};
|
||||
} catch {
|
||||
return { channel_name: null, thread_name: null };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fire-and-forget: embed a captured message and upsert it into the persistent
|
||||
* archive collection so the public web can semantic-search the corpus.
|
||||
*
|
||||
* Failures are swallowed — searching is a nice-to-have, never a precondition
|
||||
* for capture or moderation. The message text is kept in the payload so the
|
||||
* search endpoint can return results even for deleted messages.
|
||||
*
|
||||
* NSFW/age-restricted messages are skipped (never embedded) — they are stored
|
||||
* in the database for the dashboard but kept out of the public search archive.
|
||||
*/
|
||||
export function archiveMessageEmbedded(message: ArchiveMessage): void {
|
||||
if (message.isAgeRestricted) return; // never surface NSFW in public archive
|
||||
if (!config.AI_LLM_EMBEDDING_MODEL) return; // embeddings disabled → skip
|
||||
const text = message.content?.trim();
|
||||
if (!text || text.length < 3) return;
|
||||
|
||||
void (async () => {
|
||||
try {
|
||||
// Normalize once: the vector AND the stored payload both use the clean
|
||||
// text so the public search returns readable content and the vector
|
||||
// isn't diluted by @mentions/URLs/emoji (see normalizeEmbeddingContent).
|
||||
const normalized = normalizeEmbeddingContent(text);
|
||||
if (!normalized) return; // nothing meaningful left after cleanup
|
||||
const vector = await embedText(normalized);
|
||||
if (!vector) return;
|
||||
const { channel_name, thread_name } = extractChannelLabel(
|
||||
message.metadata,
|
||||
);
|
||||
const ok = await upsertQdrantPointV2(
|
||||
ARCHIVE_COLLECTION,
|
||||
qdrantPointId(`archive:${message.id}`),
|
||||
vector,
|
||||
{
|
||||
text: normalized.slice(0, 4000),
|
||||
flags: "",
|
||||
// Rich metadata so public semantic search results can be shown in
|
||||
// context (who said it, where, when) instead of a bare text blob.
|
||||
username: message.username ?? "",
|
||||
channel_id: message.channel_id ?? "",
|
||||
guild_id: message.guild_id ?? "",
|
||||
thread_id: message.thread_id ?? null,
|
||||
channel_name: channel_name ?? null,
|
||||
thread_name: thread_name ?? null,
|
||||
created_at: message.created_at ?? Date.now(),
|
||||
analyzed_at: Date.now(),
|
||||
// 5-year persistent window (archive is NOT a TTL cache).
|
||||
expires_at: Date.now() + 1000 * 60 * 60 * 24 * 365 * 5,
|
||||
content_hash: message.id,
|
||||
},
|
||||
);
|
||||
if (!ok) return;
|
||||
log.debug({ messageId: message.id }, "Archived message embedding");
|
||||
} catch (err) {
|
||||
log.debug(
|
||||
{
|
||||
messageId: message.id,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"archive embed skipped",
|
||||
);
|
||||
}
|
||||
})();
|
||||
}
|
||||
@@ -4,12 +4,10 @@ import { config } from "../../shared/config/index.js";
|
||||
import { queueMessageAnalysis } from "../ai-moderation/aiAnalyzer.js";
|
||||
import { processAttachmentUpload } from "../attachment-upload/attachmentUploader.js";
|
||||
import type { EventBroadcaster } from "../event-broadcaster/eventBroadcaster.js";
|
||||
import { archiveMessageEmbedded } from "../message-capture/archiveEmbedder.js";
|
||||
import {
|
||||
getDisplayContent,
|
||||
getMessageLocation,
|
||||
getMessageMetadata,
|
||||
isAgeRestrictedMessage,
|
||||
} from "../message-capture/messageMetadata.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type {
|
||||
@@ -213,16 +211,6 @@ export async function captureMessage(
|
||||
return;
|
||||
}
|
||||
|
||||
// Fire-and-forget: make the captured message searchable in the persistent
|
||||
// archive (public semantic search). Never blocks capture/moderation.
|
||||
// NSFW/age-restricted messages are kept OUT of the public archive.
|
||||
if (!isBacklog && messageRecord.content) {
|
||||
archiveMessageEmbedded({
|
||||
...messageRecord,
|
||||
isAgeRestricted: isAgeRestrictedMessage(message),
|
||||
});
|
||||
}
|
||||
|
||||
if (_eventBroadcaster && !isBacklog) {
|
||||
_eventBroadcaster.messageCreated(messageRecord);
|
||||
}
|
||||
|
||||
@@ -167,34 +167,6 @@ export const configSchema = z
|
||||
.describe(
|
||||
"Disable LLM chain-of-thought (reasoning/thinking) to speed up AI analysis. Set false to restore thinking.",
|
||||
),
|
||||
AI_LLM_EMBEDDING_MODEL: z.string().optional(),
|
||||
AI_LLM_EMBEDDING_MIN_SIMILARITY: z.coerce
|
||||
.number()
|
||||
.min(0)
|
||||
.max(1)
|
||||
.default(0.97),
|
||||
// Two-band semantic acceptance (2026-08-24): non-actionable verdicts
|
||||
// (clean, no flags, action=none) may be reused from a LOOSER similarity
|
||||
// band than actionable ones (warn/flagged). Actionable verdicts keep the
|
||||
// strict gate above; anything between the two bands falls through to the
|
||||
// LLM (fail-open toward accuracy).
|
||||
AI_LLM_EMBEDDING_MIN_SIMILARITY_CLEAN: z.coerce
|
||||
.number()
|
||||
.min(0)
|
||||
.max(1)
|
||||
.default(0.92),
|
||||
AI_LLM_EMBEDDING_MAX_CANDIDATES: z.coerce
|
||||
.number()
|
||||
.int()
|
||||
.positive()
|
||||
.default(50),
|
||||
// Qdrant vector store for the semantic moderation cache. When
|
||||
// QDRANT_URL is set, embeddings are stored/searched there (Postgres
|
||||
// embedding column remains as a legacy fallback).
|
||||
QDRANT_URL: z.string().optional(),
|
||||
QDRANT_COLLECTION: z.string().default("gmw_text_moderation"),
|
||||
QDRANT_ARCHIVE_COLLECTION: z.string().default("gmw_message_archive"),
|
||||
QDRANT_API_KEY: z.string().optional(),
|
||||
AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(8),
|
||||
// Media-lane LLM concurrency cap (2026-09-24): vision + media-batch calls
|
||||
// use their OWN semaphore instead of sharing AI_LLM_MAX_CONCURRENT, so a
|
||||
|
||||
@@ -395,9 +395,6 @@ export const pgTextAnalysisCacheTable = pgTable(
|
||||
analyzed_at: pgBigint("analyzed_at", { mode: "number" }).notNull(),
|
||||
expires_at: pgBigint("expires_at", { mode: "number" }).notNull(),
|
||||
hit_count: pgInteger("hit_count").notNull().default(0),
|
||||
// JSON-encoded embedding vector for semantic moderation cache lookups.
|
||||
// Null for entries stored before embeddings were enabled.
|
||||
embedding: pgText("embedding"),
|
||||
},
|
||||
(table) => ({
|
||||
expiresAtIdx: pgIndex("idx_text_analysis_cache_expires_at").on(
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import {
|
||||
deriveRecommendedAction,
|
||||
isEligibleForAutoDelete,
|
||||
} from "../src/modules/ai-moderation/autoDeleteEligibility.js";
|
||||
import type { MessageRecord } from "../src/shared/moderation-types.js";
|
||||
|
||||
const baseMessage: MessageRecord = {
|
||||
id: "test-msg-1",
|
||||
channel_id: "chan-1",
|
||||
thread_id: null,
|
||||
guild_id: "guild-1",
|
||||
user_id: "user-1",
|
||||
username: "tester",
|
||||
content: "kau ngehina aku hitam kah ?",
|
||||
created_at: Date.now(),
|
||||
ai_status: "flagged",
|
||||
ai_severity: "high",
|
||||
ai_recommended_action: "review",
|
||||
ai_moderation_flags: '["harassment"]',
|
||||
ai_confidence: 0.95,
|
||||
ai_analysis: "konfrontatif",
|
||||
ai_categories: null,
|
||||
ai_moderation_score: null,
|
||||
ai_analyzed_at: Date.now(),
|
||||
deleted_at: null,
|
||||
} as unknown as MessageRecord;
|
||||
|
||||
describe("autoDeleteEligibility — flagged high severity with conservative LLM action", () => {
|
||||
it("flagged + high severity is eligible even when LLM said review", () => {
|
||||
const eligible = isEligibleForAutoDelete(baseMessage);
|
||||
expect(eligible).toBe(true);
|
||||
});
|
||||
|
||||
it("flagged + high severity derives delete regardless of stored review action", () => {
|
||||
expect(deriveRecommendedAction(baseMessage)).toBe("delete");
|
||||
});
|
||||
|
||||
it("flagged + medium severity with review action is NOT eligible", () => {
|
||||
const medium = {
|
||||
...baseMessage,
|
||||
ai_severity: "medium",
|
||||
} as unknown as MessageRecord;
|
||||
const eligible = isEligibleForAutoDelete(medium);
|
||||
expect(eligible).toBe(false);
|
||||
});
|
||||
|
||||
it("warn status with review action is NOT eligible", () => {
|
||||
const warn = {
|
||||
...baseMessage,
|
||||
ai_status: "warn",
|
||||
ai_severity: "low",
|
||||
ai_recommended_action: "warn",
|
||||
} as unknown as MessageRecord;
|
||||
const eligible = isEligibleForAutoDelete(warn);
|
||||
expect(eligible).toBe(true); // warn action is allowed
|
||||
});
|
||||
|
||||
it("clean status is never eligible", () => {
|
||||
const clean = {
|
||||
...baseMessage,
|
||||
ai_status: "clean",
|
||||
ai_severity: "none",
|
||||
ai_recommended_action: "none",
|
||||
} as unknown as MessageRecord;
|
||||
expect(isEligibleForAutoDelete(clean)).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -6,12 +6,12 @@
|
||||
// single-key getter: unexpired rows only, malformed rows skipped, verdicts
|
||||
// normalized through the shared parser. The DB layer is mocked — no live
|
||||
// Postgres in unit tests.
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { beforeEach, describe, expect, it, jest, mock } from "bun:test";
|
||||
|
||||
const executeAll = vi.fn();
|
||||
const executeGet = vi.fn();
|
||||
const executeAll = jest.fn();
|
||||
const executeGet = jest.fn();
|
||||
|
||||
vi.mock("../src/shared/database/drizzle.js", () => ({
|
||||
mock.module("../src/shared/database/drizzle.js", () => ({
|
||||
executeAll: (...args: unknown[]) => executeAll(...args),
|
||||
executeGet: (...args: unknown[]) => executeGet(...args),
|
||||
}));
|
||||
|
||||
@@ -1,11 +1,9 @@
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// Semantic two-band acceptance + global exact-cache reuse guard
|
||||
// Global exact-cache reuse guard
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// Design (2026-08-24): cache hits may be served MORE aggressively for
|
||||
// verdicts that cannot trigger enforcement actions, and NEVER more
|
||||
// aggressively for actionable ones. Two layers enforce this:
|
||||
// - isSemanticBandAccepted: similarity thresholds differ by verdict class
|
||||
// (clean band 0.92 default vs strict actionable band 0.97 default).
|
||||
// aggressively for actionable ones:
|
||||
// - isGloballyReusableCleanVerdict: context-free (cross-channel) reuse of
|
||||
// the legacy bare key only for clean / flagless / action=none verdicts
|
||||
// with high confidence and bounded age.
|
||||
@@ -31,64 +29,6 @@ function makeVerdict(
|
||||
};
|
||||
}
|
||||
|
||||
describe("isSemanticBandAccepted", () => {
|
||||
it("accepts a non-actionable clean verdict at the loose clean band", () => {
|
||||
// Default AI_LLM_EMBEDDING_MIN_SIMILARITY_CLEAN = 0.92.
|
||||
expect(isBandAccept(makeVerdict(), 0.93)).toBe(true);
|
||||
});
|
||||
|
||||
it("accepts a clean verdict exactly at the clean band boundary", () => {
|
||||
expect(isBandAccept(makeVerdict({ confidence: 0.99 }), 0.92)).toBe(true);
|
||||
});
|
||||
|
||||
it("rejects a clean verdict below the clean band", () => {
|
||||
expect(isBandAccept(makeVerdict(), 0.91)).toBe(false);
|
||||
});
|
||||
|
||||
it("rejects an actionable flagged verdict between the bands", () => {
|
||||
// 0.93 >= clean band BUT < strict band → must NOT be served.
|
||||
expect(
|
||||
isBandAccept(
|
||||
makeVerdict({ status: "flagged", flags: ["hate_speech"] }),
|
||||
0.93,
|
||||
),
|
||||
).toBe(false);
|
||||
});
|
||||
|
||||
it("accepts a flagged verdict at the strict band", () => {
|
||||
expect(
|
||||
isBandAccept(
|
||||
makeVerdict({ status: "flagged", flags: ["hate_speech"] }),
|
||||
0.98,
|
||||
),
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("rejects a warn verdict below the strict band", () => {
|
||||
expect(
|
||||
isBandAccept(
|
||||
makeVerdict({ status: "warn", recommendedAction: "warn" }),
|
||||
0.96,
|
||||
),
|
||||
).toBe(false);
|
||||
});
|
||||
|
||||
it("treats a clean verdict WITH flags as actionable (strict band)", () => {
|
||||
expect(isBandAccept(makeVerdict({ flags: ["borderline"] }), 0.93)).toBe(
|
||||
false,
|
||||
);
|
||||
});
|
||||
|
||||
it("treats a clean verdict with a non-none action as actionable", () => {
|
||||
expect(
|
||||
isBandAccept(makeVerdict({ recommendedAction: "review" }), 0.93),
|
||||
).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
// Import indirection so the describe block reads cleanly.
|
||||
import { isSemanticBandAccepted as isBandAccept } from "../src/modules/ai-moderation/textCacheStore.js";
|
||||
|
||||
describe("isGloballyReusableCleanVerdict", () => {
|
||||
it("accepts a fresh, confident, flagless clean verdict", () => {
|
||||
const v = makeVerdict({ confidence: 0.9 });
|
||||
|
||||
@@ -1,82 +0,0 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
// The qdrant client reads config at import time — we only need the
|
||||
// reset-on-failure behavior, so mock global fetch and import fresh.
|
||||
import {
|
||||
ensureQdrantCollection,
|
||||
resetQdrantCollectionCache,
|
||||
} from "../src/modules/ai-moderation/qdrantClient.js";
|
||||
|
||||
function jsonResponse(body: unknown, ok = true, status = 200): Response {
|
||||
return {
|
||||
ok,
|
||||
status,
|
||||
text: async () => JSON.stringify(body),
|
||||
} as unknown as Response;
|
||||
}
|
||||
|
||||
describe("ensureQdrantCollection retry-on-failure", () => {
|
||||
beforeEach(() => {
|
||||
resetQdrantCollectionCache();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
resetQdrantCollectionCache();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("returns false on failure and does NOT get stuck — retries on next call", async () => {
|
||||
// First call: GET collection fails hard (not a 404) → ensure rejects.
|
||||
const fetchMock = vi
|
||||
.spyOn(globalThis, "fetch")
|
||||
.mockRejectedValueOnce(new Error("network down"))
|
||||
// Second call: GET returns a matching-size collection → success.
|
||||
.mockResolvedValueOnce(
|
||||
jsonResponse({
|
||||
result: { config: { params: { vectors: { size: 3072 } } } },
|
||||
}),
|
||||
);
|
||||
|
||||
const first = await ensureQdrantCollection(3072);
|
||||
expect(first).toBe(false);
|
||||
|
||||
// Without the fix, the second call returns the memoised rejected promise
|
||||
// and fetch is never called again. With the fix, it retries.
|
||||
const second = await ensureQdrantCollection(3072);
|
||||
expect(second).toBe(true);
|
||||
expect(fetchMock).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it("recreates collection when vector size changes (DELETE + PUT)", async () => {
|
||||
const fetchMock = vi
|
||||
.spyOn(globalThis, "fetch")
|
||||
.mockResolvedValueOnce(
|
||||
jsonResponse({
|
||||
result: { config: { params: { vectors: { size: 2048 } } } },
|
||||
}),
|
||||
) // GET old size
|
||||
.mockResolvedValueOnce(jsonResponse({ result: true })) // DELETE
|
||||
.mockResolvedValueOnce(jsonResponse({ result: true })); // PUT
|
||||
|
||||
const ok = await ensureQdrantCollection(3072);
|
||||
expect(ok).toBe(true);
|
||||
|
||||
const methods = fetchMock.mock.calls.map(
|
||||
(c) => (c[1] as RequestInit).method,
|
||||
);
|
||||
expect(methods).toEqual(["GET", "DELETE", "PUT"]);
|
||||
});
|
||||
|
||||
it("is idempotent when the collection already matches", async () => {
|
||||
const fetchMock = vi.spyOn(globalThis, "fetch").mockResolvedValueOnce(
|
||||
jsonResponse({
|
||||
result: { config: { params: { vectors: { size: 3072 } } } },
|
||||
}),
|
||||
);
|
||||
|
||||
const ok = await ensureQdrantCollection(3072);
|
||||
expect(ok).toBe(true);
|
||||
expect(fetchMock).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,7 @@
|
||||
// bun test preload — replicates the env block from the old vitest.config.ts.
|
||||
// Runs before any test module imports, so the config singleton (which reads
|
||||
// process.env at import time) gets the same test values it had under vitest.
|
||||
process.env.DISCORD_TOKEN = "test-discord-token";
|
||||
process.env.DATABASE_URL = "postgres://localhost:6432/test";
|
||||
process.env.AI_ANALYSIS_ENABLED = "true";
|
||||
process.env.AI_LLM_API_KEY = "sk-test";
|
||||
@@ -6,10 +6,10 @@
|
||||
// ["conflict_instigation"]) fell into the legacy `flags.length === 0 ?
|
||||
// clean : flagged` branch and was read back as FLAGGED. Downstream this
|
||||
// broke auto-delete eligibility gating and mislabelled warnings on the
|
||||
// dashboard. parseQdrantVerdict had the same narrowing (warn → clean).
|
||||
// dashboard.
|
||||
//
|
||||
// Fix: normalizeStoredStatus() accepts the full clean/warn/flagged union in
|
||||
// BOTH readers; unknown/legacy values still derive from flags.
|
||||
// Fix: normalizeStoredStatus() accepts the full clean/warn/flagged union;
|
||||
// unknown/legacy values still derive from flags.
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { normalizeStoredStatus } from "../src/modules/ai-moderation/textCacheStore.js";
|
||||
|
||||
|
||||
@@ -45,7 +45,7 @@ describe("tinyFishSearch fallback", () => {
|
||||
// Pin the live config object to a known-disabled state: the shell may
|
||||
// export a real TINYFISH_API_KEY (dev box), which would flip
|
||||
// isTinyFishEnabled() and let tests hit the network.
|
||||
const { config } = await import("../../src/shared/config/index.js");
|
||||
const { config } = await import("../src/shared/config/index.js");
|
||||
prevKey = config.TINYFISH_API_KEY;
|
||||
prevEnabled = config.TINYFISH_SEARCH_ENABLED;
|
||||
(config as Record<string, unknown>).TINYFISH_API_KEY = "";
|
||||
@@ -54,7 +54,7 @@ describe("tinyFishSearch fallback", () => {
|
||||
|
||||
afterEach(async () => {
|
||||
vi.restoreAllMocks();
|
||||
const { config } = await import("../../src/shared/config/index.js");
|
||||
const { config } = await import("../src/shared/config/index.js");
|
||||
(config as Record<string, unknown>).TINYFISH_API_KEY = prevKey;
|
||||
(config as Record<string, unknown>).TINYFISH_SEARCH_ENABLED = prevEnabled;
|
||||
});
|
||||
@@ -63,7 +63,7 @@ describe("tinyFishSearch fallback", () => {
|
||||
config: Record<string, unknown>;
|
||||
prev: string;
|
||||
}> {
|
||||
const { config } = await import("../../src/shared/config/index.js");
|
||||
const { config } = await import("../src/shared/config/index.js");
|
||||
const prev = config.TINYFISH_API_KEY;
|
||||
(config as Record<string, unknown>).TINYFISH_API_KEY = "sk-test-key";
|
||||
return { config: config as unknown as Record<string, unknown>, prev };
|
||||
|
||||
@@ -1,23 +0,0 @@
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { defineConfig } from "vitest/config";
|
||||
|
||||
// Resolves the "@/*" tsconfig path alias so vitest can import src modules
|
||||
// (the pre-existing test suite was broken without this).
|
||||
export default defineConfig({
|
||||
resolve: {
|
||||
alias: {
|
||||
"@": fileURLToPath(new URL("./src", import.meta.url)),
|
||||
},
|
||||
},
|
||||
test: {
|
||||
include: ["tests/**/*.test.ts"],
|
||||
// Loaded before module imports — satisfies the config singleton
|
||||
// (DISCORD_TOKEN required) and DB-agnostic pure-function tests.
|
||||
env: {
|
||||
DISCORD_TOKEN: "test-discord-token",
|
||||
DATABASE_URL: "postgres://localhost:6432/test",
|
||||
AI_ANALYSIS_ENABLED: "true",
|
||||
AI_LLM_API_KEY: "sk-test",
|
||||
},
|
||||
},
|
||||
});
|
||||
@@ -17,7 +17,7 @@
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@orpc/client": "1.15.2",
|
||||
"@orpc/client": "1.15.3",
|
||||
"clsx": "^2.1.1",
|
||||
"lucide-react": "^1.47.0",
|
||||
"next": "16.3.5",
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user