Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions packages/feishu-meeting-intake/manifest.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"pluginId": "official.feishu-meeting-intake",
"version": "0.1.0-alpha.0",
"version": "0.1.0-alpha.1",
"contractVersion": "0.1.0",
"name": "Feishu Meeting Intake",
"features": [
Expand Down Expand Up @@ -28,6 +28,7 @@
]
},
"runtime": {
"transport": "builtin"
"transport": "stdio",
"entrypoint": "dist/entrypoint.js"
}
}
5 changes: 3 additions & 2 deletions packages/feishu-meeting-intake/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@clowder-ai/feishu-meeting-intake",
"version": "0.1.0-alpha.0",
"version": "0.1.0-alpha.1",
"description": "Official Feishu/Lark generated-meeting-artifact intake adapter for Clowder AI hosts",
"license": "MIT",
"type": "module",
Expand Down Expand Up @@ -32,7 +32,8 @@
},
"dependencies": {
"@clowder-ai/plugin-contract": "0.1.0-beta.9",
"@clowder-ai/plugin-sdk": "0.1.0-beta.5"
"@clowder-ai/plugin-sdk": "0.1.0-beta.5",
"@larksuite/cli": "1.0.85"
},
"devDependencies": {
"@types/node": "^20.17.0",
Expand Down
15 changes: 14 additions & 1 deletion packages/feishu-meeting-intake/src/artifact.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,26 @@ const DESCRIPTOR = {
meetingId: 'meeting-42',
} as const;

test('ships a contract-valid builtin manifest with one declared signal', async () => {
test('ships a contract-valid stdio manifest with one declared signal', async () => {
const manifest: unknown = JSON.parse(
await readFile(new URL('../manifest.json', import.meta.url), 'utf8'),
);
const packageMetadata: unknown = JSON.parse(
await readFile(new URL('../package.json', import.meta.url), 'utf8'),
);
const result = validateManifest(manifest);
assert.equal(result.valid, true, result.valid ? undefined : JSON.stringify(result.errors));
if (result.valid) {
assert.equal(result.manifest.version, '0.1.0-alpha.1');
assert.equal(
(packageMetadata as { readonly version?: unknown }).version,
result.manifest.version,
'manifest and immutable npm artifact versions must match',
);
assert.deepEqual(result.manifest.runtime, {
transport: 'stdio',
entrypoint: 'dist/entrypoint.js',
});
assert.deepEqual(result.manifest.signals?.provides, [
{
type: 'feishu.meeting_artifact.generated.v1',
Expand Down
3 changes: 3 additions & 0 deletions packages/feishu-meeting-intake/src/entrypoint.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
import { runFeishuMeetingIntakeEntrypoint } from './stdio-entrypoint.js';

runFeishuMeetingIntakeEntrypoint();
19 changes: 19 additions & 0 deletions packages/feishu-meeting-intake/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,3 +37,22 @@ export {
type FeishuMeetingIntakeCycleResult,
type FeishuMeetingIntakeRuntime,
} from './runtime.js';

export {
createLarkCliFeishuEventGateway,
larkCliChildEnvironment,
resolveBundledLarkCliEntrypoint,
type LarkCliEventConsumer,
type LarkCliFeishuEventGateway,
type LarkCliFeishuEventGatewayOptions,
} from './lark-cli-gateway.js';

export {
meetingIntakeStatePath,
readRuntimeClaims,
runFeishuMeetingIntakeEntrypoint,
startFeishuMeetingIntakeStdio,
type FeishuMeetingIntakeStdioController,
type FeishuMeetingIntakeStdioOptions,
type FeishuStdioRuntimeContext,
} from './stdio-entrypoint.js';
198 changes: 198 additions & 0 deletions packages/feishu-meeting-intake/src/lark-cli-consumer.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,198 @@
import { spawn, type ChildProcessWithoutNullStreams } from 'node:child_process';

import { FeishuGatewayError } from './gateway.js';
import {
larkCliChildEnvironment,
resolveBundledLarkCliEntrypoint,
} from './lark-cli-runner.js';

export const LARK_EVENT_KEYS = [
'minutes.minute.generated_v1',
'vc.note.generated_v1',
] as const;

const MAX_EVENT_LINE_BYTES = 256 * 1024;
const MAX_DIAGNOSTIC_BYTES = 32 * 1024;

export interface LarkCliEventConsumer {
readonly events: AsyncIterable<unknown>;
close(): Promise<void>;
}

function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value);
}

class AsyncEventQueue implements AsyncIterable<unknown> {
private readonly values: unknown[] = [];
private readonly waiters: Array<{
readonly resolve: (result: IteratorResult<unknown>) => void;
readonly reject: (error: unknown) => void;
}> = [];
private ended = false;
private failure: unknown;

push(value: unknown): void {
if (this.ended) return;
const waiter = this.waiters.shift();
if (waiter !== undefined) waiter.resolve({ done: false, value });
else this.values.push(value);
}

fail(error: unknown): void {
if (this.ended) return;
this.ended = true;
this.failure = error;
for (const waiter of this.waiters.splice(0)) waiter.reject(error);
}

end(): void {
if (this.ended) return;
this.ended = true;
for (const waiter of this.waiters.splice(0)) waiter.resolve({ done: true, value: undefined });
}

[Symbol.asyncIterator](): AsyncIterator<unknown> {
return {
next: (): Promise<IteratorResult<unknown>> => {
const value = this.values.shift();
if (value !== undefined) return Promise.resolve({ done: false, value });
if (this.failure !== undefined) return Promise.reject(this.failure);
if (this.ended) return Promise.resolve({ done: true, value: undefined });
return new Promise((resolve, reject) => this.waiters.push({ resolve, reject }));
},
};
}
}

function mappedCliFailure(detail: string): FeishuGatewayError {
let errorType = '';
let errorSubtype = '';
for (const line of detail.split('\n')) {
try {
const candidate: unknown = JSON.parse(line);
if (isRecord(candidate) && isRecord(candidate.error)) {
errorType = String(candidate.error.type ?? '');
errorSubtype = String(candidate.error.subtype ?? '');
}
} catch {
// Readiness and exit diagnostics are intentionally non-JSON.
}
}
const classification = `${errorType}:${errorSubtype}`.toLowerCase();
if (/auth|login|token|not_configured/u.test(classification)) {
return new FeishuGatewayError('AUTH_EXPIRED', 'lark-cli user authorization is unavailable');
}
if (/permission|scope|forbidden/u.test(classification)) {
return new FeishuGatewayError('PERMISSION_DENIED', 'lark-cli lacks generated-event permission');
}
if (/rate|429/u.test(classification)) {
return new FeishuGatewayError('RATE_LIMITED', 'lark-cli generated-event source is rate limited');
}
return new FeishuGatewayError('UNAVAILABLE', 'lark-cli generated-event source stopped');
}

function startBoundedJsonLines(
child: ChildProcessWithoutNullStreams,
queue: AsyncEventQueue,
fail: (error: unknown) => void,
): void {
let pending = Buffer.alloc(0);
child.stdout.on('data', (chunk: Buffer) => {
pending = Buffer.concat([pending, chunk]);
if (pending.byteLength > MAX_EVENT_LINE_BYTES && !pending.includes(0x0a)) {
fail(new FeishuGatewayError('UNAVAILABLE', 'lark-cli emitted an oversized event'));
return;
}
while (true) {
const newline = pending.indexOf(0x0a);
if (newline < 0) break;
const line = pending.subarray(0, newline);
pending = pending.subarray(newline + 1);
if (line.byteLength === 0) continue;
if (line.byteLength > MAX_EVENT_LINE_BYTES) {
fail(new FeishuGatewayError('UNAVAILABLE', 'lark-cli emitted an oversized event'));
return;
}
try {
queue.push(JSON.parse(line.toString('utf8')) as unknown);
} catch {
fail(new FeishuGatewayError('UNAVAILABLE', 'lark-cli emitted invalid event JSON'));
return;
}
}
});
}

export async function startDefaultLarkCliConsumer(
eventKey: typeof LARK_EVENT_KEYS[number],
signal: AbortSignal,
homeDirectory: string,
): Promise<LarkCliEventConsumer> {
const child = spawn(
process.execPath,
[resolveBundledLarkCliEntrypoint(), 'event', 'consume', eventKey, '--as', 'user'],
{
cwd: process.cwd(),
env: larkCliChildEnvironment(homeDirectory),
shell: false,
stdio: ['pipe', 'pipe', 'pipe'],
},
);
const queue = new AsyncEventQueue();
let diagnostic = '';
let closed = false;
let readySeen = false;
let readyResolve!: () => void;
let readyReject!: (error: unknown) => void;
const ready = new Promise<void>((resolve, reject) => {
readyResolve = resolve;
readyReject = reject;
});
const fail = (error: unknown): void => {
queue.fail(error);
readyReject(error);
child.stdin.end();
};
startBoundedJsonLines(child, queue, fail);
child.stderr.on('data', (chunk: Buffer) => {
diagnostic = `${diagnostic}${chunk.toString('utf8')}`.slice(-MAX_DIAGNOSTIC_BYTES);
if (diagnostic.includes(`[event] ready event_key=${eventKey}`)) {
readySeen = true;
readyResolve();
}
});
child.once('error', fail);
child.once('exit', (code) => {
closed = true;
if (!readySeen) fail(mappedCliFailure(diagnostic));
else if (code === 0) queue.end();
else fail(mappedCliFailure(diagnostic));
});
const onAbort = (): void => {
if (!readySeen) readyReject(signal.reason);
child.stdin.end();
};
signal.addEventListener('abort', onAbort, { once: true });
try {
await ready;
} catch (error) {
signal.removeEventListener('abort', onAbort);
throw error;
}
return {
events: queue,
close: async () => {
signal.removeEventListener('abort', onAbort);
if (closed) return;
child.stdin.end();
await new Promise<void>(resolve => {
const timer = setTimeout(() => child.kill('SIGTERM'), 2_000);
child.once('exit', () => {
clearTimeout(timer);
resolve();
});
});
},
};
}
Loading
Loading