mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 03:50:45 +08:00
feat: add Slack bot integration
This commit is contained in:
parent
bee4c2927f
commit
7d44861a07
28 changed files with 3021 additions and 661 deletions
173
src/channels/slack/config-store.mjs
Normal file
173
src/channels/slack/config-store.mjs
Normal file
|
|
@ -0,0 +1,173 @@
|
|||
import { createHash } from 'node:crypto';
|
||||
import { mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname } from 'node:path';
|
||||
|
||||
const EMPTY_DOCUMENT = Object.freeze({ version: 1, bots: Object.freeze([]) });
|
||||
const BOT_ID_PATTERN = /^slack_[a-f0-9]{24}$/;
|
||||
const BOT_TOKEN_REF_PATTERN = /^DSH_SLACK_BOT_TOKEN_[A-F0-9]{24}$/;
|
||||
const APP_TOKEN_REF_PATTERN = /^DSH_SLACK_APP_TOKEN_[A-F0-9]{24}$/;
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
export function deriveSlackBotIdentity(platformId) {
|
||||
const raw = cleanString(platformId);
|
||||
if (!raw) throw new TypeError('platformId is required');
|
||||
const digest = createHash('sha256').update(raw).digest('hex').slice(0, 24);
|
||||
const suffix = digest.toUpperCase();
|
||||
return {
|
||||
botId: `slack_${digest}`,
|
||||
botTokenRef: `DSH_SLACK_BOT_TOKEN_${suffix}`,
|
||||
appTokenRef: `DSH_SLACK_APP_TOKEN_${suffix}`,
|
||||
};
|
||||
}
|
||||
|
||||
export function maskSlackBotId(platformId) {
|
||||
const value = cleanString(platformId) ?? '';
|
||||
const [teamId, userId] = value.split(':');
|
||||
if (teamId && userId) return `${teamId.slice(0, 5)}••• · ${userId.slice(0, 5)}•••`;
|
||||
return value ? `${value.slice(0, 6)}••••` : 'Slack机器人';
|
||||
}
|
||||
|
||||
export class SlackConfigStore {
|
||||
#path;
|
||||
#value = EMPTY_DOCUMENT;
|
||||
#writeQueue = Promise.resolve();
|
||||
|
||||
constructor(path) {
|
||||
this.#path = path;
|
||||
}
|
||||
|
||||
async load() {
|
||||
try {
|
||||
const normalized = this.#normalizeDocument(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
if (!normalized) throw new Error('dsh-im Slack config contains invalid bot data');
|
||||
this.#value = normalized;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
list() {
|
||||
return structuredClone(this.#value.bots);
|
||||
}
|
||||
|
||||
get(botId) {
|
||||
const bot = this.#value.bots.find((candidate) => candidate.botId === botId);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
getByPlatformId(platformId) {
|
||||
const bot = this.#value.bots.find((candidate) => candidate.platformId === platformId);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
async save(value) {
|
||||
const normalized = this.#normalizeBot(value);
|
||||
if (!normalized) throw new Error('Refusing to persist incomplete Slack bot data');
|
||||
return this.#mutate((bots) => {
|
||||
const collision = bots.find((bot) => (
|
||||
bot.botId !== normalized.botId
|
||||
&& (bot.platformId === normalized.platformId
|
||||
|| bot.botTokenRef === normalized.botTokenRef
|
||||
|| bot.appTokenRef === normalized.appTokenRef)
|
||||
));
|
||||
if (collision) throw new Error('Duplicate Slack bot identity');
|
||||
const index = bots.findIndex((bot) => bot.botId === normalized.botId);
|
||||
if (index === -1) bots.push(normalized);
|
||||
else bots[index] = normalized;
|
||||
return structuredClone(normalized);
|
||||
});
|
||||
}
|
||||
|
||||
async remove(botId) {
|
||||
if (!BOT_ID_PATTERN.test(botId)) throw new TypeError('Invalid Slack bot id');
|
||||
return this.#mutate((bots) => {
|
||||
const index = bots.findIndex((bot) => bot.botId === botId);
|
||||
if (index === -1) return null;
|
||||
return structuredClone(bots.splice(index, 1)[0]);
|
||||
});
|
||||
}
|
||||
|
||||
async clear() {
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
try {
|
||||
await unlink(this.#path);
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
}
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
}
|
||||
|
||||
#normalizeBot(value) {
|
||||
if (!value || typeof value !== 'object') return null;
|
||||
const botId = cleanString(value.botId);
|
||||
const platformId = cleanString(value.platformId);
|
||||
const botTokenRef = cleanString(value.botTokenRef);
|
||||
const appTokenRef = cleanString(value.appTokenRef);
|
||||
const name = cleanString(value.name);
|
||||
if (!botId || !platformId || !botTokenRef || !appTokenRef || !name
|
||||
|| !BOT_ID_PATTERN.test(botId)
|
||||
|| !BOT_TOKEN_REF_PATTERN.test(botTokenRef)
|
||||
|| !APP_TOKEN_REF_PATTERN.test(appTokenRef)) return null;
|
||||
const derived = deriveSlackBotIdentity(platformId);
|
||||
if (derived.botId !== botId
|
||||
|| derived.botTokenRef !== botTokenRef
|
||||
|| derived.appTokenRef !== appTokenRef) return null;
|
||||
return Object.freeze({
|
||||
botId,
|
||||
platformId,
|
||||
botTokenRef,
|
||||
appTokenRef,
|
||||
name,
|
||||
username: cleanString(value.username),
|
||||
teamId: cleanString(value.teamId),
|
||||
teamName: cleanString(value.teamName),
|
||||
createdAt: cleanString(value.createdAt) ?? new Date().toISOString(),
|
||||
connectedAt: cleanString(value.connectedAt),
|
||||
});
|
||||
}
|
||||
|
||||
#normalizeDocument(value) {
|
||||
if (!value || value.version !== 1 || !Array.isArray(value.bots)) return null;
|
||||
const bots = value.bots.map((bot) => this.#normalizeBot(bot));
|
||||
if (bots.some((bot) => bot === null)) return null;
|
||||
const ids = new Set();
|
||||
const platformIds = new Set();
|
||||
const refs = new Set();
|
||||
for (const bot of bots) {
|
||||
if (ids.has(bot.botId) || platformIds.has(bot.platformId)
|
||||
|| refs.has(bot.botTokenRef) || refs.has(bot.appTokenRef)) return null;
|
||||
ids.add(bot.botId);
|
||||
platformIds.add(bot.platformId);
|
||||
refs.add(bot.botTokenRef);
|
||||
refs.add(bot.appTokenRef);
|
||||
}
|
||||
return Object.freeze({ version: 1, bots: Object.freeze(bots) });
|
||||
}
|
||||
|
||||
async #mutate(mutator) {
|
||||
let result;
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
const bots = [...this.#value.bots];
|
||||
result = mutator(bots);
|
||||
const document = Object.freeze({ version: 1, bots: Object.freeze(bots) });
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(document, null, 2)}\n`, {
|
||||
encoding: 'utf8', mode: 0o600,
|
||||
});
|
||||
await rename(temporary, this.#path);
|
||||
this.#value = document;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
return result;
|
||||
}
|
||||
}
|
||||
3
src/channels/slack/harness-client.mjs
Normal file
3
src/channels/slack/harness-client.mjs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
import { HarnessClient } from '../weixin/harness-client.mjs';
|
||||
|
||||
export class SlackHarnessClient extends HarnessClient {}
|
||||
31
src/channels/slack/manifest.mjs
Normal file
31
src/channels/slack/manifest.mjs
Normal file
|
|
@ -0,0 +1,31 @@
|
|||
export const SLACK_APP_MANIFEST_YAML = `_metadata:
|
||||
major_version: 1
|
||||
display_information:
|
||||
name: DeepSeek Harness
|
||||
description: Connect Slack conversations to a local DeepSeek Harness agent.
|
||||
background_color: "#4A154B"
|
||||
features:
|
||||
app_home:
|
||||
home_tab_enabled: false
|
||||
messages_tab_enabled: true
|
||||
messages_tab_read_only_enabled: false
|
||||
bot_user:
|
||||
display_name: DeepSeek Harness
|
||||
always_online: false
|
||||
oauth_config:
|
||||
scopes:
|
||||
bot:
|
||||
- app_mentions:read
|
||||
- chat:write
|
||||
- im:history
|
||||
settings:
|
||||
event_subscriptions:
|
||||
bot_events:
|
||||
- app_mention
|
||||
- message.im
|
||||
org_deploy_enabled: false
|
||||
socket_mode_enabled: true
|
||||
token_rotation_enabled: false
|
||||
`;
|
||||
|
||||
export const SLACK_CREATE_APP_URL = 'https://api.slack.com/apps?new_app=1';
|
||||
259
src/channels/slack/slack-api.mjs
Normal file
259
src/channels/slack/slack-api.mjs
Normal file
|
|
@ -0,0 +1,259 @@
|
|||
const DEFAULT_BASE_URL = 'https://slack.com/api/';
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function requestSignal(signal, timeoutMs) {
|
||||
const timeout = AbortSignal.timeout(timeoutMs);
|
||||
return signal ? AbortSignal.any([signal, timeout]) : timeout;
|
||||
}
|
||||
|
||||
function delay(ms, signal) {
|
||||
return new Promise((resolve, reject) => {
|
||||
if (signal?.aborted) {
|
||||
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
|
||||
return;
|
||||
}
|
||||
const timer = setTimeout(resolve, ms);
|
||||
timer?.unref?.();
|
||||
signal?.addEventListener('abort', () => {
|
||||
clearTimeout(timer);
|
||||
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
|
||||
}, { once: true });
|
||||
});
|
||||
}
|
||||
|
||||
function slackId(value, name) {
|
||||
const result = cleanString(value);
|
||||
if (!result || !/^[A-Z][A-Z0-9]{4,30}$/i.test(result)) {
|
||||
throw new TypeError(`Invalid Slack ${name}`);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
function requiredString(value, name) {
|
||||
const result = cleanString(value);
|
||||
if (!result) throw new TypeError(`Slack ${name} is required`);
|
||||
return result;
|
||||
}
|
||||
|
||||
function safeOutgoingText(value, { trim = true } = {}) {
|
||||
const raw = typeof value === 'string' ? value : '';
|
||||
const text = trim ? raw.trim() : raw;
|
||||
if (!text) throw new TypeError('Slack message text is required');
|
||||
return text
|
||||
.replace(/<@([A-Z0-9]+)>/gi, '@$1')
|
||||
.replace(/<!(channel|here|everyone)(?:\^[^>]*)?>/gi, '@$1');
|
||||
}
|
||||
|
||||
function apiFailure(method, payload, tokenKind) {
|
||||
const reason = cleanString(payload?.error) ?? 'unknown_error';
|
||||
const error = new Error(`Slack ${method} failed: ${reason.replaceAll('_', ' ')}`);
|
||||
if (['invalid_auth', 'not_authed', 'token_revoked', 'account_inactive'].includes(reason)) {
|
||||
error.code = tokenKind === 'app' ? 'slack-invalid-app-token' : 'slack-invalid-bot-token';
|
||||
} else if (reason === 'missing_scope') {
|
||||
error.code = 'slack-missing-scope';
|
||||
} else if (reason === 'method_not_supported_for_channel_type'
|
||||
|| reason === 'channel_type_not_supported'
|
||||
|| reason === 'deprecated_endpoint') {
|
||||
error.code = 'slack-stream-unavailable';
|
||||
} else {
|
||||
error.code = `slack-${reason}`;
|
||||
}
|
||||
return error;
|
||||
}
|
||||
|
||||
export function validSlackBotToken(value) {
|
||||
return typeof value === 'string'
|
||||
&& /^xoxb-[A-Za-z0-9-]{16,}$/.test(value.trim());
|
||||
}
|
||||
|
||||
export function validSlackAppToken(value) {
|
||||
return typeof value === 'string'
|
||||
&& /^xapp-[A-Za-z0-9-]{16,}$/.test(value.trim());
|
||||
}
|
||||
|
||||
export class SlackApi {
|
||||
#botToken;
|
||||
#appToken;
|
||||
#fetch;
|
||||
#baseUrl;
|
||||
|
||||
constructor({ botToken, appToken, fetchImpl = fetch, baseUrl = DEFAULT_BASE_URL }) {
|
||||
if (botToken !== undefined && !validSlackBotToken(botToken)) {
|
||||
throw new TypeError('Slack Bot Token is invalid');
|
||||
}
|
||||
if (appToken !== undefined && !validSlackAppToken(appToken)) {
|
||||
throw new TypeError('Slack App Token is invalid');
|
||||
}
|
||||
if (!botToken && !appToken) throw new TypeError('SlackApi requires a token');
|
||||
if (typeof fetchImpl !== 'function') throw new TypeError('SlackApi requires fetch');
|
||||
this.#botToken = botToken?.trim();
|
||||
this.#appToken = appToken?.trim();
|
||||
this.#fetch = fetchImpl;
|
||||
this.#baseUrl = new URL(baseUrl);
|
||||
}
|
||||
|
||||
authTest(options = {}) {
|
||||
return this.#request('auth.test', { ...options, tokenKind: 'bot' });
|
||||
}
|
||||
|
||||
openConnection(options = {}) {
|
||||
return this.#request('apps.connections.open', {
|
||||
...options,
|
||||
tokenKind: 'app',
|
||||
body: undefined,
|
||||
});
|
||||
}
|
||||
|
||||
postMessage({ channelId, text, threadTs, signal }) {
|
||||
return this.#request('chat.postMessage', {
|
||||
tokenKind: 'bot',
|
||||
signal,
|
||||
body: {
|
||||
channel: slackId(channelId, 'channel id'),
|
||||
text: safeOutgoingText(text),
|
||||
...(threadTs ? { thread_ts: cleanString(threadTs) } : {}),
|
||||
mrkdwn: true,
|
||||
link_names: false,
|
||||
unfurl_links: false,
|
||||
unfurl_media: false,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
updateMessage({ channelId, ts, text, signal }) {
|
||||
return this.#request('chat.update', {
|
||||
tokenKind: 'bot',
|
||||
signal,
|
||||
body: {
|
||||
channel: slackId(channelId, 'channel id'),
|
||||
ts: requiredString(ts, 'message timestamp'),
|
||||
text: safeOutgoingText(text),
|
||||
parse: 'none',
|
||||
link_names: false,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
startStream({ channelId, threadTs, recipientTeamId, recipientUserId, markdownText, signal }) {
|
||||
return this.#request('chat.startStream', {
|
||||
tokenKind: 'bot',
|
||||
signal,
|
||||
body: {
|
||||
channel: slackId(channelId, 'channel id'),
|
||||
thread_ts: requiredString(threadTs, 'thread timestamp'),
|
||||
...(recipientTeamId ? { recipient_team_id: slackId(recipientTeamId, 'team id') } : {}),
|
||||
...(recipientUserId ? { recipient_user_id: slackId(recipientUserId, 'user id') } : {}),
|
||||
...(cleanString(markdownText) ? { markdown_text: safeOutgoingText(markdownText) } : {}),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
appendStream({ channelId, ts, markdownText, signal }) {
|
||||
return this.#request('chat.appendStream', {
|
||||
tokenKind: 'bot',
|
||||
signal,
|
||||
body: {
|
||||
channel: slackId(channelId, 'channel id'),
|
||||
ts: requiredString(ts, 'stream timestamp'),
|
||||
markdown_text: safeOutgoingText(markdownText, { trim: false }),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
stopStream({ channelId, ts, markdownText, signal }) {
|
||||
return this.#request('chat.stopStream', {
|
||||
tokenKind: 'bot',
|
||||
signal,
|
||||
body: {
|
||||
channel: slackId(channelId, 'channel id'),
|
||||
ts: requiredString(ts, 'stream timestamp'),
|
||||
...(cleanString(markdownText) ? {
|
||||
markdown_text: safeOutgoingText(markdownText, { trim: false }),
|
||||
} : {}),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
async #request(method, {
|
||||
tokenKind,
|
||||
body,
|
||||
signal,
|
||||
timeoutMs = 15_000,
|
||||
retry = true,
|
||||
}) {
|
||||
const token = tokenKind === 'app' ? this.#appToken : this.#botToken;
|
||||
if (!token) throw new TypeError(`Slack ${tokenKind} token is required for ${method}`);
|
||||
let response;
|
||||
try {
|
||||
response = await this.#fetch(new URL(method, this.#baseUrl), {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
authorization: `Bearer ${token}`,
|
||||
'content-type': body === undefined
|
||||
? 'application/x-www-form-urlencoded;charset=utf-8'
|
||||
: 'application/json;charset=utf-8',
|
||||
'user-agent': 'DeepSeek-Harness-dsh-im (https://github.com/xmanrui/dsh-im, 0.2.2)',
|
||||
},
|
||||
...(body === undefined ? {} : { body: JSON.stringify(body) }),
|
||||
signal: requestSignal(signal, timeoutMs),
|
||||
redirect: 'error',
|
||||
});
|
||||
} catch (error) {
|
||||
if (error?.name === 'AbortError' || error?.name === 'TimeoutError') throw error;
|
||||
throw new Error(`Slack ${method} transport failed`);
|
||||
}
|
||||
|
||||
let payload;
|
||||
try {
|
||||
payload = await response.json();
|
||||
} catch {
|
||||
throw new Error(`Slack ${method} returned invalid JSON`);
|
||||
}
|
||||
if (response.status === 429 && retry) {
|
||||
const seconds = Number(response.headers.get('retry-after')) || 1;
|
||||
await delay(Math.min(10_000, Math.max(100, seconds * 1_000)), signal);
|
||||
return this.#request(method, { tokenKind, body, signal, timeoutMs, retry: false });
|
||||
}
|
||||
if (!response.ok || payload?.ok !== true) throw apiFailure(method, payload, tokenKind);
|
||||
return payload;
|
||||
}
|
||||
}
|
||||
|
||||
export async function inspectSlackCredentials({ botToken, appToken }, options = {}) {
|
||||
if (!validSlackBotToken(botToken)) {
|
||||
const error = new TypeError('Slack Bot Token 必须以 xoxb- 开头。');
|
||||
error.code = 'slack-invalid-bot-token';
|
||||
throw error;
|
||||
}
|
||||
if (!validSlackAppToken(appToken)) {
|
||||
const error = new TypeError('Slack App Token 必须以 xapp- 开头。');
|
||||
error.code = 'slack-invalid-app-token';
|
||||
throw error;
|
||||
}
|
||||
const api = new SlackApi({ botToken, appToken, ...options });
|
||||
const [identity, connection] = await Promise.all([api.authTest(), api.openConnection()]);
|
||||
if (!identity?.team_id || !identity?.user_id || !identity?.bot_id) {
|
||||
throw new Error('Slack Bot Token 没有返回完整的机器人身份。');
|
||||
}
|
||||
let socketUrl;
|
||||
try {
|
||||
socketUrl = new URL(connection?.url);
|
||||
} catch {
|
||||
socketUrl = null;
|
||||
}
|
||||
if (!socketUrl || socketUrl.protocol !== 'wss:') {
|
||||
const error = new Error('Slack App Token 无法创建 Socket Mode 连接,请确认已启用 Socket Mode 和 connections:write。');
|
||||
error.code = 'slack-socket-mode';
|
||||
throw error;
|
||||
}
|
||||
return {
|
||||
platformId: `${identity.team_id}:${identity.user_id}`,
|
||||
name: cleanString(identity.user) ?? 'DeepSeek Harness',
|
||||
username: cleanString(identity.user),
|
||||
teamId: String(identity.team_id),
|
||||
teamName: cleanString(identity.team),
|
||||
};
|
||||
}
|
||||
15
src/channels/slack/slack-bridge.mjs
Normal file
15
src/channels/slack/slack-bridge.mjs
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
import { TextHarnessBridge, createTextBridgeStatus } from '../shared/text-harness-bridge.mjs';
|
||||
|
||||
export const SLACK_DESCRIPTOR = Object.freeze({
|
||||
key: 'slack',
|
||||
label: 'Slack',
|
||||
connectionLabel: ' Socket Mode 长连接',
|
||||
});
|
||||
|
||||
export class SlackHarnessBridge extends TextHarnessBridge {
|
||||
constructor(options) {
|
||||
super({ descriptor: SLACK_DESCRIPTOR, ...options });
|
||||
}
|
||||
}
|
||||
|
||||
export { createTextBridgeStatus as createSlackBridgeStatus };
|
||||
318
src/channels/slack/slack-controller.mjs
Normal file
318
src/channels/slack/slack-controller.mjs
Normal file
|
|
@ -0,0 +1,318 @@
|
|||
import { deriveSlackBotIdentity, maskSlackBotId } from './config-store.mjs';
|
||||
import { inspectSlackCredentials } from './slack-api.mjs';
|
||||
import { SLACK_DESCRIPTOR } from './slack-bridge.mjs';
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function safeError(code, message) {
|
||||
return Object.freeze({ code, message });
|
||||
}
|
||||
|
||||
export class SlackController {
|
||||
#credentials;
|
||||
#configStore;
|
||||
#inspectCredentials;
|
||||
#createRuntime;
|
||||
#deleteState;
|
||||
#logger;
|
||||
#runtimes = new Map();
|
||||
#errors = new Map();
|
||||
#transitions = new Map();
|
||||
#revision = 0;
|
||||
#closed = false;
|
||||
|
||||
constructor({
|
||||
credentials,
|
||||
configStore,
|
||||
inspectCredentials = inspectSlackCredentials,
|
||||
createRuntime,
|
||||
deleteState = async () => {},
|
||||
logger = console,
|
||||
}) {
|
||||
if (!credentials || typeof credentials.resolve !== 'function'
|
||||
|| typeof credentials.set !== 'function' || typeof credentials.unset !== 'function') {
|
||||
throw new TypeError('Slack requires the DSH credential provider');
|
||||
}
|
||||
if (!configStore || typeof configStore.list !== 'function'
|
||||
|| typeof configStore.save !== 'function' || typeof configStore.remove !== 'function') {
|
||||
throw new TypeError('Slack requires a config store');
|
||||
}
|
||||
if (typeof inspectCredentials !== 'function' || typeof createRuntime !== 'function') {
|
||||
throw new TypeError('Slack controller dependencies are incomplete');
|
||||
}
|
||||
this.#credentials = credentials;
|
||||
this.#configStore = configStore;
|
||||
this.#inspectCredentials = inspectCredentials;
|
||||
this.#createRuntime = createRuntime;
|
||||
this.#deleteState = deleteState;
|
||||
this.#logger = logger;
|
||||
}
|
||||
|
||||
async initialize() {
|
||||
if (this.#closed) return this.status();
|
||||
for (const config of this.#configStore.list()) {
|
||||
await this.#withBotTransition(config.botId, async () => {
|
||||
if (this.#closed || this.#runtimes.get(config.botId)?.status?.ready) return;
|
||||
const resolved = await this.#resolveCredentials(config);
|
||||
if (!resolved) {
|
||||
this.#errors.set(config.botId, safeError(
|
||||
'missing-token',
|
||||
'Slack机器人凭据缺失,请移除后重新接入。',
|
||||
));
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await this.#startRuntime(config, resolved);
|
||||
this.#errors.delete(config.botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(config.botId, safeError(
|
||||
'connection-failed',
|
||||
'Slack Socket Mode 连接未就绪,插件会自动重试。',
|
||||
));
|
||||
this.#logger.warn?.(
|
||||
`[dsh-im:slack] bot ${config.botId} failed to initialize:`,
|
||||
error,
|
||||
);
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
}
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async bindCredentials({ botToken, appToken } = {}) {
|
||||
if (this.#closed) throw new Error('Slack controller is closed');
|
||||
const normalizedBotToken = cleanString(botToken);
|
||||
const normalizedAppToken = cleanString(appToken);
|
||||
if (!normalizedBotToken || !normalizedAppToken) {
|
||||
throw new TypeError('Slack Bot Token and App Token are required');
|
||||
}
|
||||
const inspected = await this.#inspectCredentials({
|
||||
botToken: normalizedBotToken,
|
||||
appToken: normalizedAppToken,
|
||||
});
|
||||
const platformId = cleanString(inspected?.platformId);
|
||||
const name = cleanString(inspected?.name);
|
||||
if (!platformId || !name) throw new Error('Slack returned an invalid bot identity');
|
||||
const identity = deriveSlackBotIdentity(platformId);
|
||||
|
||||
await this.#withBotTransition(identity.botId, async () => {
|
||||
if (this.#closed) throw new Error('Slack controller is closed');
|
||||
const previousConfig = this.#configStore.getByPlatformId(platformId);
|
||||
const previousBotToken = await this.#credentials.resolve(identity.botTokenRef).catch(() => undefined);
|
||||
const previousAppToken = await this.#credentials.resolve(identity.appTokenRef).catch(() => undefined);
|
||||
const config = {
|
||||
...identity,
|
||||
platformId,
|
||||
name,
|
||||
username: cleanString(inspected.username),
|
||||
teamId: cleanString(inspected.teamId),
|
||||
teamName: cleanString(inspected.teamName),
|
||||
createdAt: previousConfig?.createdAt ?? new Date().toISOString(),
|
||||
connectedAt: new Date().toISOString(),
|
||||
};
|
||||
try {
|
||||
await this.#credentials.set(identity.botTokenRef, normalizedBotToken);
|
||||
await this.#credentials.set(identity.appTokenRef, normalizedAppToken);
|
||||
await this.#configStore.save(config);
|
||||
} catch (error) {
|
||||
await Promise.all([
|
||||
this.#restoreCredential(identity.botTokenRef, previousBotToken),
|
||||
this.#restoreCredential(identity.appTokenRef, previousAppToken),
|
||||
]);
|
||||
throw error;
|
||||
}
|
||||
try {
|
||||
await this.#startRuntime(config, {
|
||||
botToken: normalizedBotToken,
|
||||
appToken: normalizedAppToken,
|
||||
});
|
||||
this.#errors.delete(identity.botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(identity.botId, safeError(
|
||||
'connection-failed',
|
||||
'Slack机器人已接入,Socket Mode 连接暂未就绪。',
|
||||
));
|
||||
this.#logger.warn?.(
|
||||
`[dsh-im:slack] bot ${identity.botId} credential connection failed:`,
|
||||
error,
|
||||
);
|
||||
}
|
||||
this.#touch();
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async reconnectBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown Slack bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
const resolved = await this.#resolveCredentials(config);
|
||||
if (!resolved) throw new Error('Slack bot credentials are missing');
|
||||
try {
|
||||
await this.#startRuntime(config, resolved);
|
||||
this.#errors.delete(botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(botId, safeError(
|
||||
'connection-failed',
|
||||
'Slack Socket Mode 连接仍未就绪,请检查两个 Token。',
|
||||
));
|
||||
throw error;
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async deleteBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown Slack bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
const previousBotToken = await this.#credentials.resolve(config.botTokenRef).catch(() => undefined);
|
||||
const previousAppToken = await this.#credentials.resolve(config.appTokenRef).catch(() => undefined);
|
||||
await this.#stopRuntime(botId);
|
||||
try {
|
||||
await this.#credentials.unset(config.botTokenRef);
|
||||
await this.#credentials.unset(config.appTokenRef);
|
||||
await this.#configStore.remove(botId);
|
||||
} catch (error) {
|
||||
await Promise.all([
|
||||
this.#restoreCredential(config.botTokenRef, previousBotToken),
|
||||
this.#restoreCredential(config.appTokenRef, previousAppToken),
|
||||
]);
|
||||
if (previousBotToken?.value && previousAppToken?.value) {
|
||||
await this.#startRuntime(config, {
|
||||
botToken: previousBotToken.value,
|
||||
appToken: previousAppToken.value,
|
||||
}).catch(() => undefined);
|
||||
}
|
||||
throw new Error('Unable to remove the Slack bot safely.', { cause: error });
|
||||
}
|
||||
await this.#deleteState({ botId, config }).catch((error) => {
|
||||
this.#logger.warn?.(`[dsh-im:slack] bot ${botId} state cleanup failed:`, error);
|
||||
});
|
||||
this.#errors.delete(botId);
|
||||
this.#touch();
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
status() {
|
||||
const bots = this.#configStore.list().map((config) => {
|
||||
const runtimeStatus = this.#runtimes.get(config.botId)?.status ?? null;
|
||||
const connected = runtimeStatus?.ready === true
|
||||
&& runtimeStatus.connectionState === 'connected'
|
||||
&& runtimeStatus.harnessReachable === true;
|
||||
const state = connected ? 'connected'
|
||||
: runtimeStatus?.connectionState === 'connecting' ? 'connecting'
|
||||
: this.#errors.has(config.botId) || runtimeStatus?.connectionState === 'failed'
|
||||
? 'error' : 'offline';
|
||||
return {
|
||||
botId: config.botId,
|
||||
state,
|
||||
connected,
|
||||
configured: true,
|
||||
bot: {
|
||||
name: config.name,
|
||||
username: config.username,
|
||||
teamName: config.teamName,
|
||||
idMasked: maskSlackBotId(config.platformId),
|
||||
},
|
||||
health: {
|
||||
status: connected ? 'healthy' : state === 'error' ? 'error' : 'offline',
|
||||
summary: connected ? `Slack${SLACK_DESCRIPTOR.connectionLabel}运行正常`
|
||||
: state === 'error' ? 'Slack连接未就绪,插件会自动重试'
|
||||
: 'Slack连接当前离线',
|
||||
lastCheckedAt: runtimeStatus?.lastCheckedAt ?? null,
|
||||
lastConnectedAt: runtimeStatus?.lastConnectedAt ?? null,
|
||||
},
|
||||
stats: {
|
||||
messagesReceived: runtimeStatus?.messagesReceived ?? 0,
|
||||
messagesReplied: runtimeStatus?.messagesReplied ?? 0,
|
||||
},
|
||||
error: structuredClone(this.#errors.get(config.botId) ?? null),
|
||||
};
|
||||
});
|
||||
const connected = bots.filter((bot) => bot.connected).length;
|
||||
return {
|
||||
schemaVersion: 1,
|
||||
revision: this.#revision,
|
||||
state: bots.length === 0 ? 'disconnected'
|
||||
: connected === bots.length ? 'connected'
|
||||
: connected > 0 ? 'degraded' : 'offline',
|
||||
bots,
|
||||
totals: { configured: bots.length, connected },
|
||||
};
|
||||
}
|
||||
|
||||
async close() {
|
||||
if (this.#closed) return;
|
||||
this.#closed = true;
|
||||
await Promise.allSettled([...this.#transitions.values()]);
|
||||
await Promise.allSettled([...this.#runtimes.keys()].map((botId) => this.#stopRuntime(botId)));
|
||||
}
|
||||
|
||||
async #startRuntime(config, { botToken, appToken }) {
|
||||
if (this.#closed) throw new Error('Slack controller is closed');
|
||||
await this.#stopRuntime(config.botId);
|
||||
if (this.#closed) throw new Error('Slack controller is closed');
|
||||
const runtime = await this.#createRuntime({
|
||||
botId: config.botId,
|
||||
config,
|
||||
botToken,
|
||||
appToken,
|
||||
});
|
||||
if (!runtime || typeof runtime.start !== 'function' || typeof runtime.stop !== 'function') {
|
||||
throw new TypeError('createRuntime returned an invalid Slack runtime');
|
||||
}
|
||||
this.#runtimes.set(config.botId, runtime);
|
||||
try {
|
||||
await runtime.start();
|
||||
} catch (error) {
|
||||
await runtime.stop().catch(() => undefined);
|
||||
this.#runtimes.delete(config.botId);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async #stopRuntime(botId) {
|
||||
const runtime = this.#runtimes.get(botId);
|
||||
this.#runtimes.delete(botId);
|
||||
await runtime?.stop().catch((error) => {
|
||||
this.#logger.warn?.(`[dsh-im:slack] bot ${botId} failed to stop cleanly:`, error);
|
||||
});
|
||||
}
|
||||
|
||||
async #resolveCredentials(config) {
|
||||
const [bot, app] = await Promise.all([
|
||||
this.#credentials.resolve(config.botTokenRef).catch(() => undefined),
|
||||
this.#credentials.resolve(config.appTokenRef).catch(() => undefined),
|
||||
]);
|
||||
const botToken = cleanString(bot?.value);
|
||||
const appToken = cleanString(app?.value);
|
||||
return botToken && appToken ? { botToken, appToken } : null;
|
||||
}
|
||||
|
||||
async #restoreCredential(ref, previous) {
|
||||
if (previous?.value) await this.#credentials.set(ref, previous.value).catch(() => undefined);
|
||||
else await this.#credentials.unset(ref).catch(() => undefined);
|
||||
}
|
||||
|
||||
#withBotTransition(botId, operation) {
|
||||
const previous = this.#transitions.get(botId) ?? Promise.resolve();
|
||||
const current = previous.catch(() => undefined).then(operation);
|
||||
const settled = current.finally(() => {
|
||||
if (this.#transitions.get(botId) === settled) this.#transitions.delete(botId);
|
||||
});
|
||||
this.#transitions.set(botId, settled);
|
||||
return settled;
|
||||
}
|
||||
|
||||
#touch() {
|
||||
this.#revision += 1;
|
||||
}
|
||||
}
|
||||
462
src/channels/slack/slack-runtime.mjs
Normal file
462
src/channels/slack/slack-runtime.mjs
Normal file
|
|
@ -0,0 +1,462 @@
|
|||
import { splitMessageText } from '../shared/editable-message-stream.mjs';
|
||||
import { SlackApi } from './slack-api.mjs';
|
||||
import { createSlackBridgeStatus, SlackHarnessBridge } from './slack-bridge.mjs';
|
||||
|
||||
const RECONNECT_DELAYS_MS = Object.freeze([1_000, 3_000, 5_000, 10_000, 30_000]);
|
||||
const SLACK_MESSAGE_LIMIT = 38_000;
|
||||
const SLACK_STREAM_CHUNK_LIMIT = 11_000;
|
||||
|
||||
function addSocketListener(socket, event, listener) {
|
||||
if (typeof socket.addEventListener === 'function') socket.addEventListener(event, listener);
|
||||
else if (typeof socket.on === 'function') socket.on(event, listener);
|
||||
else throw new TypeError('Slack WebSocket does not support events');
|
||||
}
|
||||
|
||||
function eventData(event) {
|
||||
const value = event?.data ?? event;
|
||||
if (typeof value === 'string') return value;
|
||||
if (Buffer.isBuffer(value)) return value.toString('utf8');
|
||||
if (value instanceof ArrayBuffer) return Buffer.from(value).toString('utf8');
|
||||
if (ArrayBuffer.isView(value)) {
|
||||
return Buffer.from(value.buffer, value.byteOffset, value.byteLength).toString('utf8');
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function socketUrl(value) {
|
||||
const url = new URL(value);
|
||||
if (url.protocol !== 'wss:') throw new Error('Slack returned an insecure Socket Mode URL');
|
||||
return url.href;
|
||||
}
|
||||
|
||||
function decodeSlackText(value) {
|
||||
return typeof value === 'string' ? value
|
||||
.replaceAll('&', '&')
|
||||
.replaceAll('<', '<')
|
||||
.replaceAll('>', '>') : '';
|
||||
}
|
||||
|
||||
function stripBotMention(value, botUserId) {
|
||||
return decodeSlackText(value)
|
||||
.replace(new RegExp(`<@${botUserId}>`, 'gi'), '')
|
||||
.trim();
|
||||
}
|
||||
|
||||
export function normalizeSlackEvent(payload, botUserId) {
|
||||
const event = payload?.event;
|
||||
if (!event || !payload?.event_id || !event.channel || !event.user || !event.ts) return null;
|
||||
const direct = event.type === 'message' && event.channel_type === 'im';
|
||||
const mentioned = event.type === 'app_mention';
|
||||
if (!direct && !mentioned) return null;
|
||||
if (event.subtype || event.bot_id || event.app_id) return null;
|
||||
const threadTs = String(event.thread_ts ?? event.ts);
|
||||
return {
|
||||
messageId: String(payload.event_id),
|
||||
senderId: String(event.user),
|
||||
senderIsBot: String(event.user) === String(botUserId),
|
||||
kind: direct ? 'direct' : 'group',
|
||||
conversationId: direct ? String(event.channel) : `${event.channel}:${threadTs}`,
|
||||
content: stripBotMention(event.text ?? '', botUserId),
|
||||
addressed: direct || mentioned,
|
||||
replyTarget: {
|
||||
channelId: String(event.channel),
|
||||
threadTs,
|
||||
recipientUserId: String(event.user),
|
||||
recipientTeamId: String(event.user_team ?? payload.team_id ?? ''),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function isToolProgress(text) {
|
||||
return typeof text === 'string' && /^正在使用.+…$/.test(text.trim());
|
||||
}
|
||||
|
||||
async function appendInChunks(api, target, ts, text, signal) {
|
||||
if (!text) return;
|
||||
for (let offset = 0; offset < text.length; offset += SLACK_STREAM_CHUNK_LIMIT) {
|
||||
await api.appendStream({
|
||||
channelId: target.channelId,
|
||||
ts,
|
||||
markdownText: text.slice(offset, offset + SLACK_STREAM_CHUNK_LIMIT),
|
||||
signal,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async function createSlackMessageStream({ api, target, signal, logger }) {
|
||||
const started = await api.startStream({
|
||||
channelId: target.channelId,
|
||||
threadTs: target.threadTs,
|
||||
recipientTeamId: target.recipientTeamId || undefined,
|
||||
recipientUserId: target.recipientUserId || undefined,
|
||||
signal,
|
||||
});
|
||||
const ts = typeof started?.ts === 'string' ? started.ts : null;
|
||||
if (!ts) throw new Error('Slack did not return a streaming message timestamp');
|
||||
|
||||
let appended = '';
|
||||
let pending = '';
|
||||
let timer = null;
|
||||
let inFlight = null;
|
||||
let broken = false;
|
||||
let closed = false;
|
||||
|
||||
const appendLatest = async (text) => {
|
||||
const next = splitMessageText(text, SLACK_MESSAGE_LIMIT)[0] ?? '';
|
||||
if (!next || !next.startsWith(appended)) return;
|
||||
const delta = next.slice(appended.length);
|
||||
if (!delta) return;
|
||||
await appendInChunks(api, target, ts, delta, signal);
|
||||
appended = next;
|
||||
};
|
||||
|
||||
const schedule = () => {
|
||||
if (closed || broken || timer !== null || inFlight || !pending) return;
|
||||
timer = setTimeout(() => {
|
||||
timer = null;
|
||||
const text = pending;
|
||||
pending = '';
|
||||
inFlight = appendLatest(text)
|
||||
.catch((error) => {
|
||||
broken = true;
|
||||
logger.warn?.('[dsh-im:slack] streaming append failed:', error);
|
||||
})
|
||||
.finally(() => {
|
||||
inFlight = null;
|
||||
schedule();
|
||||
});
|
||||
}, 350);
|
||||
timer?.unref?.();
|
||||
};
|
||||
|
||||
return {
|
||||
update(text) {
|
||||
if (closed || broken || typeof text !== 'string' || !text.trim() || isToolProgress(text)) return;
|
||||
pending = text;
|
||||
schedule();
|
||||
},
|
||||
async finish(text) {
|
||||
if (closed) throw new Error('Slack message stream is already closed');
|
||||
closed = true;
|
||||
if (timer !== null) clearTimeout(timer);
|
||||
timer = null;
|
||||
pending = '';
|
||||
await inFlight?.catch(() => undefined);
|
||||
|
||||
const chunks = splitMessageText(text, SLACK_MESSAGE_LIMIT);
|
||||
const first = chunks[0] ?? '处理完成。';
|
||||
if (!broken && first.startsWith(appended)) {
|
||||
await appendInChunks(api, target, ts, first.slice(appended.length), signal);
|
||||
await api.stopStream({ channelId: target.channelId, ts, signal });
|
||||
} else {
|
||||
await api.stopStream({ channelId: target.channelId, ts, signal }).catch(() => undefined);
|
||||
await api.updateMessage({ channelId: target.channelId, ts, text: first, signal });
|
||||
}
|
||||
for (const chunk of chunks.slice(1)) {
|
||||
await api.postMessage({
|
||||
channelId: target.channelId,
|
||||
threadTs: target.threadTs,
|
||||
text: chunk,
|
||||
signal,
|
||||
});
|
||||
}
|
||||
},
|
||||
cancel() {
|
||||
closed = true;
|
||||
pending = '';
|
||||
if (timer !== null) clearTimeout(timer);
|
||||
timer = null;
|
||||
void api.stopStream({ channelId: target.channelId, ts, signal }).catch(() => undefined);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
class SlackBotClient {
|
||||
#api;
|
||||
#signal;
|
||||
#logger;
|
||||
|
||||
constructor({ api, signal, logger }) {
|
||||
this.#api = api;
|
||||
this.#signal = signal;
|
||||
this.#logger = logger;
|
||||
}
|
||||
|
||||
async sendText(target, text) {
|
||||
const chunks = splitMessageText(text, SLACK_MESSAGE_LIMIT);
|
||||
let result = null;
|
||||
for (const chunk of chunks) {
|
||||
result = await this.#api.postMessage({
|
||||
channelId: target.channelId,
|
||||
threadTs: target.threadTs,
|
||||
text: chunk,
|
||||
signal: this.#signal,
|
||||
});
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
openStream(target) {
|
||||
return createSlackMessageStream({
|
||||
api: this.#api,
|
||||
target,
|
||||
signal: this.#signal,
|
||||
logger: this.#logger,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
export function createSlackRuntimeStatus() {
|
||||
return {
|
||||
startedAt: null,
|
||||
ready: false,
|
||||
connectionState: 'idle',
|
||||
harnessReachable: false,
|
||||
lastCheckedAt: null,
|
||||
lastConnectedAt: null,
|
||||
lastError: null,
|
||||
...createSlackBridgeStatus(),
|
||||
};
|
||||
}
|
||||
|
||||
export class SlackRuntime {
|
||||
#config;
|
||||
#botToken;
|
||||
#appToken;
|
||||
#harness;
|
||||
#state;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
#createApi;
|
||||
#createWebSocket;
|
||||
#status = createSlackRuntimeStatus();
|
||||
#api = null;
|
||||
#bridge = null;
|
||||
#abortController = null;
|
||||
#socket = null;
|
||||
#appId = null;
|
||||
#reconnectTimer = null;
|
||||
#reconnectAttempt = 0;
|
||||
#generation = 0;
|
||||
#stopped = true;
|
||||
#starting = null;
|
||||
|
||||
constructor({
|
||||
config,
|
||||
botToken,
|
||||
appToken,
|
||||
harness,
|
||||
state,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
createApi = (options) => new SlackApi(options),
|
||||
createWebSocket = (url) => new WebSocket(url),
|
||||
}) {
|
||||
if (!config || !botToken || !appToken || !harness || !state) {
|
||||
throw new TypeError('SlackRuntime requires config, both tokens, Harness, and state');
|
||||
}
|
||||
if (typeof createWebSocket !== 'function') throw new TypeError('SlackRuntime requires WebSocket');
|
||||
this.#config = config;
|
||||
this.#botToken = botToken;
|
||||
this.#appToken = appToken;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
this.#createApi = createApi;
|
||||
this.#createWebSocket = createWebSocket;
|
||||
}
|
||||
|
||||
get status() {
|
||||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
async start() {
|
||||
if (this.#status.ready && this.#socket) return this.status;
|
||||
if (this.#starting) return this.#starting;
|
||||
this.#starting = this.#start().finally(() => {
|
||||
this.#starting = null;
|
||||
});
|
||||
return this.#starting;
|
||||
}
|
||||
|
||||
async #start() {
|
||||
await this.stop();
|
||||
this.#stopped = false;
|
||||
this.#reconnectAttempt = 0;
|
||||
this.#status.startedAt = new Date().toISOString();
|
||||
this.#status.connectionState = 'connecting';
|
||||
this.#status.lastError = null;
|
||||
await this.#harness.ensureRunning();
|
||||
this.#status.harnessReachable = true;
|
||||
const controller = new AbortController();
|
||||
this.#abortController = controller;
|
||||
const api = this.#createApi({ botToken: this.#botToken, appToken: this.#appToken });
|
||||
this.#api = api;
|
||||
try {
|
||||
const identity = await api.authTest({ signal: controller.signal });
|
||||
if (`${identity?.team_id}:${identity?.user_id}` !== this.#config.platformId) {
|
||||
throw new Error('Slack Bot Token identity does not match the saved bot');
|
||||
}
|
||||
const client = new SlackBotClient({ api, signal: controller.signal, logger: this.#logger });
|
||||
this.#bridge = new SlackHarnessBridge({
|
||||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
});
|
||||
let timer;
|
||||
try {
|
||||
await Promise.race([
|
||||
this.#connect(),
|
||||
new Promise((_, reject) => {
|
||||
timer = setTimeout(
|
||||
() => reject(new Error('Slack Socket Mode did not become ready in time')),
|
||||
this.#connectTimeoutMs,
|
||||
);
|
||||
timer?.unref?.();
|
||||
}),
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
return this.status;
|
||||
} catch (error) {
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'failed';
|
||||
this.#status.lastError = error?.message ?? String(error);
|
||||
await this.stop();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async #connect() {
|
||||
if (this.#stopped) throw new Error('Slack runtime is stopped');
|
||||
const connection = await this.#api.openConnection({ signal: this.#abortController?.signal });
|
||||
return this.#openSocket(connection?.url);
|
||||
}
|
||||
|
||||
#openSocket(value) {
|
||||
if (this.#stopped) return Promise.reject(new Error('Slack runtime is stopped'));
|
||||
const generation = ++this.#generation;
|
||||
const socket = this.#createWebSocket(socketUrl(value));
|
||||
this.#socket = socket;
|
||||
let settled = false;
|
||||
return new Promise((resolve, reject) => {
|
||||
const markReady = (packet) => {
|
||||
if (settled || generation !== this.#generation) return;
|
||||
settled = true;
|
||||
this.#appId = packet?.connection_info?.app_id ?? null;
|
||||
this.#reconnectAttempt = 0;
|
||||
const now = Date.now();
|
||||
this.#status.ready = true;
|
||||
this.#status.connectionState = 'connected';
|
||||
this.#status.lastCheckedAt = now;
|
||||
this.#status.lastConnectedAt = now;
|
||||
this.#status.lastError = null;
|
||||
resolve();
|
||||
};
|
||||
|
||||
addSocketListener(socket, 'message', (event) => {
|
||||
if (generation !== this.#generation || this.#stopped) return;
|
||||
const raw = eventData(event);
|
||||
if (!raw) return;
|
||||
let packet;
|
||||
try {
|
||||
packet = JSON.parse(raw);
|
||||
} catch {
|
||||
this.#logger.warn?.('[dsh-im:slack] ignored malformed Socket Mode JSON');
|
||||
return;
|
||||
}
|
||||
if (packet.type === 'hello') {
|
||||
markReady(packet);
|
||||
return;
|
||||
}
|
||||
if (packet.envelope_id && socket.readyState === 1) {
|
||||
socket.send(JSON.stringify({ envelope_id: packet.envelope_id }));
|
||||
this.#status.lastCheckedAt = Date.now();
|
||||
}
|
||||
if (packet.type === 'disconnect') {
|
||||
socket.close(4000, 'Slack requested reconnect');
|
||||
return;
|
||||
}
|
||||
if (packet.type !== 'events_api' || packet.payload?.type !== 'event_callback') return;
|
||||
if (this.#appId && packet.payload.api_app_id
|
||||
&& packet.payload.api_app_id !== this.#appId) return;
|
||||
const message = normalizeSlackEvent(packet.payload, this.#config.platformId.split(':')[1]);
|
||||
if (message) void this.#bridge?.accept(message);
|
||||
});
|
||||
|
||||
addSocketListener(socket, 'close', (event = {}) => {
|
||||
if (generation !== this.#generation) return;
|
||||
if (this.#socket === socket) this.#socket = null;
|
||||
if (this.#stopped) {
|
||||
if (!settled) reject(new DOMException('Stopped', 'AbortError'));
|
||||
return;
|
||||
}
|
||||
const code = Number(event.code) || 0;
|
||||
const error = new Error(`Slack Socket Mode closed (${code || 'unknown'})`);
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'connecting';
|
||||
this.#status.lastError = error.message;
|
||||
if (!settled) {
|
||||
settled = true;
|
||||
reject(error);
|
||||
}
|
||||
this.#scheduleReconnect();
|
||||
});
|
||||
|
||||
addSocketListener(socket, 'error', () => {
|
||||
if (generation !== this.#generation || this.#stopped) return;
|
||||
this.#status.lastError = 'Slack Socket Mode WebSocket error';
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
#scheduleReconnect() {
|
||||
if (this.#stopped || this.#reconnectTimer !== null) return;
|
||||
const delay = RECONNECT_DELAYS_MS[Math.min(this.#reconnectAttempt, RECONNECT_DELAYS_MS.length - 1)];
|
||||
this.#reconnectAttempt += 1;
|
||||
this.#reconnectTimer = setTimeout(() => {
|
||||
this.#reconnectTimer = null;
|
||||
void this.#connect().catch((error) => {
|
||||
if (this.#stopped) return;
|
||||
this.#logger.warn?.('[dsh-im:slack] Socket Mode reconnect failed:', error);
|
||||
this.#scheduleReconnect();
|
||||
});
|
||||
}, delay);
|
||||
this.#reconnectTimer?.unref?.();
|
||||
}
|
||||
|
||||
async stop() {
|
||||
this.#stopped = true;
|
||||
this.#generation += 1;
|
||||
this.#abortController?.abort();
|
||||
this.#abortController = null;
|
||||
if (this.#reconnectTimer !== null) clearTimeout(this.#reconnectTimer);
|
||||
this.#reconnectTimer = null;
|
||||
const socket = this.#socket;
|
||||
const bridge = this.#bridge;
|
||||
this.#socket = null;
|
||||
this.#bridge = null;
|
||||
this.#api = null;
|
||||
this.#appId = null;
|
||||
try {
|
||||
if (socket && socket.readyState < 2) socket.close(1000, 'Plugin stopped');
|
||||
} catch (error) {
|
||||
this.#logger.warn?.(`[dsh-im:slack] bot ${this.#config.botId} failed to close Socket Mode:`, error);
|
||||
}
|
||||
await Promise.race([
|
||||
bridge?.waitForIdle() ?? Promise.resolve(),
|
||||
new Promise((resolve) => setTimeout(resolve, 2_000)),
|
||||
]);
|
||||
this.#status.ready = false;
|
||||
this.#status.connectionState = 'idle';
|
||||
return this.status;
|
||||
}
|
||||
}
|
||||
3
src/channels/slack/state-store.mjs
Normal file
3
src/channels/slack/state-store.mjs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
import { ConversationStateStore } from '../shared/conversation-state-store.mjs';
|
||||
|
||||
export class SlackStateStore extends ConversationStateStore {}
|
||||
Loading…
Add table
Add a link
Reference in a new issue