import { constants } from "node:fs"; import { createHash } from "node:crypto"; import fs from "node:fs/promises"; import path from "node:path"; import { TextDecoder } from "node:util"; import { z } from "zod"; import { canonicalJson, sha256, type JsonObject } from "../core/json.js"; import { newId } from "../core/ids.js"; import type { EventCandidate, ThoughtEvent } from "../events/types.js"; import type { JazzThoughtStore } from "../jazz/store.js"; import type { SourceCursor } from "../store/types.js"; const SPOOL_REVISION = "telegram-spool-ndjson-v1"; const utf8 = new TextDecoder("utf-8", { fatal: true }); const attachmentSchema = z.object({ id: z.string().min(1).optional(), name: z.string().min(1).optional(), mimeType: z.string().min(1).optional(), sizeBytes: z.number().int().nonnegative().optional(), kind: z.enum(["image", "file", "audio", "video"]), reference: z.string().min(1), }).strict(); const routeSchema = z.object({ agentId: z.string().min(1), conversationId: z.string().min(1), }).strict(); const telegramSpoolRecordSchema = z.object({ schemaVersion: z.literal(1), accountId: z.string().min(1), chatId: z.string().min(1), messageId: z.string().min(1), senderId: z.string().min(1), senderName: z.string().min(1).optional(), chatType: z.enum(["direct", "channel"]).optional(), occurredAt: z.iso.datetime({ offset: true }), editedAt: z.iso.datetime({ offset: true }).optional(), text: z.string(), threadId: z.string().min(1).optional(), replyToMessageId: z.string().min(1).optional(), route: routeSchema, attachments: z.array(attachmentSchema).max(32).optional(), }).strict(); export type TelegramSpoolRecord = z.infer; export interface TelegramSpoolConnectorOptions { id: string; file: string; maxRecords?: number; maxReadBytes?: number; } export interface TelegramSpoolIngestResult { status: "updated" | "unchanged"; inserted: number; unchanged: number; consumedRecords: number; partialTrailingBytes: number; events: ThoughtEvent[]; cursor: SourceCursor; } export class TelegramSpoolConnector { readonly kind = "telegram" as const; readonly id: string; private readonly file: string; private readonly filePathHash: string; private readonly maxRecords: number; private readonly maxReadBytes: number; constructor(options: TelegramSpoolConnectorOptions) { this.id = required(options.id, "Telegram spool connector id"); this.file = path.resolve(options.file); this.filePathHash = sha256(this.file); this.maxRecords = boundedPositiveInteger(options.maxRecords ?? 1_000, "maxRecords", 100_000); this.maxReadBytes = boundedPositiveInteger(options.maxReadBytes ?? 8 * 1024 * 1024, "maxReadBytes", 64 * 1024 * 1024); } describe(): JsonObject { return { id: this.id, kind: this.kind, transport: "append-only-ndjson-spool", spoolRevision: SPOOL_REVISION, filePathHash: this.filePathHash, maxRecords: this.maxRecords, maxReadBytes: this.maxReadBytes, authority: "read-only", }; } async ingest(store: JazzThoughtStore): Promise { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const startedAt = new Date().toISOString(); const correlationId = newId("telegram_spool"); await this.appendConnectorEvent(store, "started", correlationId, startedAt, { status: "started", transport: "append-only-ndjson-spool", filePathHash: this.filePathHash, }); let handle: fs.FileHandle | undefined; try { this.assertCursorConfiguration(prior); handle = await fs.open(this.file, constants.O_RDONLY | constants.O_NOFOLLOW); const stat = await handle.stat(); if (!stat.isFile()) throw new Error("Telegram spool path must be a regular file"); const offset = cursorOffset(prior); this.assertFileContinuity(prior, stat, offset); await this.assertBoundary(handle, prior, offset); const unreadBytes = stat.size - offset; const readBytes = Math.min(unreadBytes, this.maxReadBytes); const buffer = Buffer.alloc(readBytes); if (readBytes > 0) { const result = await handle.read(buffer, 0, readBytes, offset); if (result.bytesRead !== readBytes) throw new Error("Telegram spool changed while it was being read"); } const { records, consumedBytes, partialTrailingBytes } = parseRecordWindow(buffer, this.maxRecords, offset); if (records.length === 0 && unreadBytes > this.maxReadBytes && consumedBytes === 0) { throw new Error(`Telegram spool has no complete record within maxReadBytes=${this.maxReadBytes}`); } const candidates: EventCandidate[] = records.map((record) => this.eventCandidate(record, correlationId)); const completedAt = new Date().toISOString(); const newOffset = offset + consumedBytes; const cursor: SourceCursor = { id: cursorId, source: this.id, cursor: { spoolRevision: SPOOL_REVISION, filePathHash: this.filePathHash, device: String(stat.dev), inode: String(stat.ino), byteOffset: newOffset, recordCount: cursorRecordCount(prior) + records.length, prefixHash: await prefixHash(handle, newOffset), }, lastSuccessAt: completedAt, ...(prior?.lastFailureAt ? { lastFailureAt: prior.lastFailureAt } : {}), updatedAt: completedAt, }; const sourceCount = candidates.length; candidates.push(this.cursorEvent(correlationId, completedAt, cursor)); if (prior?.lastFailureAt && (!prior.lastSuccessAt || prior.lastFailureAt > prior.lastSuccessAt)) { candidates.push(this.connectorEvent("recovered", correlationId, completedAt, { status: "recovered", previousFailureAt: prior.lastFailureAt, })); } const status = records.length > 0 ? "updated" : "unchanged"; candidates.push(this.connectorEvent("completed", correlationId, completedAt, { status, consumedRecords: records.length, partialTrailingBytes, })); const batch = await store.appendProducerBatch(candidates, cursor); const offeredEvents = batch.events.slice(0, sourceCount); const insertedIds = new Set(batch.inserted.map((event) => event.id)); const events = offeredEvents.filter((event) => insertedIds.has(event.id)); const inserted = events.length; const unchanged = sourceCount - inserted; return { status, inserted, unchanged, consumedRecords: records.length, partialTrailingBytes, events, cursor, }; } catch (error) { const failedAt = new Date().toISOString(); const message = error instanceof Error ? error.message : String(error); const failedCursor: SourceCursor = { id: cursorId, source: this.id, cursor: prior?.cursor ?? { spoolRevision: SPOOL_REVISION, filePathHash: this.filePathHash, byteOffset: 0, recordCount: 0, }, ...(prior?.lastSuccessAt ? { lastSuccessAt: prior.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: message, updatedAt: failedAt, }; await store.appendProducerBatch([this.connectorEvent("failed", correlationId, failedAt, { status: "failed", error: message, transport: "append-only-ndjson-spool", filePathHash: this.filePathHash, })], failedCursor); throw error; } finally { await handle?.close(); } } private eventCandidate(record: TelegramSpoolRecord, correlationId: string) { const revision = record.editedAt ?? "original"; const externalId = `${record.accountId}:${record.chatId}:${record.messageId}`; const payload = compact({ accountId: record.accountId, chatId: record.chatId, messageId: record.messageId, senderId: record.senderId, senderName: record.senderName, chatType: record.chatType, occurredAt: record.occurredAt, editedAt: record.editedAt, text: record.text, threadId: record.threadId, replyToMessageId: record.replyToMessageId, route: record.route, attachments: record.attachments, spoolSchemaVersion: record.schemaVersion, }); return { type: "stream.thought.source.telegram.message", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId, idempotencyKey: sha256(canonicalJson({ accountId: record.accountId, chatId: record.chatId, messageId: record.messageId, revision })), occurredAt: record.editedAt ?? record.occurredAt, actor: record.senderId, correlationId, privacy: "sensitive" as const, payload, }; } private assertCursorConfiguration(cursor: SourceCursor | undefined): void { const revision = cursor?.cursor.spoolRevision; if (revision !== undefined && revision !== SPOOL_REVISION) throw new Error("Telegram spool cursor revision does not match connector configuration"); const pathHash = cursor?.cursor.filePathHash; if (pathHash !== undefined && pathHash !== this.filePathHash) throw new Error("Telegram spool cursor belongs to a different file path"); } private assertFileContinuity(cursor: SourceCursor | undefined, stat: Awaited>, offset: number): void { if (stat.size < offset) throw new Error("Telegram spool was truncated behind its durable cursor"); const device = cursor?.cursor.device; const inode = cursor?.cursor.inode; if (device !== undefined && device !== String(stat.dev)) throw new Error("Telegram spool file identity changed behind its durable cursor"); if (inode !== undefined && inode !== String(stat.ino)) throw new Error("Telegram spool file identity changed behind its durable cursor"); } private async assertBoundary(handle: fs.FileHandle, cursor: SourceCursor | undefined, offset: number): Promise { const expected = cursor?.cursor.prefixHash; if (expected !== undefined && expected !== await prefixHash(handle, offset)) { throw new Error("Telegram spool content changed behind its durable cursor"); } } private cursorEvent(correlationId: string, at: string, cursor: SourceCursor): EventCandidate { return { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: cursor.id, idempotencyKey: `${correlationId}:cursor`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload: { source: this.id, cursorId: cursor.id, cursor: cursor.cursor }, }; } private connectorEvent( phase: "started" | "completed" | "failed" | "recovered", correlationId: string, at: string, payload: JsonObject, ): EventCandidate { const type = phase === "failed" ? "stream.thought.connector.failed" : phase === "recovered" ? "stream.thought.connector.recovered" : `stream.thought.connector.poll.${phase}`; return { type, schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: correlationId, idempotencyKey: `${correlationId}:${phase}`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload, }; } private async appendConnectorEvent( store: JazzThoughtStore, phase: "started" | "completed" | "failed" | "recovered", correlationId: string, at: string, payload: JsonObject, ): Promise { await store.appendEvent(this.connectorEvent(phase, correlationId, at, payload)); } } export function parseTelegramSpoolRecord(value: string, lineNumber?: number): TelegramSpoolRecord { try { return telegramSpoolRecordSchema.parse(JSON.parse(value) as unknown); } catch (error) { const prefix = lineNumber === undefined ? "Invalid Telegram spool record" : `Invalid Telegram spool record at line ${lineNumber}`; throw new Error(`${prefix}: ${error instanceof Error ? error.message : String(error)}`); } } function parseRecordWindow(buffer: Buffer, maxRecords: number, priorOffset: number): { records: TelegramSpoolRecord[]; consumedBytes: number; partialTrailingBytes: number; } { const records: TelegramSpoolRecord[] = []; let start = 0; let consumedBytes = 0; let physicalLine = 0; while (records.length < maxRecords) { const newline = buffer.indexOf(0x0a, start); if (newline < 0) break; physicalLine += 1; let end = newline; if (end > start && buffer[end - 1] === 0x0d) end -= 1; const lineBuffer = buffer.subarray(start, end); if (lineBuffer.length === 0) throw new Error(`Invalid Telegram spool record after byte ${priorOffset + start}: blank lines are not allowed`); const line = utf8.decode(lineBuffer); records.push(parseTelegramSpoolRecord(line, physicalLine)); consumedBytes = newline + 1; start = newline + 1; } return { records, consumedBytes, partialTrailingBytes: buffer.length - consumedBytes, }; } async function prefixHash(handle: fs.FileHandle, offset: number): Promise { const hash = createHash("sha256"); const buffer = Buffer.alloc(Math.min(1024 * 1024, Math.max(offset, 1))); let position = 0; while (position < offset) { const length = Math.min(buffer.length, offset - position); const result = await handle.read(buffer, 0, length, position); if (result.bytesRead !== length) throw new Error("Telegram spool changed while its consumed prefix was verified"); hash.update(buffer.subarray(0, length)); position += length; } return hash.digest("hex"); } function cursorOffset(cursor: SourceCursor | undefined): number { const value = cursor?.cursor.byteOffset; if (value === undefined) return 0; if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) throw new Error("Telegram spool cursor byteOffset must be a nonnegative safe integer"); return value; } function cursorRecordCount(cursor: SourceCursor | undefined): number { const value = cursor?.cursor.recordCount; if (value === undefined) return 0; if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) throw new Error("Telegram spool cursor recordCount must be a nonnegative safe integer"); return value; } function boundedPositiveInteger(value: number, label: string, maximum: number): number { if (!Number.isSafeInteger(value) || value <= 0 || value > maximum) throw new Error(`${label} must be a positive integer <= ${maximum}`); return value; } function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized; } function compact(value: Record): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject; }