Author SHA1 Message Date
mytheclipsebotreview 2b6ec59286 refactor(ai): remove semantic embedding cache + Qdrant vector store
Hapus seluruh fitur embedding/Qdrant (tidak dipakai lagi):

- gateway: drop embeddingClient.ts, qdrantClient.ts, archiveEmbedder.ts
  dan tes qdrantEnsure.test.ts; moderationOrchestrator kembali ke
  exact-hash cache -> LLM (tanpa phase-2 semantic lookup); textCacheStore
  kehilangan findSimilarTextModeration / parseQdrantVerdict /
  isSemanticBandAccepted / upsertBareKeyToQdrant; cache-prune hanya
  menyapu Postgres.
- backend: drop embed.ts + qdrant.ts, endpoint messages.semanticSearch
  dan schema/type terkait; kolom embedding dilepas dari schema
  text_analysis_cache.
- frontend: hapus toggle EXACT/SEMANTIC, hook useSemanticSearch,
  API client + tipe SemanticSearchResult.
- config: buang AI_LLM_EMBEDDING_* dan QDRANT_* (env + .env.example).
- docs: ARCHITECTURE.md / AGENTS.md / README.md / diagram arsitektur
  disesuaikan (LLM caller - vision, cache = exact-hash saja).

Verifikasi: tsc 0 (backend, gateway, frontend); bun test 135 pass +
37 pass, 0 fail; biome 0 error.
2026-09-25 01:38:51 +07:00
mytheclipsebotreview fdc7f01c26 fix(ci): replace leftover pnpm reference in gateway nativeBuildInputs with bun (build failure) 2026-09-25 00:47:56 +07:00
mytheclipsebotreview f31b1f62d4 Merge remote-tracking branch 'origin/main'
# Conflicts:
#	services/frontend/pnpm-lock.yaml
2026-09-25 00:38:10 +07:00
mytheclipsebotreviewandgit-migration[bot] 7be069d73f chore: migrate monorepo toolchain pnpm+vitest → bun (bun test, bunfig preload, bun.lock)
- Gateway + backend + frontend: pnpm/vitest fully removed → bun 1.3.14
  (bun install, bun test tests/, bunfig.toml [test] preload, bun.lock).
- vitest configs deleted; vitest→bun facade (jest/mock/spyOn/waitForCompat)
  keeps the vitest-style assertions working under bun:test.
- flake.nix: bunInstall switch; pruned prod-pass now removes post-pnpm
  dev-toolchain trees; frontend builds Next standalone via bun's next.
- CI: deploy.yml installs with bun + runs bun test tests/ per service.
- 190 tests green (146 gateway + 37 backend pass, 14 skip), 0 fail;
  tsc + biome across all 3 services clean.

