Skip to content
Draft
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -31,3 +31,6 @@ data/
plans/
*.env
.worktrees/

# WhatsApp auth state (contains session credentials)
whatsapp-auth/
3 changes: 3 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -18,18 +18,21 @@
"@earendil-works/pi-agent-core": "^0.82.1",
"@earendil-works/pi-ai": "^0.82.1",
"@earendil-works/pi-coding-agent": "^0.82.1",
"@hapi/boom": "^10.0.1",
"@linear/sdk": "^88.3.0",
"@shikijs/markdown-it": "^4.3.1",
"@silvia-odwyer/photon-node": "^0.3.4",
"@slack/socket-mode": "^3.0.0",
"@slack/web-api": "^8.0.0",
"@tobilu/qmd": "^2.5.3",
"baileys": "^7.0.0-rc13",
"croner": "^10.0.1",
"diff": "^9.0.0",
"file-type": "^22.0.1",
"jiti": "^2.7.0",
"markdown-it": "^14.3.0",
"pi-mcp-adapter": "^2.15.0",
"pino": "^9.7.0",
"quick-lru": "^7.3.0",
"typebox": "^1.3.8"
},
Expand Down
3 changes: 3 additions & 0 deletions src/agent/setup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -661,6 +661,9 @@ function defaultLogContextScope(event: BotEvent): LogContextScope {
if (event.source === "linear") {
return { source: "linear", kind: "chronological" };
}
if (event.source === "whatsapp") {
return { source: "whatsapp", kind: "chronological" };
}
if (event.threadTs) {
return { source: "slack", kind: "thread", rootTs: event.threadTs };
}
Expand Down
12 changes: 12 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,14 @@ export interface DigbyConfig {
*/
replyBehaviour?: Record<string, "mention" | "channel" | "thread">;
};
whatsapp?: {
/**
* JIDs or phone numbers allowed to trigger the bot.
* Phone numbers in E.164 format (with or without +) are normalised to JID format.
* Examples: "+447700900000", "447700900000@s.whatsapp.net", "123456789@g.us"
*/
allowFrom?: string[];
};
/** Post tool calls/thinking to thread under bot's message (default: false) */
debugThreading?: boolean;
/** Maximum time (seconds) a single run can take before being aborted (default: 600) */
Expand Down Expand Up @@ -83,3 +91,7 @@ export function getRunTimeout(): number {
export function getRunTimeoutWarnBeforeS(): number {
return loadConfig().runTimeoutWarnBeforeS ?? 60;
}

export function getWhatsAppAllowFrom(): string[] {
return loadConfig().whatsapp?.allowFrom ?? [];
}
144 changes: 144 additions & 0 deletions src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -601,6 +601,150 @@ if (LINEAR_API_KEY && LINEAR_WEBHOOK_SECRET) {
log.info("Linear agent disabled (missing LINEAR_API_KEY/LINEAR_WEBHOOK_SECRET)");
}

// ============================================================================
// WhatsApp agent (optional — only if enabled and auth dir configured)
// ============================================================================

const DIGBY_WHATSAPP_ENABLED = process.env.DIGBY_WHATSAPP_ENABLED === "true";
const DIGBY_WHATSAPP_AUTH_DIR = process.env.DIGBY_WHATSAPP_AUTH_DIR || join(workingDir, "whatsapp-auth");

if (DIGBY_WHATSAPP_ENABLED) {
const { getWhatsAppAllowFrom } = await import("./config.js");
const { WhatsAppClient } = await import("./whatsapp/client.js");
const { getWhatsAppConversationTarget } = await import("./whatsapp/conversation.js");
const { setupWhatsAppRouter, type WhatsAppRouterHandler } = await import("./whatsapp/router.js");
const { WhatsAppSurface } = await import("./surface/whatsapp.js");
type WhatsAppEvent = BotEvent;

const whatsappClient = new WhatsAppClient({
authDir: DIGBY_WHATSAPP_AUTH_DIR,
allowFrom: getWhatsAppAllowFrom(),
});

async function runWhatsAppEvent(event: WhatsAppEvent, options: { logTrigger?: boolean } = {}): Promise<void> {
const state = getChannelRunState(event.channel);
const conversation = getWhatsAppConversationTarget(event, state.channelState.channelDir);
const runnerId = conversation.runnerId;
const lane = getLaneRunState(state, runnerId);

lane.running = true;
lane.acceptingFollowUps = false;
lane.activeRunner = undefined;

const stats = createRunStats();
const ctx = new WhatsAppSurface(whatsappClient, event.channel, stats);

try {
if (options.logTrigger !== false) {
state.channelState.logUserMessage(event);
}

const runner = await getOrCreateRunner({
runnerId,
channelId: event.channel,
channelDir: state.channelState.channelDir,
sessionDir: conversation.sessionDir,
workingDir,
});

lane.activeRunner = runner;
if (lane.stopRequested) {
runner.abort();
}
lane.acceptingFollowUps = true;

log.info(`[${event.channel}] Starting WhatsApp run: ${event.text.substring(0, 50)}`);

ctx.emitThinking();

const result = await runner.run(ctx, event, state.channelState, [], [], undefined, stats, conversation.logContextScope);

if (result.stopReason === "aborted" && lane.stopRequested) {
log.info(`[${event.channel}] WhatsApp run stopped`);
}
} catch (err) {
const errMsg = err instanceof Error ? err.message : String(err);
log.warn(`[${event.channel}] WhatsApp run error`, errMsg);
ctx.reject(`Something went wrong: ${errMsg.substring(0, 500)}`);
} finally {
ctx.resolve();
await ctx.flush();

const finalText = ctx.finalText;
if (finalText && finalText !== THINKING_PLACEHOLDER && !ctx.wasDeleted) {
state.channelState.logBotResponse(finalText, String(Date.now() / 1000));
}

lane.acceptingFollowUps = false;
const queuedFollowUps = lane.followUps.drain(runnerId);
if (queuedFollowUps.length > 0) {
const trigger = createQueuedFollowUpTrigger(queuedFollowUps);
log.info(`[${event.channel}] Scheduling ${queuedFollowUps.length} queued WhatsApp follow-up message(s)`);
lane.queue.enqueue(() => runWhatsAppEvent(trigger, { logTrigger: false }));
}

lane.running = false;
lane.activeRunner = undefined;
lane.stopRequested = false;
lane.stopMessageTs = undefined;

await evictRunner(runnerId);
}
}

async function enqueueWhatsAppEvent(event: WhatsAppEvent): Promise<void> {
const state = getChannelRunState(event.channel);
const conversation = getWhatsAppConversationTarget(event, state.channelState.channelDir);
const runnerId = conversation.runnerId;
const lane = getLaneRunState(state, runnerId);

if (lane.running && lane.acceptingFollowUps) {
state.channelState.logUserMessage(event);
const count = lane.followUps.enqueue(runnerId, event);
log.info(`[${event.channel}] Queued WhatsApp follow-up ${count}`);
return;
}

lane.queue.enqueue(() => runWhatsAppEvent(event));
}

const whatsappHandler: WhatsAppRouterHandler = {
isBusy(event: WhatsAppEvent): boolean {
const state = getChannelRunState(event.channel);
const conversation = getWhatsAppConversationTarget(event, state.channelState.channelDir);
return isLaneBusy(state.lanes.get(conversation.runnerId));
},

async handleEvent(event: WhatsAppEvent): Promise<void> {
await enqueueWhatsAppEvent(event);
},

async handleStop(event: WhatsAppEvent): Promise<void> {
const state = getChannelRunState(event.channel);
const conversation = getWhatsAppConversationTarget(event, state.channelState.channelDir);
const lane = getLaneRunState(state, conversation.runnerId);
if (isLaneBusy(lane)) {
lane.stopRequested = true;
lane.activeRunner?.abort();
await whatsappClient.sendMessage(conversation.jid, "_Stopping..._");
} else {
await whatsappClient.sendMessage(conversation.jid, "_Nothing running_");
}
},

logMessage(event: WhatsAppEvent): void {
const state = getChannelRunState(event.channel);
state.channelState.logUserMessage(event);
},
};

await whatsappClient.start();

@cubic-dev-ai cubic-dev-ai Bot Jul 28, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A transient Baileys version or auth initialization failure prevents the whole process from reaching Ready, taking Slack down even though WhatsApp is optional. Isolating WhatsApp startup from the main boot path or adding an initial retry/backoff around initialization would keep the existing channel available.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At src/main.ts, line 741:

<comment>A transient Baileys version or auth initialization failure prevents the whole process from reaching `Ready`, taking Slack down even though WhatsApp is optional. Isolating WhatsApp startup from the main boot path or adding an initial retry/backoff around initialization would keep the existing channel available.</comment>

<file context>
@@ -601,6 +601,150 @@ if (LINEAR_API_KEY && LINEAR_WEBHOOK_SECRET) {
+		},
+	};
+
+	await whatsappClient.start();
+	setupWhatsAppRouter(whatsappClient, whatsappHandler, startupTs);
+	log.info("WhatsApp agent enabled");
</file context>
Fix with cubic

setupWhatsAppRouter(whatsappClient, whatsappHandler, startupTs);
log.info("WhatsApp agent enabled");
} else {
log.info("WhatsApp agent disabled (DIGBY_WHATSAPP_ENABLED not set)");
}

log.info("Ready");

// ============================================================================
Expand Down
21 changes: 17 additions & 4 deletions src/persistence/log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,8 @@ export interface LogMessage {
export type LogContextScope =
| { source: "linear"; kind: "chronological" }
| { source: "slack"; kind: "channel" }
| { source: "slack"; kind: "thread"; rootTs: string };
| { source: "slack"; kind: "thread"; rootTs: string }
| { source: "whatsapp"; kind: "chronological" };

export interface SelectedContextMessage {
id: string;
Expand Down Expand Up @@ -71,13 +72,13 @@ function displayName(logMsg: LogMessage): string {
return logMsg.userName || logMsg.user || "unknown";
}

export function formatLogMessageForContext(source: "linear" | "slack", logMsg: LogMessage): string {
export function formatLogMessageForContext(source: "linear" | "slack" | "whatsapp", logMsg: LogMessage): string {
const ts = logMsg.ts || "unknown";
const attrs = [`ts="${escapeAttribute(ts)}"`];
if (source === "slack" && logMsg.threadTs) attrs.push(`thread_ts="${escapeAttribute(logMsg.threadTs)}"`);
const user = displayName(logMsg);
attrs.push(`user="${escapeAttribute(user)}"`);
const tag = source === "slack" ? "slack_message" : "linear_message";
const tag = source === "slack" ? "slack_message" : source === "whatsapp" ? "whatsapp_message" : "linear_message";
return `<${tag} ${attrs.join(" ")}>\n[${user}]: ${logMsg.text || ""}\n</${tag}>`;
}

Expand Down Expand Up @@ -249,13 +250,22 @@ function selectLinearMessages(messages: LogMessage[], currentTs: string): Select
.map((logMsg) => contextMessageFromLog("linear", logMsg));
}

function selectWhatsAppMessages(messages: LogMessage[], currentTs: string): SelectedContextMessage[] {
return sortLogMessages(messages)
.filter((logMsg) => isBeforeCurrent(logMsg, currentTs))
.map((logMsg) => contextMessageFromLog("whatsapp", logMsg));
}

export function selectLogMessagesForContext(
messages: LogMessage[],
options: SelectLogMessagesOptions,
): SelectedContextMessage[] {
if (options.scope.source === "linear") {
return selectLinearMessages(messages, options.currentTs);
}
if (options.scope.source === "whatsapp") {
return selectWhatsAppMessages(messages, options.currentTs);
}
const dedupedMessages = dedupeLogMessagesByTs(messages);
if (options.scope.kind === "channel") {
return selectSlackChannelMessages(dedupedMessages, options.currentTs);
Expand Down Expand Up @@ -378,7 +388,7 @@ export function syncLogToContext(
*/
function normalizeMessageText(text: string): string {
let normalized = text.replace(/^\[\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}[+-]\d{2}:\d{2}\] /, "");
const wrappedMessage = normalized.match(/^<((?:slack|linear)_message)\b[^>]*>\n([\s\S]*?)\n<\/\1>/);
const wrappedMessage = normalized.match(/^<((?:slack|linear|whatsapp)_message)\b[^>]*>\n([\s\S]*?)\n<\/\1>/);
if (wrappedMessage) {
normalized = wrappedMessage[2];
}
Expand All @@ -402,6 +412,9 @@ function extractContextIds(text: string): string[] {
for (const match of text.matchAll(/<linear_message\s+[^>]*ts="([^"]+)"/g)) {
ids.push(`linear:${match[1]}`);
}
for (const match of text.matchAll(/<whatsapp_message\s+[^>]*ts="([^"]+)"/g)) {
ids.push(`whatsapp:${match[1]}`);
}
for (const match of text.matchAll(/<slack_thread_boundary\s+root_ts="([^"]+)"/g)) {
ids.push(`slack-thread-boundary:${match[1]}`);
}
Expand Down
Loading
Loading