This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
seahaven-slack-bot/lambda/slack-processor/index.ts

236 lines
6.8 KiB
TypeScript
Raw Permalink Normal View History

import { SecretsManagerClient, GetSecretValueCommand } from '@aws-sdk/client-secrets-manager';
import {
BedrockAgentRuntimeClient,
InvokeAgentCommand,
} from '@aws-sdk/client-bedrock-agent-runtime';
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import { DynamoDBDocumentClient, PutCommand, QueryCommand } from '@aws-sdk/lib-dynamodb';
const secretsClient = new SecretsManagerClient({});
const bedrockClient = new BedrockAgentRuntimeClient({ region: process.env.REGION! });
const ddb = DynamoDBDocumentClient.from(new DynamoDBClient({}));
interface SlackCredentials {
botToken: string;
signingSecret: string;
}
interface ProcessorEvent {
userId: string;
channelId: string;
text: string;
ts: string;
threadTs?: string;
eventType?: string;
}
// Cache the secret across warm invocations
let cachedCredentials: SlackCredentials | undefined;
async function getSlackCredentials(): Promise<SlackCredentials> {
if (!cachedCredentials) {
const res = await secretsClient.send(
new GetSecretValueCommand({ SecretId: process.env.SLACK_SECRET_ARN! }),
);
cachedCredentials = JSON.parse(res.SecretString!) as SlackCredentials;
}
return cachedCredentials;
}
async function postToSlack(
token: string,
channel: string,
text: string,
threadTs?: string,
): Promise<string> {
const body: Record<string, string> = { channel, text };
if (threadTs) body.thread_ts = threadTs;
const res = await fetch('https://slack.com/api/chat.postMessage', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${token}`,
},
body: JSON.stringify(body),
});
if (!res.ok) throw new Error(`Slack HTTP error ${res.status}`);
const data = (await res.json()) as { ok: boolean; ts?: string; error?: string };
if (!data.ok) throw new Error(`Slack API error: ${data.error}`);
return data.ts!;
}
async function updateSlack(
token: string,
channel: string,
ts: string,
text: string,
): Promise<void> {
const res = await fetch('https://slack.com/api/chat.update', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${token}`,
},
body: JSON.stringify({ channel, ts, text }),
});
if (!res.ok) throw new Error(`Slack HTTP error ${res.status}`);
const data = (await res.json()) as { ok: boolean; error?: string };
if (!data.ok) throw new Error(`Slack update error: ${data.error}`);
}
// Reuse the most recent Bedrock session if there's been activity in the last 30 minutes.
// This enables multi-turn conversation without re-introducing context each message.
async function resolveSessionId(userId: string): Promise<string> {
const cutoff = new Date(Date.now() - 30 * 60 * 1000).toISOString();
const result = await ddb.send(
new QueryCommand({
TableName: process.env.CONVERSATION_TABLE!,
KeyConditionExpression: 'userId = :uid AND #ts > :cutoff',
ExpressionAttributeNames: { '#ts': 'timestamp' },
ExpressionAttributeValues: { ':uid': userId, ':cutoff': cutoff },
Limit: 1,
ScanIndexForward: false,
ProjectionExpression: 'sessionId',
}),
);
if (result.Items?.length && result.Items[0].sessionId) {
return result.Items[0].sessionId as string;
}
return `${userId}-${Date.now()}`;
}
async function invokeAgent(sessionId: string, inputText: string): Promise<string> {
const res = await bedrockClient.send(
new InvokeAgentCommand({
agentId: process.env.AGENT_ID!,
agentAliasId: process.env.AGENT_ALIAS_ID!,
sessionId,
inputText,
}),
);
let answer = '';
if (res.completion) {
for await (const event of res.completion) {
if (event.chunk?.bytes) {
answer += Buffer.from(event.chunk.bytes).toString('utf-8');
}
}
}
return answer.trim() || 'I was unable to generate a response. Please try again.';
}
/** Convert Markdown from Bedrock Agent output to Slack mrkdwn. */
function toSlackMrkdwn(md: string): string {
return md
// Headers → bold text
.replace(/^#{1,3}\s+(.+)$/gm, '*$1*')
// Bold: **text** or __text__ → *text*
.replace(/\*\*(.+?)\*\*/g, '*$1*')
.replace(/__(.+?)__/g, '*$1*')
// Italic: _text_ is the same in Slack, leave as-is
// Inline code and code blocks are the same in Slack, leave as-is
// Links: [text](url) → <url|text>
.replace(/\[([^\]]+)\]\(([^)]+)\)/g, '<$2|$1>')
// Collapse 3+ consecutive newlines to 2
.replace(/\n{3,}/g, '\n\n');
}
const UNANSWERED_SIGNALS = [
"i don't have",
"not in the knowledge base",
"i couldn't find",
"i could not find",
"no information available",
"outside my capabilities",
"no results found",
"i'm unable to",
"i am unable to",
"i don't have access",
"not available in",
"not something i can",
"beyond what i can",
"i don't currently have",
];
function isUnanswered(response: string): boolean {
const lower = response.toLowerCase();
return UNANSWERED_SIGNALS.some((signal) => lower.includes(signal));
}
export const handler = async (event: ProcessorEvent): Promise<void> => {
const { userId, channelId, text, ts, threadTs, eventType } = event;
const [credentials, sessionId] = await Promise.all([
getSlackCredentials(),
resolveSessionId(userId),
]);
// In channels, reply in a thread; in DMs, reply flat (no thread_ts)
const replyThreadTs = eventType === 'app_mention' ? (threadTs ?? ts) : threadTs;
const placeholderTs = await postToSlack(
credentials.botToken,
channelId,
'_Alex is thinking..._',
replyThreadTs,
);
let answer: string;
try {
answer = await invokeAgent(sessionId, text);
} catch (err) {
console.error('Bedrock agent invocation failed:', err);
answer = 'Sorry, I encountered an error processing your request. Please try again in a moment.';
}
await updateSlack(credentials.botToken, channelId, placeholderTs, toSlackMrkdwn(answer));
const now = new Date().toISOString();
const ttl = Math.floor(Date.now() / 1000) + 90 * 24 * 60 * 60; // 90-day TTL
await ddb.send(
new PutCommand({
TableName: process.env.CONVERSATION_TABLE!,
Item: {
userId,
timestamp: now,
sessionId,
question: text,
answer,
channelId,
messageTs: ts,
eventType: eventType ?? 'dm',
ttl,
},
}),
);
if (process.env.UNANSWERED_TABLE && isUnanswered(answer)) {
const qTtl = Math.floor(Date.now() / 1000) + 180 * 24 * 60 * 60;
await ddb.send(
new PutCommand({
TableName: process.env.UNANSWERED_TABLE,
Item: {
questionId: crypto.randomUUID(),
userId,
question: text,
agentResponse: answer,
timestamp: now,
status: 'pending',
channelId,
eventType: eventType ?? 'dm',
ttl: qTtl,
},
}),
);
}
};