import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test, vi } from "vitest"; import { XActivityConnector } from "../src/connectors/x-activity.js"; import { xCrcResponseToken, startXWebhookServer, xWebhookSignature, type XCrcDiagnostic, } from "../src/connectors/x-webhook.js"; import { sha256 } from "../src/core/json.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import { buildRootActivity, presentActivityEvent } from "../src/projections/activity.js"; import { temporaryProject, testStore } from "./helpers.js"; const roots: string[] = []; const stores: JazzThoughtStore[] = []; afterEach(async () => { vi.restoreAllMocks(); await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); }); describe("X Activity webhook", () => { test("normalizes the captured official post create and delete fixtures", async () => { const root = await temporaryProject("thoughtstream-x-captured-fixture-"); roots.push(root); const store = await testStore(root); stores.push(store); const connector = new XActivityConnector({ id: "x:official-fixture", lane: "public-watchlist", expectedSubscriptions: [ { eventType: "post.create", userId: "1111111111111111111", tag: "posts" }, { eventType: "post.delete", userId: "1111111111111111111", tag: "post-deletes" }, ], }); for (const name of ["post-create.json", "post-delete.json"]) { const text = await fs.readFile(path.join(process.cwd(), "fixtures", "x-activity", name), "utf8"); await connector.ingest(store, JSON.parse(text) as unknown, sha256(text)); } const events = await store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(events.map((event) => event.payload.eventType)).toEqual(["post.create", "post.delete"]); expect(events[0]?.payload.post).toMatchObject({ postId: "2080761390344937796", text: "Hello from the Activity API", }); expect(events[1]?.payload.post).toEqual({ postId: "2080761269309813134", authorId: "1111111111111111111", }); const activity = await buildRootActivity(store); expect(activity.items.find((item) => item.id === events[0]?.id)?.presentation).toEqual({ renderer: "x-post", title: "Posted on X", body: "Hello from the Activity API", objectLabel: "X post", url: "https://x.com/i/web/status/2080761390344937796", urlLabel: "Open post", x: { postId: "2080761390344937796", createdAt: "2026-07-24T21:05:48.000Z", language: "en", attachmentCount: 0, possiblySensitive: false, }, }); expect(activity.items.find((item) => item.id === events[1]?.id)?.presentation).toEqual({ renderer: "generic", title: "Deleted an X post", body: "", objectLabel: "X post", }); expect(presentActivityEvent({ ...events[0]!, payload: { ...events[0]!.payload, subscriptionTag: "thoughtstream:x:official-fixture:alice_fixture:post-create", post: { ...(events[0]!.payload.post as Record), text: "@bob_fixture A direct response", inReplyToUserId: "2222222222222222222", referencedPosts: [{ type: "replied_to", postId: "2080761000000000000" }], entities: { mentions: [{ start: 0, end: 12, username: "bob_fixture", userId: "2222222222222222222" }] }, }, }, })).toEqual({ renderer: "x-post", title: "Replied on X", body: "A direct response", objectLabel: "X post", url: "https://x.com/i/web/status/2080761390344937796", urlLabel: "Open reply", parentUrl: "https://x.com/i/web/status/2080761000000000000", parentLabel: "Open parent post", x: { postId: "2080761390344937796", authorHandle: "alice_fixture", createdAt: "2026-07-24T21:05:48.000Z", replyToHandle: "bob_fixture", language: "en", attachmentCount: 0, possiblySensitive: false, }, }); for (const event of events) { expect(JSON.stringify(event.payload)).not.toContain("public_metrics"); expect(JSON.stringify(event.payload)).not.toContain("followers_count"); } }); test("normalizes an outbound private like with sensitive custody and the action timestamp", async () => { const root = await temporaryProject("thoughtstream-x-private-like-"); roots.push(root); const store = await testStore(root); stores.push(store); const connector = new XActivityConnector({ id: "x:cameron-private", lane: "personal-private", expectedSubscriptions: [{ eventType: "like.create", userId: "1111111111111111111", direction: "outbound", tag: "likes", }], }); const text = await fs.readFile(path.join(process.cwd(), "fixtures", "x-activity", "like-create.json"), "utf8"); await connector.ingest(store, JSON.parse(text) as unknown, sha256(text)); const [event] = await store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(event).toMatchObject({ externalId: "-3402274206530057851", actor: "1111111111111111111", privacy: "sensitive", occurredAt: "2026-07-24T21:12:45.236Z", payload: { eventType: "like.create", matchedUserId: "1111111111111111111", direction: "outbound", lane: "personal-private", like: { likeId: "87719f50ee17bdfa06a3089098a2b9ed", actorId: "1111111111111111111", postId: "2079814480427442556", postAuthorId: "3333333333333333333", likedPostCreatedAt: "2026-07-24T21:12:45.000Z", eventTimestampMs: "1784927565236", eventAt: "2026-07-24T21:12:45.236Z", }, }, }); expect(JSON.stringify(event?.payload)).not.toContain("includes"); expect((await buildRootActivity(store)).items.find((item) => item.id === event?.id)?.presentation).toEqual({ renderer: "generic", title: "Liked a post on X", body: "", objectLabel: "X post", url: "https://x.com/i/web/status/2079814480427442556", urlLabel: "Open liked post", }); const lifecycle = await store.listEvents({ source: "connector:x:cameron-private" }); expect(lifecycle.every((item) => item.privacy === "sensitive")).toBe(true); const wrongDirection = JSON.parse(text) as { data: { event_uuid: string; filter: { direction: string } } }; wrongDirection.data.event_uuid = "wrong-direction"; wrongDirection.data.filter.direction = "inbound"; await connector.ingest(store, wrongDirection, sha256(JSON.stringify(wrongDirection))); expect(await store.listEvents({ types: ["stream.thought.source.x.activity"] })).toHaveLength(1); }); test("binds public and private X subscription lanes to different event contracts", () => { expect(() => new XActivityConnector({ id: "x:public-invalid", lane: "personal-public", expectedSubscriptions: [{ eventType: "like.create", userId: "1111111111111111111", direction: "outbound", tag: "likes" }], })).toThrow("public sources admit only directionless post subscriptions"); expect(() => new XActivityConnector({ id: "x:private-invalid", lane: "personal-private", expectedSubscriptions: [{ eventType: "like.create", userId: "1111111111111111111", tag: "likes" }], })).toThrow("personal-private sources admit only outbound like.create subscriptions"); }); test("answers CRC without opening or mutating the event store", async () => { const diagnostics: XCrcDiagnostic[] = []; const fixture = await setup({ onCrcDiagnostic: (diagnostic) => diagnostics.push(diagnostic) }); const append = vi.spyOn(fixture.store, "appendEvent"); const response = await fetch(`${fixture.endpoint}?crc_token=fixture-token`); expect(response.status).toBe(200); expect(xCrcResponseToken("fixture-token", "fixture-secret")).toBe("sha256=iFUrz8TmBK5ze3siT/pRI6oBU8cd0cXwoCw+7/1/oMc="); expect(await response.json()).toEqual({ response_token: "sha256=iFUrz8TmBK5ze3siT/pRI6oBU8cd0cXwoCw+7/1/oMc=" }); expect(response.headers.get("cache-control")).toBe("no-store"); expect(append).not.toHaveBeenCalled(); expect((await fetch(`${fixture.endpoint}?crc_token=one&nonce=cache-buster`)).status).toBe(200); expect((await fetch(`${fixture.endpoint}?crc_token=one&extra=two`)).status).toBe(400); expect((await fetch(`${fixture.endpoint}?crc_token=one&crc_token=two`)).status).toBe(400); expect((await fetch(`${fixture.endpoint}?crc_token=one&nonce=first&nonce=second`)).status).toBe(400); expect(diagnostics).toEqual([ { accepted: true, queryKeys: ["crc_token"], tokenCount: 1, tokenLength: 13, nonceCount: 0, nonceLength: 0 }, { accepted: true, queryKeys: ["crc_token", "nonce"], tokenCount: 1, tokenLength: 3, nonceCount: 1, nonceLength: 12 }, { accepted: false, queryKeys: ["crc_token", "extra"], tokenCount: 1, tokenLength: 3, nonceCount: 0, nonceLength: 0 }, { accepted: false, queryKeys: ["crc_token", "crc_token"], tokenCount: 2, tokenLength: 3, nonceCount: 0, nonceLength: 0 }, { accepted: false, queryKeys: ["crc_token", "nonce", "nonce"], tokenCount: 1, tokenLength: 3, nonceCount: 2, nonceLength: 5 }, ]); expect(JSON.stringify(diagnostics)).not.toContain("fixture-token"); await fixture.receiver.close(); }); test("verifies the raw body signature and durably admits one allowlisted post", async () => { const fixture = await setup(); const envelope = postCreateEnvelope("event-1", "1001", "first observation"); const body = JSON.stringify(envelope); expect((await fetch(fixture.endpoint, { method: "POST", headers: { "content-type": "application/json" }, body, })).status).toBe(401); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("wrong", body), body, })).status).toBe(401); expect((await fetch(`${fixture.endpoint}?forbidden=true`, { method: "POST", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(404); expect((await fetch(fixture.endpoint, { method: "PUT", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(405); expect((await fetch(fixture.endpoint, { method: "POST", headers: { ...signedHeaders("fixture-secret", body), "content-type": "text/plain" }, body, })).status).toBe(415); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(200); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(200); const ignoredReplayChanges = structuredClone(envelope); ignoredReplayChanges.data.payload.public_metrics.like_count = 1_000; ignoredReplayChanges.data.includes.users[0]!.public_metrics.followers_count = 2; const changedRawBody = JSON.stringify(ignoredReplayChanges, null, 2); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", changedRawBody), body: changedRawBody, })).status).toBe(200); const events = await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(events).toHaveLength(1); expect(events[0]).toMatchObject({ externalId: "event-1", actor: "123456789", privacy: "public-source", payload: { eventType: "post.create", matchedUserId: "123456789", subscriptionTag: "thoughtstream:x:fixture:post-create", lane: "public-watchlist", post: { postId: "1001", authorId: "123456789", text: "first observation", language: "en", }, transport: { signatureVerified: true, kind: "x-v2-webhook" }, }, }); expect(events[0]!.payload).not.toHaveProperty("includes"); expect(events[0]!.payload.post).not.toHaveProperty("public_metrics"); const malformed = JSON.stringify({ ...envelope, data: { ...envelope.data, payload: { id: "1001" } } }); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", malformed), body: malformed, })).status).toBe(400); const emptyEventUuid = JSON.stringify({ ...envelope, data: { ...envelope.data, event_uuid: "" } }); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", emptyEventUuid), body: emptyEventUuid, })).status).toBe(400); const numericEventUuid = JSON.stringify({ ...envelope, data: { ...envelope.data, event_uuid: 2080761390344937796 } }); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", numericEventUuid), body: numericEventUuid, })).status).toBe(400); await fixture.receiver.close(); }); test("acknowledges signed but unconfigured activity without admitting source content", async () => { const fixture = await setup(); const wrongTag = postCreateEnvelope("event-wrong-tag", "1002", "ignore wrong tag"); wrongTag.data.tag = "thoughtstream:x:someone-else:post-create"; const wrongUser = postCreateEnvelope("event-wrong-user", "1003", "ignore wrong user"); wrongUser.data.filter.user_id = "999999999"; const wrongType = postCreateEnvelope("event-wrong-type", "1004", "ignore wrong type"); wrongType.data.event_type = "profile.update.bio"; for (const envelope of [wrongTag, wrongUser, wrongType]) { const body = JSON.stringify(envelope); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body })).status).toBe(200); } expect(await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] })).toHaveLength(0); const completions = await fixture.store.listEvents({ types: ["stream.thought.connector.ingest.completed"] }); expect(completions.at(-1)?.payload).toMatchObject({ accepted: 0, ignored: 1, eventType: "unconfigured" }); await fixture.receiver.close(); }); test("refuses divergent replay of one event UUID", async () => { const fixture = await setup(); const original = JSON.stringify(postCreateEnvelope("event-divergent", "1003", "original")); const changed = JSON.stringify(postCreateEnvelope("event-divergent", "1003", "changed")); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", original), body: original })).status).toBe(200); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", changed), body: changed })).status).toBe(503); const events = await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(events).toHaveLength(1); expect(events[0]!.payload.post).toMatchObject({ text: "original" }); await fixture.receiver.close(); }); test("admits out-of-order post timestamps in durable receipt order", async () => { const fixture = await setup(); const newer = JSON.stringify(postCreateEnvelope("event-new", "1004", "newer", "2026-08-10T12:00:00.000Z")); const older = JSON.stringify(postCreateEnvelope("event-old", "1005", "older", "2026-08-09T12:00:00.000Z")); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", newer), body: newer })).status).toBe(200); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", older), body: older })).status).toBe(200); const events = await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(events.map((event) => [event.externalId, event.sourceSequence, event.occurredAt])).toEqual([ ["event-new", events[0]!.sourceSequence, "2026-08-10T12:00:00.000Z"], ["event-old", events[1]!.sourceSequence, "2026-08-09T12:00:00.000Z"], ]); expect(events[1]!.sourceSequence).toBeGreaterThan(events[0]!.sourceSequence); await fixture.receiver.close(); }); test("serializes authenticated deliveries before touching Jazz", async () => { const fixture = await setup(); const actualIngest = fixture.connector.ingest.bind(fixture.connector); let active = 0; let maximumActive = 0; vi.spyOn(fixture.connector, "ingest").mockImplementation(async (...arguments_) => { active += 1; maximumActive = Math.max(maximumActive, active); await new Promise((resolve) => setTimeout(resolve, 20)); try { return await actualIngest(...arguments_); } finally { active -= 1; } }); const one = JSON.stringify(postCreateEnvelope("event-serial-1", "1006", "one")); const two = JSON.stringify(postCreateEnvelope("event-serial-2", "1007", "two")); const responses = await Promise.all([ fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", one), body: one }), fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", two), body: two }), ]); expect(responses.map((response) => response.status)).toEqual([200, 200]); expect(maximumActive).toBe(1); await fixture.receiver.close(); }); test("bounds the serial ingestion queue and leaves overflow for X to retry", async () => { const fixture = await setup({ pendingRequestLimit: 1 }); const actualIngest = fixture.connector.ingest.bind(fixture.connector); vi.spyOn(fixture.connector, "ingest").mockImplementationOnce(async (...arguments_) => { await new Promise((resolve) => setTimeout(resolve, 150)); return actualIngest(...arguments_); }); const first = JSON.stringify(postCreateEnvelope("event-queue-1", "1010", "first queued")); const overflow = JSON.stringify(postCreateEnvelope("event-queue-2", "1011", "retry later")); const firstRequest = fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", first), body: first, }); await new Promise((resolve) => setTimeout(resolve, 30)); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", overflow), body: overflow, })).status).toBe(503); expect((await firstRequest).status).toBe(200); const events = await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(events.map((event) => event.externalId)).toEqual(["event-queue-1"]); await fixture.receiver.close(); }); test("returns 5xx when durable settlement fails and enforces the body bound", async () => { const fixture = await setup({ maxBodyBytes: 1_024 }); const envelope = JSON.stringify(postCreateEnvelope("event-failure", "1008", "failure")); const producerBatch = vi.spyOn(fixture.store, "appendProducerBatch").mockRejectedValueOnce(new Error("fixture durable failure")); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", envelope), body: envelope, })).status).toBe(503); producerBatch.mockRestore(); expect(await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] })).toHaveLength(0); const oversized = JSON.stringify({ data: { padding: "x".repeat(2_000) } }); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", oversized), body: oversized, })).status).toBe(413); await fixture.receiver.close(); }); test("times out before X's acknowledgement deadline while allowing the durable attempt to finish", async () => { const fixture = await setup({ requestTimeoutMs: 1_000 }); const actualIngest = fixture.connector.ingest.bind(fixture.connector); vi.spyOn(fixture.connector, "ingest").mockImplementationOnce(async (...arguments_) => { await new Promise((resolve) => setTimeout(resolve, 1_100)); return actualIngest(...arguments_); }); const body = JSON.stringify(postCreateEnvelope("event-timeout", "1009", "late but durable")); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(503); await fixture.receiver.drain(); expect(await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] })).toHaveLength(1); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body, })).status).toBe(200); expect(await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] })).toHaveLength(1); await fixture.receiver.close(); }); test("normalizes allowlisted delete events", async () => { const fixture = await setup(); const envelope = { data: { event_uuid: "event-delete", event_type: "post.delete", filter: { user_id: "123456789" }, tag: "thoughtstream:x:fixture:post-delete", payload: { id: "999999999", author_id: "123456789", public_metrics: { likes: 5 } }, }, }; const body = JSON.stringify(envelope); expect((await fetch(fixture.endpoint, { method: "POST", headers: signedHeaders("fixture-secret", body), body })).status).toBe(200); const [event] = await fixture.store.listEvents({ types: ["stream.thought.source.x.activity"] }); expect(event?.payload).toMatchObject({ eventType: "post.delete", post: { postId: "999999999", authorId: "123456789" }, }); expect(event?.payload.post).not.toHaveProperty("public_metrics"); await fixture.receiver.close(); }); }); async function setup(options: { maxBodyBytes?: number; requestTimeoutMs?: number; pendingRequestLimit?: number; onCrcDiagnostic?: (diagnostic: XCrcDiagnostic) => void; } = {}) { const root = await temporaryProject("thoughtstream-x-webhook-"); roots.push(root); const store = await testStore(root); stores.push(store); const connector = new XActivityConnector({ id: "x:fixture", lane: "public-watchlist", expectedSubscriptions: [ { eventType: "post.create", userId: "123456789", tag: "thoughtstream:x:fixture:post-create" }, { eventType: "post.delete", userId: "123456789", tag: "thoughtstream:x:fixture:post-delete" }, ], }); const receiver = await startXWebhookServer({ connector, store, consumerSecret: "fixture-secret", path: "/webhooks/x", host: "127.0.0.1", port: 0, maxBodyBytes: options.maxBodyBytes ?? 2 * 1024 * 1024, requestTimeoutMs: options.requestTimeoutMs ?? 8_000, pendingRequestLimit: options.pendingRequestLimit ?? 100, ...(options.onCrcDiagnostic ? { onCrcDiagnostic: options.onCrcDiagnostic } : {}), }); return { store, connector, receiver, endpoint: `http://${receiver.host}:${receiver.port}${receiver.path}` }; } function signedHeaders(secret: string, body: string): Record { return { "content-type": "application/json", "x-twitter-webhooks-signature": xWebhookSignature(body, secret), }; } function postCreateEnvelope( eventUuid: string, postId: string, text: string, createdAt = "2026-08-10T10:00:00.000Z", ) { return { data: { event_uuid: eventUuid, event_type: "post.create", filter: { user_id: "123456789" }, tag: "thoughtstream:x:fixture:post-create", payload: { id: postId, author_id: "123456789", text, created_at: createdAt, conversation_id: postId, edit_history_tweet_ids: [postId], lang: "en", public_metrics: { like_count: 999 }, entities: { hashtags: [{ start: 0, end: 4, tag: "test", mutable_count: 1 }], }, }, includes: { users: [{ id: "123456789", username: "fixture", public_metrics: { followers_count: 1 } }] }, }, }; }