feat: リアルタイム議事録システム初期リリース
Made-with: Cursor
This commit is contained in:
@@ -0,0 +1,155 @@
|
||||
import fs from "node:fs";
|
||||
import { parse as parseYaml } from "yaml";
|
||||
import { getEnv } from "./env.js";
|
||||
import { logger } from "./logger.js";
|
||||
|
||||
export interface Replacement {
|
||||
find: string;
|
||||
replace: string;
|
||||
}
|
||||
|
||||
export interface Config {
|
||||
server: {
|
||||
port: number;
|
||||
};
|
||||
deepgram: {
|
||||
model: string;
|
||||
language: string;
|
||||
diarize: boolean;
|
||||
interimResults: boolean;
|
||||
endpointing: number | false;
|
||||
utteranceEndMs: number;
|
||||
smartFormat: boolean;
|
||||
punctuate: boolean;
|
||||
keyterms: string[];
|
||||
replacements: Replacement[];
|
||||
};
|
||||
keepAlive: {
|
||||
intervalMs: number;
|
||||
idleThresholdMs: number;
|
||||
};
|
||||
taskExtraction: {
|
||||
model: string;
|
||||
intervalMs: number;
|
||||
llmMaxRetries: number;
|
||||
maxConsecutiveFailures: number;
|
||||
maxBackoffMs: number;
|
||||
};
|
||||
reconnect: {
|
||||
maxAttempts: number;
|
||||
maxBufferBytes: number;
|
||||
};
|
||||
output: {
|
||||
dir: string;
|
||||
};
|
||||
participants: string[];
|
||||
}
|
||||
|
||||
// --- config.yaml 読み込み ---
|
||||
|
||||
const CONFIG_YAML_PATH = "config.yaml";
|
||||
|
||||
interface MeetingConfig {
|
||||
participants: string[];
|
||||
keyterms: string[];
|
||||
replacements: Replacement[];
|
||||
}
|
||||
|
||||
const DEFAULT_MEETING_CONFIG: MeetingConfig = {
|
||||
participants: [],
|
||||
keyterms: [],
|
||||
replacements: [],
|
||||
};
|
||||
|
||||
function filterStrings(value: unknown): string[] {
|
||||
return Array.isArray(value)
|
||||
? value.filter((v): v is string => typeof v === "string")
|
||||
: [];
|
||||
}
|
||||
|
||||
function isReplacement(v: unknown): v is Replacement {
|
||||
if (typeof v !== "object" || v === null) return false;
|
||||
const r = v as Record<string, unknown>;
|
||||
return typeof r.find === "string" && typeof r.replace === "string";
|
||||
}
|
||||
|
||||
function loadMeetingConfig(): MeetingConfig {
|
||||
if (!fs.existsSync(CONFIG_YAML_PATH)) {
|
||||
logger.info("config.yaml が見つかりません。デフォルト設定で起動します");
|
||||
return DEFAULT_MEETING_CONFIG;
|
||||
}
|
||||
|
||||
try {
|
||||
const raw = fs.readFileSync(CONFIG_YAML_PATH, "utf-8");
|
||||
const parsed = parseYaml(raw) as Record<string, unknown>;
|
||||
|
||||
const config: MeetingConfig = {
|
||||
participants: filterStrings(parsed.participants),
|
||||
keyterms: filterStrings(parsed.keyterms),
|
||||
replacements: Array.isArray(parsed.replacements)
|
||||
? parsed.replacements.filter(isReplacement)
|
||||
: [],
|
||||
};
|
||||
|
||||
logger.info(
|
||||
`config.yaml: 参加者 ${config.participants.length} 名, キーターム ${config.keyterms.length} 件, 置換ルール ${config.replacements.length} 件`,
|
||||
);
|
||||
|
||||
return config;
|
||||
} catch (e) {
|
||||
logger.warn(
|
||||
"config.yaml の読み込みに失敗しました。デフォルト設定で起動します",
|
||||
e,
|
||||
);
|
||||
return DEFAULT_MEETING_CONFIG;
|
||||
}
|
||||
}
|
||||
|
||||
// プロセスライフタイムで1回だけ初期化される設定値(録音セッション状態とは無関係)
|
||||
let _config: Config | null = null;
|
||||
|
||||
export function getConfig(): Config {
|
||||
if (!_config) {
|
||||
const env = getEnv();
|
||||
const meeting = loadMeetingConfig();
|
||||
_config = {
|
||||
server: {
|
||||
port: env.PORT,
|
||||
},
|
||||
deepgram: {
|
||||
model: "nova-3",
|
||||
language: "ja",
|
||||
diarize: false,
|
||||
interimResults: true,
|
||||
endpointing: 500,
|
||||
utteranceEndMs: 1000,
|
||||
smartFormat: true,
|
||||
punctuate: true,
|
||||
keyterms: [
|
||||
...new Set([...meeting.participants, ...meeting.keyterms]),
|
||||
],
|
||||
replacements: meeting.replacements,
|
||||
},
|
||||
keepAlive: {
|
||||
intervalMs: 3_000,
|
||||
idleThresholdMs: 5_000,
|
||||
},
|
||||
taskExtraction: {
|
||||
model: "gemini-3-flash-preview",
|
||||
intervalMs: 30_000,
|
||||
llmMaxRetries: 0,
|
||||
maxConsecutiveFailures: 3,
|
||||
maxBackoffMs: 180_000,
|
||||
},
|
||||
reconnect: {
|
||||
maxAttempts: 3,
|
||||
maxBufferBytes: 2 * 1024 * 1024,
|
||||
},
|
||||
output: {
|
||||
dir: "output",
|
||||
},
|
||||
participants: meeting.participants,
|
||||
};
|
||||
}
|
||||
return _config;
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
// preflight.ts と env.ts の両方から参照される定数。
|
||||
// 外部ライブラリに依存しないため、preflight.ts を単独実行しても安全。
|
||||
|
||||
export const REQUIRED_KEYS = [
|
||||
"DEEPGRAM_API_KEY",
|
||||
"GOOGLE_GENERATIVE_AI_API_KEY",
|
||||
] as const;
|
||||
|
||||
export const DEFAULT_PORT = 3001;
|
||||
@@ -0,0 +1,278 @@
|
||||
import { DeepgramClient } from "@deepgram/sdk";
|
||||
import type { WebSocket } from "ws";
|
||||
import { getConfig, type Config, type Replacement } from "./config.js";
|
||||
import { TranscriptWriter } from "./transcript-writer.js";
|
||||
import type { ServerMessage } from "./types.js";
|
||||
import { logger } from "./logger.js";
|
||||
|
||||
type V1Socket = Awaited<
|
||||
ReturnType<InstanceType<typeof DeepgramClient>["listen"]["v1"]["connect"]>
|
||||
>;
|
||||
|
||||
function formatReplacement(r: Replacement): string {
|
||||
return `${r.find}:${r.replace}`;
|
||||
}
|
||||
|
||||
export class DeepgramRelay {
|
||||
readonly timestamp: string;
|
||||
readonly llmInputBuffer: string[] = [];
|
||||
readonly outputPath: string;
|
||||
|
||||
private readonly config: Config;
|
||||
private readonly writer: TranscriptWriter;
|
||||
private readonly meetingStartTime: Date;
|
||||
private readonly connectArgs: Parameters<InstanceType<typeof DeepgramClient>["listen"]["v1"]["connect"]>[0];
|
||||
|
||||
private connection: V1Socket | null = null;
|
||||
private keepAliveTimer: ReturnType<typeof setInterval> | null = null;
|
||||
private lastDataSentAt = Date.now();
|
||||
private timestampOffsetMs = 0;
|
||||
|
||||
private reconnecting = false;
|
||||
private readonly audioBuffer: Buffer[] = [];
|
||||
private audioBufferBytes = 0;
|
||||
|
||||
constructor(
|
||||
private readonly apiKey: string,
|
||||
private readonly browserWs: WebSocket,
|
||||
) {
|
||||
this.config = getConfig();
|
||||
this.meetingStartTime = new Date();
|
||||
this.timestamp = this.meetingStartTime
|
||||
.toISOString()
|
||||
.replace(/[:.]/g, "-")
|
||||
.slice(0, 19);
|
||||
this.writer = new TranscriptWriter(this.timestamp);
|
||||
this.outputPath = this.writer.getOutputPath();
|
||||
|
||||
this.connectArgs = {
|
||||
model: this.config.deepgram.model,
|
||||
language: this.config.deepgram.language,
|
||||
encoding: "linear16",
|
||||
sample_rate: "16000",
|
||||
channels: "1",
|
||||
...(this.config.deepgram.diarize ? { diarize: "true" } : {}),
|
||||
interim_results: String(this.config.deepgram.interimResults),
|
||||
endpointing: String(this.config.deepgram.endpointing),
|
||||
utterance_end_ms: String(this.config.deepgram.utteranceEndMs),
|
||||
smart_format: String(this.config.deepgram.smartFormat),
|
||||
punctuate: String(this.config.deepgram.punctuate),
|
||||
Authorization: this.apiKey,
|
||||
queryParams: this.buildQueryParams(),
|
||||
};
|
||||
}
|
||||
|
||||
async start(participants: string[]): Promise<void> {
|
||||
await this.writer.init(participants);
|
||||
this.connection = await this.openConnection();
|
||||
this.startKeepAlive();
|
||||
}
|
||||
|
||||
sendAudio(chunk: Buffer): void {
|
||||
if (chunk.byteLength === 0) return;
|
||||
this.lastDataSentAt = Date.now();
|
||||
|
||||
if (this.reconnecting) {
|
||||
this.bufferChunk(chunk);
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
this.connection?.sendMedia(chunk);
|
||||
} catch {
|
||||
logger.warn("音声送信失敗 — 再接続を開始");
|
||||
this.bufferChunk(chunk);
|
||||
void this.reconnect();
|
||||
}
|
||||
}
|
||||
|
||||
async stop(): Promise<void> {
|
||||
this.sendToClient({ type: "stopping" });
|
||||
|
||||
try { this.connection?.sendFinalize({ type: "Finalize" }); } catch {}
|
||||
await this.wait(5000);
|
||||
try { this.connection?.sendCloseStream({ type: "CloseStream" }); } catch {}
|
||||
|
||||
this.clearKeepAlive();
|
||||
await this.writer.flush();
|
||||
|
||||
try { this.connection?.close(); } catch {}
|
||||
this.connection = null;
|
||||
}
|
||||
|
||||
private buildQueryParams(): Record<string, string[]> | undefined {
|
||||
const params: Record<string, string[]> = {};
|
||||
const { keyterms, replacements } = this.config.deepgram;
|
||||
|
||||
if (keyterms.length > 0) {
|
||||
params.keyterm = keyterms;
|
||||
logger.info(`Deepgram keyterms: ${keyterms.join(", ")}`);
|
||||
}
|
||||
|
||||
if (replacements.length > 0) {
|
||||
params.replace = replacements.map(formatReplacement);
|
||||
logger.info(
|
||||
`Deepgram replace: ${replacements.map(formatReplacement).join(", ")}`,
|
||||
);
|
||||
}
|
||||
|
||||
return Object.keys(params).length > 0 ? params : undefined;
|
||||
}
|
||||
|
||||
// --- Private: Connection lifecycle ---
|
||||
|
||||
private async openConnection(): Promise<V1Socket> {
|
||||
const client = new DeepgramClient({ apiKey: this.apiKey });
|
||||
const conn = await client.listen.v1.connect(this.connectArgs);
|
||||
|
||||
conn.on("open", () => {
|
||||
logger.info("Deepgram 接続確立");
|
||||
if (!this.reconnecting) {
|
||||
this.sendToClient({ type: "ready", outputPath: this.outputPath });
|
||||
}
|
||||
});
|
||||
|
||||
conn.on("message", (data) => this.handleTranscriptMessage(data));
|
||||
|
||||
conn.on("error", (error) => {
|
||||
logger.error("Deepgram エラー", error);
|
||||
this.sendToClient({ type: "error", code: "DEEPGRAM_ERROR", message: error.message });
|
||||
});
|
||||
|
||||
conn.on("close", () => {
|
||||
logger.info("Deepgram 接続終了");
|
||||
});
|
||||
|
||||
conn.connect();
|
||||
await conn.waitForOpen();
|
||||
return conn;
|
||||
}
|
||||
|
||||
private handleTranscriptMessage(data: { type: string }): void {
|
||||
if (data.type !== "Results") return;
|
||||
|
||||
const result = data as unknown as {
|
||||
is_final?: boolean;
|
||||
speech_final?: boolean;
|
||||
start: number;
|
||||
channel: { alternatives: { transcript: string; words: { speaker?: number }[] }[] };
|
||||
};
|
||||
|
||||
const alt = result.channel.alternatives[0];
|
||||
if (!alt) return;
|
||||
|
||||
const { transcript } = alt;
|
||||
const isFinal = result.is_final ?? false;
|
||||
const speechFinal = result.speech_final ?? false;
|
||||
|
||||
this.sendToClient({ type: "transcript", text: transcript, isFinal });
|
||||
|
||||
if (!isFinal || !transcript.trim()) return;
|
||||
|
||||
const ts = this.formatTimestamp(result.start + this.timestampOffsetMs / 1000);
|
||||
this.writer.appendUtterance(transcript, ts, speechFinal);
|
||||
this.llmInputBuffer.push(transcript);
|
||||
}
|
||||
|
||||
// --- Private: Reconnection ---
|
||||
|
||||
private async reconnect(): Promise<void> {
|
||||
if (this.reconnecting) return;
|
||||
this.reconnecting = true;
|
||||
this.sendToClient({ type: "reconnecting" });
|
||||
this.clearKeepAlive();
|
||||
|
||||
const elapsedMs = Date.now() - this.meetingStartTime.getTime();
|
||||
|
||||
for (let attempt = 1; attempt <= this.config.reconnect.maxAttempts; attempt++) {
|
||||
const backoff = Math.pow(2, attempt - 1) * 1000;
|
||||
logger.info(`Deepgram 再接続 試行 ${attempt}/${this.config.reconnect.maxAttempts} (${backoff}ms後)`);
|
||||
await this.wait(backoff);
|
||||
|
||||
try {
|
||||
this.connection = await this.openConnection();
|
||||
this.timestampOffsetMs = elapsedMs;
|
||||
this.drainAudioBuffer();
|
||||
this.startKeepAlive();
|
||||
this.reconnecting = false;
|
||||
this.sendToClient({ type: "reconnected" });
|
||||
logger.info("Deepgram 再接続成功");
|
||||
return;
|
||||
} catch (e) {
|
||||
logger.error(`再接続試行 ${attempt} 失敗`, e);
|
||||
}
|
||||
}
|
||||
|
||||
this.reconnecting = false;
|
||||
this.sendToClient({
|
||||
type: "error",
|
||||
code: "DEEPGRAM_RECONNECT_FAILED",
|
||||
message: `Deepgram への再接続に ${this.config.reconnect.maxAttempts} 回失敗しました`,
|
||||
});
|
||||
}
|
||||
|
||||
// --- Private: KeepAlive ---
|
||||
|
||||
private startKeepAlive(): void {
|
||||
this.clearKeepAlive();
|
||||
this.keepAliveTimer = setInterval(() => {
|
||||
if (Date.now() - this.lastDataSentAt > this.config.keepAlive.idleThresholdMs) {
|
||||
try { this.connection?.sendKeepAlive({ type: "KeepAlive" }); } catch {
|
||||
logger.warn("KeepAlive 送信失敗");
|
||||
}
|
||||
}
|
||||
}, this.config.keepAlive.intervalMs);
|
||||
}
|
||||
|
||||
private clearKeepAlive(): void {
|
||||
if (this.keepAliveTimer) {
|
||||
clearInterval(this.keepAliveTimer);
|
||||
this.keepAliveTimer = null;
|
||||
}
|
||||
}
|
||||
|
||||
// --- Private: Audio buffer (reconnection) ---
|
||||
|
||||
private bufferChunk(chunk: Buffer): void {
|
||||
while (
|
||||
this.audioBufferBytes + chunk.byteLength > this.config.reconnect.maxBufferBytes &&
|
||||
this.audioBuffer.length > 0
|
||||
) {
|
||||
const dropped = this.audioBuffer.shift()!;
|
||||
this.audioBufferBytes -= dropped.byteLength;
|
||||
logger.warn("再接続バッファ上限超過: 古いチャンクを破棄");
|
||||
}
|
||||
this.audioBuffer.push(chunk);
|
||||
this.audioBufferBytes += chunk.byteLength;
|
||||
}
|
||||
|
||||
private drainAudioBuffer(): void {
|
||||
const chunks = this.audioBuffer.splice(0);
|
||||
this.audioBufferBytes = 0;
|
||||
for (const chunk of chunks) {
|
||||
try { this.connection?.sendMedia(chunk); } catch { break; }
|
||||
}
|
||||
}
|
||||
|
||||
// --- Private: Utilities ---
|
||||
|
||||
private sendToClient(msg: ServerMessage): void {
|
||||
if (this.browserWs.readyState === this.browserWs.OPEN) {
|
||||
this.browserWs.send(JSON.stringify(msg));
|
||||
}
|
||||
}
|
||||
|
||||
private formatTimestamp(startSeconds: number): string {
|
||||
const d = new Date(this.meetingStartTime.getTime() + startSeconds * 1000);
|
||||
return d.toLocaleTimeString("ja-JP", {
|
||||
hour: "2-digit",
|
||||
minute: "2-digit",
|
||||
second: "2-digit",
|
||||
hour12: false,
|
||||
});
|
||||
}
|
||||
|
||||
private wait(ms: number): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
}
|
||||
+133
@@ -0,0 +1,133 @@
|
||||
import { config as dotenvConfig } from "dotenv";
|
||||
import { execSync } from "node:child_process";
|
||||
import fs from "node:fs";
|
||||
import { z } from "zod";
|
||||
import { REQUIRED_KEYS } from "./constants.js";
|
||||
import { logger } from "./logger.js";
|
||||
|
||||
const envSchema = z.object({
|
||||
DEEPGRAM_API_KEY: z.string().min(1),
|
||||
GOOGLE_GENERATIVE_AI_API_KEY: z.string().min(1),
|
||||
PORT: z
|
||||
.string()
|
||||
.default("3001")
|
||||
.transform(Number)
|
||||
.pipe(z.number().int().min(1).max(65535)),
|
||||
});
|
||||
|
||||
export type Env = z.infer<typeof envSchema>;
|
||||
|
||||
export class EnvError extends Error {
|
||||
constructor(message: string) {
|
||||
super(message);
|
||||
this.name = "EnvError";
|
||||
}
|
||||
}
|
||||
|
||||
// プロセスライフタイムで1回だけ初期化される設定値(録音セッション状態とは無関係)
|
||||
let _env: Env | null = null;
|
||||
|
||||
export function getEnv(): Env {
|
||||
if (!_env) _env = loadEnv();
|
||||
return _env;
|
||||
}
|
||||
|
||||
function loadEnv(): Env {
|
||||
dotenvConfig();
|
||||
|
||||
const first = envSchema.safeParse(process.env);
|
||||
if (first.success) return first.data;
|
||||
|
||||
if (!fs.existsSync(".env") && tryOpInject()) {
|
||||
dotenvConfig({ override: true });
|
||||
const retry = envSchema.safeParse(process.env);
|
||||
if (retry.success) {
|
||||
logger.info("1Password CLI で .env を自動生成しました");
|
||||
return retry.data;
|
||||
}
|
||||
}
|
||||
|
||||
const { fieldErrors } = z.flattenError(first.error);
|
||||
throw new EnvError(formatDiagnostics(fieldErrors));
|
||||
}
|
||||
|
||||
function tryOpInject(): boolean {
|
||||
try {
|
||||
execSync("op --version", { stdio: "ignore", timeout: 3_000 });
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
try {
|
||||
logger.info("1Password CLI で .env を生成中...");
|
||||
execSync("op inject -i .env.example -o .env", {
|
||||
stdio: "inherit",
|
||||
timeout: 30_000,
|
||||
});
|
||||
return true;
|
||||
} catch {
|
||||
logger.warn("1Password CLI での .env 生成に失敗しました");
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
function hasOpCli(): boolean {
|
||||
try {
|
||||
execSync("op --version", { stdio: "ignore", timeout: 3_000 });
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
function formatDiagnostics(
|
||||
fieldErrors: Record<string, string[] | undefined>,
|
||||
): string {
|
||||
const issues = Object.keys(fieldErrors)
|
||||
.map((key) => {
|
||||
const val = process.env[key];
|
||||
const status =
|
||||
val === undefined ? "未設定" : val === "" ? "空です" : "値が不正です";
|
||||
return ` ${key}: ${status}`;
|
||||
})
|
||||
.join("\n");
|
||||
|
||||
const opAvailable = hasOpCli();
|
||||
|
||||
const lines = [
|
||||
"",
|
||||
"============================================================",
|
||||
" 環境変数が設定されていません",
|
||||
"============================================================",
|
||||
"",
|
||||
" 不足している変数:",
|
||||
issues,
|
||||
"",
|
||||
" -- 解決方法 -----------------------------------------------",
|
||||
"",
|
||||
];
|
||||
|
||||
if (opAvailable) {
|
||||
lines.push(
|
||||
" [方法1] 1Password CLI で自動生成(推奨)",
|
||||
" $ op inject -i .env.example -o .env",
|
||||
"",
|
||||
);
|
||||
}
|
||||
|
||||
lines.push(
|
||||
opAvailable
|
||||
? " [方法2] 手動で .env ファイルを作成"
|
||||
: " .env ファイルを手動で作成してください",
|
||||
"",
|
||||
...REQUIRED_KEYS.map((key) => ` ${key}=your-key-here`),
|
||||
"",
|
||||
" API キーの取得先:",
|
||||
" Deepgram : https://console.deepgram.com/",
|
||||
" Google AI : https://aistudio.google.com/apikey",
|
||||
"",
|
||||
"============================================================",
|
||||
"",
|
||||
);
|
||||
|
||||
return lines.join("\n");
|
||||
}
|
||||
+195
@@ -0,0 +1,195 @@
|
||||
import http from "node:http";
|
||||
import express from "express";
|
||||
import { WebSocketServer, type WebSocket } from "ws";
|
||||
import { DeepgramClient } from "@deepgram/sdk";
|
||||
import { getEnv, EnvError } from "./env.js";
|
||||
import { logger } from "./logger.js";
|
||||
import { getConfig } from "./config.js";
|
||||
import { DeepgramRelay } from "./deepgram-relay.js";
|
||||
import { taskExtractionLoop } from "./task-extractor.js";
|
||||
import type { ClientMessage, ServerMessage, TaskStatus } from "./types.js";
|
||||
|
||||
async function validateDeepgramKey(apiKey: string): Promise<void> {
|
||||
try {
|
||||
const client = new DeepgramClient({ apiKey });
|
||||
await client.manage.v1.projects.list();
|
||||
logger.info("Deepgram: OK");
|
||||
} catch (e) {
|
||||
if (hasHttpStatus(e, 401, 403)) {
|
||||
logger.error("Deepgram: APIキーが無効です。.env を確認してください");
|
||||
process.exit(1);
|
||||
}
|
||||
logger.warn("Deepgram: 検証リクエストに失敗しました(起動は続行)", e);
|
||||
}
|
||||
}
|
||||
|
||||
async function validateGoogleAiKey(apiKey: string): Promise<void> {
|
||||
const url = `https://generativelanguage.googleapis.com/v1beta/models?key=${apiKey}&pageSize=1`;
|
||||
try {
|
||||
const res = await fetch(url);
|
||||
if (res.status === 401 || res.status === 403) {
|
||||
logger.error("Google AI: APIキーが無効です。.env を確認してください");
|
||||
process.exit(1);
|
||||
}
|
||||
if (!res.ok) {
|
||||
logger.warn(`Google AI: 検証リクエストが ${res.status} を返しました(起動は続行)`);
|
||||
return;
|
||||
}
|
||||
logger.info("Google AI: OK");
|
||||
} catch (e) {
|
||||
logger.warn("Google AI: 検証リクエストに失敗しました(起動は続行)", e);
|
||||
}
|
||||
}
|
||||
|
||||
function hasHttpStatus(e: unknown, ...codes: number[]): boolean {
|
||||
if (!(e instanceof Error)) return false;
|
||||
const record = e as Record<string, unknown>;
|
||||
const status = record["statusCode"] ?? record["status"];
|
||||
if (typeof status === "number") return codes.includes(status);
|
||||
const msg = e.message.toLowerCase();
|
||||
return codes.some(
|
||||
(code) => msg.includes(String(code)) || msg.includes(HTTP_STATUS_NAMES[code] ?? ""),
|
||||
);
|
||||
}
|
||||
|
||||
const HTTP_STATUS_NAMES: Record<number, string> = {
|
||||
401: "unauthorized",
|
||||
403: "forbidden",
|
||||
};
|
||||
|
||||
function sendToWs(ws: WebSocket, msg: ServerMessage): void {
|
||||
if (ws.readyState === ws.OPEN) {
|
||||
ws.send(JSON.stringify(msg));
|
||||
}
|
||||
}
|
||||
|
||||
function tryParseClientMessage(text: string): ClientMessage | null {
|
||||
try {
|
||||
const parsed = JSON.parse(text) as Record<string, unknown>;
|
||||
if (parsed.type === "stop") return { type: "stop" };
|
||||
return null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async function main() {
|
||||
const env = getEnv();
|
||||
const config = getConfig();
|
||||
|
||||
const participants = config.participants;
|
||||
logger.info(`参加者: ${participants.length > 0 ? participants.join(", ") : "(未指定)"}`);
|
||||
|
||||
await validateDeepgramKey(env.DEEPGRAM_API_KEY);
|
||||
await validateGoogleAiKey(env.GOOGLE_GENERATIVE_AI_API_KEY);
|
||||
|
||||
const app = express();
|
||||
const server = http.createServer(app);
|
||||
const wss = new WebSocketServer({ noServer: true });
|
||||
|
||||
server.on("upgrade", (req, socket, head) => {
|
||||
wss.handleUpgrade(req, socket, head, (ws) => {
|
||||
wss.emit("connection", ws, req);
|
||||
});
|
||||
});
|
||||
|
||||
wss.on("connection", (ws: WebSocket) => {
|
||||
void handleSession(ws, env.DEEPGRAM_API_KEY, participants);
|
||||
});
|
||||
|
||||
const port = config.server.port;
|
||||
server.listen(port, () => {
|
||||
logger.info(`サーバー起動: http://localhost:${port}`);
|
||||
});
|
||||
}
|
||||
|
||||
async function handleSession(
|
||||
ws: WebSocket,
|
||||
apiKey: string,
|
||||
participants: string[],
|
||||
): Promise<void> {
|
||||
logger.info("ブラウザ WebSocket 接続");
|
||||
|
||||
const config = getConfig();
|
||||
const relay = new DeepgramRelay(apiKey, ws);
|
||||
|
||||
let stopping = false;
|
||||
let lastTaskStatus: TaskStatus = { taskCount: 0, failing: false };
|
||||
let taskAbortController: AbortController | null = null;
|
||||
let taskExtractionPromise: Promise<TaskStatus> | null = null;
|
||||
|
||||
try {
|
||||
await relay.start(participants);
|
||||
} catch (e) {
|
||||
logger.error("Deepgram リレー起動失敗", e);
|
||||
sendToWs(ws, {
|
||||
type: "error",
|
||||
code: "DEEPGRAM_CONNECT_FAILED",
|
||||
message: "Deepgram への接続に失敗しました",
|
||||
});
|
||||
ws.close();
|
||||
return;
|
||||
}
|
||||
|
||||
taskAbortController = new AbortController();
|
||||
taskExtractionPromise = taskExtractionLoop(
|
||||
{
|
||||
participants,
|
||||
outputDir: config.output.dir,
|
||||
timestamp: relay.timestamp,
|
||||
llmInputBuffer: relay.llmInputBuffer,
|
||||
onStatusChange: (status) => {
|
||||
lastTaskStatus = status;
|
||||
sendToWs(ws, { type: "task_status", ...status });
|
||||
},
|
||||
},
|
||||
taskAbortController.signal,
|
||||
);
|
||||
|
||||
async function gracefulStop(): Promise<void> {
|
||||
await relay.stop();
|
||||
|
||||
sendToWs(ws, { type: "stopped", ...lastTaskStatus });
|
||||
ws.close();
|
||||
|
||||
taskAbortController?.abort();
|
||||
if (taskExtractionPromise) {
|
||||
const result = await taskExtractionPromise.catch((): null => null);
|
||||
if (result && result.taskCount > 0) {
|
||||
logger.info(`タスク抽出完了: ${result.taskCount} 件`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ws.on("message", (data: Buffer | string, isBinary: boolean) => {
|
||||
if (isBinary) {
|
||||
relay.sendAudio(data as Buffer);
|
||||
return;
|
||||
}
|
||||
|
||||
const msg = tryParseClientMessage(String(data));
|
||||
if (msg?.type === "stop" && !stopping) {
|
||||
stopping = true;
|
||||
logger.info("停止要求を受信");
|
||||
void gracefulStop();
|
||||
}
|
||||
});
|
||||
|
||||
ws.on("close", () => {
|
||||
logger.info("ブラウザ WebSocket 切断");
|
||||
if (!stopping) {
|
||||
taskAbortController?.abort();
|
||||
taskExtractionPromise?.catch(() => {});
|
||||
void relay.stop().catch(() => {});
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
if (e instanceof EnvError) {
|
||||
console.error(e.message);
|
||||
} else {
|
||||
logger.error("起動エラー", e);
|
||||
}
|
||||
process.exit(1);
|
||||
});
|
||||
@@ -0,0 +1,10 @@
|
||||
const timestamp = () => new Date().toISOString();
|
||||
|
||||
export const logger = {
|
||||
info: (msg: string, ...args: unknown[]) =>
|
||||
console.log(`[${timestamp()}] INFO ${msg}`, ...args),
|
||||
warn: (msg: string, ...args: unknown[]) =>
|
||||
console.warn(`[${timestamp()}] WARN ${msg}`, ...args),
|
||||
error: (msg: string, ...args: unknown[]) =>
|
||||
console.error(`[${timestamp()}] ERROR ${msg}`, ...args),
|
||||
};
|
||||
@@ -0,0 +1,114 @@
|
||||
import { config as dotenvConfig } from "dotenv";
|
||||
import fs from "node:fs";
|
||||
import net from "node:net";
|
||||
import { REQUIRED_KEYS, DEFAULT_PORT } from "./constants.js";
|
||||
|
||||
const MIN_NODE_VERSION = 20;
|
||||
|
||||
function fail(label: string, ...lines: string[]): never {
|
||||
console.error(`\n[check] ${label} ... NG\n`);
|
||||
for (const line of lines) console.error(line);
|
||||
console.error();
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
function checkNodeVersion(): void {
|
||||
const major = parseInt(process.versions.node.split(".")[0]!, 10);
|
||||
if (major < MIN_NODE_VERSION) {
|
||||
fail(
|
||||
`Node.js v${process.versions.node}`,
|
||||
` Node.js v${MIN_NODE_VERSION} 以上が必要です(現在 v${process.versions.node})。`,
|
||||
" https://nodejs.org/ から最新版をインストールしてください。",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
function checkEnvFile(): void {
|
||||
if (!fs.existsSync(".env")) {
|
||||
fail(
|
||||
".env",
|
||||
" .env ファイルが見つかりません。",
|
||||
" 以下の手順で作成してください:",
|
||||
"",
|
||||
" 1. cp .env.example .env",
|
||||
' 2. テキストエディタで .env を開く(例: code .env)',
|
||||
" 3. API キーを貼り付ける",
|
||||
"",
|
||||
" API キーの取得先:",
|
||||
" Deepgram : https://console.deepgram.com/",
|
||||
" Google AI : https://aistudio.google.com/apikey",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
function checkRequiredKeys(): void {
|
||||
for (const key of REQUIRED_KEYS) {
|
||||
const val = process.env[key];
|
||||
if (!val || val.trim() === "") {
|
||||
fail(
|
||||
key,
|
||||
` ${key} が設定されていません。`,
|
||||
" .env ファイルを開いて値を設定してください。",
|
||||
"",
|
||||
" API キーの取得先:",
|
||||
" Deepgram : https://console.deepgram.com/",
|
||||
" Google AI : https://aistudio.google.com/apikey",
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function resolvePort(): number {
|
||||
const portStr = process.env["PORT"] ?? String(DEFAULT_PORT);
|
||||
const port = parseInt(portStr, 10);
|
||||
if (isNaN(port) || port < 1 || port > 65535) {
|
||||
fail("PORT", ` PORT の値が不正です: ${portStr}`);
|
||||
}
|
||||
return port;
|
||||
}
|
||||
|
||||
function checkPortAvailable(port: number): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
const server = net.createServer();
|
||||
server.once("error", (err: NodeJS.ErrnoException) => {
|
||||
if (err.code === "EADDRINUSE") {
|
||||
fail(
|
||||
`ポート ${port}`,
|
||||
` ポート ${port} が使用中です。`,
|
||||
" 前回のサーバーが残っている場合: just stop を実行",
|
||||
" 他のアプリが使用中の場合: そのアプリを終了してください",
|
||||
);
|
||||
} else if (err.code === "EACCES") {
|
||||
fail(
|
||||
`ポート ${port}`,
|
||||
` ポート ${port} にアクセス権がありません。`,
|
||||
" 1024 以上のポート番号を .env の PORT に設定してください。",
|
||||
);
|
||||
} else {
|
||||
fail(`ポート ${port}`, ` ポートチェックでエラーが発生しました: ${err.message}`);
|
||||
}
|
||||
});
|
||||
server.listen({ port, host: "0.0.0.0" }, () => {
|
||||
server.close(() => resolve());
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
async function main(): Promise<void> {
|
||||
checkNodeVersion();
|
||||
checkEnvFile();
|
||||
dotenvConfig();
|
||||
checkRequiredKeys();
|
||||
|
||||
const port = resolvePort();
|
||||
await checkPortAvailable(port);
|
||||
|
||||
console.log(
|
||||
`[check] 環境チェック完了 (Node.js v${process.versions.node}, .env, APIキー, ポート${port})`,
|
||||
);
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
console.error("プリフライトチェックで予期しないエラーが発生しました:", e);
|
||||
process.exit(1);
|
||||
});
|
||||
@@ -0,0 +1,15 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const ExtractedTaskSchema = z.object({
|
||||
summary: z.string().describe("タスクの1行要約"),
|
||||
assignee: z.string().nullable().describe("担当者名(推定)"),
|
||||
deadline: z.string().nullable().describe("期限(言及があれば)"),
|
||||
evidence: z.string().describe("根拠となる発話の引用"),
|
||||
});
|
||||
|
||||
export const TaskExtractionResultSchema = z.object({
|
||||
tasks: z.array(ExtractedTaskSchema),
|
||||
});
|
||||
|
||||
export type ExtractedTask = z.infer<typeof ExtractedTaskSchema>;
|
||||
export type TaskExtractionResult = z.infer<typeof TaskExtractionResultSchema>;
|
||||
@@ -0,0 +1,190 @@
|
||||
import { setTimeout as delay } from "node:timers/promises";
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { generateObject, APICallError } from "ai";
|
||||
import { google } from "@ai-sdk/google";
|
||||
import { getConfig } from "./config.js";
|
||||
import { TaskExtractionResultSchema, type ExtractedTask } from "./schemas.js";
|
||||
import type { TaskStatus } from "./types.js";
|
||||
import { logger } from "./logger.js";
|
||||
|
||||
export interface TaskExtractorOptions {
|
||||
participants: string[];
|
||||
outputDir: string;
|
||||
timestamp: string;
|
||||
llmInputBuffer: string[];
|
||||
onStatusChange: (status: TaskStatus) => void;
|
||||
}
|
||||
|
||||
type TaskWithId = ExtractedTask & { id: string };
|
||||
|
||||
function isQuotaOrRateLimitError(error: unknown): boolean {
|
||||
if (!APICallError.isInstance(error)) return false;
|
||||
if (error.statusCode === 429) return true;
|
||||
if (error.statusCode === 503) {
|
||||
const body = error.responseBody ?? "";
|
||||
return body.includes("RESOURCE_EXHAUSTED") || body.includes("quota");
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
function buildPrompt(participants: string[], transcriptText: string): string {
|
||||
const participantLine = participants.length > 0
|
||||
? `参加者は ${participants.join("、")} です。`
|
||||
: "参加者リストは未提供です。assignee は常に null にしてください。";
|
||||
|
||||
const assigneeRule = participants.length > 0
|
||||
? `- assignee は参加者リスト(${participants.join("、")})の中からのみ選ぶ。リストにない名前を担当者にしてはならない
|
||||
- 会話中で「〇〇さんお願い」「〇〇がやります」のように名前が明示されている場合のみ assignee を設定する
|
||||
- 名前が明示されていない場合は assignee を null にする(推測しない)`
|
||||
: "- assignee は常に null にする";
|
||||
|
||||
return `以下は会議の文字起こしの一部です。${participantLine}
|
||||
話者の区別はありません(すべての発話が話者ラベルなしで記録されています)。
|
||||
|
||||
---
|
||||
${transcriptText}
|
||||
---
|
||||
|
||||
上記の会話から以下を抽出してください。
|
||||
|
||||
タスク:
|
||||
- 誰かが何かをやると約束した、または依頼された内容を抽出する
|
||||
- 「検討します」「考えておきます」などの曖昧な表現もタスクとして抽出する
|
||||
- 明確な行動(「作る」「送る」「確認する」「レビューする」等)だけでなく、検討・調査系の意思表示も含める
|
||||
- evidence には、そのタスクの根拠となる発話を原文のまま引用する
|
||||
|
||||
担当者の推定ルール:
|
||||
${assigneeRule}
|
||||
- 期限の言及がない場合は deadline を null にする
|
||||
|
||||
エッジケース:
|
||||
- テキストが2文以下の場合は tasks を空配列で返す
|
||||
- タスクが見当たらない場合は tasks を空配列で返す(無理に抽出しない)`;
|
||||
}
|
||||
|
||||
function formatTasksMarkdown(tasks: TaskWithId[]): string {
|
||||
const rows = tasks.map(
|
||||
(t) =>
|
||||
`| ${t.id} | ${t.summary} | ${t.assignee ?? "-"} | ${t.deadline ?? "-"} | ${t.evidence} |`,
|
||||
);
|
||||
|
||||
return [
|
||||
"# 抽出タスク",
|
||||
"",
|
||||
"| # | タスク | 担当 | 期限 | 根拠 |",
|
||||
"|---|--------|------|------|------|",
|
||||
...rows,
|
||||
"",
|
||||
"---",
|
||||
`*最終更新: ${new Date().toLocaleTimeString("ja-JP", { hour: "2-digit", minute: "2-digit", second: "2-digit", hour12: false })}*`,
|
||||
"",
|
||||
].join("\n");
|
||||
}
|
||||
|
||||
async function abortableSleep(ms: number, signal: AbortSignal): Promise<void> {
|
||||
try {
|
||||
await delay(ms, undefined, { signal });
|
||||
} catch {
|
||||
// AbortError on cancellation — expected
|
||||
}
|
||||
}
|
||||
|
||||
export async function taskExtractionLoop(
|
||||
options: TaskExtractorOptions,
|
||||
signal: AbortSignal,
|
||||
): Promise<TaskStatus> {
|
||||
const { participants, outputDir, timestamp, llmInputBuffer, onStatusChange } = options;
|
||||
const cfg = getConfig().taskExtraction;
|
||||
const tasksFilePath = path.join(outputDir, `meeting-${timestamp}-tasks.md`);
|
||||
const allTasks: TaskWithId[] = [];
|
||||
let taskCounter = 0;
|
||||
let consecutiveFailures = 0;
|
||||
|
||||
function computeBackoffMs(): number {
|
||||
if (consecutiveFailures === 0) return cfg.intervalMs;
|
||||
return Math.min(
|
||||
cfg.intervalMs * Math.pow(2, consecutiveFailures - 1),
|
||||
cfg.maxBackoffMs,
|
||||
);
|
||||
}
|
||||
|
||||
async function extractAndWrite(lines: string[]): Promise<void> {
|
||||
const transcriptText = lines.join("\n");
|
||||
logger.info(`タスク抽出開始 (${lines.length} 行, ${transcriptText.length} 文字)`);
|
||||
|
||||
const { object } = await generateObject({
|
||||
model: google(cfg.model),
|
||||
schema: TaskExtractionResultSchema,
|
||||
prompt: buildPrompt(participants, transcriptText),
|
||||
maxRetries: cfg.llmMaxRetries,
|
||||
abortSignal: signal,
|
||||
});
|
||||
|
||||
if (object.tasks.length > 0) {
|
||||
const newTasks = object.tasks.map((t) => ({
|
||||
...t,
|
||||
id: String(++taskCounter),
|
||||
}));
|
||||
allTasks.push(...newTasks);
|
||||
await fs.writeFile(tasksFilePath, formatTasksMarkdown(allTasks), "utf-8");
|
||||
logger.info(`タスク ${newTasks.length} 件抽出 → ${tasksFilePath}`);
|
||||
} else {
|
||||
logger.info("タスクなし(今回のサイクル)");
|
||||
}
|
||||
}
|
||||
|
||||
function currentStatus(overrides?: Partial<TaskStatus>): TaskStatus {
|
||||
return {
|
||||
taskCount: allTasks.length,
|
||||
failing: consecutiveFailures > 0,
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
while (!signal.aborted) {
|
||||
if (llmInputBuffer.length > 0) {
|
||||
const snapshot = llmInputBuffer.splice(0);
|
||||
try {
|
||||
await extractAndWrite(snapshot);
|
||||
consecutiveFailures = 0;
|
||||
onStatusChange(currentStatus());
|
||||
} catch (e) {
|
||||
consecutiveFailures++;
|
||||
const isQuota = isQuotaOrRateLimitError(e);
|
||||
|
||||
if (isQuota) {
|
||||
logger.warn(
|
||||
`タスク抽出失敗: API クォータ超過 (連続 ${consecutiveFailures} 回)。Google AI Studio でプランを確認してください`,
|
||||
);
|
||||
} else {
|
||||
logger.error(`タスク抽出失敗 (連続 ${consecutiveFailures} 回)`, e);
|
||||
}
|
||||
|
||||
logger.warn(`失敗した ${snapshot.length} 行を破棄(次サイクルの新規バッファで再試行)`);
|
||||
|
||||
const message = isQuota
|
||||
? "タスク抽出が API 制限で停止中"
|
||||
: "タスク抽出でエラーが発生中";
|
||||
onStatusChange(currentStatus({ message }));
|
||||
}
|
||||
}
|
||||
|
||||
await abortableSleep(computeBackoffMs(), signal);
|
||||
}
|
||||
|
||||
if (llmInputBuffer.length > 0 && consecutiveFailures < cfg.maxConsecutiveFailures) {
|
||||
try {
|
||||
logger.info("最終 flush: タスク抽出");
|
||||
await extractAndWrite(llmInputBuffer.splice(0));
|
||||
} catch {
|
||||
logger.warn("最終 flush 失敗");
|
||||
}
|
||||
} else if (llmInputBuffer.length > 0) {
|
||||
logger.warn(
|
||||
`最終 flush スキップ: API が ${consecutiveFailures} 回連続で失敗中のため`,
|
||||
);
|
||||
}
|
||||
|
||||
return currentStatus();
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { getConfig } from "./config.js";
|
||||
|
||||
export class TranscriptWriter {
|
||||
private writeQueue = Promise.resolve();
|
||||
private readonly filePath: string;
|
||||
|
||||
constructor(timestamp: string) {
|
||||
this.filePath = path.join(
|
||||
getConfig().output.dir,
|
||||
`meeting-${timestamp}-transcript.md`,
|
||||
);
|
||||
}
|
||||
|
||||
async init(participants: string[]): Promise<void> {
|
||||
await fs.mkdir(path.dirname(this.filePath), { recursive: true });
|
||||
|
||||
const now = new Date();
|
||||
const dateStr = now.toLocaleDateString("ja-JP", {
|
||||
year: "numeric",
|
||||
month: "2-digit",
|
||||
day: "2-digit",
|
||||
});
|
||||
const timeStr = now.toLocaleTimeString("ja-JP", {
|
||||
hour: "2-digit",
|
||||
minute: "2-digit",
|
||||
hour12: false,
|
||||
});
|
||||
|
||||
const header = [
|
||||
`# 会議メモ ${dateStr} ${timeStr}`,
|
||||
"",
|
||||
"## 参加者",
|
||||
participants.length > 0 ? participants.join("、") : "(未指定)",
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
].join("\n");
|
||||
|
||||
await fs.writeFile(this.filePath, header, "utf-8");
|
||||
}
|
||||
|
||||
appendUtterance(
|
||||
text: string,
|
||||
timestamp: string,
|
||||
isSpeechFinal: boolean,
|
||||
): void {
|
||||
if (!text.trim()) return;
|
||||
|
||||
let content = `**[${timestamp}]** ${text}\n`;
|
||||
if (isSpeechFinal) {
|
||||
content += "\n";
|
||||
}
|
||||
|
||||
this.append(content);
|
||||
}
|
||||
|
||||
private append(text: string): void {
|
||||
this.writeQueue = this.writeQueue.then(() =>
|
||||
fs.appendFile(this.filePath, text, "utf-8"),
|
||||
);
|
||||
}
|
||||
|
||||
async flush(): Promise<void> {
|
||||
await this.writeQueue;
|
||||
}
|
||||
|
||||
getOutputPath(): string {
|
||||
return this.filePath;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
export interface TaskStatus {
|
||||
taskCount: number;
|
||||
failing: boolean;
|
||||
message?: string;
|
||||
}
|
||||
|
||||
export type ServerMessage =
|
||||
| { type: "ready"; outputPath: string }
|
||||
| { type: "transcript"; text: string; isFinal: boolean }
|
||||
| { type: "reconnecting" }
|
||||
| { type: "reconnected" }
|
||||
| { type: "stopping" }
|
||||
| { type: "stopped" } & TaskStatus
|
||||
| { type: "task_status" } & TaskStatus
|
||||
| { type: "error"; code: string; message: string };
|
||||
|
||||
export type ClientMessage = { type: "stop" };
|
||||
Reference in New Issue
Block a user