import { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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; 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): 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, right: ReturnType): 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; }