remove the telegram and whatsapp channels, reduce discord to notifications
All three existed to drive the platform from a chat app. The phone app does that now, so they are dead weight — three bot gateways, three command parsers, account pairing, admin config screens and bot tokens sitting in the database. Gone entirely: telegram/ and whatsapp/, discord/'s bot and command handler, the shared channel plumbing they were the only users of (pairing.ts, send-and-await.ts, types.ts, routes.ts and its 20 config/pairing/status endpoints), the six settings components, and their sections in the integrations screen. What Discord keeps is the one piece worth keeping — pushing a message out — as notify/discord.ts, configured by DISCORD_WEBHOOK_URL in the env. No UI, no pairing, no stored credential, and it never throws: a notification that fails to send is logged and dropped. Unset means notifications are silently skipped, which is the default state. Kept deliberately: send-claude-code.ts and send-opencode.ts. They live under channels/ but have nothing to do with chat apps — they are how /chat and the pipeline executor drive an agent turn. Also drops four dependencies with no remaining importer (discord.js, node-telegram-bot-api, whatsapp-web.js, qrcode), and the /sync-now route added to the email sidecar an hour ago, whose only consumer was the channel handlers. The "Channel Models" settings section stays, with its description corrected — it is keyed off the general access policy rather than anything channel-specific, so it governs non-owner accounts, not chat apps. Whether that whole class still earns its place is the open question already noted against origin-validation. Not touched: the telegram, whatsapp and discord rows in server_integrations, which still hold their bot tokens. Deleting rows is a different kind of decision and the SQL is in the handover. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -2,9 +2,6 @@ import { mkdirSync } from 'node:fs';
|
||||
import { DATA_PATH, ensureItemDirs } from './data-path';
|
||||
import { ensureToolLoader } from './ensure-tool-loader';
|
||||
// Queue is now owned by the sidecar process
|
||||
import { startDiscordBotIfConfigured } from './channels/discord/bot';
|
||||
import { startTelegramBotIfConfigured } from './channels/telegram/bot';
|
||||
import { startWhatsAppBotIfConfigured } from './channels/whatsapp/bot';
|
||||
import { startChatEventRetention } from './api/chat/retention';
|
||||
|
||||
mkdirSync(DATA_PATH, { recursive: true });
|
||||
@@ -16,15 +13,6 @@ ensureItemDirs();
|
||||
|
||||
startChatEventRetention();
|
||||
|
||||
await startDiscordBotIfConfigured().catch((err) => {
|
||||
console.error('[channels] Failed to start Discord bot:', err);
|
||||
});
|
||||
|
||||
await startTelegramBotIfConfigured().catch((err) => {
|
||||
console.error('[channels] Failed to start Telegram bot:', err);
|
||||
});
|
||||
|
||||
await startWhatsAppBotIfConfigured().catch((err) => {
|
||||
console.error('[channels] Failed to start WhatsApp bot:', err);
|
||||
});
|
||||
})();
|
||||
|
||||
@@ -1,65 +0,0 @@
|
||||
import { Client, GatewayIntentBits, Partials, Events } from 'discord.js';
|
||||
import { getServerIntegration } from 'officerdb';
|
||||
import { handleDiscordMessage } from './handler';
|
||||
|
||||
let client: Client | null = null;
|
||||
|
||||
export async function startDiscordBot(token: string): Promise<void> {
|
||||
if (client) {
|
||||
await stopDiscordBot();
|
||||
}
|
||||
|
||||
client = new Client({
|
||||
intents: [
|
||||
GatewayIntentBits.Guilds,
|
||||
GatewayIntentBits.GuildMembers,
|
||||
GatewayIntentBits.GuildPresences,
|
||||
GatewayIntentBits.DirectMessages,
|
||||
GatewayIntentBits.MessageContent,
|
||||
],
|
||||
partials: [Partials.Channel],
|
||||
});
|
||||
|
||||
client.on(Events.MessageCreate, (message) => {
|
||||
handleDiscordMessage(message).catch((err) => {
|
||||
console.error('[discord] Unhandled error in message handler:', err);
|
||||
});
|
||||
});
|
||||
|
||||
client.once(Events.ClientReady, (c) => {
|
||||
console.log(`[discord] Bot logged in as ${c.user.tag}`);
|
||||
});
|
||||
|
||||
await client.login(token);
|
||||
}
|
||||
|
||||
export async function stopDiscordBot(): Promise<void> {
|
||||
if (client) {
|
||||
client.destroy();
|
||||
client = null;
|
||||
console.log('[discord] Bot stopped');
|
||||
}
|
||||
}
|
||||
|
||||
export function isDiscordBotRunning(): boolean {
|
||||
return client !== null && client.isReady();
|
||||
}
|
||||
|
||||
export function getDiscordBotUsername(): string | null {
|
||||
return client?.user?.tag ?? null;
|
||||
}
|
||||
|
||||
type DiscordConfig = {
|
||||
botToken: string;
|
||||
};
|
||||
|
||||
export async function startDiscordBotIfConfigured(): Promise<void> {
|
||||
const integration = await getServerIntegration('discord');
|
||||
if (!integration?.enabled) return;
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
const botToken = config.botToken as string | undefined;
|
||||
if (!botToken) return;
|
||||
|
||||
await startDiscordBot(botToken);
|
||||
}
|
||||
@@ -1,51 +0,0 @@
|
||||
const MAX_LENGTH = 2000;
|
||||
|
||||
export function chunkMessage(text: string): string[] {
|
||||
if (text.length <= MAX_LENGTH) return [text];
|
||||
|
||||
const chunks: string[] = [];
|
||||
const paragraphs = text.split('\n\n');
|
||||
|
||||
let current = '';
|
||||
|
||||
for (const paragraph of paragraphs) {
|
||||
if (paragraph.length > MAX_LENGTH) {
|
||||
// Flush current chunk
|
||||
if (current) {
|
||||
chunks.push(current.trim());
|
||||
current = '';
|
||||
}
|
||||
// Split long paragraph on newlines
|
||||
const lines = paragraph.split('\n');
|
||||
for (const line of lines) {
|
||||
if (line.length > MAX_LENGTH) {
|
||||
// Flush current
|
||||
if (current) {
|
||||
chunks.push(current.trim());
|
||||
current = '';
|
||||
}
|
||||
// Hard-split long line
|
||||
for (let i = 0; i < line.length; i += MAX_LENGTH) {
|
||||
chunks.push(line.slice(i, i + MAX_LENGTH));
|
||||
}
|
||||
} else if (current.length + 1 + line.length > MAX_LENGTH) {
|
||||
chunks.push(current.trim());
|
||||
current = line;
|
||||
} else {
|
||||
current += (current ? '\n' : '') + line;
|
||||
}
|
||||
}
|
||||
} else if (current.length + 2 + paragraph.length > MAX_LENGTH) {
|
||||
chunks.push(current.trim());
|
||||
current = paragraph;
|
||||
} else {
|
||||
current += (current ? '\n\n' : '') + paragraph;
|
||||
}
|
||||
}
|
||||
|
||||
if (current.trim()) {
|
||||
chunks.push(current.trim());
|
||||
}
|
||||
|
||||
return chunks;
|
||||
}
|
||||
@@ -1,223 +0,0 @@
|
||||
import type { Message as DiscordMessage } from 'discord.js';
|
||||
import { findUserByIntegrationConfig, readConfigValue } from 'officerdb';
|
||||
import { sendAndAwait, getSessionModel, setSessionModel } from '../send-and-await';
|
||||
import { consumePairingCode } from '../pairing';
|
||||
import { chunkMessage } from './chunker';
|
||||
import { listChatModels } from '@@/api/chat/list-models';
|
||||
import { enqueueJob } from '../../queue/init';
|
||||
import { readJob } from '@@/queue/storage';
|
||||
import type { ModelInfo } from '@@/api/chat/types';
|
||||
import { toShellUsername } from '@@/data-path';
|
||||
|
||||
const PAIRING_CODE_PATTERN = /^[A-Z0-9]{6}$/;
|
||||
const TYPING_INTERVAL_MS = 8_000;
|
||||
|
||||
type SendableChannel = { send: (content: string) => Promise<unknown> };
|
||||
|
||||
type AccessPolicy = { allowedModels: string[] };
|
||||
const ACCESS_POLICY_KEY = 'chat-access-policy';
|
||||
|
||||
async function getVisibleModels(): Promise<ModelInfo[]> {
|
||||
const allModels = await listChatModels();
|
||||
const policy = await readConfigValue<AccessPolicy>(ACCESS_POLICY_KEY, { allowedModels: [] });
|
||||
const allowed = policy.allowedModels;
|
||||
|
||||
if (allowed.length === 0) return allModels;
|
||||
|
||||
const allowedSet = new Set(allowed);
|
||||
const allowedProviderSet = new Set(allowed.map((key) => key.split(':')[0]));
|
||||
|
||||
return allModels.filter((m) => {
|
||||
const key = `${m.provider}:${m.id}`;
|
||||
const isExplicitlyAllowed = allowedSet.has(key);
|
||||
const isFromNewProvider = !allowedProviderSet.has(m.provider);
|
||||
return isExplicitlyAllowed || isFromNewProvider;
|
||||
});
|
||||
}
|
||||
|
||||
import { runEmailSyncCommand } from '../email-sync-command';
|
||||
|
||||
type CommandContext = {
|
||||
content: string;
|
||||
channel: SendableChannel;
|
||||
userId: number;
|
||||
email: string;
|
||||
discordId: string;
|
||||
};
|
||||
|
||||
async function handleEmailSync(ctx: CommandContext): Promise<void> {
|
||||
const { channel, userId } = ctx;
|
||||
await channel.send('Syncing emails...');
|
||||
const text = await runEmailSyncCommand(userId);
|
||||
for (const chunk of chunkMessage(text)) await channel.send(chunk);
|
||||
}
|
||||
|
||||
async function handleCommand(ctx: CommandContext): Promise<boolean> {
|
||||
const { content, channel, discordId } = ctx;
|
||||
const lower = content.toLowerCase();
|
||||
|
||||
const helpSections: Record<string, string> = {
|
||||
models:
|
||||
'**Models:**\n' +
|
||||
'`!model` — show current model\n' +
|
||||
'`!model <id>` — switch model\n' +
|
||||
'`!models` — list available models',
|
||||
email: '**Email:**\n' + '`!email sync` — sync Gmail and show new emails',
|
||||
};
|
||||
|
||||
if (lower === '!help' || lower.startsWith('!help ')) {
|
||||
const topic = content.slice('!help'.length).trim().toLowerCase();
|
||||
if (topic && topic in helpSections) {
|
||||
await channel.send(helpSections[topic]!);
|
||||
return true;
|
||||
}
|
||||
if (topic) {
|
||||
await channel.send(
|
||||
`Unknown topic: \`${topic}\`\nAvailable: ${Object.keys(helpSections)
|
||||
.map((k) => `\`${k}\``)
|
||||
.join(', ')}`,
|
||||
);
|
||||
return true;
|
||||
}
|
||||
const full = Object.values(helpSections).join('\n\n');
|
||||
await channel.send(full + '\n\n`!help <topic>` — show commands for a topic');
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!models') {
|
||||
const models = await getVisibleModels();
|
||||
if (models.length === 0) {
|
||||
await channel.send('No models available.');
|
||||
return true;
|
||||
}
|
||||
const current = getSessionModel('discord', ctx.userId, discordId);
|
||||
const grouped = new Map<string, string[]>();
|
||||
for (const m of models) {
|
||||
const list = grouped.get(m.provider) ?? [];
|
||||
list.push(m.id === current ? `**${m.id}** (current)` : m.id);
|
||||
grouped.set(m.provider, list);
|
||||
}
|
||||
let text = '**Available models:**\n';
|
||||
for (const [provider, ids] of grouped) {
|
||||
text += `\n__${provider}__\n${ids.map((id) => ` ${id}`).join('\n')}\n`;
|
||||
}
|
||||
text += '\nUse `!model <id>` to switch.';
|
||||
await channel.send(text);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!model') {
|
||||
const current = getSessionModel('discord', ctx.userId, discordId);
|
||||
await channel.send(
|
||||
current
|
||||
? `Current model: **${current}**`
|
||||
: 'No active session yet — the default model will be used on your next message.',
|
||||
);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower.startsWith('!model ')) {
|
||||
const requested = content.slice('!model '.length).trim();
|
||||
if (!requested) {
|
||||
const current = getSessionModel('discord', ctx.userId, discordId);
|
||||
await channel.send(current ? `Current model: **${current}**` : 'No active session yet.');
|
||||
return true;
|
||||
}
|
||||
const models = await getVisibleModels();
|
||||
const match = models.find((m) => m.id === requested || m.name === requested);
|
||||
if (!match) {
|
||||
await channel.send(`Model not found: \`${requested}\`\nUse \`!models\` to see available models.`);
|
||||
return true;
|
||||
}
|
||||
setSessionModel('discord', ctx.userId, discordId, match.id);
|
||||
await channel.send(`Model switched to **${match.id}**. The new model will be used on your next message.`);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!email sync') {
|
||||
await handleEmailSync(ctx);
|
||||
return true;
|
||||
}
|
||||
|
||||
// Not a recognized command — pass through to PI
|
||||
return false;
|
||||
}
|
||||
|
||||
export async function handleDiscordMessage(message: DiscordMessage): Promise<void> {
|
||||
// Ignore bots and non-DM messages
|
||||
if (message.author.bot) return;
|
||||
if (!message.channel.isDMBased() || !('send' in message.channel)) return;
|
||||
|
||||
const channel = message.channel;
|
||||
const discordId = message.author.id;
|
||||
const content = message.content.trim();
|
||||
if (!content) return;
|
||||
|
||||
// Look up linked Officer user
|
||||
const linked = await findUserByIntegrationConfig('discord', 'discordId', discordId);
|
||||
|
||||
if (!linked) {
|
||||
// Check if this is a pairing code
|
||||
if (PAIRING_CODE_PATTERN.test(content.toUpperCase())) {
|
||||
const result = await consumePairingCode(content.toUpperCase(), discordId);
|
||||
if (result) {
|
||||
await channel.send('Account linked! You can now chat with me.');
|
||||
return;
|
||||
}
|
||||
await channel.send('Invalid or expired pairing code. Please generate a new one from Officer Settings.');
|
||||
return;
|
||||
}
|
||||
|
||||
await channel.send(
|
||||
"I don't recognize your Discord account. To link it:\n" +
|
||||
'1. Go to Officer Settings → Integrations → Discord\n' +
|
||||
'2. Click "Link Discord" to get a pairing code\n' +
|
||||
'3. Send the 6-character code to me here',
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle commands
|
||||
if (content.startsWith('!')) {
|
||||
const handled = await handleCommand({
|
||||
content,
|
||||
channel,
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
discordId,
|
||||
});
|
||||
if (handled) return;
|
||||
}
|
||||
|
||||
// Start typing indicator with keep-alive
|
||||
const sendTyping = () => {
|
||||
if ('sendTyping' in channel) {
|
||||
(channel as { sendTyping: () => Promise<void> }).sendTyping().catch(() => {});
|
||||
}
|
||||
};
|
||||
const typingInterval = setInterval(sendTyping, TYPING_INTERVAL_MS);
|
||||
sendTyping();
|
||||
|
||||
try {
|
||||
const result = await sendAndAwait({
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
username: toShellUsername(linked.user.username ?? '', linked.user.email),
|
||||
prompt: content,
|
||||
context: 'discord',
|
||||
contextId: discordId,
|
||||
});
|
||||
|
||||
clearInterval(typingInterval);
|
||||
|
||||
const signature = `\`${result.model}\`\n`;
|
||||
const chunks = chunkMessage(result.text);
|
||||
for (let i = 0; i < chunks.length; i++) {
|
||||
await channel.send(i === 0 ? signature + chunks[i]! : chunks[i]!);
|
||||
}
|
||||
} catch (err) {
|
||||
clearInterval(typingInterval);
|
||||
console.error('[discord] Error handling message:', err);
|
||||
await channel.send('Sorry, something went wrong processing your message.').catch(() => {});
|
||||
}
|
||||
}
|
||||
@@ -1,45 +0,0 @@
|
||||
import { getEmailServerUrl } from '../api/email/router';
|
||||
|
||||
// The "sync my email" chat command, once, for all three channels.
|
||||
//
|
||||
// Telegram, Discord and WhatsApp each carried their own copy: open the mail store directly, count rows,
|
||||
// enqueue a `gmail-sync` job, poll it, count again, diff. That coupled three chat bridges to the mail
|
||||
// schema, and the job type was hardcoded to the OAuth path even though the account syncs over IMAP — so
|
||||
// the command was already broken before sync moved into the sidecar. Now it is one HTTP call to the
|
||||
// sidecar, which does the sync and reports what arrived.
|
||||
|
||||
type SyncNowResponse = {
|
||||
saved: number;
|
||||
skipped?: number;
|
||||
errors?: number;
|
||||
newest: Array<{ from_name: string | null; from_address: string; subject: string }>;
|
||||
error?: string;
|
||||
};
|
||||
|
||||
/** Runs the sync and returns the message to send back to the user. */
|
||||
export async function runEmailSyncCommand(userId: number): Promise<string> {
|
||||
const base = getEmailServerUrl();
|
||||
if (!base) return 'Email is not available right now — the mail service is starting up.';
|
||||
|
||||
let res: Response;
|
||||
try {
|
||||
res = await fetch(`${base}/sync-now`, {
|
||||
method: 'POST',
|
||||
// Loopback-only, same trust as the platform's own proxy.
|
||||
headers: { 'X-Officer-User': String(userId) },
|
||||
});
|
||||
} catch {
|
||||
return 'Email sync failed — the mail service is unreachable.';
|
||||
}
|
||||
|
||||
if (!res.ok) return `Email sync failed (${res.status}).`;
|
||||
|
||||
const body = (await res.json()) as SyncNowResponse;
|
||||
if (body.error) return body.error;
|
||||
if (body.saved <= 0) return 'Sync complete — no new emails.';
|
||||
|
||||
const lines = body.newest.map((e) => `- *${e.from_name || e.from_address}*: ${e.subject}`);
|
||||
let text = `Sync complete — *${body.saved}* new email${body.saved !== 1 ? 's' : ''}`;
|
||||
if (body.saved > 20) text += ' (showing latest 20)';
|
||||
return `${text}:\n\n${lines.join('\n')}`;
|
||||
}
|
||||
@@ -1,84 +0,0 @@
|
||||
import type { ChannelProvider } from './types';
|
||||
import { upsertUserIntegration } from 'officerdb';
|
||||
|
||||
const configKeyMap: Record<ChannelProvider, string> = {
|
||||
discord: 'discordId',
|
||||
telegram: 'telegramId',
|
||||
whatsapp: 'whatsappId',
|
||||
};
|
||||
|
||||
type PairingEntry = {
|
||||
userId: number;
|
||||
email: string;
|
||||
provider: ChannelProvider;
|
||||
expiresAt: number;
|
||||
};
|
||||
|
||||
const pairingCodes = new Map<string, PairingEntry>();
|
||||
|
||||
const CODE_TTL_MS = 10 * 60 * 1000; // 10 minutes
|
||||
const CODE_LENGTH = 6;
|
||||
const CODE_CHARS = 'ABCDEFGHJKLMNPQRSTUVWXYZ23456789'; // no 0/O/1/I ambiguity
|
||||
|
||||
function generateCode(): string {
|
||||
let code = '';
|
||||
for (let i = 0; i < CODE_LENGTH; i++) {
|
||||
code += CODE_CHARS[Math.floor(Math.random() * CODE_CHARS.length)]!;
|
||||
}
|
||||
return code;
|
||||
}
|
||||
|
||||
function cleanupExpiredCodes(): void {
|
||||
const now = Date.now();
|
||||
for (const [code, entry] of pairingCodes) {
|
||||
if (entry.expiresAt <= now) {
|
||||
pairingCodes.delete(code);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function generatePairingCode(userId: number, email: string, provider: ChannelProvider): string {
|
||||
cleanupExpiredCodes();
|
||||
|
||||
// Revoke any existing code for this user+provider
|
||||
for (const [code, entry] of pairingCodes) {
|
||||
if (entry.userId === userId && entry.provider === provider) {
|
||||
pairingCodes.delete(code);
|
||||
}
|
||||
}
|
||||
|
||||
let code: string;
|
||||
do {
|
||||
code = generateCode();
|
||||
} while (pairingCodes.has(code));
|
||||
|
||||
pairingCodes.set(code, {
|
||||
userId,
|
||||
email,
|
||||
provider,
|
||||
expiresAt: Date.now() + CODE_TTL_MS,
|
||||
});
|
||||
|
||||
return code;
|
||||
}
|
||||
|
||||
type PairingResult = {
|
||||
userId: number;
|
||||
email: string;
|
||||
provider: ChannelProvider;
|
||||
};
|
||||
|
||||
export async function consumePairingCode(code: string, channelUserId: string): Promise<PairingResult | null> {
|
||||
const entry = pairingCodes.get(code.toUpperCase());
|
||||
if (!entry || entry.expiresAt <= Date.now()) return null;
|
||||
|
||||
pairingCodes.delete(code.toUpperCase());
|
||||
|
||||
await upsertUserIntegration({
|
||||
userId: entry.userId,
|
||||
provider: entry.provider,
|
||||
config: { [configKeyMap[entry.provider]]: channelUserId },
|
||||
});
|
||||
|
||||
return { userId: entry.userId, email: entry.email, provider: entry.provider };
|
||||
}
|
||||
@@ -1,312 +0,0 @@
|
||||
import { createRouter } from '@@/create-router';
|
||||
import { getServerIntegration, upsertServerIntegration, getUserIntegration, deleteUserIntegration } from 'officerdb';
|
||||
import { startDiscordBot, stopDiscordBot, isDiscordBotRunning, getDiscordBotUsername } from './discord/bot';
|
||||
import { startTelegramBot, stopTelegramBot, isTelegramBotRunning, getTelegramBotUsername } from './telegram/bot';
|
||||
import {
|
||||
startWhatsAppBot,
|
||||
stopWhatsAppBot,
|
||||
isWhatsAppBotRunning,
|
||||
getWhatsAppBotPhone,
|
||||
getWhatsAppQR,
|
||||
disconnectWhatsApp,
|
||||
} from './whatsapp/bot';
|
||||
import { generatePairingCode } from './pairing';
|
||||
|
||||
export const channelsRouter = createRouter();
|
||||
|
||||
// ── Admin: Discord config ──
|
||||
|
||||
channelsRouter.get('/discord/config', async (ctx) => {
|
||||
const integration = await getServerIntegration('discord');
|
||||
if (!integration) return ctx.json({ configured: false });
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
const botToken = config.botToken as string | undefined;
|
||||
|
||||
return ctx.json({
|
||||
configured: !!botToken,
|
||||
enabled: integration.enabled,
|
||||
botToken: botToken ? `${botToken.slice(0, 8)}...${botToken.slice(-4)}` : null,
|
||||
serverInvite: (config.serverInvite as string) ?? null,
|
||||
botHandle: (config.botHandle as string) ?? null,
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.put('/discord/config', async (ctx) => {
|
||||
const body = ctx.get('body') as Record<string, unknown>;
|
||||
const botToken = body.botToken as string | undefined;
|
||||
const enabled = body.enabled as boolean | undefined;
|
||||
const serverInvite = body.serverInvite as string | undefined;
|
||||
const botHandle = body.botHandle as string | undefined;
|
||||
|
||||
if (!botToken && enabled === undefined && serverInvite === undefined && botHandle === undefined) {
|
||||
return ctx.json({ error: 'At least one field required' }, 400);
|
||||
}
|
||||
|
||||
const existing = await getServerIntegration('discord');
|
||||
const existingConfig = (existing?.config ?? {}) as Record<string, unknown>;
|
||||
const newConfig = { ...existingConfig };
|
||||
if (botToken) newConfig.botToken = botToken;
|
||||
if (serverInvite !== undefined) newConfig.serverInvite = serverInvite;
|
||||
if (botHandle !== undefined) newConfig.botHandle = botHandle;
|
||||
|
||||
const shouldRun = enabled ?? existing?.enabled ?? true;
|
||||
const token = (botToken ?? existingConfig.botToken) as string | undefined;
|
||||
|
||||
// If a new token is provided, validate it by starting the bot before saving
|
||||
if (botToken && shouldRun) {
|
||||
try {
|
||||
await startDiscordBot(botToken);
|
||||
} catch (err) {
|
||||
console.error('[channels] Discord bot token validation failed:', err);
|
||||
return ctx.json({ error: 'Invalid bot token — connection failed' }, 400);
|
||||
}
|
||||
}
|
||||
|
||||
await upsertServerIntegration('discord', newConfig, shouldRun);
|
||||
|
||||
// Start/restart with existing token (already validated on initial save)
|
||||
if (!botToken && token && shouldRun) {
|
||||
try {
|
||||
await startDiscordBot(token);
|
||||
} catch (err) {
|
||||
console.error('[channels] Failed to start Discord bot:', err);
|
||||
return ctx.json({ success: true, botStarted: false, error: String(err) });
|
||||
}
|
||||
}
|
||||
|
||||
if (!shouldRun) {
|
||||
await stopDiscordBot();
|
||||
}
|
||||
|
||||
return ctx.json({ success: true, botStarted: token && shouldRun });
|
||||
});
|
||||
|
||||
channelsRouter.get('/discord/status', async (ctx) => {
|
||||
const integration = await getServerIntegration('discord');
|
||||
const config = (integration?.config ?? {}) as Record<string, unknown>;
|
||||
|
||||
return ctx.json({
|
||||
configured: !!config.botToken,
|
||||
enabled: integration?.enabled ?? false,
|
||||
running: isDiscordBotRunning(),
|
||||
botUsername: getDiscordBotUsername(),
|
||||
serverInvite: (config.serverInvite as string) ?? null,
|
||||
botHandle: (config.botHandle as string) ?? null,
|
||||
});
|
||||
});
|
||||
|
||||
// ── User: Discord pairing ──
|
||||
|
||||
channelsRouter.post('/discord/pair', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const code = generatePairingCode(user.id, user.email, 'discord');
|
||||
return ctx.json({ code, expiresIn: 600 });
|
||||
});
|
||||
|
||||
channelsRouter.get('/discord/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const integration = await getUserIntegration(user.id, 'discord');
|
||||
if (!integration) return ctx.json({ linked: false });
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
return ctx.json({
|
||||
linked: true,
|
||||
discordId: config.discordId,
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.delete('/discord/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const deleted = await deleteUserIntegration(user.id, 'discord');
|
||||
return ctx.json({ success: deleted });
|
||||
});
|
||||
|
||||
// ── Admin: Telegram config ──
|
||||
|
||||
channelsRouter.get('/telegram/config', async (ctx) => {
|
||||
const integration = await getServerIntegration('telegram');
|
||||
if (!integration) return ctx.json({ configured: false });
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
const botToken = config.botToken as string | undefined;
|
||||
|
||||
return ctx.json({
|
||||
configured: !!botToken,
|
||||
enabled: integration.enabled,
|
||||
botToken: botToken ? `${botToken.slice(0, 8)}...${botToken.slice(-4)}` : null,
|
||||
serverInvite: (config.serverInvite as string) ?? null,
|
||||
botHandle: (config.botHandle as string) ?? null,
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.put('/telegram/config', async (ctx) => {
|
||||
const body = ctx.get('body') as Record<string, unknown>;
|
||||
const botToken = body.botToken as string | undefined;
|
||||
const enabled = body.enabled as boolean | undefined;
|
||||
const serverInvite = body.serverInvite as string | undefined;
|
||||
const botHandle = body.botHandle as string | undefined;
|
||||
|
||||
if (!botToken && enabled === undefined && serverInvite === undefined && botHandle === undefined) {
|
||||
return ctx.json({ error: 'At least one field required' }, 400);
|
||||
}
|
||||
|
||||
const existing = await getServerIntegration('telegram');
|
||||
const existingConfig = (existing?.config ?? {}) as Record<string, unknown>;
|
||||
const newConfig = { ...existingConfig };
|
||||
if (botToken) newConfig.botToken = botToken;
|
||||
if (serverInvite !== undefined) newConfig.serverInvite = serverInvite;
|
||||
if (botHandle !== undefined) newConfig.botHandle = botHandle;
|
||||
|
||||
const shouldRun = enabled ?? existing?.enabled ?? true;
|
||||
const token = (botToken ?? existingConfig.botToken) as string | undefined;
|
||||
|
||||
// If a new token is provided, validate it by starting the bot before saving
|
||||
if (botToken && shouldRun) {
|
||||
try {
|
||||
await startTelegramBot(botToken);
|
||||
} catch (err) {
|
||||
console.error('[channels] Telegram bot token validation failed:', err);
|
||||
return ctx.json({ error: 'Invalid bot token — connection failed' }, 400);
|
||||
}
|
||||
}
|
||||
|
||||
await upsertServerIntegration('telegram', newConfig, shouldRun);
|
||||
|
||||
// Start/restart with existing token (already validated on initial save)
|
||||
if (!botToken && token && shouldRun) {
|
||||
try {
|
||||
await startTelegramBot(token);
|
||||
} catch (err) {
|
||||
console.error('[channels] Failed to start Telegram bot:', err);
|
||||
return ctx.json({ success: true, botStarted: false, error: String(err) });
|
||||
}
|
||||
}
|
||||
|
||||
if (!shouldRun) {
|
||||
await stopTelegramBot();
|
||||
}
|
||||
|
||||
return ctx.json({ success: true, botStarted: token && shouldRun });
|
||||
});
|
||||
|
||||
channelsRouter.get('/telegram/status', async (ctx) => {
|
||||
const integration = await getServerIntegration('telegram');
|
||||
const config = (integration?.config ?? {}) as Record<string, unknown>;
|
||||
|
||||
return ctx.json({
|
||||
configured: !!config.botToken,
|
||||
enabled: integration?.enabled ?? false,
|
||||
running: isTelegramBotRunning(),
|
||||
botUsername: getTelegramBotUsername(),
|
||||
serverInvite: (config.serverInvite as string) ?? null,
|
||||
botHandle: (config.botHandle as string) ?? null,
|
||||
});
|
||||
});
|
||||
|
||||
// ── User: Telegram pairing ──
|
||||
|
||||
channelsRouter.post('/telegram/pair', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const code = generatePairingCode(user.id, user.email, 'telegram');
|
||||
return ctx.json({ code, expiresIn: 600 });
|
||||
});
|
||||
|
||||
channelsRouter.get('/telegram/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const integration = await getUserIntegration(user.id, 'telegram');
|
||||
if (!integration) return ctx.json({ linked: false });
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
return ctx.json({
|
||||
linked: true,
|
||||
telegramId: config.telegramId,
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.delete('/telegram/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const deleted = await deleteUserIntegration(user.id, 'telegram');
|
||||
return ctx.json({ success: deleted });
|
||||
});
|
||||
|
||||
// ── Admin: WhatsApp config ──
|
||||
|
||||
channelsRouter.get('/whatsapp/config', async (ctx) => {
|
||||
const integration = await getServerIntegration('whatsapp');
|
||||
|
||||
return ctx.json({
|
||||
configured: integration?.enabled ?? false,
|
||||
enabled: integration?.enabled ?? false,
|
||||
running: isWhatsAppBotRunning(),
|
||||
phone: getWhatsAppBotPhone(),
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.put('/whatsapp/config', async (ctx) => {
|
||||
const body = ctx.get('body') as Record<string, unknown>;
|
||||
const enabled = body.enabled as boolean | undefined;
|
||||
|
||||
if (enabled === true) {
|
||||
await upsertServerIntegration('whatsapp', {}, true);
|
||||
// Don't await — initialization is slow (launches Chromium) and QR events
|
||||
// are delivered via SSE. Return immediately so the client can connect SSE.
|
||||
startWhatsAppBot().catch((err) => {
|
||||
console.error('[channels] Failed to start WhatsApp bot:', err);
|
||||
});
|
||||
return ctx.json({ success: true, botStarted: false });
|
||||
}
|
||||
|
||||
if (enabled === false) {
|
||||
await disconnectWhatsApp();
|
||||
return ctx.json({ success: true, botStarted: false });
|
||||
}
|
||||
|
||||
return ctx.json({ error: 'enabled field required' }, 400);
|
||||
});
|
||||
|
||||
channelsRouter.get('/whatsapp/status', async (ctx) => {
|
||||
const integration = await getServerIntegration('whatsapp');
|
||||
|
||||
return ctx.json({
|
||||
configured: integration?.enabled ?? false,
|
||||
enabled: integration?.enabled ?? false,
|
||||
running: isWhatsAppBotRunning(),
|
||||
phone: getWhatsAppBotPhone(),
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.get('/whatsapp/qr', async (ctx) => {
|
||||
const qr = getWhatsAppQR();
|
||||
return ctx.json({
|
||||
qr,
|
||||
running: isWhatsAppBotRunning(),
|
||||
phone: getWhatsAppBotPhone(),
|
||||
});
|
||||
});
|
||||
|
||||
// ── User: WhatsApp pairing ──
|
||||
|
||||
channelsRouter.post('/whatsapp/pair', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const code = generatePairingCode(user.id, user.email, 'whatsapp');
|
||||
return ctx.json({ code, expiresIn: 600 });
|
||||
});
|
||||
|
||||
channelsRouter.get('/whatsapp/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const integration = await getUserIntegration(user.id, 'whatsapp');
|
||||
if (!integration) return ctx.json({ linked: false });
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
return ctx.json({
|
||||
linked: true,
|
||||
whatsappId: config.whatsappId,
|
||||
});
|
||||
});
|
||||
|
||||
channelsRouter.delete('/whatsapp/connection', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const deleted = await deleteUserIntegration(user.id, 'whatsapp');
|
||||
return ctx.json({ success: deleted });
|
||||
});
|
||||
@@ -1,90 +0,0 @@
|
||||
import type { MessageCost } from '@@/api/chat/types';
|
||||
import { getUserSettings } from 'officerdb';
|
||||
import { logger } from '@@/api/chat/logger';
|
||||
import { sendClaudeCode, clearClaudeCodeSession } from './send-claude-code';
|
||||
|
||||
const DEFAULT_MODEL = 'claude-code';
|
||||
|
||||
type SendAndAwaitParams = {
|
||||
userId: number;
|
||||
email: string;
|
||||
username: string;
|
||||
prompt: string;
|
||||
context: string;
|
||||
contextId: string;
|
||||
model?: string;
|
||||
};
|
||||
|
||||
type SendAndAwaitResult = {
|
||||
text: string;
|
||||
sessionId: string;
|
||||
model: string;
|
||||
cost: MessageCost;
|
||||
};
|
||||
|
||||
// Per-session mutex to serialize concurrent prompts
|
||||
const sessionLocks = new Map<string, Promise<void>>();
|
||||
|
||||
// Channel model overrides — survive session eviction/recreation
|
||||
const channelModelOverrides = new Map<string, string>();
|
||||
|
||||
function buildSessionId(context: string, userId: number, contextId: string): string {
|
||||
return `channel-${context}-${userId}-${contextId}`;
|
||||
}
|
||||
|
||||
async function getUserDefaultModel(userId: number): Promise<string | null> {
|
||||
try {
|
||||
const settings = await getUserSettings(userId);
|
||||
const chat = settings?.chat as Record<string, unknown> | undefined;
|
||||
return (chat?.defaultModel as string) || null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export function getSessionModel(context: string, userId: number, contextId: string): string | null {
|
||||
return channelModelOverrides.get(buildSessionId(context, userId, contextId)) ?? null;
|
||||
}
|
||||
|
||||
export function setSessionModel(context: string, userId: number, contextId: string, model: string): void {
|
||||
const sessionId = buildSessionId(context, userId, contextId);
|
||||
channelModelOverrides.set(sessionId, model);
|
||||
// Reset the Claude session so the next prompt starts fresh under the new model.
|
||||
clearClaudeCodeSession(sessionId);
|
||||
logger.info('Channel model override stored', { sessionId, model });
|
||||
}
|
||||
|
||||
export async function sendAndAwait(params: SendAndAwaitParams): Promise<SendAndAwaitResult> {
|
||||
const { userId, context, contextId } = params;
|
||||
const sessionId = buildSessionId(context, userId, contextId);
|
||||
|
||||
// Serialize per session — if two messages arrive at once, the second waits for the first.
|
||||
const existing = sessionLocks.get(sessionId) ?? Promise.resolve();
|
||||
let releaseLock: () => void;
|
||||
const lockPromise = new Promise<void>((resolve) => {
|
||||
releaseLock = resolve;
|
||||
});
|
||||
const chained = existing.then(() => lockPromise);
|
||||
sessionLocks.set(sessionId, chained);
|
||||
|
||||
await existing;
|
||||
|
||||
try {
|
||||
const override = channelModelOverrides.get(sessionId);
|
||||
let model = params.model ?? override ?? (await getUserDefaultModel(userId)) ?? DEFAULT_MODEL;
|
||||
// Claude-only: coerce any legacy non-Claude model preference to the Claude default.
|
||||
if (!model.startsWith('claude-code')) model = DEFAULT_MODEL;
|
||||
|
||||
return await sendClaudeCode({
|
||||
userId: params.userId,
|
||||
email: params.email,
|
||||
username: params.username,
|
||||
prompt: params.prompt,
|
||||
sessionKey: sessionId,
|
||||
model,
|
||||
});
|
||||
} finally {
|
||||
releaseLock!();
|
||||
if (sessionLocks.get(sessionId) === chained) sessionLocks.delete(sessionId);
|
||||
}
|
||||
}
|
||||
@@ -1,56 +0,0 @@
|
||||
import TelegramBot from 'node-telegram-bot-api';
|
||||
import { getServerIntegration } from 'officerdb';
|
||||
import { handleTelegramMessage } from './handler';
|
||||
|
||||
let bot: TelegramBot | null = null;
|
||||
let cachedUsername: string | null = null;
|
||||
|
||||
export async function startTelegramBot(token: string): Promise<void> {
|
||||
if (bot) {
|
||||
await stopTelegramBot();
|
||||
}
|
||||
|
||||
bot = new TelegramBot(token, { polling: true });
|
||||
|
||||
bot.on('message', (msg) => {
|
||||
handleTelegramMessage(msg).catch((err) => {
|
||||
console.error('[telegram] Unhandled error in message handler:', err);
|
||||
});
|
||||
});
|
||||
|
||||
const me = await bot.getMe();
|
||||
cachedUsername = me.username ?? null;
|
||||
console.log(`[telegram] Bot logged in as @${cachedUsername}`);
|
||||
}
|
||||
|
||||
export async function stopTelegramBot(): Promise<void> {
|
||||
if (bot) {
|
||||
await bot.stopPolling();
|
||||
bot = null;
|
||||
cachedUsername = null;
|
||||
console.log('[telegram] Bot stopped');
|
||||
}
|
||||
}
|
||||
|
||||
export function isTelegramBotRunning(): boolean {
|
||||
return bot !== null && bot.isPolling();
|
||||
}
|
||||
|
||||
export function getTelegramBotUsername(): string | null {
|
||||
return cachedUsername;
|
||||
}
|
||||
|
||||
export function getTelegramBot(): TelegramBot | null {
|
||||
return bot;
|
||||
}
|
||||
|
||||
export async function startTelegramBotIfConfigured(): Promise<void> {
|
||||
const integration = await getServerIntegration('telegram');
|
||||
if (!integration?.enabled) return;
|
||||
|
||||
const config = integration.config as Record<string, unknown>;
|
||||
const botToken = config.botToken as string | undefined;
|
||||
if (!botToken) return;
|
||||
|
||||
await startTelegramBot(botToken);
|
||||
}
|
||||
@@ -1,51 +0,0 @@
|
||||
const MAX_LENGTH = 4096;
|
||||
|
||||
export function chunkMessage(text: string): string[] {
|
||||
if (text.length <= MAX_LENGTH) return [text];
|
||||
|
||||
const chunks: string[] = [];
|
||||
const paragraphs = text.split('\n\n');
|
||||
|
||||
let current = '';
|
||||
|
||||
for (const paragraph of paragraphs) {
|
||||
if (paragraph.length > MAX_LENGTH) {
|
||||
// Flush current chunk
|
||||
if (current) {
|
||||
chunks.push(current.trim());
|
||||
current = '';
|
||||
}
|
||||
// Split long paragraph on newlines
|
||||
const lines = paragraph.split('\n');
|
||||
for (const line of lines) {
|
||||
if (line.length > MAX_LENGTH) {
|
||||
// Flush current
|
||||
if (current) {
|
||||
chunks.push(current.trim());
|
||||
current = '';
|
||||
}
|
||||
// Hard-split long line
|
||||
for (let i = 0; i < line.length; i += MAX_LENGTH) {
|
||||
chunks.push(line.slice(i, i + MAX_LENGTH));
|
||||
}
|
||||
} else if (current.length + 1 + line.length > MAX_LENGTH) {
|
||||
chunks.push(current.trim());
|
||||
current = line;
|
||||
} else {
|
||||
current += (current ? '\n' : '') + line;
|
||||
}
|
||||
}
|
||||
} else if (current.length + 2 + paragraph.length > MAX_LENGTH) {
|
||||
chunks.push(current.trim());
|
||||
current = paragraph;
|
||||
} else {
|
||||
current += (current ? '\n\n' : '') + paragraph;
|
||||
}
|
||||
}
|
||||
|
||||
if (current.trim()) {
|
||||
chunks.push(current.trim());
|
||||
}
|
||||
|
||||
return chunks;
|
||||
}
|
||||
@@ -1,226 +0,0 @@
|
||||
import type TelegramBot from 'node-telegram-bot-api';
|
||||
import { findUserByIntegrationConfig, readConfigValue } from 'officerdb';
|
||||
import { sendAndAwait, getSessionModel, setSessionModel } from '../send-and-await';
|
||||
import { consumePairingCode } from '../pairing';
|
||||
import { chunkMessage } from './chunker';
|
||||
import { getTelegramBot } from './bot';
|
||||
import { listChatModels } from '@@/api/chat/list-models';
|
||||
import { enqueueJob } from '../../queue/init';
|
||||
import { readJob } from '@@/queue/storage';
|
||||
import type { ModelInfo } from '@@/api/chat/types';
|
||||
import { toShellUsername } from '@@/data-path';
|
||||
|
||||
const PAIRING_CODE_PATTERN = /^[A-Z0-9]{6}$/;
|
||||
const TYPING_INTERVAL_MS = 5_000;
|
||||
|
||||
type SendFn = (text: string) => Promise<unknown>;
|
||||
|
||||
type AccessPolicy = { allowedModels: string[] };
|
||||
const ACCESS_POLICY_KEY = 'chat-access-policy';
|
||||
|
||||
async function getVisibleModels(): Promise<ModelInfo[]> {
|
||||
const allModels = await listChatModels();
|
||||
const policy = await readConfigValue<AccessPolicy>(ACCESS_POLICY_KEY, { allowedModels: [] });
|
||||
const allowed = policy.allowedModels;
|
||||
|
||||
if (allowed.length === 0) return allModels;
|
||||
|
||||
const allowedSet = new Set(allowed);
|
||||
const allowedProviderSet = new Set(allowed.map((key) => key.split(':')[0]));
|
||||
|
||||
return allModels.filter((m) => {
|
||||
const key = `${m.provider}:${m.id}`;
|
||||
const isExplicitlyAllowed = allowedSet.has(key);
|
||||
const isFromNewProvider = !allowedProviderSet.has(m.provider);
|
||||
return isExplicitlyAllowed || isFromNewProvider;
|
||||
});
|
||||
}
|
||||
|
||||
import { runEmailSyncCommand } from '../email-sync-command';
|
||||
|
||||
type CommandContext = {
|
||||
content: string;
|
||||
send: SendFn;
|
||||
userId: number;
|
||||
email: string;
|
||||
telegramId: string;
|
||||
};
|
||||
|
||||
async function handleEmailSync(ctx: CommandContext): Promise<void> {
|
||||
const { send, userId } = ctx;
|
||||
await send('Syncing emails...');
|
||||
const text = await runEmailSyncCommand(userId);
|
||||
for (const chunk of chunkMessage(text)) await send(chunk);
|
||||
}
|
||||
|
||||
async function handleCommand(ctx: CommandContext): Promise<boolean> {
|
||||
const { content, send, telegramId } = ctx;
|
||||
const lower = content.toLowerCase();
|
||||
|
||||
const helpSections: Record<string, string> = {
|
||||
models:
|
||||
'*Models:*\n' +
|
||||
'`!model` — show current model\n' +
|
||||
'`!model <id>` — switch model\n' +
|
||||
'`!models` — list available models',
|
||||
email: '*Email:*\n' + '`!email sync` — sync Gmail and show new emails',
|
||||
};
|
||||
|
||||
if (lower === '!help' || lower.startsWith('!help ')) {
|
||||
const topic = content.slice('!help'.length).trim().toLowerCase();
|
||||
if (topic && topic in helpSections) {
|
||||
await send(helpSections[topic]!);
|
||||
return true;
|
||||
}
|
||||
if (topic) {
|
||||
await send(
|
||||
`Unknown topic: \`${topic}\`\nAvailable: ${Object.keys(helpSections)
|
||||
.map((k) => `\`${k}\``)
|
||||
.join(', ')}`,
|
||||
);
|
||||
return true;
|
||||
}
|
||||
const full = Object.values(helpSections).join('\n\n');
|
||||
await send(full + '\n\n`!help <topic>` — show commands for a topic');
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!models') {
|
||||
const models = await getVisibleModels();
|
||||
if (models.length === 0) {
|
||||
await send('No models available.');
|
||||
return true;
|
||||
}
|
||||
const current = getSessionModel('telegram', ctx.userId, telegramId);
|
||||
const grouped = new Map<string, string[]>();
|
||||
for (const m of models) {
|
||||
const list = grouped.get(m.provider) ?? [];
|
||||
list.push(m.id === current ? `*${m.id}* (current)` : m.id);
|
||||
grouped.set(m.provider, list);
|
||||
}
|
||||
let text = '*Available models:*\n';
|
||||
for (const [provider, ids] of grouped) {
|
||||
text += `\n_${provider}_\n${ids.map((id) => ` ${id}`).join('\n')}\n`;
|
||||
}
|
||||
text += '\nUse `!model <id>` to switch.';
|
||||
await send(text);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!model') {
|
||||
const current = getSessionModel('telegram', ctx.userId, telegramId);
|
||||
await send(
|
||||
current
|
||||
? `Current model: *${current}*`
|
||||
: 'No active session yet — the default model will be used on your next message.',
|
||||
);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower.startsWith('!model ')) {
|
||||
const requested = content.slice('!model '.length).trim();
|
||||
if (!requested) {
|
||||
const current = getSessionModel('telegram', ctx.userId, telegramId);
|
||||
await send(current ? `Current model: *${current}*` : 'No active session yet.');
|
||||
return true;
|
||||
}
|
||||
const models = await getVisibleModels();
|
||||
const match = models.find((m) => m.id === requested || m.name === requested);
|
||||
if (!match) {
|
||||
await send(`Model not found: \`${requested}\`\nUse \`!models\` to see available models.`);
|
||||
return true;
|
||||
}
|
||||
setSessionModel('telegram', ctx.userId, telegramId, match.id);
|
||||
await send(`Model switched to *${match.id}*. The new model will be used on your next message.`);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!email sync') {
|
||||
await handleEmailSync(ctx);
|
||||
return true;
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
export async function handleTelegramMessage(msg: TelegramBot.Message): Promise<void> {
|
||||
const bot = getTelegramBot();
|
||||
if (!bot) return;
|
||||
|
||||
// Ignore non-private chats, bot messages, non-text
|
||||
if (msg.chat.type !== 'private') return;
|
||||
if (msg.from?.is_bot) return;
|
||||
if (!msg.text) return;
|
||||
|
||||
const chatId = msg.chat.id;
|
||||
const telegramId = String(msg.from!.id);
|
||||
const content = msg.text.trim();
|
||||
if (!content) return;
|
||||
|
||||
const send: SendFn = (text: string) => bot.sendMessage(chatId, text);
|
||||
|
||||
// Look up linked Officer user
|
||||
const linked = await findUserByIntegrationConfig('telegram', 'telegramId', telegramId);
|
||||
|
||||
if (!linked) {
|
||||
if (PAIRING_CODE_PATTERN.test(content.toUpperCase())) {
|
||||
const result = await consumePairingCode(content.toUpperCase(), telegramId);
|
||||
if (result) {
|
||||
await send('Account linked! You can now chat with me.');
|
||||
return;
|
||||
}
|
||||
await send('Invalid or expired pairing code. Please generate a new one from Officer Settings.');
|
||||
return;
|
||||
}
|
||||
|
||||
await send(
|
||||
"I don't recognize your Telegram account. To link it:\n" +
|
||||
'1. Go to Officer Settings → Integrations → Telegram\n' +
|
||||
'2. Click "Link Telegram" to get a pairing code\n' +
|
||||
'3. Send the 6-character code to me here',
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle commands
|
||||
if (content.startsWith('!')) {
|
||||
const handled = await handleCommand({
|
||||
content,
|
||||
send,
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
telegramId,
|
||||
});
|
||||
if (handled) return;
|
||||
}
|
||||
|
||||
// Typing indicator
|
||||
const sendTyping = () => {
|
||||
bot.sendChatAction(chatId, 'typing').catch(() => {});
|
||||
};
|
||||
const typingInterval = setInterval(sendTyping, TYPING_INTERVAL_MS);
|
||||
sendTyping();
|
||||
|
||||
try {
|
||||
const result = await sendAndAwait({
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
username: toShellUsername(linked.user.username ?? '', linked.user.email),
|
||||
prompt: content,
|
||||
context: 'telegram',
|
||||
contextId: telegramId,
|
||||
});
|
||||
|
||||
clearInterval(typingInterval);
|
||||
|
||||
const signature = `\`${result.model}\`\n`;
|
||||
const chunks = chunkMessage(result.text);
|
||||
for (let i = 0; i < chunks.length; i++) {
|
||||
await send(i === 0 ? signature + chunks[i]! : chunks[i]!);
|
||||
}
|
||||
} catch (err) {
|
||||
clearInterval(typingInterval);
|
||||
console.error('[telegram] Error handling message:', err);
|
||||
await send('Sorry, something went wrong processing your message.').catch(() => {});
|
||||
}
|
||||
}
|
||||
@@ -1,8 +0,0 @@
|
||||
export type ChannelProvider = 'discord' | 'telegram' | 'whatsapp';
|
||||
|
||||
export type ChannelBot = {
|
||||
provider: ChannelProvider;
|
||||
start: (token: string) => Promise<void>;
|
||||
stop: () => Promise<void>;
|
||||
isRunning: () => boolean;
|
||||
};
|
||||
@@ -1,135 +0,0 @@
|
||||
import { join } from 'path';
|
||||
import { rmSync } from 'node:fs';
|
||||
import { Client, LocalAuth } from 'whatsapp-web.js';
|
||||
import { getServerIntegration, upsertServerIntegration } from 'officerdb';
|
||||
import { DATA_PATH } from '@@/data-path';
|
||||
import { handleWhatsAppMessage } from './handler';
|
||||
|
||||
let client: Client | null = null;
|
||||
let currentQR: string | null = null;
|
||||
let clientReady = false;
|
||||
|
||||
type QRListener = (qr: string | null, event: 'qr' | 'authenticated' | 'disconnected') => void;
|
||||
const qrListeners = new Set<QRListener>();
|
||||
|
||||
export async function startWhatsAppBot(): Promise<void> {
|
||||
if (client) {
|
||||
await stopWhatsAppBot();
|
||||
}
|
||||
|
||||
clientReady = false;
|
||||
currentQR = null;
|
||||
|
||||
client = new Client({
|
||||
authStrategy: new LocalAuth({ dataPath: join(DATA_PATH, '.wwebjs_auth') }),
|
||||
puppeteer: {
|
||||
headless: true,
|
||||
args: ['--no-sandbox', '--disable-setuid-sandbox', '--disable-dev-shm-usage', '--disable-gpu'],
|
||||
},
|
||||
});
|
||||
|
||||
client.on('qr', (qr) => {
|
||||
currentQR = qr;
|
||||
for (const listener of qrListeners) {
|
||||
listener(qr, 'qr');
|
||||
}
|
||||
console.log('[whatsapp] QR code received — scan with your phone');
|
||||
});
|
||||
|
||||
client.on('ready', () => {
|
||||
clientReady = true;
|
||||
currentQR = null;
|
||||
for (const listener of qrListeners) {
|
||||
listener(null, 'authenticated');
|
||||
}
|
||||
const phone = client?.info?.wid?.user ?? 'unknown';
|
||||
console.log(`[whatsapp] Bot ready — phone: ${phone}`);
|
||||
});
|
||||
|
||||
client.on('authenticated', () => {
|
||||
console.log('[whatsapp] Authenticated');
|
||||
});
|
||||
|
||||
client.on('auth_failure', (msg) => {
|
||||
console.error('[whatsapp] Auth failure:', msg);
|
||||
});
|
||||
|
||||
client.on('disconnected', (reason) => {
|
||||
clientReady = false;
|
||||
currentQR = null;
|
||||
for (const listener of qrListeners) {
|
||||
listener(null, 'disconnected');
|
||||
}
|
||||
console.log('[whatsapp] Disconnected:', reason);
|
||||
});
|
||||
|
||||
client.on('message', (msg) => {
|
||||
handleWhatsAppMessage(msg).catch((err) => {
|
||||
console.error('[whatsapp] Unhandled error in message handler:', err);
|
||||
});
|
||||
});
|
||||
|
||||
await client.initialize();
|
||||
}
|
||||
|
||||
export async function stopWhatsAppBot(): Promise<void> {
|
||||
if (client) {
|
||||
try {
|
||||
await client.destroy();
|
||||
} catch {
|
||||
// May fail if not connected
|
||||
}
|
||||
client = null;
|
||||
clientReady = false;
|
||||
currentQR = null;
|
||||
console.log('[whatsapp] Bot stopped');
|
||||
}
|
||||
}
|
||||
|
||||
export function isWhatsAppBotRunning(): boolean {
|
||||
return client !== null && clientReady;
|
||||
}
|
||||
|
||||
export function getWhatsAppBotPhone(): string | null {
|
||||
if (!client || !clientReady) return null;
|
||||
return client.info?.wid?.user ?? null;
|
||||
}
|
||||
|
||||
export function getWhatsAppQR(): string | null {
|
||||
return currentQR;
|
||||
}
|
||||
|
||||
export function subscribeQR(listener: QRListener): () => void {
|
||||
qrListeners.add(listener);
|
||||
return () => {
|
||||
qrListeners.delete(listener);
|
||||
};
|
||||
}
|
||||
|
||||
export function getWhatsAppClient(): Client | null {
|
||||
return client;
|
||||
}
|
||||
|
||||
export async function startWhatsAppBotIfConfigured(): Promise<void> {
|
||||
const integration = await getServerIntegration('whatsapp');
|
||||
if (!integration?.enabled) return;
|
||||
|
||||
await startWhatsAppBot();
|
||||
}
|
||||
|
||||
export async function disconnectWhatsApp(): Promise<void> {
|
||||
if (client) {
|
||||
try {
|
||||
await client.logout();
|
||||
} catch {
|
||||
// May fail if not authenticated
|
||||
}
|
||||
}
|
||||
await stopWhatsAppBot();
|
||||
|
||||
// Remove cached session so a new QR is shown on next connect
|
||||
const authPath = join(DATA_PATH, '.wwebjs_auth');
|
||||
rmSync(authPath, { recursive: true, force: true });
|
||||
|
||||
await upsertServerIntegration('whatsapp', {}, false);
|
||||
}
|
||||
@@ -1,233 +0,0 @@
|
||||
import type { Message as WAMessage } from 'whatsapp-web.js';
|
||||
import { findUserByIntegrationConfig, readConfigValue } from 'officerdb';
|
||||
import { sendAndAwait, getSessionModel, setSessionModel } from '../send-and-await';
|
||||
import { consumePairingCode } from '../pairing';
|
||||
import { getWhatsAppClient } from './bot';
|
||||
import { listChatModels } from '@@/api/chat/list-models';
|
||||
import { enqueueJob } from '../../queue/init';
|
||||
import { readJob } from '@@/queue/storage';
|
||||
import type { ModelInfo } from '@@/api/chat/types';
|
||||
import { toShellUsername } from '@@/data-path';
|
||||
|
||||
const PAIRING_CODE_PATTERN = /^[A-Z0-9]{6}$/;
|
||||
const TYPING_INTERVAL_MS = 5_000;
|
||||
|
||||
type SendFn = (text: string) => Promise<unknown>;
|
||||
|
||||
type AccessPolicy = { allowedModels: string[] };
|
||||
const ACCESS_POLICY_KEY = 'chat-access-policy';
|
||||
|
||||
function extractPhone(waId: string): string {
|
||||
// WhatsApp ID format: 5511999999999@c.us → 5511999999999
|
||||
return waId.split('@')[0]!;
|
||||
}
|
||||
|
||||
async function getVisibleModels(): Promise<ModelInfo[]> {
|
||||
const allModels = await listChatModels();
|
||||
const policy = await readConfigValue<AccessPolicy>(ACCESS_POLICY_KEY, { allowedModels: [] });
|
||||
const allowed = policy.allowedModels;
|
||||
|
||||
if (allowed.length === 0) return allModels;
|
||||
|
||||
const allowedSet = new Set(allowed);
|
||||
const allowedProviderSet = new Set(allowed.map((key) => key.split(':')[0]));
|
||||
|
||||
return allModels.filter((m) => {
|
||||
const key = `${m.provider}:${m.id}`;
|
||||
const isExplicitlyAllowed = allowedSet.has(key);
|
||||
const isFromNewProvider = !allowedProviderSet.has(m.provider);
|
||||
return isExplicitlyAllowed || isFromNewProvider;
|
||||
});
|
||||
}
|
||||
|
||||
import { runEmailSyncCommand } from '../email-sync-command';
|
||||
|
||||
type CommandContext = {
|
||||
content: string;
|
||||
send: SendFn;
|
||||
userId: number;
|
||||
email: string;
|
||||
whatsappId: string;
|
||||
};
|
||||
|
||||
async function handleEmailSync(ctx: CommandContext): Promise<void> {
|
||||
const { send, userId } = ctx;
|
||||
await send('Syncing emails...');
|
||||
const text = await runEmailSyncCommand(userId);
|
||||
await send(text); // WhatsApp's 65k limit means no chunking is needed
|
||||
}
|
||||
|
||||
async function handleCommand(ctx: CommandContext): Promise<boolean> {
|
||||
const { content, send, whatsappId } = ctx;
|
||||
const lower = content.toLowerCase();
|
||||
|
||||
const helpSections: Record<string, string> = {
|
||||
models:
|
||||
'*Models:*\n' +
|
||||
'`!model` — show current model\n' +
|
||||
'`!model <id>` — switch model\n' +
|
||||
'`!models` — list available models',
|
||||
email: '*Email:*\n' + '`!email sync` — sync Gmail and show new emails',
|
||||
};
|
||||
|
||||
if (lower === '!help' || lower.startsWith('!help ')) {
|
||||
const topic = content.slice('!help'.length).trim().toLowerCase();
|
||||
if (topic && topic in helpSections) {
|
||||
await send(helpSections[topic]!);
|
||||
return true;
|
||||
}
|
||||
if (topic) {
|
||||
await send(
|
||||
`Unknown topic: \`${topic}\`\nAvailable: ${Object.keys(helpSections)
|
||||
.map((k) => `\`${k}\``)
|
||||
.join(', ')}`,
|
||||
);
|
||||
return true;
|
||||
}
|
||||
const full = Object.values(helpSections).join('\n\n');
|
||||
await send(full + '\n\n`!help <topic>` — show commands for a topic');
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!models') {
|
||||
const models = await getVisibleModels();
|
||||
if (models.length === 0) {
|
||||
await send('No models available.');
|
||||
return true;
|
||||
}
|
||||
const current = getSessionModel('whatsapp', ctx.userId, whatsappId);
|
||||
const grouped = new Map<string, string[]>();
|
||||
for (const m of models) {
|
||||
const list = grouped.get(m.provider) ?? [];
|
||||
list.push(m.id === current ? `*${m.id}* (current)` : m.id);
|
||||
grouped.set(m.provider, list);
|
||||
}
|
||||
let text = '*Available models:*\n';
|
||||
for (const [provider, ids] of grouped) {
|
||||
text += `\n_${provider}_\n${ids.map((id) => ` ${id}`).join('\n')}\n`;
|
||||
}
|
||||
text += '\nUse `!model <id>` to switch.';
|
||||
await send(text);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!model') {
|
||||
const current = getSessionModel('whatsapp', ctx.userId, whatsappId);
|
||||
await send(
|
||||
current
|
||||
? `Current model: *${current}*`
|
||||
: 'No active session yet — the default model will be used on your next message.',
|
||||
);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower.startsWith('!model ')) {
|
||||
const requested = content.slice('!model '.length).trim();
|
||||
if (!requested) {
|
||||
const current = getSessionModel('whatsapp', ctx.userId, whatsappId);
|
||||
await send(current ? `Current model: *${current}*` : 'No active session yet.');
|
||||
return true;
|
||||
}
|
||||
const models = await getVisibleModels();
|
||||
const match = models.find((m) => m.id === requested || m.name === requested);
|
||||
if (!match) {
|
||||
await send(`Model not found: \`${requested}\`\nUse \`!models\` to see available models.`);
|
||||
return true;
|
||||
}
|
||||
setSessionModel('whatsapp', ctx.userId, whatsappId, match.id);
|
||||
await send(`Model switched to *${match.id}*. The new model will be used on your next message.`);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (lower === '!email sync') {
|
||||
await handleEmailSync(ctx);
|
||||
return true;
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
export async function handleWhatsAppMessage(msg: WAMessage): Promise<void> {
|
||||
const waClient = getWhatsAppClient();
|
||||
if (!waClient) return;
|
||||
|
||||
// Ignore group chats, status broadcasts, own messages
|
||||
if (msg.from.endsWith('@g.us')) return;
|
||||
if (msg.from === 'status@broadcast') return;
|
||||
if (msg.fromMe) return;
|
||||
if (!msg.body) return;
|
||||
|
||||
const phone = extractPhone(msg.from);
|
||||
const content = msg.body.trim();
|
||||
if (!content) return;
|
||||
|
||||
const send: SendFn = (text: string) => waClient.sendMessage(msg.from, text);
|
||||
|
||||
// Look up linked Officer user
|
||||
const linked = await findUserByIntegrationConfig('whatsapp', 'whatsappId', phone);
|
||||
|
||||
if (!linked) {
|
||||
if (PAIRING_CODE_PATTERN.test(content.toUpperCase())) {
|
||||
const result = await consumePairingCode(content.toUpperCase(), phone);
|
||||
if (result) {
|
||||
await send('Account linked! You can now chat with me.');
|
||||
return;
|
||||
}
|
||||
await send('Invalid or expired pairing code. Please generate a new one from Officer Settings.');
|
||||
return;
|
||||
}
|
||||
|
||||
await send(
|
||||
"I don't recognize your WhatsApp number. To link it:\n" +
|
||||
'1. Go to Officer Settings → Integrations → WhatsApp\n' +
|
||||
'2. Click "Link WhatsApp" to get a pairing code\n' +
|
||||
'3. Send the 6-character code to me here',
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle commands
|
||||
if (content.startsWith('!')) {
|
||||
const handled = await handleCommand({
|
||||
content,
|
||||
send,
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
whatsappId: phone,
|
||||
});
|
||||
if (handled) return;
|
||||
}
|
||||
|
||||
// Typing indicator
|
||||
const sendTyping = async () => {
|
||||
try {
|
||||
const chat = await msg.getChat();
|
||||
await chat.sendStateTyping();
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
};
|
||||
const typingInterval = setInterval(sendTyping, TYPING_INTERVAL_MS);
|
||||
sendTyping();
|
||||
|
||||
try {
|
||||
const result = await sendAndAwait({
|
||||
userId: linked.user.id,
|
||||
email: linked.user.email,
|
||||
username: toShellUsername(linked.user.username ?? '', linked.user.email),
|
||||
prompt: content,
|
||||
context: 'whatsapp',
|
||||
contextId: phone,
|
||||
});
|
||||
|
||||
clearInterval(typingInterval);
|
||||
|
||||
// WhatsApp has 65k char limit — no chunking needed
|
||||
const signature = `\`${result.model}\`\n`;
|
||||
await send(signature + result.text);
|
||||
} catch (err) {
|
||||
clearInterval(typingInterval);
|
||||
console.error('[whatsapp] Error handling message:', err);
|
||||
await send('Sorry, something went wrong processing your message.').catch(() => {});
|
||||
}
|
||||
}
|
||||
@@ -37,7 +37,6 @@ import { dockRouter } from './api/dock/dock';
|
||||
import { integrationsRouter, googleCallbackHandler } from './api/integrations/integrations';
|
||||
import { queueRouter } from './api/queue/queue';
|
||||
import { emailRouter } from './api/email/router';
|
||||
import { channelsRouter } from './channels/routes';
|
||||
import { browserRouter } from './api/browser/router';
|
||||
import { desktopRouter } from './api/desktop/rest';
|
||||
import { bugReportRouter } from './api/bug-report/bug-report';
|
||||
@@ -123,7 +122,6 @@ protectedRouter.route('/dock', dockRouter);
|
||||
protectedRouter.route('/integrations', integrationsRouter);
|
||||
protectedRouter.route('/queue', queueRouter);
|
||||
protectedRouter.route('/email', emailRouter);
|
||||
protectedRouter.route('/channels', channelsRouter);
|
||||
protectedRouter.route('/browser', browserRouter);
|
||||
protectedRouter.route('/bug-report', bugReportRouter);
|
||||
protectedRouter.route('/chat', chatRouter);
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
// Discord notifications, and nothing else.
|
||||
//
|
||||
// Officer used to run a full Discord bot: a gateway connection, command parsing, account pairing, an admin
|
||||
// config screen and a token in the database — and the same again for Telegram and WhatsApp. All three
|
||||
// existed to drive the platform from a chat app, which the phone app does now. What is left is the one
|
||||
// piece worth keeping: the ability to push a message out.
|
||||
//
|
||||
// Configured by env, so there is no UI, no pairing and no stored credential:
|
||||
// DISCORD_WEBHOOK_URL a channel webhook. Unset = notifications are silently skipped.
|
||||
|
||||
const WEBHOOK_URL = process.env.DISCORD_WEBHOOK_URL;
|
||||
|
||||
/** True when a webhook is configured; callers can skip building a message otherwise. */
|
||||
export const isDiscordNotifyConfigured = (): boolean => Boolean(WEBHOOK_URL);
|
||||
|
||||
/**
|
||||
* Post a message to the configured Discord channel. Never throws and never blocks anything important — a
|
||||
* notification that fails to send is logged and dropped, not retried.
|
||||
*/
|
||||
export async function notifyDiscord(content: string): Promise<void> {
|
||||
if (!WEBHOOK_URL) return;
|
||||
try {
|
||||
const res = await fetch(WEBHOOK_URL, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
// Discord rejects anything over 2000 characters outright.
|
||||
body: JSON.stringify({ content: content.slice(0, 2000) }),
|
||||
});
|
||||
if (!res.ok) console.error(`[notify:discord] webhook returned ${res.status}`);
|
||||
} catch (err) {
|
||||
console.error('[notify:discord] failed to send:', err instanceof Error ? err.message : err);
|
||||
}
|
||||
}
|
||||
@@ -8,7 +8,6 @@ import { getEmailAttachmentCacheDir } from '@@/data-path';
|
||||
import { getEmailAccounts } from 'officerdb';
|
||||
import { openEmailDb, openUserEmailDb, rowToSummary, getSyncMeta, searchEmails } from './store';
|
||||
import { accountsRouter } from './accounts';
|
||||
import { performResync } from './resync';
|
||||
|
||||
export const emailRouter = createRouter();
|
||||
|
||||
@@ -365,33 +364,6 @@ emailRouter.delete('/messages/:id', async (ctx) => {
|
||||
}
|
||||
});
|
||||
|
||||
// POST /sync-now — run a resync and report what arrived.
|
||||
//
|
||||
// For the chat channels ("sync my email" from Telegram/Discord/WhatsApp). They used to open the mail
|
||||
// store directly and enqueue a `gmail-sync` job, which stopped existing when sync moved in here; and the
|
||||
// job type was hardcoded to the OAuth path even though an app-password account syncs over IMAP, so that
|
||||
// command had been failing regardless. One call now: sync, then say what is new.
|
||||
emailRouter.post('/sync-now', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
const accounts = await getEmailAccounts(user.id);
|
||||
const account = accounts.find((a) => a.enabled) ?? accounts[0];
|
||||
if (!account) return ctx.json({ saved: 0, newest: [], error: 'No email account configured' });
|
||||
|
||||
const result = await performResync({ accountId: account.id, userEmail: user.email, userId: user.id });
|
||||
|
||||
const db = openEmailDb(user.email, account.email);
|
||||
try {
|
||||
const newest =
|
||||
result.saved > 0
|
||||
? (db
|
||||
.query('SELECT from_name, from_address, subject FROM emails WHERE deleted = 0 ORDER BY date DESC LIMIT ?')
|
||||
.all(Math.min(result.saved, 20)) as Array<{ from_name: string | null; from_address: string; subject: string }>)
|
||||
: [];
|
||||
return ctx.json({ saved: result.saved, skipped: result.skipped, errors: result.errors, newest });
|
||||
} finally {
|
||||
db.close();
|
||||
}
|
||||
});
|
||||
|
||||
emailRouter.get('/sync-status', async (ctx) => {
|
||||
const user = ctx.get('user');
|
||||
|
||||
Reference in New Issue
Block a user