Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
15 kB · 394 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395import { 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<typeof telegramSpoolRecordSchema>;
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<TelegramSpoolIngestResult> { 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<ReturnType<fs.FileHandle["stat"]>>, 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<void> { 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<void> { 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<string> { 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<string, unknown>): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject;}