Co-authored-by: git-migration[bot] <noreply@gmw.local>
2026-09-24 23:50:30 +07:00
mytheclipsebotreview[bot] 88373985ec Auto-merge PR #89
build(deps): bump @orpc/client from 1.15.2 to 1.15.3 in /services/frontend in the production group
2026-09-24 14:48:31 +00:00
mytheclipsebotreview[bot] 7acc49e1fb Auto-merge PR #88
build(deps-dev): bump the development group in /services/discord-gateway with 2 updates
2026-09-24 14:45:56 +00:00
mytheclipsebotreview[bot] 81b3128bb0 Auto-merge PR #87
build(deps): bump the production group in /services/backend with 2 updates
2026-09-24 14:39:35 +00:00
dependabot[bot] 1d0b3f422d build(deps): bump @orpc/client
Bumps the production group in /services/frontend with 1 update: [@orpc/client](https://github.com/middleapi/orpc/tree/HEAD/packages/client).


Updates `@orpc/client` from 1.15.2 to 1.15.3
- [Release notes](https://github.com/middleapi/orpc/releases)
- [Commits](https://github.com/middleapi/orpc/commits/v1.15.3/packages/client)

---
updated-dependencies:
- dependency-name: "@orpc/client"
  dependency-version: 1.15.3
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: production
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-24 14:25:59 +00:00
dependabot[bot] 0aef0b6c08 build(deps-dev): bump the development group
Bumps the development group in /services/discord-gateway with 2 updates: [@types/node](https://github.com/DefinitelyTyped/DefinitelyTyped/tree/HEAD/types/node) and [drizzle-kit](https://github.com/drizzle-team/drizzle-orm).


Updates `@types/node` from 26.4.0 to 26.6.2
- [Release notes](https://github.com/DefinitelyTyped/DefinitelyTyped/releases)
- [Commits](https://github.com/DefinitelyTyped/DefinitelyTyped/commits/HEAD/types/node)

Updates `drizzle-kit` from 0.31.10 to 0.31.11
- [Release notes](https://github.com/drizzle-team/drizzle-orm/releases)
- [Commits](https://github.com/drizzle-team/drizzle-orm/compare/drizzle-kit@0.31.10...drizzle-kit@0.31.11)

---
updated-dependencies:
- dependency-name: "@types/node"
  dependency-version: 26.6.2
  dependency-type: direct:development
  update-type: version-update:semver-minor
  dependency-group: development
- dependency-name: drizzle-kit
  dependency-version: 0.31.11
  dependency-type: direct:development
  update-type: version-update:semver-patch
  dependency-group: development
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-24 14:25:56 +00:00
dependabot[bot] cd1d6b2d0d build(deps): bump the production group
Bumps the production group in /services/backend with 2 updates: [@orpc/server](https://github.com/middleapi/orpc/tree/HEAD/packages/server) and [drizzle-orm](https://github.com/drizzle-team/drizzle-orm).


Updates `@orpc/server` from 1.15.2 to 1.15.3
- [Release notes](https://github.com/middleapi/orpc/releases)
- [Commits](https://github.com/middleapi/orpc/commits/v1.15.3/packages/server)

Updates `drizzle-orm` from 0.45.2 to 0.45.3
- [Release notes](https://github.com/drizzle-team/drizzle-orm/releases)
- [Commits](https://github.com/drizzle-team/drizzle-orm/compare/0.45.2...0.45.3)

---
updated-dependencies:
- dependency-name: "@orpc/server"
  dependency-version: 1.15.3
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: production
- dependency-name: drizzle-orm
  dependency-version: 0.45.3
  dependency-type: direct:production
  update-type: version-update:semver-patch
  dependency-group: production
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-24 14:25:31 +00:00
asepharyana c888f23901 fix(gateway): flagged high/critical severity always eligible for auto-delete
The LLM's recommended_action is conservative — for a flagged message at
high/critical severity it frequently emits 'review' (screenshot/context
ambiguity) even when the violation is severe (harassment, SARA,
threats). autoDeleteEligibility trusted that value, so serious violations
slipped through undeleted (e.g. harassment flagged high but
review → not eligible).

Fix: flagged + high/critical severity bypasses the recommended-action
check entirely (always eligible); the action check now only gates
warn/flagged-medium. deriveRecommendedAction also returns delete for
flagged high/critical BEFORE consulting the stored LLM action. +5
regression tests.
2026-09-24 18:57:41 +07:00
asepharyana 54d02098c8 fix(frontend): rename stale 'AI heuristic reasoning' label to 'AI moderation analysis'
The label predates the LLM pipeline — all analysis (including LLM output)
was mislabeled as heuristic. The analysis field is LLM-generated
moderation justification.
2026-09-24 18:08:55 +07:00
asepharyana 5fc0e8b3cb fix(gateway): surface failed URL fetches to moderation LLM
Links that fail to fetch (Facebook share 400, login walls, anti-bot)
previously vanished from the prompt entirely — the LLM only saw the raw
URL in the message body, so its analysis degraded to a template like
"Pengguna membagikan tautan Facebook". Now the text batch tracks
failed-fetch URLs and injects <web_content fetch_error="true"> telling
the LLM the page content is unverifiable and to NOT invent it, plus a
prompt rule: describe what the user actually expressed from the text +
conversation context, or admit the subject is unclear — never make the
link itself the subject.
2026-09-24 18:03:41 +07:00
asepharyana 045cdf1f75 fix(gateway): align stuck-recovery threshold with batch timeout (300s→120s)
recovery-worker used a hardcoded STUCK_PROCESSING_AGE_MS=300_000 while
messagesCleanup's default and AI_ANALYSIS_PROCESSING_TIMEOUT_MS are both
120s. Rows stuck between 2 and 5 minutes were never reverted by the
recovery worker — they looked permanently stuck (and accumulated under
load) even though the batch budget had long passed. Now derives the
threshold from config so the two knobs can never drift again.
2026-09-24 17:52:17 +07:00
asepharyana deed7bdfb0 fix(gateway): destroy AI worker pools on graceful shutdown
Piscina worker threads outlive process.exit() and linger as orphaned
processes holding DB connections/locks after a deploy restart. Two live
gateways then fight over the same messages table rows (one claims
processing, the other reverts), which left messages stuck in
ai_status='processing' forever.

Destroy both worker pools before closing the DB, with a 5s fallback so
a hung vision job cannot block shutdown indefinitely.
2026-09-24 17:26:54 +07:00
asepharyana c4d9ade85e debug(gateway): log lane-lock skip + dispatch in scheduleLaneTimer 2026-09-24 17:11:04 +07:00
asepharyana 80daa9f045 fix(gateway): un-claim budget-overflow messages stuck in processing
getPendingMessagesByConversation() flips every fetched pending row to
'processing', then pickBatchWithinBudget() may stop early on the token
budget. The tail rows that did NOT make the batch were never un-claimed,
so they stayed 'processing' forever — the recovery worker reverted them
(120s) only for the next wave to re-claim them, an infinite loop of
stuck messages that never get analyzed (saw 22 rows, some recycled for
40+ minutes).

- add computeBudgetOverflowMessages() pure helper (batchBudget.ts)
- batchScheduler un-claims overflow rows back to 'pending' before
  dispatching the trimmed batch
- recovery-worker now reverts stuck processing unconditionally (the old
  'conversationProcessing.size > 0' guard skipped the revert when the
  in-memory lock map was empty, e.g. fresh boot — exactly when stranded
  rows from a previous process need rescuing)
- 3 regression tests for the overflow helper
2026-09-24 16:49:27 +07:00
asepharyana ef7708bf7d feat(gmw): route all LLM traffic through 9router
GMW moves off omniroute (100.121.180.82:20128) and off the direct NVIDIA
vision endpoint onto 9router, which runs on the same host as both services
(127.0.0.1:4014) — loopback avoids the TLS/proxy hop and localhost calls
bypass 9router's remote-key guard.

- gateway + backend: AI_LLM_BASE_URL default -> http://127.0.0.1:4014/v1
- drop stale 'omniroute' router references from comments/docs now that the
  active router is 9router (llmClient, llmCaller, ARCHITECTURE, AGENTS)

Verified against 9router before wiring: model 'text' -> gemini-3.5-flash-lite
(SSE, as the pipeline expects), 'multimodal' -> nemotron-3-nano-omni answers
image input, and gemini/gemini-embedding-001 returns 3072 dims — matching the
existing Qdrant collections (no reindex needed). The GMW key is already
registered in 9router's apiKeys table.

typecheck + lint + tests green (gateway 138, backend 37 excluding e2e).
2026-09-24 16:05:36 +07:00
asepharyana 750f3aa598 docs(gateway): correct metrics port — 4016 was wrong
The metrics server binds METRICS_PORT (code default 9090; this host runs it
on 4018 — 4016 is occupied by another process). Docs said 4016 in three
places, which is not what the service does.
2026-09-24 15:50:15 +07:00
asepharyana c57ee12da1 docs(gateway): consolidate README/ARCHITECTURE, drop stale MODULE_STRUCTURE
README.md was the extraction-era document (referenced winston, mock-crc.ts,
llmModerationClient.ts, indonesianTextNormalizer.ts — all long gone) and
duplicated ARCHITECTURE.md. Rewritten as a short run-the-service guide;
layout/design lives only in ARCHITECTURE.md.

MODULE_STRUCTURE.md deleted: it was a stale duplicate of ARCHITECTURE.md,
referenced by nothing but itself.

ARCHITECTURE.md updated to the post-refactor reality: app/ lifecycle split
(bootstrap/lifecycle/process-guards/metrics-collector), ai-moderation
recovery-worker + cache-prune, per-module index.ts facades, one-way
dependency rule, corrected init/shutdown/observability sections.
2026-09-24 15:09:13 +07:00
asepharyana 494e16b3b3 refactor(gateway): split bootstrap + aiAnalyzer, add module barrels
app/:
- bootstrap.ts 277 -> 145 lines: config guard, DB connect, client debug
  logging and startup order are now named steps with a comment header
- lifecycle.ts (new): everything wired on the Discord 'ready' hook, in
  explicit order (inject broadcaster -> register listeners -> start workers)
- process-guards.ts (new): SIGINT/SIGTERM/uncaughtException/unhandledRejection
  in ONE place, using isTransientStreamError() instead of two duplicated
  inline code lists
- metrics-collector.ts (new): AI pipeline Prometheus gauges

modules/:
- ai-moderation/index.ts + message-capture/index.ts (new): public facades so
  app/ never reaches into internal files
- aiAnalyzer.ts 317 -> 146 lines: pure entry API; skip-verdict recording
  extracted into recordSkip()
- recovery-worker.ts (new): stranded-message recovery + stale lane/CB pruning
- cache-prune.ts (new): 6h expired-verdict sweep, throttled + resettable
- drop 3 dead re-exports (pickBatchWithinBudget/onCircuitBreakerAlert/
  getConversationKey) whose consumers import the origin files directly

No behavior change. typecheck + lint + 138 tests green; nix build OK.
2026-09-24 15:04:53 +07:00
asepharyana 6bf3b40cc7 refactor(gateway): drop shared barrel, migrate to granular imports + shared error helpers
- delete src/shared/index.ts fat barrel; point 10 importers at the exact
  module they use (redis-channels, moderation-types, utils/pagination)
- message-capture/types.ts re-exports from shared/moderation-types directly
- shared/errors: add errorMessage() + isTransientStreamError() helpers,
  replacing the repeated err-message and transient-code checks
- drop unused imports flagged by biome

No behavior change. typecheck + lint + 138 tests green.
2026-09-24 14:58:19 +07:00
asepharyana 8743fcc0b5 fix(ci): recover attic client binary in VPS-hop fallback path 2026-09-24 14:31:26 +07:00
asepharyana e34dcd6bc6 fix(gateway): separate text and media lanes in AI analysis queue (#86)
Image messages previously blocked the whole analysis pipeline:
- conversationProcessing was a single lock per conversation; processBatch
  awaited BOTH text and media worker jobs before releasing it, so a fast
  text verdict sat unused until the slow vision/media batch finished
- one global LLM semaphore (AI_LLM_MAX_CONCURRENT) was shared by text and
  media, so a vision backlog could starve text inference
- recovery worker gated on conversationProcessing.size

Now the queue is split into independent text/media lanes:
- conversationProcessing maps key -> Partial<Record<lane, startedAt>>;
  each lane holds its own lock and frees it the moment ITS worker job
  resolves (ownership-guarded clear prevents stale timers clearing newer
  slots)
- two LLM semaphores: AI_LLM_MAX_CONCURRENT (text, default 8) and
  AI_LLM_MEDIA_MAX_CONCURRENT (media, default 4) via
  withLlmConcurrency(fn, { lane })
- batchScheduler schedules per conversation+lane (timer keys
  '<key>::<lane>'); splitMessagesByLane/laneOfMessage moved to pure
  analysisLanes.ts (unit-testable without Piscina)
- ai-analysis-worker batch jobs carry a lane field; per-lane active
  request gauges (active_text_requests / active_media_requests)
- added tests/analysisLaneLock.test.ts (7 tests: independent lane locks,
  preserving other-lane lock, clear-all, ownership guard, lane split)

Docs: ARCHITECTURE.md + AGENTS.md concurrency model updated.
typecheck/lint/test(138)/build all green.
2026-09-24 14:09:53 +07:00
asepharyana f9fecfc144 chore: clean up dead barrels, duplicate config, and orphaned frontend components
Gateway:
- Remove dead barrels (ai-moderation/index, attachment-upload/index, message-capture/index) — all consumers import files directly
- Remove orphaned schema/ split dir (analytics/cache/messages/meta) — schema.ts is monolithic
- Merge duplicate config singleton: delete shared/config/config.ts, point all 44 imports at shared/config/index

Backend:
- Remove dead commandHelper.ts (voice-era fallback), ws/index.ts barrel, health.schema.ts, moderationMetrics.ts, analysis.schema.ts (0 importers; metrics/handlers route directly)

Frontend:
- Remove orphaned CategoryDrilldown/CoverageTiles/TopicTrends, primitives/slot, use-mobile, use-mounted
- Remove unused charts donut/sparkline (TopicTrends was only consumer)

Kept (verified active): shared/database/index.ts facade (11 importers), hooks/index + lib/api/index barrels (10 importers), orpc/ws.ts, charts/index.ts barrel.
Verified: tsc + biome + vitest per service (backend e2e 3 failures pre-existing on main); frontend next build 8 routes.
2026-09-24 13:23:38 +07:00
asepharyana a59f3132ee fix: add findings and progress documentation to .gitignore 2026-09-24 12:11:21 +07:00
asepharyana c7f53e4f7e feat(gateway): remove Jev (System One) analyzer, restore LLM-only text moderation
Jev (oc/jev-1.13-free via 9router /v1/systemone) added as primary text
analyzer was underperforming. Delete the whole feature:
- jevAnalyzer.ts + its unit & live-smoke tests
- Jev-first branch in textBatchProcessor, restore pure callModerationLLM path
- AI_LLM_JEV_* config vars (zod) and .env.example entries
- @typesafe-ai/sdk dependency (+ lockfile)

Behavior: text moderation is LLM-only again, exactly as before the
Jev feature; AGENTS.md invariant 'LLM is the only judge' holds.
2026-09-24 12:06:36 +07:00
asepharyana 8583bcdf17 fe(fe): enrich dashboard Top Reacted Messages
Show relative time, always-on username with channel context, and
collapse top emojis to two; add VIEW ALL toggle to reveal all 20
fetched reactions instead of the top 4.
2026-09-23 23:56:02 +07:00
asepharyana b12eb0a038 fe(fe): surface moderation explainability in live feed
Show confidence bar, flag chips, status dot, and error text per action
in the Live Stream Audit Log; drop stale media comment on SkeletonHero.
2026-09-23 23:46:22 +07:00
asepharyana b72423c64d fix(fe): remove remaining voice/recording/media leftovers from navigation, dashboard, and ws layer 2026-09-23 23:21:05 +07:00
asepharyana 6b34fc97ec Merge branch 'fix/remove-voice-recording' 2026-09-23 22:13:00 +07:00
dependabot[bot] 8ce8978755 build(deps-dev): bump tsx (#84)
Bumps the development group in /services/discord-gateway with 1 update: [tsx](https://github.com/privatenumber/tsx).


Updates `tsx` from 4.23.13 to 4.23.15
- [Release notes](https://github.com/privatenumber/tsx/releases)
- [Changelog](https://github.com/privatenumber/tsx/blob/master/release.config.cjs)
- [Commits](https://github.com/privatenumber/tsx/compare/v4.23.13...v4.23.15)

---
updated-dependencies:
- dependency-name: tsx
  dependency-version: 4.23.15
  dependency-type: direct:development
  update-type: version-update:semver-patch
  dependency-group: development
...

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-09-23 22:10:04 +07:00
mytheclipsebotreview[bot] 869cad88e4 Auto-merge PR #83
build(deps-dev): bump tsx from 4.23.13 to 4.23.15 in /services/backend in the development group
2026-09-23 14:41:14 +00:00
mytheclipsebotreview[bot] f136cf3f6b Auto-merge PR #82
build(deps): bump openai from 7.19.0 to 7.20.0 in /services/discord-gateway in the production group
2026-09-23 14:31:27 +00:00
dependabot[bot] cd3ee5b5d8 build(deps-dev): bump tsx in /services/backend in the development group
Bumps the development group in /services/backend with 1 update: [tsx](https://github.com/privatenumber/tsx).


Updates `tsx` from 4.23.13 to 4.23.15
- [Release notes](https://github.com/privatenumber/tsx/releases)
- [Changelog](https://github.com/privatenumber/tsx/blob/master/release.config.cjs)
- [Commits](https://github.com/privatenumber/tsx/compare/v4.23.13...v4.23.15)

---
updated-dependencies:
- dependency-name: tsx
  dependency-version: 4.23.15
  dependency-type: direct:development
  update-type: version-update:semver-patch
  dependency-group: development
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-23 14:24:58 +00:00
dependabot[bot] 93288afa0b build(deps): bump openai
Bumps the production group in /services/discord-gateway with 1 update: [openai](https://github.com/openai/openai-node).


Updates `openai` from 7.19.0 to 7.20.0
- [Release notes](https://github.com/openai/openai-node/releases)
- [Changelog](https://github.com/openai/openai-node/blob/main/CHANGELOG.md)
- [Commits](https://github.com/openai/openai-node/compare/v7.19.0...v7.20.0)

---
updated-dependencies:
- dependency-name: openai
  dependency-version: 7.20.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: production
...

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

Some files were not shown because too many files have changed in this diff Show More