Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
16 kB · 458 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459import { 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<string, XExpectedSubscription>; 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<XActivityIngestResult> { 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<typeof upstreamEntitiesSchema> | 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<typeof upstreamAttachmentsSchema> | undefined) { if (!value) return undefined; return compact({ mediaKeys: value.media_keys, pollIds: value.poll_ids });}
function compact<T extends Record<string, unknown>>(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;}