diff --git a/src/cli/ingest.ts b/src/cli/ingest.ts index fe2fb98..1d85f0b 100644 --- a/src/cli/ingest.ts +++ b/src/cli/ingest.ts @@ -4,9 +4,9 @@ import { readFileSync, writeFileSync, mkdirSync, existsSync, readdirSync, renameSync } from 'node:fs'; import { join } from 'node:path'; import matter from 'gray-matter'; -import chokidar from 'chokidar'; import { callLLM, loadLLMConfig, type ChatRequest, type LLMConfigFromVault } from '../llm.js'; import { preCheckConflict } from '../ingest.js'; +import { startInboxWatcher, type ActiveWatcher } from '../watcher.js'; interface IngestOptions { mock?: boolean; @@ -98,8 +98,7 @@ export async function ingestCommand(vaultPath: string, args: string[]): Promise< /** * Lauscht auf neue Dateien in 00_Inbox/ und ingestiert sie automatisch. - * Nutzt chokidar (robuster als node:fs.watch auf Windows). - * Debounce 1s, um mehrfache Events pro Datei zu vermeiden. + * Wrapper um startInboxWatcher mit Ctrl+C-Handling. */ async function watchInbox( vaultPath: string, @@ -110,40 +109,23 @@ async function watchInbox( dryRun: boolean, inboxDir: string, ): Promise { - const watcher = chokidar.watch(inboxDir, { - ignored: /(^|[/\\])(\..*|Verarbeitet|Problemfaelle|.*\.tmp$)/, - persistent: true, - ignoreInitial: true, - awaitWriteFinish: { stabilityThreshold: 500, pollInterval: 100 }, - }); - - let debounceTimer: NodeJS.Timeout | null = null; - const pending = new Set(); - - watcher.on('add', (filePath) => { - if (!filePath.endsWith('.md')) return; - pending.add(filePath); - if (debounceTimer) clearTimeout(debounceTimer); - debounceTimer = setTimeout(async () => { - const toProcess = [...pending]; - pending.clear(); - for (const notePath of toProcess) { - console.log(`[watch] Neue Notiz: ${notePath}`); - try { - await processNote(vaultPath, notePath, schema, sysIndex, wikiIndex, llmConfig, dryRun); - } catch (err) { - console.error(`[watch] FEHLER: ${(err as Error).message}`); - moveToProblemfaelle(vaultPath, notePath, (err as Error).message); - } - } - }, 1000); + let activeWatcher: ActiveWatcher | null = null; + activeWatcher = await startInboxWatcher(inboxDir, { + processNote: async (notePath: string) => { + console.log(`[watch] Neue Notiz: ${notePath}`); + await processNote(vaultPath, notePath, schema, sysIndex, wikiIndex, llmConfig, dryRun); + }, + onError: (err: Error, notePath: string) => { + console.error(`[watch] FEHLER: ${err.message}`); + moveToProblemfaelle(vaultPath, notePath, err.message); + }, }); // Auf Ctrl+C warten await new Promise((resolve) => { process.on('SIGINT', () => { console.log('\n=== Watch-Mode beendet. ==='); - watcher.close().then(() => resolve()); + activeWatcher?.close().then(() => resolve()); }); }); } diff --git a/src/watcher.ts b/src/watcher.ts new file mode 100644 index 0000000..f37b4e2 --- /dev/null +++ b/src/watcher.ts @@ -0,0 +1,60 @@ +// Watcher-Logik (testbar) +// Bonus 21: Watcher in separate Datei, mit Dependency-Injection +// fuer processNote (testbar ohne chokidar) + +import chokidar from 'chokidar'; + +export interface WatcherDeps { + processNote: (notePath: string) => Promise; + onError?: (err: Error, notePath: string) => void; + debounceMs?: number; +} + +export interface ActiveWatcher { + close: () => Promise; +} + +/** + * Startet einen Watcher auf `inboxDir` und ruft `processNote` fuer jede + * neue .md-Datei (mit Debounce). + * + * Testbar: processNote wird als Dependency uebergeben. + */ +export async function startInboxWatcher(inboxDir: string, deps: WatcherDeps): Promise { + const debounceMs = deps.debounceMs ?? 1000; + const onError = deps.onError ?? ((err) => console.error(`[watch] FEHLER: ${err.message}`)); + + const watcher = chokidar.watch(inboxDir, { + ignored: /(^|[/\\])(\..*|Verarbeitet|Problemfaelle|.*\.tmp$)/, + persistent: true, + ignoreInitial: true, + awaitWriteFinish: { stabilityThreshold: 500, pollInterval: 100 }, + }); + + let debounceTimer: ReturnType | null = null; + const pending = new Set(); + + watcher.on('add', (filePath) => { + if (!filePath.endsWith('.md')) return; + pending.add(filePath); + if (debounceTimer) clearTimeout(debounceTimer); + debounceTimer = setTimeout(async () => { + const toProcess = [...pending]; + pending.clear(); + for (const notePath of toProcess) { + try { + await deps.processNote(notePath); + } catch (err) { + onError(err as Error, notePath); + } + } + }, debounceMs); + }); + + return { + close: async () => { + if (debounceTimer) clearTimeout(debounceTimer); + await watcher.close(); + }, + }; +} diff --git a/tests/watcher.test.ts b/tests/watcher.test.ts new file mode 100644 index 0000000..5f175c1 --- /dev/null +++ b/tests/watcher.test.ts @@ -0,0 +1,155 @@ +// Tests fuer startInboxWatcher +// Bonus 21: Watcher-Edge-Cases + +import { describe, it, expect, afterEach } from 'vitest'; +import { mkdtempSync, mkdirSync, writeFileSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { startInboxWatcher, type ActiveWatcher } from '../src/watcher.js'; + +const watchers: ActiveWatcher[] = []; +afterEach(async () => { + for (const w of watchers) await w.close(); + watchers.length = 0; +}); + +function setupInbox(): string { + const dir = mkdtempSync(join(tmpdir(), 'mindomat-watcher-test-')); + // Inbox-Struktur + for (const d of ['00_Inbox', '00_Inbox/Verarbeitet', '00_Inbox/Problemfaelle']) { + mkdirSync(join(dir, d), { recursive: true }); + } + return dir; +} + +describe('startInboxWatcher', () => { + it('startet und schliesst sauber', async () => { + const inbox = setupInbox(); + const processed: string[] = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async (p) => { processed.push(p); }, + debounceMs: 50, + }); + watchers.push(watcher); + expect(watcher).toBeDefined(); + expect(typeof watcher.close).toBe('function'); + await watcher.close(); + }); + + it('verarbeitet neue .md-Dateien', async () => { + const inbox = setupInbox(); + const processed: string[] = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async (p) => { processed.push(p); }, + debounceMs: 100, + }); + watchers.push(watcher); + + // Warte, bis Watcher bereit ist + await new Promise((r) => setTimeout(r, 500)); + + // Schreibe neue .md-Datei + const notePath = join(inbox, '00_Inbox', 'test.md'); + writeFileSync(notePath, '# Test'); + + // Warte auf Debounce + Verarbeitung + await new Promise((r) => setTimeout(r, 1500)); + + expect(processed.length).toBe(1); + expect(processed[0]).toContain('test.md'); + }); + + it('ignoriert Nicht-.md-Dateien', async () => { + const inbox = setupInbox(); + const processed: string[] = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async (p) => { processed.push(p); }, + debounceMs: 100, + }); + watchers.push(watcher); + + await new Promise((r) => setTimeout(r, 500)); + + // Schreibe .txt-Datei (sollte ignoriert werden) + writeFileSync(join(inbox, '00_Inbox', 'test.txt'), 'not markdown'); + await new Promise((r) => setTimeout(r, 800)); + + expect(processed.length).toBe(0); + }); + + it('ignoriert Verarbeitet/-Unterordner', async () => { + const inbox = setupInbox(); + const processed: string[] = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async (p) => { processed.push(p); }, + debounceMs: 100, + }); + watchers.push(watcher); + + await new Promise((r) => setTimeout(r, 500)); + + // In Verarbeitet/ schreiben (sollte ignoriert werden) + writeFileSync(join(inbox, '00_Inbox', 'Verarbeitet', 'old.md'), '# alt'); + await new Promise((r) => setTimeout(r, 800)); + + expect(processed.length).toBe(0); + }); + + it('Debounce: mehrere schnelle Writes -> eine Verarbeitung', async () => { + const inbox = setupInbox(); + const processed: string[] = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async (p) => { processed.push(p); }, + debounceMs: 500, + }); + watchers.push(watcher); + + await new Promise((r) => setTimeout(r, 500)); + + // 3 schnelle Writes + for (let i = 0; i < 3; i++) { + writeFileSync(join(inbox, '00_Inbox', `note-${i}.md`), `# Note ${i}`); + await new Promise((r) => setTimeout(r, 50)); + } + // Warte auf Debounce + await new Promise((r) => setTimeout(r, 1500)); + + // Alle 3 Notizen sollten verarbeitet sein + expect(processed.length).toBeGreaterThanOrEqual(1); + expect(processed.length).toBeLessThanOrEqual(3); + }); + + it('Fehler in processNote werden via onError-Callback gemeldet', async () => { + const inbox = setupInbox(); + const errors: Array<{ err: Error; path: string }> = []; + const watcher = await startInboxWatcher(inbox, { + processNote: async () => { throw new Error('Test-Fehler'); }, + onError: (err, path) => { errors.push({ err, path }); }, + debounceMs: 100, + }); + watchers.push(watcher); + + await new Promise((r) => setTimeout(r, 500)); + + writeFileSync(join(inbox, '00_Inbox', 'fail.md'), '# fail'); + await new Promise((r) => setTimeout(r, 1500)); + + expect(errors.length).toBeGreaterThanOrEqual(1); + expect(errors[0].err.message).toBe('Test-Fehler'); + expect(errors[0].path).toContain('fail.md'); + }); + + it('schliesst sauber ohne haengende Prozesse', async () => { + const inbox = setupInbox(); + const watcher = await startInboxWatcher(inbox, { + processNote: async () => {}, + debounceMs: 50, + }); + + // Vor close: Watcher aktiv + expect(watcher).toBeDefined(); + + // Close sollte schnell gehen und keine Errors werfen + await expect(watcher.close()).resolves.not.toThrow(); + }); +});