feat: add Redis Stream handoff for AI dev + monitor all messages
- Publish every non-spam message to tower:stream:messages immediately after DB insert in ingest.processor — before any admin approval or routing decision - Remove early tags.length === 0 drop in main.ts so no-rule-match messages are now stored and streamed (monitor-all behaviour) - Wire classify worker, review/p1/dnc queues into main.ts - Add streams/message-stream.ts with XADD wrapper (capped at ~50k entries) - Add docs/DATA_MODELS.md — full data model reference for all 29 models Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+31
-5
@@ -1,13 +1,19 @@
|
||||
import Redis from 'ioredis';
|
||||
import { PrismaClient } from '@prisma/client';
|
||||
import { createLogger } from '@tower/logger';
|
||||
import { validateEnv } from '@tower/config';
|
||||
import { createMeiliClient, configureIndex } from '@tower/search';
|
||||
import { createIngestQueue } from './queues/ingest.queue';
|
||||
import { createIngestWorker } from './queues/ingest.processor';
|
||||
import { createClassifyQueue } from './queues/classify.queue';
|
||||
import { createClassifyWorker } from './queues/classify.processor';
|
||||
import { createForwardQueue } from './queues/forward.queue';
|
||||
import { createForwardWorker } from './queues/forward.processor';
|
||||
import { createIndexQueue } from './queues/index.queue';
|
||||
import { createIndexWorker } from './queues/index.processor';
|
||||
import { createReviewQueue } from './queues/review.queue';
|
||||
import { createP1Queue } from './queues/p1.queue';
|
||||
import { createDncQueue } from './queues/dnc.queue';
|
||||
import { createSpamQueue } from './queues/spam.queue';
|
||||
import { createSpamWorker } from './queues/spam.processor';
|
||||
import { createEventExtractQueue } from './queues/event-extract.queue';
|
||||
@@ -35,14 +41,28 @@ async function bootstrap() {
|
||||
);
|
||||
|
||||
const ingestQueue = createIngestQueue(env.REDIS_URL);
|
||||
const classifyQueue = createClassifyQueue(env.REDIS_URL);
|
||||
const forwardQueue = createForwardQueue(env.REDIS_URL);
|
||||
const indexQueue = createIndexQueue(env.REDIS_URL);
|
||||
const reviewQueue = createReviewQueue(env.REDIS_URL);
|
||||
const p1Queue = createP1Queue(env.REDIS_URL);
|
||||
const dncQueue = createDncQueue(env.REDIS_URL);
|
||||
const spamQueue = createSpamQueue(env.REDIS_URL);
|
||||
const eventExtractQueue = createEventExtractQueue(env.REDIS_URL);
|
||||
const faqCandidateQueue = createFaqCandidateQueue(env.REDIS_URL);
|
||||
const pool = new WhatsAppSessionPool();
|
||||
|
||||
const ingestWorker = createIngestWorker(env.REDIS_URL, prisma, pool, forwardQueue, indexQueue);
|
||||
// Dedicated Redis client for XADD stream writes (separate from BullMQ connections)
|
||||
const { hostname, port } = new URL(env.REDIS_URL);
|
||||
const streamRedis = new Redis({ host: hostname, port: parseInt(port || '6379', 10), maxRetriesPerRequest: null, lazyConnect: true });
|
||||
await streamRedis.connect().catch((err) => logger.warn({ err }, 'Stream Redis connect failed — stream publish will be degraded'));
|
||||
|
||||
const ingestWorker = createIngestWorker(env.REDIS_URL, prisma, pool, forwardQueue, indexQueue, streamRedis);
|
||||
const classifyWorker = createClassifyWorker(
|
||||
env.REDIS_URL, prisma, pool,
|
||||
forwardQueue, indexQueue, reviewQueue, p1Queue, dncQueue, spamQueue,
|
||||
eventExtractQueue, faqCandidateQueue,
|
||||
);
|
||||
const forwardWorker = createForwardWorker(env.REDIS_URL, pool);
|
||||
const indexWorker = createIndexWorker(env.REDIS_URL, meiliClient);
|
||||
const spamWorker = createSpamWorker(env.REDIS_URL, prisma);
|
||||
@@ -51,6 +71,7 @@ async function bootstrap() {
|
||||
|
||||
ingestWorker.on('completed', (job) => logger.info({ jobId: job.id }, 'Ingest job completed'));
|
||||
ingestWorker.on('failed', (job, err) => logger.error({ jobId: job?.id, err }, 'Ingest job failed'));
|
||||
classifyWorker.on('failed', (job, err) => logger.error({ jobId: job?.id, err }, 'Classify job failed'));
|
||||
forwardWorker.on('completed', (job) => logger.info({ jobId: job.id }, 'Forward job completed'));
|
||||
forwardWorker.on('failed', (job, err) => logger.error({ jobId: job?.id, err }, 'Forward job failed'));
|
||||
indexWorker.on('completed', (job) => logger.info({ jobId: job.id }, 'Index job completed'));
|
||||
@@ -139,14 +160,13 @@ async function bootstrap() {
|
||||
const { tags, effectiveAction } = matchContentRules(msg.content, ruleRows);
|
||||
logger.info({ tags, effectiveAction }, 'Rule match result');
|
||||
|
||||
// No matching rules — drop the message
|
||||
if (tags.length === 0) return;
|
||||
|
||||
// SKIP action — silently drop
|
||||
// SKIP action — silently drop (explicit admin instruction to ignore)
|
||||
if (effectiveAction === 'SKIP') {
|
||||
logger.info({ platformMsgId: msg.platformMsgId }, 'Message skipped by rule');
|
||||
return;
|
||||
}
|
||||
// No rule match → ingest with null effectiveAction so it's stored and streamed
|
||||
// but not forwarded. This is the "monitor all" behaviour.
|
||||
|
||||
// For AUTO_APPROVE, check if the sender is a group admin
|
||||
let finalAction = effectiveAction;
|
||||
@@ -312,17 +332,23 @@ async function bootstrap() {
|
||||
logger.info('Shutting down...');
|
||||
await pool.closeAll();
|
||||
await ingestWorker.close();
|
||||
await classifyWorker.close();
|
||||
await forwardWorker.close();
|
||||
await indexWorker.close();
|
||||
await spamWorker.close();
|
||||
await eventExtractWorker.close();
|
||||
await faqCuratorWorker.close();
|
||||
await ingestQueue.close();
|
||||
await classifyQueue.close();
|
||||
await forwardQueue.close();
|
||||
await indexQueue.close();
|
||||
await reviewQueue.close();
|
||||
await p1Queue.close();
|
||||
await dncQueue.close();
|
||||
await spamQueue.close();
|
||||
await eventExtractQueue.close();
|
||||
await faqCandidateQueue.close();
|
||||
await streamRedis.quit();
|
||||
await prisma.$disconnect();
|
||||
process.exit(0);
|
||||
};
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
import { Worker, Queue } from 'bullmq';
|
||||
import Redis from 'ioredis';
|
||||
import { IngestJobData, ForwardJobData, IndexJobData } from '@tower/types';
|
||||
import { parseRedisUrl } from './redis-connection';
|
||||
import { approveMessage } from '../core/approve-message';
|
||||
import { publishToMessageStream } from '../streams/message-stream';
|
||||
import { createLogger } from '@tower/logger';
|
||||
|
||||
const logger = createLogger('ingest-processor');
|
||||
@@ -12,6 +14,7 @@ export async function processIngestJob(
|
||||
pool?: any,
|
||||
forwardQueue?: Queue<ForwardJobData>,
|
||||
indexQueue?: Queue<IndexJobData>,
|
||||
redis?: Redis,
|
||||
): Promise<void> {
|
||||
// Defensive: drop messages from non-CLAIMED groups
|
||||
const group = await prisma.group.findUnique({
|
||||
@@ -87,6 +90,22 @@ export async function processIngestJob(
|
||||
update: {},
|
||||
});
|
||||
|
||||
// Publish to AI-dev stream immediately — before any routing decision.
|
||||
// Admin approval controls forwarding only; it must never block the stream.
|
||||
if (redis) {
|
||||
publishToMessageStream(redis, {
|
||||
messageId: msg.id,
|
||||
tenantId: job.tenantId,
|
||||
content: job.content,
|
||||
senderJid: job.senderJid,
|
||||
senderName: job.senderName,
|
||||
sourceGroupId: job.sourceGroupId,
|
||||
tags: job.tags,
|
||||
effectiveAction: job.effectiveAction ?? null,
|
||||
timestamp: new Date().toISOString(),
|
||||
}).catch((err: unknown) => logger.warn({ err, messageId: msg.id }, 'Failed to publish to message stream'));
|
||||
}
|
||||
|
||||
// For REJECT, create an approval record so it's searchable as rejected
|
||||
if (job.effectiveAction === 'REJECT') {
|
||||
await prisma.approval.upsert({
|
||||
@@ -133,10 +152,11 @@ export function createIngestWorker(
|
||||
pool?: any,
|
||||
forwardQueue?: Queue<ForwardJobData>,
|
||||
indexQueue?: Queue<IndexJobData>,
|
||||
redis?: Redis,
|
||||
): Worker<IngestJobData> {
|
||||
return new Worker<IngestJobData>(
|
||||
'tower-ingest',
|
||||
async (job) => processIngestJob(job.data, prisma, pool, forwardQueue, indexQueue),
|
||||
async (job) => processIngestJob(job.data, prisma, pool, forwardQueue, indexQueue, redis),
|
||||
{ connection: parseRedisUrl(redisUrl) },
|
||||
);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
import Redis from 'ioredis';
|
||||
|
||||
export const STREAM_KEY = 'tower:stream:messages';
|
||||
|
||||
// Approximate cap — Redis trims to ~this many entries.
|
||||
// At 1000 msg/day that's ~50 days of history.
|
||||
const STREAM_MAX_LEN = 50_000;
|
||||
|
||||
export interface StreamMessagePayload {
|
||||
messageId: string;
|
||||
tenantId: string;
|
||||
content: string;
|
||||
senderJid: string;
|
||||
senderName?: string;
|
||||
sourceGroupId: string;
|
||||
tags: string[];
|
||||
effectiveAction: string | null;
|
||||
traceId?: string;
|
||||
timestamp: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Appends one entry to the Redis Stream used as the AI-dev handoff channel.
|
||||
* XADD tower:stream:messages MAXLEN ~ 50000 * field value ...
|
||||
*
|
||||
* The AI developer reads with: XREAD COUNT 100 BLOCK 0 STREAMS tower:stream:messages >
|
||||
* Consumer groups are supported for parallel AI workers.
|
||||
*/
|
||||
export async function publishToMessageStream(
|
||||
redis: Redis,
|
||||
payload: StreamMessagePayload,
|
||||
): Promise<void> {
|
||||
await redis.xadd(
|
||||
STREAM_KEY,
|
||||
'MAXLEN', '~', String(STREAM_MAX_LEN),
|
||||
'*',
|
||||
'messageId', payload.messageId,
|
||||
'tenantId', payload.tenantId,
|
||||
'content', payload.content,
|
||||
'senderJid', payload.senderJid,
|
||||
'senderName', payload.senderName ?? '',
|
||||
'sourceGroupId', payload.sourceGroupId,
|
||||
'tags', JSON.stringify(payload.tags),
|
||||
'effectiveAction', payload.effectiveAction ?? '',
|
||||
'traceId', payload.traceId ?? '',
|
||||
'timestamp', payload.timestamp,
|
||||
);
|
||||
}
|
||||
+1141
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user