Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
16 kB · 410 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411import { z } from "zod";import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import type { XExpectedSubscription } from "./x-activity.js";import { xActivityEventTypeSchema, xNumericIdSchema, xUsernameSchema } from "./x-contract.js";
const webhookSchema = z.object({ id: xNumericIdSchema, url: z.string().url(), valid: z.boolean(), created_at: z.string().optional(),}).passthrough();
const subscriptionSchema = z.object({ subscription_id: xNumericIdSchema, event_type: z.string().min(1).max(100), filter: z.object({ user_id: xNumericIdSchema.optional(), direction: z.string().max(100).optional(), keyword: z.string().max(500).optional(), }).passthrough(), tag: z.string().max(200).optional(), webhook_id: xNumericIdSchema.optional(), created_at: z.string().optional(), updated_at: z.string().optional(),}).passthrough();
const userSchema = z.object({ id: xNumericIdSchema, username: z.string().min(1).max(100), name: z.string().min(1).max(500),}).passthrough();
export interface XWebhookRecord { id: string; url: string; valid: boolean; createdAt?: string | undefined;}
export interface XSubscriptionRecord { id: string; eventType: string; userId?: string | undefined; direction?: string | undefined; keyword?: string | undefined; tag?: string | undefined; webhookId?: string | undefined; createdAt?: string | undefined; updatedAt?: string | undefined;}
export interface XUserRecord { id: string; username: string; name: string;}
export interface XApiClientOptions { bearerToken: string; baseUrl?: string | undefined; timeoutMs?: number | undefined; fetchImpl?: typeof fetch | undefined;}
export class XApiClient { private readonly bearerToken: string; private readonly baseUrl: string; private readonly timeoutMs: number; private readonly fetchImpl: typeof fetch;
constructor(options: XApiClientOptions) { this.bearerToken = required(options.bearerToken, "X bearer token"); const base = new URL(options.baseUrl ?? "https://api.x.com"); if (base.username || base.password || base.search || base.hash || base.pathname !== "/") { throw new Error("X API base URL must contain only scheme, host, and optional test port"); } if (base.protocol !== "https:" && !(base.protocol === "http:" && ["127.0.0.1", "::1", "localhost"].includes(base.hostname))) { throw new Error("X API base URL must use HTTPS or loopback HTTP"); } this.baseUrl = base.href.replace(/\/$/, ""); this.timeoutMs = boundedInteger(options.timeoutMs ?? 15_000, "X API timeoutMs", 1_000, 120_000); this.fetchImpl = options.fetchImpl ?? fetch; }
async listWebhooks(signal?: AbortSignal): Promise<XWebhookRecord[]> { const response = await this.request("GET", "/2/webhooks", undefined, signal); const data = z.object({ data: z.array(webhookSchema).default([]) }).passthrough().parse(response).data; return data.map((webhook) => ({ id: webhook.id, url: webhook.url, valid: webhook.valid, ...(webhook.created_at ? { createdAt: webhook.created_at } : {}), })); }
async createWebhook(url: string, signal?: AbortSignal): Promise<XWebhookRecord> { const response = await this.request("POST", "/2/webhooks", { url: exactHttpsUrl(url) }, signal); const webhook = z.object({ data: webhookSchema }).passthrough().parse(response).data; return { id: webhook.id, url: webhook.url, valid: webhook.valid, ...(webhook.created_at ? { createdAt: webhook.created_at } : {}), }; }
async validateWebhook(webhookId: string, signal?: AbortSignal): Promise<boolean> { const response = await this.request("PUT", `/2/webhooks/${xNumericIdSchema.parse(webhookId)}`, undefined, signal); return z.object({ data: z.object({ valid: z.boolean() }).passthrough() }).passthrough().parse(response).data.valid; }
async deleteWebhook(webhookId: string, signal?: AbortSignal): Promise<boolean> { const response = await this.request("DELETE", `/2/webhooks/${xNumericIdSchema.parse(webhookId)}`, undefined, signal); return z.object({ data: z.object({ deleted: z.boolean() }).passthrough() }).passthrough().parse(response).data.deleted; }
async listSubscriptions(signal?: AbortSignal): Promise<XSubscriptionRecord[]> { const response = await this.request("GET", "/2/activity/subscriptions", undefined, signal); const parsed = z.object({ data: z.array(subscriptionSchema).max(1_500).default([]) }).passthrough().parse(response); return parsed.data.map(normalizeSubscription); }
async createSubscription( subscription: XExpectedSubscription & { webhookId: string }, signal?: AbortSignal, ): Promise<XSubscriptionRecord> { const response = await this.request("POST", "/2/activity/subscriptions", { event_type: xActivityEventTypeSchema.parse(subscription.eventType), filter: { user_id: xNumericIdSchema.parse(subscription.userId), ...(subscription.direction ? { direction: subscription.direction } : {}), }, tag: boundedTag(subscription.tag), webhook_id: xNumericIdSchema.parse(subscription.webhookId), }, signal); const parsed = z.object({ data: z.union([subscriptionSchema, z.array(subscriptionSchema).min(1)]), }).passthrough().safeParse(response); if (parsed.success) { const data = parsed.data.data; return normalizeSubscription(Array.isArray(data) ? data[0]! : data); } const matches = (await this.listSubscriptions(signal)).filter((candidate) => ( candidate.eventType === subscription.eventType && candidate.userId === subscription.userId && candidate.direction === subscription.direction && candidate.tag === subscription.tag && candidate.webhookId === subscription.webhookId )); if (matches.length !== 1) { throw new Error("X subscription creation response was not a subscription and exact readback did not converge"); } return matches[0]!; }
async updateSubscription( subscriptionId: string, update: { tag: string; webhookId: string }, signal?: AbortSignal, ): Promise<XSubscriptionRecord> { const response = await this.request("PUT", `/2/activity/subscriptions/${xNumericIdSchema.parse(subscriptionId)}`, { tag: boundedTag(update.tag), webhook_id: xNumericIdSchema.parse(update.webhookId), }, signal); const data = z.object({ data: subscriptionSchema }).passthrough().parse(response).data; return normalizeSubscription(data); }
async deleteSubscription(subscriptionId: string, signal?: AbortSignal): Promise<boolean> { const response = await this.request("DELETE", `/2/activity/subscriptions/${xNumericIdSchema.parse(subscriptionId)}`, undefined, signal); return z.object({ data: z.object({ deleted: z.boolean() }).passthrough() }).passthrough().parse(response).data.deleted; }
async createReplay( webhookId: string, fromDate: string, toDate: string, signal?: AbortSignal, ): Promise<{ jobId: string; createdAt: string }> { const timestamp = z.string().regex(/^[0-9]{12}$/); const response = await this.request("POST", "/2/webhooks/replay", { webhook_id: xNumericIdSchema.parse(webhookId), from_date: timestamp.parse(fromDate), to_date: timestamp.parse(toDate), }, signal); const data = z.object({ data: z.object({ job_id: z.string().min(1).max(200), created_at: z.string().min(1).max(200) }).passthrough(), }).passthrough().parse(response).data; return { jobId: data.job_id, createdAt: data.created_at }; }
async lookupUser(username: string, signal?: AbortSignal): Promise<XUserRecord> { const normalized = xUsernameSchema.parse(username.trim().replace(/^@/, "")); const response = await this.request("GET", `/2/users/by/username/${encodeURIComponent(normalized)}`, undefined, signal); const user = z.object({ data: userSchema }).passthrough().parse(response).data; return { id: user.id, username: user.username, name: user.name }; }
private async request(method: string, pathname: string, body?: object, signal?: AbortSignal): Promise<unknown> { const timeout = AbortSignal.timeout(this.timeoutMs); const boundedSignal = signal ? AbortSignal.any([signal, timeout]) : timeout; const response = await this.fetchImpl(`${this.baseUrl}${pathname}`, { method, headers: { authorization: `Bearer ${this.bearerToken}`, accept: "application/json", ...(body ? { "content-type": "application/json" } : {}), }, ...(body ? { body: JSON.stringify(body) } : {}), redirect: "error", signal: boundedSignal, }); if (!response.ok) { throw new Error(`X API ${method} ${pathname.split("?", 1)[0]} failed with HTTP ${response.status}`); } const text = await response.text(); if (!text) return {}; try { return JSON.parse(text) as unknown; } catch { throw new Error(`X API ${method} ${pathname.split("?", 1)[0]} returned invalid JSON`); } }}
export interface XSubscriptionPlan { sourceId: string; webhookId: string; planHash: string; create: Array<XExpectedSubscription & { webhookId: string }>; update: Array<{ subscriptionId: string; tag: string; webhookId: string }>; delete: Array<{ subscriptionId: string }>; unchanged: string[]; conflicts: string[]; unmanagedCount: number;}
export function planXSubscriptions( sourceId: string, webhookId: string, desired: XExpectedSubscription[], live: XSubscriptionRecord[],): XSubscriptionPlan { const managedPrefix = `thoughtstream:${sourceId}:`; const normalizedDesired = [...desired].sort(compareDesired); const managed = live.filter((subscription) => subscription.tag?.startsWith(managedPrefix)); const unmanaged = live.filter((subscription) => !subscription.tag?.startsWith(managedPrefix)); const create: XSubscriptionPlan["create"] = []; const update: XSubscriptionPlan["update"] = []; const remove: XSubscriptionPlan["delete"] = []; const unchanged: string[] = []; const conflicts: string[] = []; const desiredKeys = new Set(normalizedDesired.map(subscriptionIdentity));
for (const item of normalizedDesired) { const sameIdentity = live.filter((subscription) => ( subscription.eventType === item.eventType && subscription.userId === item.userId && subscription.direction === item.direction )); const managedMatches = sameIdentity.filter((subscription) => subscription.tag?.startsWith(managedPrefix)); const unmanagedMatches = sameIdentity.filter((subscription) => !subscription.tag?.startsWith(managedPrefix)); if (unmanagedMatches.length > 0) { conflicts.push(`unmanaged-existing:${item.eventType}:${item.userId}:${item.direction ?? "none"}`); continue; } if (managedMatches.length > 1) { conflicts.push(`duplicate-managed:${item.eventType}:${item.userId}:${item.direction ?? "none"}`); continue; } const existing = managedMatches[0]; if (!existing) { create.push({ ...item, webhookId }); } else if (existing.tag === item.tag && existing.webhookId === webhookId) { unchanged.push(existing.id); } else { update.push({ subscriptionId: existing.id, tag: item.tag, webhookId }); } } for (const subscription of managed) { if (!subscription.userId || !desiredKeys.has(subscriptionRecordIdentity(subscription))) { remove.push({ subscriptionId: subscription.id }); } } const planBody = { sourceId, webhookId, desired: normalizedDesired, live: live.map(normalizePlanSubscription).sort(comparePlanSubscriptions), create, update, delete: remove, conflicts: [...conflicts].sort(), }; return { sourceId, webhookId, planHash: sha256(canonicalJson(planBody as unknown as JsonObject)), create, update, delete: remove, unchanged: unchanged.sort(), conflicts: [...conflicts].sort(), unmanagedCount: unmanaged.length, };}
export function assertBoundedXReplayWindow(fromDate: string, toDate: string, now = Date.now()): void { const from = parseXReplayTimestamp(fromDate); const to = parseXReplayTimestamp(toDate); const currentUtcMinute = Math.floor(now / 60_000) * 60_000; if (from >= to) throw new Error("X replay --from must precede --to"); if (to > currentUtcMinute) throw new Error("X replay --to cannot be in the future"); if (to - from > 24 * 60 * 60_000) throw new Error("X replay window cannot exceed 24 hours"); if (from < currentUtcMinute - 24 * 60 * 60_000) { throw new Error("X replay --from is outside the preceding 24-hour replay horizon"); }}
function parseXReplayTimestamp(value: string): number { if (!/^[0-9]{12}$/.test(value)) throw new Error("X replay timestamps must use YYYYMMDDHHmm UTC"); const year = Number(value.slice(0, 4)); const month = Number(value.slice(4, 6)); const day = Number(value.slice(6, 8)); const hour = Number(value.slice(8, 10)); const minute = Number(value.slice(10, 12)); const timestamp = Date.UTC(year, month - 1, day, hour, minute); const date = new Date(timestamp); if (date.getUTCFullYear() !== year || date.getUTCMonth() !== month - 1 || date.getUTCDate() !== day || date.getUTCHours() !== hour || date.getUTCMinutes() !== minute) { throw new Error("X replay timestamp is not a valid UTC minute"); } return timestamp;}
function normalizeSubscription(value: z.infer<typeof subscriptionSchema>): XSubscriptionRecord { return { id: value.subscription_id, eventType: value.event_type, ...(value.filter.user_id ? { userId: value.filter.user_id } : {}), ...(value.filter.direction ? { direction: value.filter.direction } : {}), ...(value.filter.keyword ? { keyword: value.filter.keyword } : {}), ...(value.tag ? { tag: value.tag } : {}), ...(value.webhook_id ? { webhookId: value.webhook_id } : {}), ...(value.created_at ? { createdAt: value.created_at } : {}), ...(value.updated_at ? { updatedAt: value.updated_at } : {}), };}
function normalizePlanSubscription(value: XSubscriptionRecord) { return { id: value.id, eventType: value.eventType, userId: value.userId ?? "", direction: value.direction ?? "", keyword: value.keyword ?? "", tag: value.tag ?? "", webhookId: value.webhookId ?? "", };}
function comparePlanSubscriptions(left: ReturnType<typeof normalizePlanSubscription>, right: ReturnType<typeof normalizePlanSubscription>): number { return left.eventType.localeCompare(right.eventType) || left.userId.localeCompare(right.userId) || left.direction.localeCompare(right.direction) || left.tag.localeCompare(right.tag) || left.id.localeCompare(right.id);}
function subscriptionIdentity(value: XExpectedSubscription): string { return `${value.eventType}\u0000${value.userId}\u0000${value.direction ?? ""}`;}
function subscriptionRecordIdentity(value: XSubscriptionRecord): string { return `${value.eventType}\u0000${value.userId ?? ""}\u0000${value.direction ?? ""}`;}
function compareDesired(left: XExpectedSubscription, right: XExpectedSubscription): number { return subscriptionIdentity(left).localeCompare(subscriptionIdentity(right)) || left.tag.localeCompare(right.tag);}
function exactHttpsUrl(value: string): string { const url = new URL(value); if (url.protocol !== "https:" || url.username || url.password || url.port || url.search || url.hash) { throw new Error("X webhook URL must be an exact HTTPS URL without credentials, port, query, or fragment"); } return url.href;}
function boundedTag(value: string): string { const tag = required(value, "X subscription tag"); if (tag.length > 200) throw new Error("X subscription tag exceeds 200 characters"); return tag;}
function boundedInteger(value: number, label: string, minimum: number, maximum: number): number { if (!Number.isSafeInteger(value) || value < minimum || value > maximum) { throw new Error(`${label} must be an integer between ${minimum} and ${maximum}`); } return value;}
function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized;}