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; } // Cache the secret across warm invocations let cachedCredentials: SlackCredentials | undefined; async function getSlackCredentials(): Promise { 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 { const body: Record = { 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; error?: string }; if (!data.ok) throw new Error(`Slack API 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 { 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 { 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.'; } export const handler = async (event: ProcessorEvent): Promise => { const { userId, channelId, text, ts, threadTs } = event; const [credentials, sessionId] = await Promise.all([ getSlackCredentials(), resolveSessionId(userId), ]); 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 postToSlack(credentials.botToken, channelId, answer, threadTs ?? ts); 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, ttl, }, }), ); };