Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 282 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283import fs from "node:fs/promises";import path from "node:path";import { assertPrivateDestination, openPrivateDirectory, type PrivateDestinationOptions,} from "../security/private-files.js";
export type ConsumerCredentialProvider = "letta" | "letta-local" | "tinker" | "openai" | "openai-compatible";export type CredentialCompartment = | "telegram-webhook" | "x-webhook" | "x-management" | "x-user-management" | "fastmail-jmap" | "consumer" | "telegram-dispatcher" | "jetstream";
export interface SplitCredentialOptions extends PrivateDestinationOptions { consumerProviders?: ConsumerCredentialProvider[] | undefined; allowedRemovedVariableNames?: string[] | undefined;}
export interface CredentialCompartmentReceipt { outputDirectory: string; files: Record<CredentialCompartment, { path: string; variableNames: string[] }>; omittedVariableNames: string[]; removedVariableNames: string[];}
const outputNames: Record<CredentialCompartment, string> = { "telegram-webhook": "telegram-webhook.env", "x-webhook": "x-webhook.env", "x-management": "x-management.env", "x-user-management": "x-user-management.env", "fastmail-jmap": "fastmail-jmap.env", consumer: "consumer.env", "telegram-dispatcher": "telegram-dispatcher.env", jetstream: "jetstream.env",};
const consumerRuntimeVariables = new Set([ "THOUGHTSTREAM_ENABLE_COIL_PUBLIC_KNOWLEDGE", "THOUGHTSTREAM_PUBLIC_KNOWLEDGE_POLICY_PATH", "THOUGHTSTREAM_PUBLIC_KNOWLEDGE_CATALOG_ROOT",]);
export async function splitServiceCredentialFile( source: string, outputDirectory: string, options: SplitCredentialOptions = {},): Promise<CredentialCompartmentReceipt> { const absoluteOutput = path.resolve(outputDirectory); const destinations = Object.fromEntries(Object.entries(outputNames).map(([compartment, name]) => [ compartment, path.join(absoluteOutput, name), ])) as Record<CredentialCompartment, string>; for (const destination of Object.values(destinations)) { await assertPrivateDestination(destination, { publicContentRoots: options.publicContentRoots }); }
const assignments = parseEnvironmentAssignments(await fs.readFile(path.resolve(source), "utf8")); const providers = options.consumerProviders ?? ["letta"]; validateProviders(providers); validateRequiredAssignments(assignments, providers); const selected = selectAssignments(assignments, providers); const selectedNames = new Set(Object.values(selected).flat().map((assignment) => assignment.name)); const omittedVariableNames = [...assignments.keys()].filter((name) => !selectedNames.has(name)).sort(); const unassignedCredentials = omittedVariableNames.filter((name) => isCredentialLikeName(name) && !isKnownCredentialName(name)); if (unassignedCredentials.length > 0) { throw new Error(`Credential-like variables are not assigned to a service compartment: ${unassignedCredentials.join(", ")}`); } const removedVariableNames = await validateExistingAssignmentRemovals( destinations, selected, options.allowedRemovedVariableNames ?? [], );
const privateDirectory = await openPrivateDirectory(absoluteOutput, { publicContentRoots: options.publicContentRoots, beforeFinalize: options.beforeFinalize, }); try { for (const compartment of Object.keys(outputNames) as CredentialCompartment[]) { const lines = selected[compartment].map((assignment) => assignment.raw); const content = [ "# Generated by thought stream credential compartment splitter.", "# Contains service-specific assignments. Do not commit or print this file.", ...lines, "", ].join("\n"); await privateDirectory.write(outputNames[compartment], content); } } finally { await privateDirectory.close(); }
return { outputDirectory: absoluteOutput, files: Object.fromEntries((Object.keys(outputNames) as CredentialCompartment[]).map((compartment) => [ compartment, { path: destinations[compartment], variableNames: selected[compartment].map((assignment) => assignment.name), }, ])) as CredentialCompartmentReceipt["files"], omittedVariableNames, removedVariableNames, };}
async function validateExistingAssignmentRemovals( destinations: Record<CredentialCompartment, string>, selected: Record<CredentialCompartment, EnvironmentAssignment[]>, allowedRemovedVariableNames: string[],): Promise<string[]> { const allowed = new Set<string>(); for (const name of allowedRemovedVariableNames) { if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(name)) { throw new Error(`Invalid explicitly removed variable name: ${name}`); } if (allowed.has(name)) throw new Error(`Duplicate explicitly removed variable name: ${name}`); allowed.add(name); }
const removals: Array<{ compartment: CredentialCompartment; name: string }> = []; for (const compartment of Object.keys(destinations) as CredentialCompartment[]) { const destination = destinations[compartment]; const existing = await fs.lstat(destination).catch((error: NodeJS.ErrnoException) => { if (error.code === "ENOENT") return undefined; throw error; }); if (!existing) continue; if (!existing.isFile() || existing.isSymbolicLink()) { throw new Error(`Existing credential compartment must be a regular non-symlink file: ${compartment}`); } const previous = parseEnvironmentAssignments(await fs.readFile(destination, "utf8")); const nextNames = new Set(selected[compartment].map((assignment) => assignment.name)); for (const name of previous.keys()) { if (!nextNames.has(name)) removals.push({ compartment, name }); } }
const removalNames = [...new Set(removals.map((removal) => removal.name))].sort(); const unauthorized = removals.filter((removal) => !allowed.has(removal.name)); if (unauthorized.length > 0) { const detail = unauthorized .map((removal) => `${removal.compartment}:${removal.name}`) .sort() .join(", "); throw new Error(`Credential split would remove existing assignments without explicit authorization: ${detail}`); } const unusedAllowances = [...allowed].filter((name) => !removalNames.includes(name)).sort(); if (unusedAllowances.length > 0) { throw new Error(`Explicitly removed variable names do not match the current removal plan: ${unusedAllowances.join(", ")}`); } return removalNames;}
interface EnvironmentAssignment { name: string; raw: string;}
export function parseEnvironmentAssignments(contents: string): Map<string, EnvironmentAssignment> { const assignments = new Map<string, EnvironmentAssignment>(); for (const [index, original] of contents.split(/\r?\n/).entries()) { const line = original.trim(); if (!line || line.startsWith("#")) continue; const match = /^([A-Za-z_][A-Za-z0-9_]*)=(.*)$/.exec(line); if (!match) throw new Error(`Invalid environment assignment at line ${index + 1}`); const name = match[1]!; if (assignments.has(name)) throw new Error(`Duplicate environment assignment: ${name}`); assignments.set(name, { name, raw: line }); } return assignments;}
function selectAssignments( assignments: Map<string, EnvironmentAssignment>, providers: ConsumerCredentialProvider[],): Record<CredentialCompartment, EnvironmentAssignment[]> { const selected: Record<CredentialCompartment, EnvironmentAssignment[]> = { "telegram-webhook": [], "x-webhook": [], "x-management": [], "x-user-management": [], "fastmail-jmap": [], consumer: [], "telegram-dispatcher": [], jetstream: [], }; for (const assignment of assignments.values()) { if (isTelegramBotToken(assignment.name) || assignment.name === "THOUGHTSTREAM_TELEGRAM_API_BASE_URL") { selected["telegram-webhook"].push(assignment); selected["telegram-dispatcher"].push(assignment); } if (assignment.name === "THOUGHTSTREAM_TELEGRAM_WEBHOOK_SECRET") { selected["telegram-webhook"].push(assignment); } if (assignment.name === "THOUGHTSTREAM_X_CONSUMER_SECRET") { selected["x-webhook"].push(assignment); } if (assignment.name === "THOUGHTSTREAM_X_MANAGEMENT_BEARER_TOKEN" || assignment.name === "THOUGHTSTREAM_X_API_BASE_URL") { selected["x-management"].push(assignment); } if (assignment.name === "THOUGHTSTREAM_X_CAMERON_USER_ACCESS_TOKEN" || assignment.name === "THOUGHTSTREAM_X_API_BASE_URL") { selected["x-user-management"].push(assignment); } if (assignment.name === "FASTMAIL_API_KEY") selected["fastmail-jmap"].push(assignment); if (isConsumerVariable(assignment.name, providers)) selected.consumer.push(assignment); if (assignment.name === "THOUGHTSTREAM_JETSTREAM_URL") selected.jetstream.push(assignment); } for (const compartment of Object.keys(selected) as CredentialCompartment[]) { selected[compartment].sort((left, right) => left.name.localeCompare(right.name)); } return selected;}
function validateRequiredAssignments( assignments: Map<string, EnvironmentAssignment>, providers: ConsumerCredentialProvider[],): void { for (const name of ["THOUGHTSTREAM_TELEGRAM_BOT_TOKEN", "THOUGHTSTREAM_TELEGRAM_WEBHOOK_SECRET"]) { if (!assignments.has(name)) throw new Error(`Required service credential assignment is missing: ${name}`); } if (providers.includes("letta")) { if (!assignments.has("LETTA_API_KEY")) throw new Error("Required Letta consumer credential assignment is missing: LETTA_API_KEY"); } if ((providers.includes("letta") || providers.includes("letta-local")) && ![...assignments.keys()].some((name) => /^THOUGHTSTREAM_LETTA_[A-Z0-9_]+_AGENT_ID$/.test(name))) { throw new Error("A Letta consumer compartment requires at least one THOUGHTSTREAM_LETTA_*_AGENT_ID assignment"); } if (providers.includes("tinker") && !assignments.has("TINKER_API_KEY")) { throw new Error("Required Tinker consumer credential assignment is missing: TINKER_API_KEY"); } if (providers.includes("openai") && !assignments.has("OPENAI_API_KEY")) { throw new Error("Required OpenAI consumer credential assignment is missing: OPENAI_API_KEY"); } if (providers.includes("openai-compatible") && !assignments.has("THOUGHTSTREAM_MODEL_API_KEY")) { throw new Error("Required OpenAI-compatible consumer credential assignment is missing: THOUGHTSTREAM_MODEL_API_KEY"); }}
function validateProviders(providers: ConsumerCredentialProvider[]): void { if (providers.length === 0) throw new Error("At least one consumer credential provider is required"); if (new Set(providers).size !== providers.length) throw new Error("Consumer credential providers must be unique");}
function isTelegramBotToken(name: string): boolean { return name === "THOUGHTSTREAM_TELEGRAM_BOT_TOKEN";}
function isConsumerVariable(name: string, providers: ConsumerCredentialProvider[]): boolean { return consumerRuntimeVariables.has(name) || (providers.includes("letta") && name === "LETTA_API_KEY") || ((providers.includes("letta") || providers.includes("letta-local")) && /^THOUGHTSTREAM_LETTA_[A-Z0-9_]+$/.test(name)) || (providers.includes("tinker") && (name === "TINKER_API_KEY" || /^THOUGHTSTREAM_TINKER_[A-Z0-9_]+$/.test(name))) || (providers.includes("openai") && name === "OPENAI_API_KEY") || (providers.includes("openai-compatible") && (name === "THOUGHTSTREAM_MODEL_API_KEY" || /^THOUGHTSTREAM_MODEL_[A-Z0-9_]+$/.test(name)));}
function isKnownCredentialName(name: string): boolean { return name === "THOUGHTSTREAM_TELEGRAM_BOT_TOKEN" || name === "THOUGHTSTREAM_TELEGRAM_WEBHOOK_SECRET" || name === "THOUGHTSTREAM_X_CONSUMER_SECRET" || name === "THOUGHTSTREAM_X_MANAGEMENT_BEARER_TOKEN" || name === "THOUGHTSTREAM_X_CAMERON_USER_ACCESS_TOKEN" || name === "FASTMAIL_API_KEY" || name === "LETTA_API_KEY" || name === "TINKER_API_KEY" || name === "OPENAI_API_KEY" || name === "THOUGHTSTREAM_MODEL_API_KEY" || /^THOUGHTSTREAM_LETTA_[A-Z0-9_]+_AGENT_ID$/.test(name);}
function isCredentialLikeName(name: string): boolean { return /(?:API_KEY|TOKEN|SECRET|PASSWORD|CREDENTIAL|AGENT_ID|PRIVATE_KEY)(?:_|$)/.test(name);}