import fs from "node:fs/promises"; import path from "node:path"; import { createTwoFilesPatch } from "diff"; import fg from "fast-glob"; import YAML from "yaml"; import { newId, stableKey } from "../core/ids.js"; import { sha256 } from "../core/json.js"; import type { PrivacyClass, ThoughtEvent } from "../events/types.js"; import type { JazzThoughtStore } from "../jazz/store.js"; export interface FilesystemConnectorConfig { id: string; root: string; include?: string[]; ignore?: string[]; extensions?: string[]; maxFileBytes?: number; maxDiffChars?: number; privacy?: PrivacyClass; storeContent?: boolean; } export interface FilesystemScanResult { correlationId: string; scanned: number; unchanged: number; added: number; changed: number; renamed: number; deleted: number; skipped: number; events: ThoughtEvent[]; } interface FileObservation { relativePath: string; content: string; contentType: string; sha256: string; sizeBytes: number; mtimeMs: number; occurredAt: string; explicitId?: string; } const defaultIgnores = [ "**/.git/**", "**/.obsidian/**", "**/.thoughtstream/**", "**/node_modules/**", "**/.DS_Store", "**/*.swp", "**/*.swo", "**/*~", "**/.env", "**/.env.*", "**/*secret*", "**/*credential*", ]; export class FilesystemConnector { readonly id: string; private readonly include: string[]; private readonly ignore: string[]; private readonly extensions: Set; private readonly maxFileBytes: number; private readonly maxDiffChars: number; private readonly privacy: PrivacyClass; private readonly storeContent: boolean; constructor(private readonly config: FilesystemConnectorConfig) { this.id = config.id; this.include = config.include ?? ["**/*"]; this.ignore = [...defaultIgnores, ...(config.ignore ?? [])]; this.extensions = new Set((config.extensions ?? [".md", ".mdx", ".txt", ".json", ".yaml", ".yml", ".toml"]).map((item) => item.toLowerCase())); this.maxFileBytes = config.maxFileBytes ?? 1_000_000; this.maxDiffChars = config.maxDiffChars ?? 32_000; this.privacy = config.privacy ?? "sensitive"; this.storeContent = config.storeContent ?? true; } async scan(store: JazzThoughtStore): Promise { const root = await fs.realpath(this.config.root); const correlationId = newId("scan"); const paths = await fg(this.include, { cwd: root, ignore: this.ignore, onlyFiles: true, dot: false, followSymbolicLinks: false, unique: true, }); const current = await store.listCurrentDocuments(this.id); const activeCurrent = current.filter((document) => !document.deleted); const currentByPath = new Map(activeCurrent.map((document) => [document.path, document])); const currentByDocumentId = new Map(activeCurrent.map((document) => [document.documentId, document])); const unmatchedByHash = new Map(); for (const document of activeCurrent) { const bucket = unmatchedByHash.get(document.sha256) ?? []; bucket.push(document); unmatchedByHash.set(document.sha256, bucket); } const result: FilesystemScanResult = { correlationId, scanned: 0, unchanged: 0, added: 0, changed: 0, renamed: 0, deleted: 0, skipped: 0, events: [], }; const seenDocumentIds = new Set(); for (const relativePath of paths.sort()) { const observation = await this.observe(root, relativePath); if (!observation) { result.skipped += 1; continue; } result.scanned += 1; const explicitDocumentId = observation.explicitId ? stableDocumentId(this.id, `frontmatter:${observation.explicitId}`) : undefined; const byPath = currentByPath.get(observation.relativePath); const byExplicitId = explicitDocumentId ? currentByDocumentId.get(explicitDocumentId) : undefined; const byContent = findUniqueUnseen(unmatchedByHash.get(observation.sha256) ?? [], seenDocumentIds); const previous = byExplicitId ?? byPath ?? byContent; const identityConfidence = observation.explicitId ? "explicit" : byPath ? "path" : byContent ? "content" : "path"; const documentId = explicitDocumentId ?? previous?.documentId ?? stableDocumentId(this.id, `path:${observation.relativePath}`); seenDocumentIds.add(documentId); if (previous?.path === observation.relativePath && previous.sha256 === observation.sha256) { result.unchanged += 1; continue; } const versionId = stableKey("version", this.id, documentId, observation.sha256); await store.appendDocumentVersion({ id: versionId, source: this.id, documentId, path: observation.relativePath, contentType: observation.contentType, sha256: observation.sha256, content: this.storeContent ? observation.content : "", sizeBytes: observation.sizeBytes, mtimeMs: observation.mtimeMs, createdAt: new Date().toISOString(), }); const previousVersion = previous?.versionId ? await store.getDocumentVersion(previous.versionId) : undefined; const renamed = Boolean(previous && previous.path !== observation.relativePath); const type = renamed ? "stream.thought.source.file.renamed" : previous ? "stream.thought.source.file.changed" : "stream.thought.source.file.added"; const diff = previousVersion && previousVersion.sha256 !== observation.sha256 ? boundedDiff(previousVersion.path, previousVersion.content, observation.relativePath, observation.content, this.maxDiffChars) : undefined; const append = await store.appendEvent({ type, schemaVersion: 1, source: this.id, sourceKind: "filesystem", externalId: documentId, idempotencyKey: `${documentId}:${type}:${observation.relativePath}:${observation.sha256}`, occurredAt: observation.occurredAt, actor: this.id, correlationId, privacy: this.privacy, payload: { documentId, path: observation.relativePath, versionId, ...(previous ? { previousVersionId: previous.versionId, previousSha256: previous.sha256 } : {}), ...(renamed && previous ? { previousPath: previous.path } : {}), sha256: observation.sha256, contentType: observation.contentType, sizeBytes: observation.sizeBytes, mtimeMs: observation.mtimeMs, ...(diff ? { diff } : {}), identityConfidence, }, }); if (append.inserted) result.events.push(append.event); await store.upsertCurrentDocument({ id: stableKey("document", this.id, documentId), source: this.id, documentId, path: observation.relativePath, versionId, sha256: observation.sha256, contentType: observation.contentType, sizeBytes: observation.sizeBytes, mtimeMs: observation.mtimeMs, deleted: false, updatedAt: new Date().toISOString(), }); if (renamed) result.renamed += 1; else if (previous) result.changed += 1; else result.added += 1; } for (const previous of activeCurrent) { if (seenDocumentIds.has(previous.documentId)) continue; const now = new Date().toISOString(); const append = await store.appendEvent({ type: "stream.thought.source.file.deleted", schemaVersion: 1, source: this.id, sourceKind: "filesystem", externalId: previous.documentId, idempotencyKey: `${previous.documentId}:deleted:${previous.path}:${previous.sha256}`, occurredAt: now, actor: this.id, correlationId, privacy: this.privacy, payload: { documentId: previous.documentId, path: previous.path, previousVersionId: previous.versionId, previousSha256: previous.sha256, contentType: previous.contentType, sizeBytes: previous.sizeBytes, mtimeMs: previous.mtimeMs, identityConfidence: "path", }, }); if (append.inserted) result.events.push(append.event); await store.upsertCurrentDocument({ ...previous, deleted: true, updatedAt: now }); result.deleted += 1; } await store.flush(); return result; } private async observe(root: string, relativePath: string): Promise { const extension = path.extname(relativePath).toLowerCase(); if (!this.extensions.has(extension)) return undefined; const absolutePath = path.resolve(root, relativePath); const realPath = await fs.realpath(absolutePath); if (!isInside(root, realPath)) throw new Error(`Filesystem source escaped root: ${relativePath}`); const stat = await fs.stat(realPath); if (!stat.isFile() || stat.size > this.maxFileBytes) return undefined; const raw = await fs.readFile(realPath, "utf8"); const content = raw.replace(/\r\n/g, "\n"); return { relativePath: toPosix(path.relative(root, realPath)), content, contentType: contentTypeForExtension(extension), sha256: sha256(content), sizeBytes: Buffer.byteLength(content), mtimeMs: stat.mtimeMs, occurredAt: stat.mtime.toISOString(), ...(extension === ".md" || extension === ".mdx" ? explicitMarkdownId(content) : {}), }; } } function explicitMarkdownId(content: string): { explicitId?: string } { if (!content.startsWith("---\n")) return {}; const end = content.indexOf("\n---\n", 4); if (end === -1) return {}; try { const frontmatter: unknown = YAML.parse(content.slice(4, end)); if (!frontmatter || typeof frontmatter !== "object" || !("id" in frontmatter)) return {}; const id = String((frontmatter as { id: unknown }).id).trim(); return id ? { explicitId: id } : {}; } catch { return {}; } } function stableDocumentId(source: string, identity: string): string { return `doc_${sha256(`${source}\0${identity}`).slice(0, 32)}`; } function findUniqueUnseen(items: T[], seen: Set): T | undefined { const unseen = items.filter((item) => !seen.has(item.documentId)); return unseen.length === 1 ? unseen[0] : undefined; } function boundedDiff(previousPath: string, previous: string, currentPath: string, current: string, maxChars: number): string { const patch = createTwoFilesPatch(previousPath, currentPath, previous, current, "previous", "current", { context: 3 }); if (patch.length <= maxChars) return patch; return `${patch.slice(0, maxChars)}\n... diff truncated at ${maxChars} characters ...\n`; } function isInside(root: string, candidate: string): boolean { const relative = path.relative(root, candidate); return relative === "" || (!relative.startsWith("..") && !path.isAbsolute(relative)); } function toPosix(value: string): string { return value.split(path.sep).join("/"); } function contentTypeForExtension(extension: string): string { switch (extension) { case ".md": return "text/markdown"; case ".mdx": return "text/mdx"; case ".json": return "application/json"; case ".yaml": case ".yml": return "application/yaml"; case ".toml": return "application/toml"; default: return "text/plain"; } }