mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 20:55:59 +08:00
feat: secure Telegram access and connect AI Office
This commit is contained in:
parent
e3ae772106
commit
b0606159e0
30 changed files with 2045 additions and 302 deletions
107
src/channels/office/config-store.mjs
Normal file
107
src/channels/office/config-store.mjs
Normal file
|
|
@ -0,0 +1,107 @@
|
|||
import { mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname, isAbsolute } from 'node:path';
|
||||
|
||||
import { normalizeOfficeBaseUrl, OFFICE_PROTOCOL_VERSION } from './protocol.mjs';
|
||||
|
||||
const ALIAS = /^[a-z][a-z0-9-]{1,63}$/;
|
||||
const DEVICE_ID = /^[a-z0-9][a-z0-9_-]{2,63}$/i;
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function normalizeMap(value, { kind }) {
|
||||
if (!value || typeof value !== 'object' || Array.isArray(value)) return null;
|
||||
const output = {};
|
||||
for (const [rawId, rawValue] of Object.entries(value)) {
|
||||
const id = cleanString(rawId);
|
||||
const item = cleanString(rawValue);
|
||||
if (!id || !ALIAS.test(id) || !item) return null;
|
||||
if (kind === 'workspace' && !isAbsolute(item)) return null;
|
||||
if (kind === 'preset' && item.length > 8_000) return null;
|
||||
output[id] = item;
|
||||
}
|
||||
return Object.freeze(output);
|
||||
}
|
||||
|
||||
export function normalizeOfficeConfig(value) {
|
||||
if (!value || typeof value !== 'object' || Array.isArray(value)) return null;
|
||||
let baseUrl;
|
||||
try {
|
||||
baseUrl = normalizeOfficeBaseUrl(value.baseUrl).origin;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
const deviceId = cleanString(value.deviceId);
|
||||
const deviceTokenRef = cleanString(value.deviceTokenRef);
|
||||
const workspaces = normalizeMap(value.workspaces ?? {}, { kind: 'workspace' });
|
||||
const instructionPresets = normalizeMap(value.instructionPresets ?? {}, { kind: 'preset' });
|
||||
const maxConcurrency = Number(value.maxConcurrency ?? 1);
|
||||
const heartbeatSeconds = Number(value.heartbeatSeconds ?? 30);
|
||||
if (!deviceId || !DEVICE_ID.test(deviceId) || !deviceTokenRef
|
||||
|| !/^DSH_OFFICE_DEVICE_TOKEN_[A-F0-9]{24}$/.test(deviceTokenRef)
|
||||
|| !workspaces || !instructionPresets
|
||||
|| !Number.isInteger(maxConcurrency) || maxConcurrency < 1 || maxConcurrency > 4
|
||||
|| !Number.isInteger(heartbeatSeconds) || heartbeatSeconds < 10 || heartbeatSeconds > 300) {
|
||||
return null;
|
||||
}
|
||||
return Object.freeze({
|
||||
version: 1,
|
||||
protocolVersion: OFFICE_PROTOCOL_VERSION,
|
||||
baseUrl,
|
||||
deviceId,
|
||||
deviceTokenRef,
|
||||
maxConcurrency,
|
||||
heartbeatSeconds,
|
||||
workspaces,
|
||||
instructionPresets,
|
||||
createdAt: cleanString(value.createdAt) ?? new Date().toISOString(),
|
||||
updatedAt: cleanString(value.updatedAt) ?? new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
|
||||
export class OfficeConfigStore {
|
||||
#path;
|
||||
#value = null;
|
||||
#queue = Promise.resolve();
|
||||
|
||||
constructor(path) { this.#path = path; }
|
||||
|
||||
async load() {
|
||||
try {
|
||||
const normalized = normalizeOfficeConfig(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
if (!normalized) throw new Error('dsh-im AI Office config is invalid');
|
||||
this.#value = normalized;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#value = null;
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
get() { return this.#value ? structuredClone(this.#value) : null; }
|
||||
|
||||
async save(value) {
|
||||
const normalized = normalizeOfficeConfig(value);
|
||||
if (!normalized) throw new Error('Refusing to persist invalid AI Office configuration');
|
||||
const operation = this.#queue.then(async () => {
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(normalized, null, 2)}\n`, { mode: 0o600 });
|
||||
await rename(temporary, this.#path);
|
||||
this.#value = normalized;
|
||||
});
|
||||
this.#queue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
return this.get();
|
||||
}
|
||||
|
||||
async clear() {
|
||||
const operation = this.#queue.then(async () => {
|
||||
try { await unlink(this.#path); } catch (error) { if (error?.code !== 'ENOENT') throw error; }
|
||||
this.#value = null;
|
||||
});
|
||||
this.#queue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
}
|
||||
}
|
||||
167
src/channels/office/office-controller.mjs
Normal file
167
src/channels/office/office-controller.mjs
Normal file
|
|
@ -0,0 +1,167 @@
|
|||
import { createHash } from 'node:crypto';
|
||||
|
||||
import { normalizeOfficeBaseUrl, officeHookUrls } from './protocol.mjs';
|
||||
import { normalizeOfficeConfig } from './config-store.mjs';
|
||||
import { OfficeRuntime } from './office-runtime.mjs';
|
||||
|
||||
function clean(value) { return typeof value === 'string' && value.trim() ? value.trim() : null; }
|
||||
|
||||
export function officeTokenRef(baseUrl, deviceId) {
|
||||
const digest = createHash('sha256').update(`${baseUrl}\n${deviceId}`).digest('hex').slice(0, 24).toUpperCase();
|
||||
return `DSH_OFFICE_DEVICE_TOKEN_${digest}`;
|
||||
}
|
||||
|
||||
export class OfficeController {
|
||||
#credentials;
|
||||
#store;
|
||||
#logger;
|
||||
#createRuntime;
|
||||
#runtime = null;
|
||||
#transition = Promise.resolve();
|
||||
|
||||
constructor({ credentials, configStore, logger = console, createRuntime }) {
|
||||
if (!credentials?.resolve || !credentials?.set || !credentials?.unset) {
|
||||
throw new TypeError('AI Office requires the Harness credential provider');
|
||||
}
|
||||
if (!configStore?.get || !configStore?.save || !configStore?.clear) {
|
||||
throw new TypeError('AI Office requires a config store');
|
||||
}
|
||||
this.#credentials = credentials;
|
||||
this.#store = configStore;
|
||||
this.#logger = logger;
|
||||
this.#createRuntime = createRuntime ?? ((options) => new OfficeRuntime(options));
|
||||
}
|
||||
|
||||
async initialize() {
|
||||
const config = this.#store.get();
|
||||
if (config) {
|
||||
const credential = await this.#credentials.resolve(config.deviceTokenRef).catch(() => undefined);
|
||||
if (credential?.value) await this.#start(config, credential.value);
|
||||
}
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async configure(input = {}) {
|
||||
return this.#serial(async () => {
|
||||
const previous = this.#store.get();
|
||||
const requestedBaseUrl = clean(input.baseUrl);
|
||||
const deviceId = clean(input.deviceId);
|
||||
if (!requestedBaseUrl || !deviceId) throw new TypeError('Office URL and Device ID are required');
|
||||
const baseUrl = normalizeOfficeBaseUrl(requestedBaseUrl).origin;
|
||||
const tokenRef = officeTokenRef(baseUrl, deviceId);
|
||||
const suppliedToken = clean(input.deviceToken);
|
||||
const priorCredential = await this.#credentials.resolve(tokenRef).catch(() => undefined);
|
||||
const token = suppliedToken ?? priorCredential?.value;
|
||||
if (!token || token.length < 32) throw new TypeError('Device Token must contain at least 32 characters');
|
||||
const now = new Date().toISOString();
|
||||
const config = normalizeOfficeConfig({
|
||||
version: 1,
|
||||
baseUrl,
|
||||
deviceId,
|
||||
deviceTokenRef: tokenRef,
|
||||
maxConcurrency: input.maxConcurrency,
|
||||
heartbeatSeconds: input.heartbeatSeconds,
|
||||
workspaces: input.workspaces,
|
||||
instructionPresets: input.instructionPresets,
|
||||
createdAt: previous?.createdAt ?? now,
|
||||
updatedAt: now,
|
||||
});
|
||||
if (!config) throw new TypeError('AI Office connector configuration is invalid');
|
||||
|
||||
await this.#credentials.set(tokenRef, token);
|
||||
try {
|
||||
await this.#store.save(config);
|
||||
} catch (error) {
|
||||
if (priorCredential?.value) await this.#credentials.set(tokenRef, priorCredential.value).catch(() => undefined);
|
||||
else await this.#credentials.unset(tokenRef).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
if (previous?.deviceTokenRef && previous.deviceTokenRef !== tokenRef) {
|
||||
await this.#credentials.unset(previous.deviceTokenRef).catch(() => undefined);
|
||||
}
|
||||
await this.#start(config, token);
|
||||
return this.status();
|
||||
});
|
||||
}
|
||||
|
||||
async reconnect() {
|
||||
return this.#serial(async () => {
|
||||
const config = this.#store.get();
|
||||
if (!config) throw new Error('AI Office connector is not configured');
|
||||
await this.#start(config);
|
||||
return this.status();
|
||||
});
|
||||
}
|
||||
|
||||
async test() {
|
||||
const config = this.#store.get();
|
||||
if (!config) throw new Error('AI Office connector is not configured');
|
||||
const token = await this.#resolveToken(config);
|
||||
const runtime = this.#createRuntime({ config, token, logger: this.#logger });
|
||||
try { await runtime.testConnection(AbortSignal.timeout(10_000)); } finally { await runtime.stop().catch(() => undefined); }
|
||||
return { tested: true, snapshot: await this.status() };
|
||||
}
|
||||
|
||||
async remove() {
|
||||
return this.#serial(async () => {
|
||||
const config = this.#store.get();
|
||||
await this.#stop();
|
||||
if (config?.deviceTokenRef) await this.#credentials.unset(config.deviceTokenRef).catch(() => undefined);
|
||||
await this.#store.clear();
|
||||
return this.status();
|
||||
});
|
||||
}
|
||||
|
||||
async status() {
|
||||
const config = this.#store.get();
|
||||
if (!config) return { schemaVersion: 1, configured: false, connected: false, state: 'unconfigured' };
|
||||
const credential = await this.#credentials.resolve(config.deviceTokenRef).catch(() => undefined);
|
||||
const runtime = this.#runtime?.status ?? null;
|
||||
return {
|
||||
schemaVersion: 1,
|
||||
configured: true,
|
||||
connected: runtime?.connected === true,
|
||||
state: runtime?.state ?? (credential?.value ? 'idle' : 'missing-token'),
|
||||
config: {
|
||||
protocolVersion: config.protocolVersion,
|
||||
baseUrl: config.baseUrl,
|
||||
deviceId: config.deviceId,
|
||||
maxConcurrency: config.maxConcurrency,
|
||||
heartbeatSeconds: config.heartbeatSeconds,
|
||||
workspaces: config.workspaces,
|
||||
instructionPresets: config.instructionPresets,
|
||||
hooks: officeHookUrls(config.baseUrl),
|
||||
},
|
||||
health: runtime,
|
||||
tokenConfigured: Boolean(credential?.value),
|
||||
};
|
||||
}
|
||||
|
||||
async close() { await this.#transition.catch(() => undefined); await this.#stop(); }
|
||||
|
||||
async #resolveToken(config) {
|
||||
const credential = await this.#credentials.resolve(config.deviceTokenRef).catch(() => undefined);
|
||||
if (!credential?.value) throw new Error('AI Office Device Token is missing');
|
||||
return credential.value;
|
||||
}
|
||||
|
||||
async #start(config, knownToken) {
|
||||
await this.#stop();
|
||||
const token = knownToken ?? await this.#resolveToken(config);
|
||||
const runtime = this.#createRuntime({ config, token, logger: this.#logger });
|
||||
this.#runtime = runtime;
|
||||
runtime.start();
|
||||
}
|
||||
|
||||
async #stop() {
|
||||
const runtime = this.#runtime;
|
||||
this.#runtime = null;
|
||||
if (runtime) await runtime.stop();
|
||||
}
|
||||
|
||||
#serial(operation) {
|
||||
const run = this.#transition.then(operation, operation);
|
||||
this.#transition = run.then(() => undefined, () => undefined);
|
||||
return run;
|
||||
}
|
||||
}
|
||||
132
src/channels/office/office-runtime.mjs
Normal file
132
src/channels/office/office-runtime.mjs
Normal file
|
|
@ -0,0 +1,132 @@
|
|||
import { setTimeout as sleep } from 'node:timers/promises';
|
||||
|
||||
import { OfficeTransport } from './office-transport.mjs';
|
||||
import { OFFICE_PROTOCOL_VERSION } from './protocol.mjs';
|
||||
|
||||
const RETRY_DELAYS = Object.freeze([1_000, 3_000, 10_000, 30_000]);
|
||||
|
||||
function safeConnectionError(error) {
|
||||
const code = typeof error?.code === 'string' ? error.code : 'office-connection-failed';
|
||||
const messages = {
|
||||
'invalid-device-token': 'AI Office 拒绝了 Device Token。',
|
||||
'office-hook-unavailable': 'AI Office Connector Hook 尚未就绪。',
|
||||
'office-protocol-mismatch': 'AI Office Connector 协议版本不兼容。',
|
||||
'office-transport-failed': '本机暂时无法访问 AI Office。',
|
||||
};
|
||||
return { code, message: messages[code] ?? 'AI Office 连接已中断。' };
|
||||
}
|
||||
|
||||
export class OfficeRuntime {
|
||||
#config;
|
||||
#token;
|
||||
#logger;
|
||||
#transport;
|
||||
#controller = null;
|
||||
#task = null;
|
||||
#status;
|
||||
|
||||
constructor({ config, token, logger = console, transport }) {
|
||||
this.#config = config;
|
||||
this.#token = token;
|
||||
this.#logger = logger;
|
||||
this.#transport = transport ?? new OfficeTransport({
|
||||
baseUrl: config.baseUrl, deviceId: config.deviceId, token,
|
||||
});
|
||||
this.#status = {
|
||||
state: 'idle', connected: false, startedAt: null, lastHeartbeatAt: null,
|
||||
lastEventAt: null, lastEventId: null, lastEventType: null, reconnects: 0,
|
||||
jobsOffered: 0, error: null,
|
||||
};
|
||||
}
|
||||
|
||||
get status() { return structuredClone(this.#status); }
|
||||
|
||||
capabilities() {
|
||||
return {
|
||||
protocolVersion: OFFICE_PROTOCOL_VERSION,
|
||||
deviceId: this.#config.deviceId,
|
||||
workspaces: Object.keys(this.#config.workspaces),
|
||||
instructionPresets: Object.keys(this.#config.instructionPresets),
|
||||
maxConcurrency: this.#config.maxConcurrency,
|
||||
};
|
||||
}
|
||||
|
||||
async testConnection(signal) {
|
||||
await this.#transport.heartbeat({ ...this.capabilities(), probe: true }, { signal });
|
||||
return { ok: true };
|
||||
}
|
||||
|
||||
start() {
|
||||
if (this.#task) return this.status;
|
||||
this.#controller = new AbortController();
|
||||
this.#status.startedAt = new Date().toISOString();
|
||||
this.#status.state = 'connecting';
|
||||
this.#task = this.#run(this.#controller.signal).finally(() => { this.#task = null; });
|
||||
this.#task.catch((error) => {
|
||||
if (this.#controller?.signal.aborted) return;
|
||||
this.#logger.error?.('[dsh-im:office] connector stopped:', error);
|
||||
});
|
||||
return this.status;
|
||||
}
|
||||
|
||||
async #run(signal) {
|
||||
let attempt = 0;
|
||||
while (!signal.aborted) {
|
||||
const attemptController = new AbortController();
|
||||
const attemptSignal = AbortSignal.any([signal, attemptController.signal]);
|
||||
try {
|
||||
await this.#transport.heartbeat(this.capabilities(), { signal: attemptSignal });
|
||||
this.#status.lastHeartbeatAt = new Date().toISOString();
|
||||
const heartbeat = this.#heartbeatLoop(attemptSignal);
|
||||
const stream = this.#transport.stream({
|
||||
signal: attemptSignal,
|
||||
lastEventId: this.#status.lastEventId,
|
||||
onOpen: () => {
|
||||
this.#status.connected = true;
|
||||
this.#status.state = 'connected';
|
||||
this.#status.error = null;
|
||||
attempt = 0;
|
||||
},
|
||||
onEvent: async (event) => {
|
||||
this.#status.lastEventAt = new Date().toISOString();
|
||||
this.#status.lastEventId = event.id ?? this.#status.lastEventId;
|
||||
this.#status.lastEventType = event.type;
|
||||
if (event.type === 'job.available') this.#status.jobsOffered += 1;
|
||||
},
|
||||
});
|
||||
await Promise.race([stream, heartbeat]);
|
||||
} catch (error) {
|
||||
if (signal.aborted) break;
|
||||
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; }
|
||||
} finally {
|
||||
attemptController.abort();
|
||||
}
|
||||
}
|
||||
this.#status.connected = false;
|
||||
this.#status.state = 'idle';
|
||||
}
|
||||
|
||||
async #heartbeatLoop(signal) {
|
||||
while (!signal.aborted) {
|
||||
await sleep(this.#config.heartbeatSeconds * 1_000, undefined, { signal });
|
||||
await this.#transport.heartbeat(this.capabilities(), { signal });
|
||||
this.#status.lastHeartbeatAt = new Date().toISOString();
|
||||
}
|
||||
}
|
||||
|
||||
async stop() {
|
||||
const task = this.#task;
|
||||
this.#controller?.abort();
|
||||
this.#controller = null;
|
||||
if (task) await task.catch(() => undefined);
|
||||
this.#status.connected = false;
|
||||
this.#status.state = 'idle';
|
||||
return this.status;
|
||||
}
|
||||
}
|
||||
118
src/channels/office/office-transport.mjs
Normal file
118
src/channels/office/office-transport.mjs
Normal file
|
|
@ -0,0 +1,118 @@
|
|||
import { OFFICE_HOOK_PATHS, OFFICE_PROTOCOL_VERSION, officeHookUrls } from './protocol.mjs';
|
||||
|
||||
function safeTransportError(operation, response) {
|
||||
const error = new Error(`AI Office ${operation} failed: HTTP ${response.status}`);
|
||||
error.code = response.status === 401 ? 'invalid-device-token'
|
||||
: response.status === 404 ? 'office-hook-unavailable' : 'office-transport-failed';
|
||||
return error;
|
||||
}
|
||||
|
||||
function transportFailure(message, code = 'office-transport-failed', cause) {
|
||||
const error = new Error(message, cause === undefined ? undefined : { cause });
|
||||
error.code = code;
|
||||
return error;
|
||||
}
|
||||
|
||||
function isAbort(error, signal) {
|
||||
return signal?.aborted || error?.name === 'AbortError';
|
||||
}
|
||||
|
||||
function parseFrame(frame) {
|
||||
let type = 'message';
|
||||
let id;
|
||||
const data = [];
|
||||
for (const line of frame.split(/\r?\n/)) {
|
||||
if (line.startsWith('event:')) type = line.slice(6).trim() || 'message';
|
||||
else if (line.startsWith('id:')) id = line.slice(3).trim() || undefined;
|
||||
else if (line.startsWith('data:')) data.push(line.slice(5).trimStart());
|
||||
}
|
||||
if (data.length === 0) return null;
|
||||
let value;
|
||||
try { value = JSON.parse(data.join('\n')); } catch { throw new Error('AI Office SSE returned invalid JSON'); }
|
||||
return { id, type: typeof value?.type === 'string' ? value.type : type, data: value };
|
||||
}
|
||||
|
||||
export class OfficeTransport {
|
||||
#baseUrl;
|
||||
#deviceId;
|
||||
#token;
|
||||
#fetch;
|
||||
|
||||
constructor({ baseUrl, deviceId, token, fetchImpl = fetch }) {
|
||||
this.#baseUrl = baseUrl;
|
||||
this.#deviceId = deviceId;
|
||||
this.#token = token;
|
||||
this.#fetch = fetchImpl;
|
||||
}
|
||||
|
||||
hooks() { return officeHookUrls(this.#baseUrl); }
|
||||
|
||||
#headers(extra = {}) {
|
||||
return {
|
||||
authorization: `Bearer ${this.#token}`,
|
||||
'x-harness-device-id': this.#deviceId,
|
||||
...extra,
|
||||
};
|
||||
}
|
||||
|
||||
async heartbeat(payload, { signal } = {}) {
|
||||
let response;
|
||||
try {
|
||||
response = await this.#fetch(new URL(OFFICE_HOOK_PATHS.heartbeat, this.#baseUrl), {
|
||||
method: 'POST',
|
||||
headers: this.#headers({ accept: 'application/json', 'content-type': 'application/json' }),
|
||||
body: JSON.stringify(payload),
|
||||
signal,
|
||||
redirect: 'error',
|
||||
});
|
||||
} catch (error) {
|
||||
if (isAbort(error, signal)) throw error;
|
||||
throw transportFailure('AI Office heartbeat request could not be completed', undefined, error);
|
||||
}
|
||||
if (!response.ok) throw safeTransportError('heartbeat', response);
|
||||
let value;
|
||||
try { value = await response.json(); }
|
||||
catch (error) { throw transportFailure('AI Office heartbeat returned invalid JSON', 'office-protocol-mismatch', error); }
|
||||
if (!value || typeof value !== 'object' || value.ok !== true
|
||||
|| value.protocolVersion !== OFFICE_PROTOCOL_VERSION) {
|
||||
throw transportFailure('AI Office heartbeat protocol does not match', 'office-protocol-mismatch');
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
async stream({ signal, lastEventId, onOpen, onEvent }) {
|
||||
const headers = this.#headers({ accept: 'text/event-stream' });
|
||||
if (lastEventId) headers['last-event-id'] = lastEventId;
|
||||
let response;
|
||||
try {
|
||||
response = await this.#fetch(new URL(OFFICE_HOOK_PATHS.stream, this.#baseUrl), {
|
||||
method: 'GET', headers, signal, redirect: 'error', cache: 'no-store',
|
||||
});
|
||||
} catch (error) {
|
||||
if (isAbort(error, signal)) throw error;
|
||||
throw transportFailure('AI Office SSE request could not be completed', undefined, error);
|
||||
}
|
||||
if (!response.ok) throw safeTransportError('stream', response);
|
||||
if (!response.headers.get('content-type')?.toLowerCase().includes('text/event-stream')) {
|
||||
throw new Error('AI Office stream did not return text/event-stream');
|
||||
}
|
||||
if (!response.body) throw new Error('AI Office stream returned no body');
|
||||
onOpen?.();
|
||||
const reader = response.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = '';
|
||||
while (true) {
|
||||
const { value, done } = await reader.read();
|
||||
if (done) break;
|
||||
buffer = `${buffer}${decoder.decode(value, { stream: true })}`.replaceAll('\r\n', '\n');
|
||||
let boundary;
|
||||
while ((boundary = buffer.indexOf('\n\n')) !== -1) {
|
||||
const raw = buffer.slice(0, boundary);
|
||||
buffer = buffer.slice(boundary + 2);
|
||||
const event = parseFrame(raw);
|
||||
if (event) await onEvent?.(event);
|
||||
}
|
||||
}
|
||||
throw new Error('AI Office SSE stream ended');
|
||||
}
|
||||
}
|
||||
43
src/channels/office/protocol.mjs
Normal file
43
src/channels/office/protocol.mjs
Normal file
|
|
@ -0,0 +1,43 @@
|
|||
export const OFFICE_PROTOCOL_VERSION = 'office-harness.v1';
|
||||
export const OFFICE_RPC_CHANNEL = '/office';
|
||||
|
||||
export const OFFICE_RPC_ENDPOINTS = Object.freeze({
|
||||
status: 'connection.status',
|
||||
configure: 'connector.configure',
|
||||
reconnect: 'connector.reconnect',
|
||||
test: 'connector.test',
|
||||
remove: 'connector.remove',
|
||||
});
|
||||
|
||||
export const OFFICE_HOOK_PATHS = Object.freeze({
|
||||
stream: '/api/harness/connector/stream',
|
||||
heartbeat: '/api/harness/connector/heartbeat',
|
||||
job: '/api/harness/connector/jobs/:id',
|
||||
accept: '/api/harness/connector/jobs/:id/accept',
|
||||
renew: '/api/harness/connector/jobs/:id/renew',
|
||||
progress: '/api/harness/connector/jobs/:id/progress',
|
||||
approval: '/api/harness/connector/jobs/:id/approval',
|
||||
result: '/api/harness/connector/jobs/:id/result',
|
||||
fail: '/api/harness/connector/jobs/:id/fail',
|
||||
});
|
||||
|
||||
export function normalizeOfficeBaseUrl(value) {
|
||||
const url = new URL(typeof value === 'string' ? value.trim() : '');
|
||||
const localHttp = url.protocol === 'http:' && ['localhost', '127.0.0.1', '[::1]'].includes(url.hostname);
|
||||
if (url.protocol !== 'https:' && !localHttp) {
|
||||
throw new TypeError('AI Office URL must use HTTPS (HTTP is allowed only for loopback testing)');
|
||||
}
|
||||
if (url.username || url.password || url.search || url.hash) {
|
||||
throw new TypeError('AI Office URL must be a bare origin');
|
||||
}
|
||||
url.pathname = '/';
|
||||
return url;
|
||||
}
|
||||
|
||||
export function officeHookUrls(baseUrl) {
|
||||
const origin = normalizeOfficeBaseUrl(baseUrl);
|
||||
return Object.fromEntries(Object.entries(OFFICE_HOOK_PATHS).map(([name, path]) => [
|
||||
name,
|
||||
new URL(path, origin).toString(),
|
||||
]));
|
||||
}
|
||||
|
|
@ -115,6 +115,12 @@ export function normalizeTelegramUpdate(update, { botId, username, loadFile = as
|
|||
};
|
||||
}
|
||||
|
||||
export function telegramInboundAllowed(message, allowedPrivateUserIds) {
|
||||
return message?.kind === 'direct'
|
||||
&& allowedPrivateUserIds instanceof Set
|
||||
&& allowedPrivateUserIds.has(String(message.senderId));
|
||||
}
|
||||
|
||||
class TelegramBotClient {
|
||||
#api;
|
||||
#signal;
|
||||
|
|
@ -198,6 +204,7 @@ export class TelegramRuntime {
|
|||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#createApi;
|
||||
#allowedPrivateUserIds;
|
||||
#status = createTelegramRuntimeStatus();
|
||||
#api = null;
|
||||
#bridge = null;
|
||||
|
|
@ -213,6 +220,7 @@ 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');
|
||||
|
|
@ -224,6 +232,9 @@ export class TelegramRuntime {
|
|||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#createApi = createApi;
|
||||
this.#allowedPrivateUserIds = new Set(
|
||||
Array.isArray(allowedPrivateUserIds) ? allowedPrivateUserIds.map(String) : [],
|
||||
);
|
||||
}
|
||||
|
||||
get status() {
|
||||
|
|
@ -323,7 +334,7 @@ export class TelegramRuntime {
|
|||
username: this.#config.username,
|
||||
loadFile: (fileId, options) => this.#api.downloadFile({ fileId, ...options }),
|
||||
});
|
||||
if (message) {
|
||||
if (message && telegramInboundAllowed(message, this.#allowedPrivateUserIds)) {
|
||||
void this.#bridge.accept(message).catch((error) => {
|
||||
if (signal.aborted) return;
|
||||
this.#logger.error?.(
|
||||
|
|
@ -331,6 +342,9 @@ export class TelegramRuntime {
|
|||
error,
|
||||
);
|
||||
});
|
||||
} else if (message) {
|
||||
this.#status.messagesRejected += 1;
|
||||
this.#status.lastRejectedAt = new Date().toISOString();
|
||||
}
|
||||
cursor = update.update_id + 1;
|
||||
await this.#state.setCursor(cursor);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue