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 { X_ACTIVITY_SOURCE_EVENT_TYPE, X_ACTIVITY_SOURCE_SCHEMA_VERSION, X_WEBHOOK_REVISION, type XActivityPayload, type XActivityDirection, type XActivityEventType, type XSourceLane, xActivityDirectionSchema, xActivityEventTypeSchema, xActivityPayloadSchema, xNumericIdSchema, } from "./x-contract.js"; const upstreamFilterSchema = z.object({ user_id: xNumericIdSchema, direction: xActivityDirectionSchema.optional(), }).passthrough(); const upstreamEnvelopeSchema = z.object({ data: z.object({ event_uuid: z.string().min(1).max(200), filter: upstreamFilterSchema, event_type: z.string().min(1).max(100), tag: z.string().min(1).max(200), payload: z.unknown(), includes: z.unknown().optional(), }).passthrough(), }).passthrough(); const upstreamReferenceSchema = z.object({ type: z.enum(["replied_to", "quoted", "retweeted"]), id: xNumericIdSchema, }).passthrough(); const upstreamPositionedTagSchema = z.object({ start: z.number().int().nonnegative(), end: z.number().int().positive(), tag: z.string().min(1).max(200), }).passthrough(); const upstreamMentionSchema = z.object({ start: z.number().int().nonnegative(), end: z.number().int().positive(), username: z.string().min(1).max(100), id: xNumericIdSchema.optional(), }).passthrough(); const upstreamUrlSchema = z.object({ start: z.number().int().nonnegative(), end: z.number().int().positive(), url: z.string().min(1).max(2_048), expanded_url: z.string().min(1).max(8_192).optional(), display_url: z.string().min(1).max(2_048).optional(), unwound_url: z.string().min(1).max(8_192).optional(), }).passthrough(); const upstreamEntitiesSchema = z.object({ hashtags: z.array(upstreamPositionedTagSchema).max(100).optional(), cashtags: z.array(upstreamPositionedTagSchema).max(100).optional(), mentions: z.array(upstreamMentionSchema).max(100).optional(), urls: z.array(upstreamUrlSchema).max(100).optional(), }).passthrough(); const upstreamAttachmentsSchema = z.object({ media_keys: z.array(z.string().min(1).max(200)).max(50).optional(), poll_ids: z.array(z.string().min(1).max(200)).max(10).optional(), }).passthrough(); const upstreamPostCreateSchema = z.object({ id: xNumericIdSchema, author_id: xNumericIdSchema, text: z.string().max(100_000), created_at: z.iso.datetime(), conversation_id: xNumericIdSchema.optional(), in_reply_to_user_id: xNumericIdSchema.optional(), edit_history_tweet_ids: z.array(xNumericIdSchema).min(1).max(100).optional(), referenced_tweets: z.array(upstreamReferenceSchema).max(100).optional(), lang: z.string().min(1).max(32).optional(), reply_settings: z.string().min(1).max(100).optional(), possibly_sensitive: z.boolean().optional(), entities: upstreamEntitiesSchema.optional(), attachments: upstreamAttachmentsSchema.optional(), }).passthrough(); const upstreamPostDeleteSchema = z.object({ id: xNumericIdSchema, author_id: xNumericIdSchema, }).passthrough(); const upstreamLikeCreateSchema = z.object({ id: z.string().min(1).max(200), liked_tweet_id: xNumericIdSchema, liked_tweet_author_id: xNumericIdSchema, created_at: z.iso.datetime().optional(), timestamp_ms: z.string().regex(/^[0-9]{1,16}$/).optional(), }).passthrough(); export interface XExpectedSubscription { eventType: XActivityEventType; userId: string; direction?: XActivityDirection | undefined; tag: string; } export interface XActivityConnectorOptions { id: string; lane: XSourceLane; expectedSubscriptions: XExpectedSubscription[]; } export interface XActivityIngestResult { status: "updated" | "unchanged"; accepted: number; ignored: number; inserted: number; unchanged: number; events: ThoughtEvent[]; } export interface ParsedXActivityEnvelope { eventUuid: string; eventType: string; matchedUserId: string; direction?: XActivityDirection | undefined; tag: string; payload: unknown; } export class XActivityConnector { readonly kind = "x-webhook" as const; readonly id: string; readonly lane: XSourceLane; private readonly expectedByTag: Map; private readonly subscriptionSetHash: string; constructor(options: XActivityConnectorOptions) { this.id = required(options.id, "X activity connector id"); this.lane = options.lane; if (options.expectedSubscriptions.length === 0) { throw new Error("X activity connector requires at least one expected subscription"); } this.expectedByTag = new Map(); for (const subscription of options.expectedSubscriptions) { const normalized = { eventType: xActivityEventTypeSchema.parse(subscription.eventType), userId: xNumericIdSchema.parse(subscription.userId), ...(subscription.direction ? { direction: xActivityDirectionSchema.parse(subscription.direction) } : {}), tag: required(subscription.tag, "X subscription tag"), }; assertSubscriptionLane(this.lane, normalized); if (this.expectedByTag.has(normalized.tag)) { throw new Error(`Duplicate X subscription tag: ${normalized.tag}`); } this.expectedByTag.set(normalized.tag, normalized); } this.subscriptionSetHash = sha256(canonicalJson( [...this.expectedByTag.values()] .sort(compareSubscriptions) .map((item) => ({ eventType: item.eventType, userId: item.userId, ...(item.direction ? { direction: item.direction } : {}), tag: item.tag, })) as JsonObject[], )); } describe(): JsonObject { return { id: this.id, kind: this.kind, transport: "x-v2-webhook", revision: X_WEBHOOK_REVISION, lane: this.lane, expectedSubscriptionCount: this.expectedByTag.size, expectedSubscriptionsSha256: this.subscriptionSetHash, authority: "ingest-allowlisted-x-activity", }; } async ingest( store: JazzThoughtStore, rawEnvelope: unknown, bodySha256: string, receivedAt = new Date().toISOString(), ): Promise { const envelope = parseXActivityEnvelope(rawEnvelope); const attemptId = newId("x_webhook"); const correlationId = `${this.id}:x-event:${envelope.eventUuid}`; await store.appendEvent(this.connectorEvent("started", attemptId, receivedAt, { status: "started", transport: "x-v2-webhook", bodySha256, eventUuidSha256: sha256(envelope.eventUuid), })); try { const expected = this.expectedByTag.get(envelope.tag); const admitted = expected !== undefined && expected.eventType === envelope.eventType && expected.userId === envelope.matchedUserId && expected.direction === envelope.direction; const candidate = admitted ? this.eventCandidate(envelope, receivedAt, correlationId) : undefined; const completedAt = new Date().toISOString(); const candidates: EventCandidate[] = candidate ? [candidate] : []; const sourceCount = candidates.length; candidates.push(this.connectorEvent("completed", attemptId, completedAt, { status: admitted ? "updated" : "unchanged", accepted: admitted ? 1 : 0, ignored: admitted ? 0 : 1, bodySha256, eventUuidSha256: sha256(envelope.eventUuid), eventType: admitted ? envelope.eventType : "unconfigured", })); const batch = await store.appendProducerBatch(candidates); const sourceEvents = batch.events.slice(0, sourceCount); const insertedIds = new Set(batch.inserted.map((event) => event.id)); const events = sourceEvents.filter((event) => insertedIds.has(event.id)); return { status: admitted ? "updated" : "unchanged", accepted: admitted ? 1 : 0, ignored: admitted ? 0 : 1, inserted: events.length, unchanged: sourceCount - events.length, events, }; } catch (error) { const failedAt = new Date().toISOString(); await store.appendEvent(this.connectorEvent("failed", attemptId, failedAt, { status: "failed", phase: "x-webhook-ingest", bodySha256, eventUuidSha256: sha256(envelope.eventUuid), errorCode: classifyXIngestError(error), })); throw error; } } private eventCandidate( envelope: ParsedXActivityEnvelope, receivedAt: string, correlationId: string, ): EventCandidate { const payload = normalizeXActivityPayload(envelope, this.lane); return { type: X_ACTIVITY_SOURCE_EVENT_TYPE, schemaVersion: X_ACTIVITY_SOURCE_SCHEMA_VERSION, source: this.id, sourceKind: this.kind, externalId: envelope.eventUuid, idempotencyKey: sha256(canonicalJson({ source: this.id, revision: X_WEBHOOK_REVISION, eventUuid: envelope.eventUuid, })), occurredAt: payload.eventType === "post.create" ? payload.post.createdAt : payload.eventType === "like.create" && payload.like.eventAt ? payload.like.eventAt : receivedAt, actor: payload.eventType === "like.create" ? payload.like.actorId : payload.post.authorId, correlationId, privacy: sourcePrivacy(this.lane), payload: payload as unknown as JsonObject, }; } private connectorEvent( phase: "started" | "completed" | "failed", attemptId: string, at: string, payload: JsonObject, ): EventCandidate { return { type: phase === "failed" ? "stream.thought.connector.failed" : `stream.thought.connector.ingest.${phase}`, schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: attemptId, idempotencyKey: `${attemptId}:${phase}`, occurredAt: at, actor: this.id, correlationId: attemptId, privacy: sourcePrivacy(this.lane), payload, }; } } export function parseXActivityEnvelope(value: unknown): ParsedXActivityEnvelope { const parsed = upstreamEnvelopeSchema.parse(value).data; return { eventUuid: parsed.event_uuid, eventType: parsed.event_type, matchedUserId: parsed.filter.user_id, ...(parsed.filter.direction ? { direction: parsed.filter.direction } : {}), tag: parsed.tag, payload: parsed.payload, }; } export function isPermanentXActivityError(error: unknown): boolean { return error instanceof z.ZodError || (error instanceof Error && error.message.startsWith("Unsupported admitted X activity")); } function normalizeXActivityPayload( envelope: ParsedXActivityEnvelope, lane: XSourceLane, ): XActivityPayload { const common = { eventUuid: envelope.eventUuid, matchedUserId: envelope.matchedUserId, subscriptionTag: envelope.tag, lane, transport: { revision: X_WEBHOOK_REVISION as typeof X_WEBHOOK_REVISION, signatureVerified: true as const, kind: "x-v2-webhook" as const, }, }; if (envelope.eventType === "post.create") { const post = upstreamPostCreateSchema.parse(envelope.payload); const normalized = { ...common, eventType: "post.create", post: compact({ postId: post.id, authorId: post.author_id, text: post.text, createdAt: post.created_at, conversationId: post.conversation_id, inReplyToUserId: post.in_reply_to_user_id, editHistoryPostIds: post.edit_history_tweet_ids, referencedPosts: post.referenced_tweets?.map((reference) => ({ type: reference.type, postId: reference.id })), language: post.lang, replySettings: post.reply_settings, possiblySensitive: post.possibly_sensitive, entities: normalizeEntities(post.entities), attachments: normalizeAttachments(post.attachments), }), }; return xActivityPayloadSchema.parse(normalized as unknown as JsonObject) as unknown as XActivityPayload; } if (envelope.eventType === "post.delete") { const post = upstreamPostDeleteSchema.parse(envelope.payload); const normalized = { ...common, eventType: "post.delete", post: { postId: post.id, authorId: post.author_id }, }; return xActivityPayloadSchema.parse(normalized as unknown as JsonObject) as unknown as XActivityPayload; } if (envelope.eventType === "like.create") { if (lane !== "personal-private" || envelope.direction !== "outbound") { throw new Error("Unsupported admitted X activity event type or direction: like.create"); } const like = upstreamLikeCreateSchema.parse(envelope.payload); const normalized = { ...common, eventType: "like.create", direction: "outbound", like: compact({ likeId: like.id, actorId: envelope.matchedUserId, postId: like.liked_tweet_id, postAuthorId: like.liked_tweet_author_id, likedPostCreatedAt: like.created_at, eventTimestampMs: like.timestamp_ms, eventAt: eventTimestamp(like.timestamp_ms), }), }; return xActivityPayloadSchema.parse(normalized as unknown as JsonObject) as unknown as XActivityPayload; } throw new Error(`Unsupported admitted X activity event type: ${envelope.eventType}`); } function normalizeEntities(value: z.infer | undefined) { if (!value) return undefined; return compact({ hashtags: value.hashtags?.map((item) => ({ start: item.start, end: item.end, tag: item.tag })), cashtags: value.cashtags?.map((item) => ({ start: item.start, end: item.end, tag: item.tag })), mentions: value.mentions?.map((item) => compact({ start: item.start, end: item.end, username: item.username, userId: item.id, })), urls: value.urls?.map((item) => compact({ start: item.start, end: item.end, url: item.url, expandedUrl: item.expanded_url, displayUrl: item.display_url, unwoundUrl: item.unwound_url, })), }); } function normalizeAttachments(value: z.infer | undefined) { if (!value) return undefined; return compact({ mediaKeys: value.media_keys, pollIds: value.poll_ids }); } function compact>(value: T): T { return Object.fromEntries(Object.entries(value).filter(([, item]) => item !== undefined)) as T; } function compareSubscriptions(left: XExpectedSubscription, right: XExpectedSubscription): number { return left.eventType.localeCompare(right.eventType) || left.userId.localeCompare(right.userId) || (left.direction ?? "").localeCompare(right.direction ?? "") || left.tag.localeCompare(right.tag); } function assertSubscriptionLane(lane: XSourceLane, subscription: XExpectedSubscription): void { if (lane === "personal-private") { if (subscription.eventType !== "like.create" || subscription.direction !== "outbound") { throw new Error("X personal-private sources admit only outbound like.create subscriptions"); } return; } if (subscription.eventType === "like.create" || subscription.direction !== undefined) { throw new Error("X public sources admit only directionless post subscriptions"); } } function sourcePrivacy(lane: XSourceLane): "public-source" | "sensitive" { return lane === "personal-private" ? "sensitive" : "public-source"; } function eventTimestamp(value: string | undefined): string | undefined { if (!value) return undefined; const timestamp = Number(value); if (!Number.isSafeInteger(timestamp) || timestamp < 0) throw new Error("X like timestamp_ms is invalid"); const date = new Date(timestamp); if (!Number.isFinite(date.getTime())) throw new Error("X like timestamp_ms is outside the supported date range"); return date.toISOString(); } function classifyXIngestError(error: unknown): string { if (isPermanentXActivityError(error)) return "x-payload-invalid"; return "x-ingest-failed"; } function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized; }