mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 03:03:24 +08:00
fix: preserve Telegram compatibility and stabilize Office
This commit is contained in:
parent
7f114e1c36
commit
24e690498b
23 changed files with 1679 additions and 314 deletions
|
|
@ -95,6 +95,7 @@ export class OfficeJobExecutor {
|
|||
#transport;
|
||||
#createHarness;
|
||||
#logger;
|
||||
#sleep;
|
||||
#active = new Map();
|
||||
#queued = new Set();
|
||||
#completed = new Set();
|
||||
|
|
@ -107,7 +108,13 @@ export class OfficeJobExecutor {
|
|||
lastJobAt: null,
|
||||
};
|
||||
|
||||
constructor({ config, transport, createHarness, logger = console }) {
|
||||
constructor({
|
||||
config,
|
||||
transport,
|
||||
createHarness,
|
||||
logger = console,
|
||||
sleepImpl = sleep,
|
||||
}) {
|
||||
if (!config || !transport || typeof createHarness !== 'function') {
|
||||
throw new TypeError('OfficeJobExecutor requires config, transport, and createHarness');
|
||||
}
|
||||
|
|
@ -115,6 +122,7 @@ export class OfficeJobExecutor {
|
|||
this.#transport = transport;
|
||||
this.#createHarness = createHarness;
|
||||
this.#logger = logger;
|
||||
this.#sleep = sleepImpl;
|
||||
}
|
||||
|
||||
get status() { return structuredClone(this.#status); }
|
||||
|
|
@ -256,18 +264,10 @@ export class OfficeJobExecutor {
|
|||
|
||||
async #renew(jobId, entry) {
|
||||
while (!entry.controller.signal.aborted) {
|
||||
try { await sleep(RENEW_MS, undefined, { signal: entry.controller.signal }); }
|
||||
try { await this.#sleep(RENEW_MS, undefined, { signal: entry.controller.signal }); }
|
||||
catch { return; }
|
||||
try {
|
||||
await this.#transport.renewJob(jobId, entry.leaseToken, { signal: entry.controller.signal });
|
||||
const snapshot = await this.#transport.getJob(jobId, { signal: entry.controller.signal });
|
||||
const approval = snapshot?.job?.approval;
|
||||
if (approval && (approval.status === 'approved' || approval.status === 'rejected')) {
|
||||
entry.approvals.get(approval.id)?.resolve({
|
||||
decision: approval.status,
|
||||
answer: clean(approval.answer),
|
||||
});
|
||||
}
|
||||
} catch (error) {
|
||||
if (entry.controller.signal.aborted) return;
|
||||
entry.cancelled = true;
|
||||
|
|
@ -280,6 +280,22 @@ export class OfficeJobExecutor {
|
|||
}
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const snapshot = await this.#transport.getJob(jobId, { signal: entry.controller.signal });
|
||||
const approval = snapshot?.job?.approval;
|
||||
if (approval && (approval.status === 'approved' || approval.status === 'rejected')) {
|
||||
entry.approvals.get(approval.id)?.resolve({
|
||||
decision: approval.status,
|
||||
answer: clean(approval.answer),
|
||||
});
|
||||
}
|
||||
} catch (error) {
|
||||
if (entry.controller.signal.aborted) return;
|
||||
this.#logger.warn?.(
|
||||
`[dsh-im:office] Job ${jobId} approval poll failed; will retry after the next renewal:`,
|
||||
error.message,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -22,15 +22,25 @@ export class OfficeRuntime {
|
|||
#token;
|
||||
#logger;
|
||||
#transport;
|
||||
#sleep;
|
||||
#controller = null;
|
||||
#task = null;
|
||||
#status;
|
||||
#jobs;
|
||||
|
||||
constructor({ config, token, logger = console, transport, createHarness, jobExecutor }) {
|
||||
constructor({
|
||||
config,
|
||||
token,
|
||||
logger = console,
|
||||
transport,
|
||||
createHarness,
|
||||
jobExecutor,
|
||||
sleepImpl = sleep,
|
||||
}) {
|
||||
this.#config = config;
|
||||
this.#token = token;
|
||||
this.#logger = logger;
|
||||
this.#sleep = sleepImpl;
|
||||
this.#transport = transport ?? new OfficeTransport({
|
||||
baseUrl: config.baseUrl, deviceId: config.deviceId, token,
|
||||
});
|
||||
|
|
@ -91,7 +101,10 @@ export class OfficeRuntime {
|
|||
const heartbeat = await this.#transport.heartbeat(this.capabilities(), { signal: attemptSignal });
|
||||
this.#offerJobs(heartbeat?.jobs);
|
||||
this.#status.lastHeartbeatAt = new Date().toISOString();
|
||||
const heartbeatTask = this.#heartbeatLoop(attemptSignal);
|
||||
let streamOpened = false;
|
||||
const heartbeatTask = this.#heartbeatLoop(attemptSignal, () => {
|
||||
if (streamOpened && !attemptSignal.aborted) attempt = 0;
|
||||
});
|
||||
const stream = this.#transport.stream({
|
||||
signal: attemptSignal,
|
||||
lastEventId: this.#status.lastEventId,
|
||||
|
|
@ -99,7 +112,7 @@ export class OfficeRuntime {
|
|||
this.#status.connected = true;
|
||||
this.#status.state = 'connected';
|
||||
this.#status.error = null;
|
||||
attempt = 0;
|
||||
streamOpened = true;
|
||||
},
|
||||
onEvent: async (event) => {
|
||||
this.#status.lastEventAt = new Date().toISOString();
|
||||
|
|
@ -112,13 +125,14 @@ export class OfficeRuntime {
|
|||
await Promise.race([stream, heartbeatTask]);
|
||||
} catch (error) {
|
||||
if (signal.aborted) break;
|
||||
attemptController.abort();
|
||||
this.#status.connected = false;
|
||||
this.#status.state = 'reconnecting';
|
||||
this.#status.error = safeConnectionError(error);
|
||||
this.#status.reconnects += 1;
|
||||
const delay = RETRY_DELAYS[Math.min(attempt, RETRY_DELAYS.length - 1)];
|
||||
attempt += 1;
|
||||
try { await sleep(delay, undefined, { signal }); } catch { break; }
|
||||
try { await this.#sleep(delay, undefined, { signal }); } catch { break; }
|
||||
} finally {
|
||||
attemptController.abort();
|
||||
}
|
||||
|
|
@ -127,12 +141,13 @@ export class OfficeRuntime {
|
|||
this.#status.state = 'idle';
|
||||
}
|
||||
|
||||
async #heartbeatLoop(signal) {
|
||||
async #heartbeatLoop(signal, onSuccess) {
|
||||
while (!signal.aborted) {
|
||||
await sleep(this.#config.heartbeatSeconds * 1_000, undefined, { signal });
|
||||
await this.#sleep(this.#config.heartbeatSeconds * 1_000, undefined, { signal });
|
||||
const heartbeat = await this.#transport.heartbeat(this.capabilities(), { signal });
|
||||
this.#offerJobs(heartbeat?.jobs);
|
||||
this.#status.lastHeartbeatAt = new Date().toISOString();
|
||||
onSuccess?.();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -163,6 +163,34 @@ export class TokenBotController {
|
|||
return this.status();
|
||||
}
|
||||
|
||||
async updateBotConfig(botId, update) {
|
||||
if (this.#closed) throw new Error(`${this.#descriptor.label} controller is closed`);
|
||||
if (typeof update !== 'function') throw new TypeError('Bot config update must be a function');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
if (this.#closed) throw new Error(`${this.#descriptor.label} controller is closed`);
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error(`Unknown ${this.#descriptor.label} bot`);
|
||||
const token = await this.#resolveToken(config.tokenRef);
|
||||
if (!token) throw new Error(`${this.#descriptor.label} bot token is missing`);
|
||||
if (this.#closed) throw new Error(`${this.#descriptor.label} controller is closed`);
|
||||
const nextConfig = update(config);
|
||||
const savedConfig = await this.#configStore.save(nextConfig);
|
||||
try {
|
||||
await this.#startRuntime(savedConfig, token);
|
||||
this.#errors.delete(botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(botId, safeError(
|
||||
'connection-failed',
|
||||
`${this.#descriptor.label}连接仍未就绪,请稍后重试。`,
|
||||
));
|
||||
throw error;
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async sendConnectionTest(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error(`Unknown ${this.#descriptor.label} bot`);
|
||||
|
|
|
|||
|
|
@ -33,16 +33,26 @@ export class TokenBotConfigStore {
|
|||
#channel;
|
||||
#botPrefix;
|
||||
#tokenRefPrefix;
|
||||
#normalizeBotExtension;
|
||||
#botIdPattern;
|
||||
#tokenRefPattern;
|
||||
#value = EMPTY_DOCUMENT;
|
||||
#writeQueue = Promise.resolve();
|
||||
|
||||
constructor(path, { channel, botPrefix, tokenRefPrefix }) {
|
||||
constructor(path, {
|
||||
channel,
|
||||
botPrefix,
|
||||
tokenRefPrefix,
|
||||
normalizeBotExtension = () => ({}),
|
||||
}) {
|
||||
if (typeof normalizeBotExtension !== 'function') {
|
||||
throw new TypeError('normalizeBotExtension must be a function');
|
||||
}
|
||||
this.#path = path;
|
||||
this.#channel = channel;
|
||||
this.#botPrefix = botPrefix;
|
||||
this.#tokenRefPrefix = tokenRefPrefix;
|
||||
this.#normalizeBotExtension = normalizeBotExtension;
|
||||
this.#botIdPattern = new RegExp(`^${escapePattern(botPrefix)}_[a-f0-9]{24}$`);
|
||||
this.#tokenRefPattern = new RegExp(`^${escapePattern(tokenRefPrefix)}_[A-F0-9]{24}$`);
|
||||
}
|
||||
|
|
@ -127,6 +137,8 @@ export class TokenBotConfigStore {
|
|||
tokenRefPrefix: this.#tokenRefPrefix,
|
||||
});
|
||||
if (derived.botId !== botId || derived.tokenRef !== tokenRef) return null;
|
||||
const extension = this.#normalizeBotExtension(value);
|
||||
if (!extension || typeof extension !== 'object' || Array.isArray(extension)) return null;
|
||||
return Object.freeze({
|
||||
botId,
|
||||
platformId,
|
||||
|
|
@ -135,6 +147,7 @@ export class TokenBotConfigStore {
|
|||
username: cleanString(value.username),
|
||||
createdAt: cleanString(value.createdAt) ?? new Date().toISOString(),
|
||||
connectedAt: cleanString(value.connectedAt),
|
||||
...extension,
|
||||
});
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -9,6 +9,58 @@ const IDENTITY_OPTIONS = Object.freeze({
|
|||
tokenRefPrefix: 'DSH_TELEGRAM_BOT_TOKEN',
|
||||
});
|
||||
|
||||
export const TELEGRAM_ACCESS_MODES = Object.freeze({
|
||||
compatible: 'compatible',
|
||||
privateAllowlist: 'private-allowlist',
|
||||
});
|
||||
|
||||
const TELEGRAM_USER_ID = /^[1-9]\d{0,15}$/;
|
||||
|
||||
export function normalizeTelegramAllowedUsers(value) {
|
||||
if (value === undefined) return Object.freeze([]);
|
||||
if (!Array.isArray(value)) {
|
||||
throw new TypeError('allowedUsers must be an array of numeric Telegram User IDs');
|
||||
}
|
||||
const normalized = value.map((entry) => {
|
||||
const userId = typeof entry === 'number' && Number.isSafeInteger(entry)
|
||||
? String(entry) : typeof entry === 'string' ? entry.trim() : '';
|
||||
if (!TELEGRAM_USER_ID.test(userId)) {
|
||||
throw new TypeError('allowedUsers contains an invalid Telegram User ID');
|
||||
}
|
||||
return userId;
|
||||
});
|
||||
return Object.freeze([...new Set(normalized)]);
|
||||
}
|
||||
|
||||
export function normalizeTelegramAccessPolicy(value = {}) {
|
||||
if (!value || typeof value !== 'object' || Array.isArray(value)) {
|
||||
throw new TypeError('Telegram access policy must be an object');
|
||||
}
|
||||
const accessMode = value.accessMode ?? TELEGRAM_ACCESS_MODES.compatible;
|
||||
if (!Object.values(TELEGRAM_ACCESS_MODES).includes(accessMode)) {
|
||||
throw new TypeError('Telegram accessMode must be compatible or private-allowlist');
|
||||
}
|
||||
return Object.freeze({
|
||||
accessMode,
|
||||
allowedUsers: normalizeTelegramAllowedUsers(value.allowedUsers),
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeTelegramBotExtension(value) {
|
||||
const hasAccessMode = Object.hasOwn(value, 'accessMode');
|
||||
const hasAllowedUsers = Object.hasOwn(value, 'allowedUsers');
|
||||
if (!hasAccessMode && !hasAllowedUsers) return {};
|
||||
try {
|
||||
const policy = normalizeTelegramAccessPolicy(value);
|
||||
return {
|
||||
...(hasAccessMode ? { accessMode: policy.accessMode } : {}),
|
||||
...(hasAllowedUsers || hasAccessMode ? { allowedUsers: policy.allowedUsers } : {}),
|
||||
};
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export function deriveTelegramBotIdentity(platformId) {
|
||||
return deriveTokenBotIdentity(platformId, IDENTITY_OPTIONS);
|
||||
}
|
||||
|
|
@ -19,6 +71,15 @@ export function maskTelegramBotId(platformId) {
|
|||
|
||||
export class TelegramConfigStore extends TokenBotConfigStore {
|
||||
constructor(path) {
|
||||
super(path, { channel: 'Telegram', ...IDENTITY_OPTIONS });
|
||||
super(path, {
|
||||
channel: 'Telegram',
|
||||
...IDENTITY_OPTIONS,
|
||||
normalizeBotExtension: normalizeTelegramBotExtension,
|
||||
});
|
||||
}
|
||||
|
||||
async save(value) {
|
||||
const previous = value?.platformId ? this.getByPlatformId(String(value.platformId)) : null;
|
||||
return super.save({ ...previous, ...value });
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,9 +1,15 @@
|
|||
import { TokenBotController } from '../shared/token-bot-controller.mjs';
|
||||
import { deriveTelegramBotIdentity, maskTelegramBotId } from './config-store.mjs';
|
||||
import {
|
||||
deriveTelegramBotIdentity,
|
||||
maskTelegramBotId,
|
||||
normalizeTelegramAccessPolicy,
|
||||
} from './config-store.mjs';
|
||||
import { inspectTelegramToken } from './telegram-api.mjs';
|
||||
import { TELEGRAM_DESCRIPTOR } from './telegram-bridge.mjs';
|
||||
|
||||
export class TelegramController extends TokenBotController {
|
||||
#configStore;
|
||||
|
||||
constructor(options) {
|
||||
super({
|
||||
...options,
|
||||
|
|
@ -12,5 +18,23 @@ export class TelegramController extends TokenBotController {
|
|||
deriveIdentity: deriveTelegramBotIdentity,
|
||||
maskPlatformId: maskTelegramBotId,
|
||||
});
|
||||
this.#configStore = options.configStore;
|
||||
}
|
||||
|
||||
status() {
|
||||
const snapshot = super.status();
|
||||
return {
|
||||
...snapshot,
|
||||
bots: snapshot.bots.map((bot) => {
|
||||
const config = this.#configStore.get(bot.botId);
|
||||
const accessPolicy = normalizeTelegramAccessPolicy(config ?? {});
|
||||
return { ...bot, accessPolicy };
|
||||
}),
|
||||
};
|
||||
}
|
||||
|
||||
async setAccessPolicy(botId, value) {
|
||||
const accessPolicy = normalizeTelegramAccessPolicy(value);
|
||||
return this.updateBotConfig(botId, (config) => ({ ...config, ...accessPolicy }));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,6 +1,10 @@
|
|||
import { createEditableMessageStream, splitMessageText } from '../shared/editable-message-stream.mjs';
|
||||
import { TelegramApi } from './telegram-api.mjs';
|
||||
import { createTelegramBridgeStatus, TelegramHarnessBridge } from './telegram-bridge.mjs';
|
||||
import {
|
||||
TELEGRAM_ACCESS_MODES,
|
||||
normalizeTelegramAccessPolicy,
|
||||
} from './config-store.mjs';
|
||||
|
||||
function escaped(value) {
|
||||
return value.replace(/[.*+?^${}()|[\]\\]/g, '\\$&');
|
||||
|
|
@ -115,7 +119,11 @@ export function normalizeTelegramUpdate(update, { botId, username, loadFile = as
|
|||
};
|
||||
}
|
||||
|
||||
export function telegramInboundAllowed(message, allowedPrivateUserIds) {
|
||||
export function telegramInboundAllowed(message, {
|
||||
accessMode = TELEGRAM_ACCESS_MODES.compatible,
|
||||
allowedPrivateUserIds = new Set(),
|
||||
} = {}) {
|
||||
if (accessMode !== TELEGRAM_ACCESS_MODES.privateAllowlist) return true;
|
||||
return message?.kind === 'direct'
|
||||
&& allowedPrivateUserIds instanceof Set
|
||||
&& allowedPrivateUserIds.has(String(message.senderId));
|
||||
|
|
@ -204,6 +212,7 @@ export class TelegramRuntime {
|
|||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#createApi;
|
||||
#accessMode;
|
||||
#allowedPrivateUserIds;
|
||||
#status = createTelegramRuntimeStatus();
|
||||
#api = null;
|
||||
|
|
@ -220,7 +229,6 @@ export class TelegramRuntime {
|
|||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
createApi = (options) => new TelegramApi(options),
|
||||
allowedPrivateUserIds = [],
|
||||
}) {
|
||||
if (!config || !token || !harness || !state) {
|
||||
throw new TypeError('TelegramRuntime requires config, token, Harness, and state');
|
||||
|
|
@ -232,9 +240,9 @@ export class TelegramRuntime {
|
|||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#createApi = createApi;
|
||||
this.#allowedPrivateUserIds = new Set(
|
||||
Array.isArray(allowedPrivateUserIds) ? allowedPrivateUserIds.map(String) : [],
|
||||
);
|
||||
const accessPolicy = normalizeTelegramAccessPolicy(config);
|
||||
this.#accessMode = accessPolicy.accessMode;
|
||||
this.#allowedPrivateUserIds = new Set(accessPolicy.allowedUsers);
|
||||
}
|
||||
|
||||
get status() {
|
||||
|
|
@ -334,7 +342,10 @@ export class TelegramRuntime {
|
|||
username: this.#config.username,
|
||||
loadFile: (fileId, options) => this.#api.downloadFile({ fileId, ...options }),
|
||||
});
|
||||
if (message && telegramInboundAllowed(message, this.#allowedPrivateUserIds)) {
|
||||
if (message && telegramInboundAllowed(message, {
|
||||
accessMode: this.#accessMode,
|
||||
allowedPrivateUserIds: this.#allowedPrivateUserIds,
|
||||
})) {
|
||||
void this.#bridge.accept(message).catch((error) => {
|
||||
if (signal.aborted) return;
|
||||
this.#logger.error?.(
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue