Files
tower/apps/worker/src/whatsapp/command-handler.ts
T
2026-06-09 02:02:40 +05:30

233 lines
7.1 KiB
TypeScript

import type { NormalizedMessage } from '@tower/types';
import { createLogger } from '@tower/logger';
import { WhatsAppSessionPool } from './session-pool';
const logger = createLogger('command-handler');
const PORTAL_BASE = process.env['TOWER_PORTAL_BASE_URL'] ?? 'http://localhost:3000';
const STOP_REGEX = /^\s*stop\s*$/i;
const START_REGEX = /^\s*start\s*$/i;
const PORTAL_REGEX = /^\s*portal\s*$/i;
const COMMANDS_REGEX = /^\s*commands\s*$/i;
export async function handleCommand(
msg: NormalizedMessage,
accountId: string,
prisma: any,
pool: WhatsAppSessionPool,
): Promise<boolean> {
if (!msg.content) return false;
const text = msg.content.trim();
if (STOP_REGEX.test(text)) {
return await handleStop(msg, accountId, prisma, pool);
}
if (START_REGEX.test(text)) {
return await handleStart(msg, accountId, prisma, pool);
}
if (PORTAL_REGEX.test(text)) {
return await handlePortal(msg, accountId, prisma, pool);
}
if (COMMANDS_REGEX.test(text)) {
return await handleCommands(msg, accountId, pool);
}
return false;
}
async function handleStop(
msg: NormalizedMessage,
accountId: string,
prisma: any,
pool: WhatsAppSessionPool,
): Promise<boolean> {
const group = await prisma.group.findUnique({
where: { platform_platformId: { platform: 'whatsapp', platformId: msg.sourceGroupJid } },
select: { id: true, tenantId: true, claimStatus: true },
});
if (!group || group.claimStatus !== 'CLAIMED' || !group.tenantId) {
return false;
}
// Find or create a TowerUser by jid (phone not yet known for STOP-only flow)
const phoneHash = `stop:${msg.senderJid}`;
const user = await prisma.towerUser.upsert({
where: { tenantId_phoneHash: { tenantId: group.tenantId, phoneHash } },
update: { jid: msg.senderJid },
create: {
tenantId: group.tenantId,
phoneHash,
jid: msg.senderJid,
displayName: msg.senderName ?? msg.senderJid,
},
});
// Revoke all active consents in this group
await prisma.consentRecord.updateMany({
where: { userId: user.id, tenantId: group.tenantId, groupId: group.id, status: 'GRANTED' },
data: { status: 'REVOKED', revokedAt: new Date() },
});
await prisma.memberOptOut.create({
data: {
tenantId: group.tenantId,
userId: user.id,
groupId: group.id,
reason: 'STOP_KEYWORD',
},
});
await prisma.auditEvent.create({
data: {
tenantId: group.tenantId,
actorType: 'MEMBER',
actorId: user.id,
action: 'MEMBER_OPT_OUT',
resourceType: 'TowerUser',
resourceId: user.id,
payload: { jid: msg.senderJid, groupId: group.id, reason: 'STOP_KEYWORD' },
},
});
try {
await pool.sendMessage(
accountId,
msg.senderJid,
"You've been opted out. Type START in this group to rejoin.",
);
} catch (err) {
logger.warn({ err, jid: msg.senderJid }, 'Failed to send STOP confirmation DM');
}
logger.info({ jid: msg.senderJid, groupId: group.id }, 'STOP processed');
return true;
}
async function handleStart(
msg: NormalizedMessage,
accountId: string,
prisma: any,
pool: WhatsAppSessionPool,
): Promise<boolean> {
const group = await prisma.group.findUnique({
where: { platform_platformId: { platform: 'whatsapp', platformId: msg.sourceGroupJid } },
select: { id: true, tenantId: true, claimStatus: true },
});
if (!group || group.claimStatus !== 'CLAIMED' || !group.tenantId) {
return false;
}
const phoneHash = `stop:${msg.senderJid}`;
const user = await prisma.towerUser.upsert({
where: { tenantId_phoneHash: { tenantId: group.tenantId, phoneHash } },
update: { jid: msg.senderJid },
create: {
tenantId: group.tenantId,
phoneHash,
jid: msg.senderJid,
displayName: msg.senderName ?? msg.senderJid,
},
});
const existing = await prisma.consentRecord.findFirst({
where: { userId: user.id, tenantId: group.tenantId, groupId: group.id },
});
if (existing) {
await prisma.consentRecord.update({
where: { id: existing.id },
data: {
status: 'GRANTED',
scopes: existing.scopes.length > 0 ? existing.scopes : ['INGEST', 'DISPLAY'],
retentionDays: existing.retentionDays,
revokedAt: null,
effectiveAt: new Date(),
},
});
} else {
await prisma.consentRecord.create({
data: {
tenantId: group.tenantId,
groupId: group.id,
userId: user.id,
scopes: ['INGEST', 'DISPLAY'],
retentionDays: 90,
policyVersion: 'v1',
status: 'GRANTED',
proofEventId: user.id,
},
});
}
try {
await pool.sendMessage(
accountId,
msg.senderJid,
"Welcome back. Your default scopes (INGEST, DISPLAY) are re-granted. Visit your portal to customize: " +
`${PORTAL_BASE}/my`,
);
} catch (err) {
logger.warn({ err, jid: msg.senderJid }, 'Failed to send START confirmation DM');
}
logger.info({ jid: msg.senderJid, groupId: group.id }, 'START processed');
return true;
}
async function handlePortal(
msg: NormalizedMessage,
accountId: string,
prisma: any,
pool: WhatsAppSessionPool,
): Promise<boolean> {
const group = await prisma.group.findUnique({
where: { platform_platformId: { platform: 'whatsapp', platformId: msg.sourceGroupJid } },
select: { id: true, tenantId: true, claimStatus: true },
});
if (!group || group.claimStatus !== 'CLAIMED' || !group.tenantId) {
return false;
}
// Reuse same phoneHash scheme as STOP/START so the same user is found
const phoneHash = `stop:${msg.senderJid}`;
const user = await prisma.towerUser.upsert({
where: { tenantId_phoneHash: { tenantId: group.tenantId, phoneHash } },
update: { jid: msg.senderJid },
create: {
tenantId: group.tenantId,
phoneHash,
jid: msg.senderJid,
displayName: msg.senderName ?? msg.senderJid,
},
});
const onboardingToken = await issueOnboardingToken(prisma, group.tenantId, group.id, msg.senderJid);
try {
await pool.sendMessage(
accountId,
msg.senderJid,
`Manage your data: ${PORTAL_BASE}/onboard?token=${onboardingToken}`,
);
} catch (err) {
logger.warn({ err, jid: msg.senderJid }, 'Failed to send PORTAL link DM');
}
return true;
}
async function handleCommands(
msg: NormalizedMessage,
accountId: string,
pool: WhatsAppSessionPool,
): Promise<boolean> {
try {
await pool.sendMessage(
accountId,
msg.senderJid,
'TOWER commands: STOP (opt out), START (rejoin), PORTAL (get your data link).',
);
} catch (err) {
logger.warn({ err, jid: msg.senderJid }, 'Failed to send COMMANDS reply');
}
return true;
}
async function issueOnboardingToken(
prisma: any,
tenantId: string,
groupId: string,
jid: string,
): Promise<string> {
// For Phase 2B the onboarding token is the base64url({groupId, jid, tenantId}) shape.
// The API re-verifies this token (decodes the payload) and trusts the OTP step
// for actual authentication. In a future phase we can sign with JWT_SECRET here.
const payload = { tenantId, groupId, jid };
return Buffer.from(JSON.stringify(payload), 'utf8').toString('base64url');
}