Skip to content
← Public packages

@kentcdodds/x

X API v2 helpers for tweets, search, legacy DMs, and encrypted X Chat via a Fly XDK sidecar.

src/chat.ts

477 lines · 14.8 KB · TypeScript
import { getMe, getUserByUsername, refreshOauthToken, xRequest, XApiError } from './client.ts'
import { resolveIntegrationName, resolveOAuthAccount } from './accounts.ts'
import { cleanObject, extractRateLimit } from './domain.ts'
import { withSendIdempotency } from './idempotency.ts'
import { parseSendChatInput, withSendTargetResult } from './send-input.ts'
import { rawXRequest } from './openapi-client.ts'
import { sidecarRequest } from './sidecar.ts'
import type {
	JsonRecord,
	QueryInput,
	XAccountParams,
	XChatConversation,
	XChatMessage,
	XChatPublicKeyRecord,
	XChatSigningKey,
	XGetChatConversationParams,
	XListChatConversationsParams,
	XSendChatParams,
} from './types.ts'

function accountSelection(params: XAccountParams): XAccountParams {
	return { account: params.account, integration: params.integration }
}

async function parseOkJson<T>(response: Response): Promise<T> {
	const text = await response.text()
	let body: unknown = null
	if (text) {
		try {
			body = JSON.parse(text)
		} catch {
			body = text
		}
	}
	const rateLimit = extractRateLimit(response)
	if (!response.ok) {
		const record = body && typeof body === 'object' ? (body as JsonRecord) : null
		const message =
			(typeof record?.detail === 'string' && record.detail) ||
			(typeof record?.title === 'string' && record.title) ||
			`X API request failed: ${response.status} ${response.statusText}`
		throw new XApiError(message, {
			status: response.status,
			statusText: response.statusText,
			details: body,
			rateLimit,
		})
	}
	if (body && typeof body === 'object' && !Array.isArray(body)) {
		return { ...(body as object), rateLimit } as T
	}
	return body as T
}

async function chatGet(input: {
	path: string
	query?: QueryInput
	params: XAccountParams
}): Promise<JsonRecord> {
	const oauth = await resolveOAuthAccount(input.params)
	let response = await rawXRequest(input.path, {
		authMode: 'oauth',
		integrationName: oauth.integrationName,
		query: input.query,
	})
	if (response.status === 401) {
		await refreshOauthToken({
			internal: true,
			account: input.params.account,
			integration: input.params.integration,
		})
		response = await rawXRequest(input.path, {
			authMode: 'oauth',
			integrationName: oauth.integrationName,
			query: input.query,
		})
	}
	return parseOkJson(response)
}

export function hyphenConversationId(id: string): string {
	return id.split(':').join('-')
}

export function colonConversationId(id: string): string {
	return id.split('-').join(':')
}

function asRecord(value: unknown): JsonRecord | null {
	return value && typeof value === 'object' && !Array.isArray(value) ? (value as JsonRecord) : null
}

function asRecordList(value: unknown): JsonRecord[] {
	if (Array.isArray(value)) {
		return value.filter((item): item is JsonRecord => Boolean(asRecord(item)))
	}
	const single = asRecord(value)
	return single ? [single] : []
}

function publicKeyRecords(payload: JsonRecord): XChatPublicKeyRecord[] {
	return asRecordList(payload.data) as XChatPublicKeyRecord[]
}

function latestPublicKey(records: XChatPublicKeyRecord[]): XChatPublicKeyRecord | null {
	if (records.length === 0) return null
	return records.reduce((latest, record) => {
		const current = Number(record.public_key_version || 0)
		const best = Number(latest.public_key_version || 0)
		return current > best ? record : latest
	})
}

export function toSigningKey(userId: string, record: XChatPublicKeyRecord): XChatSigningKey {
	const identity = String(record.identity_public_key || record.public_key || '')
	const signing = String(record.signing_public_key || record.public_key || '')
	return {
		user_id: String(userId),
		public_key_version: String(record.public_key_version ?? ''),
		public_key: signing,
		identity_public_key: identity,
		identity_public_key_signature: String(record.identity_public_key_signature || ''),
	}
}

function collectEncodedEvents(payload: JsonRecord): string[] {
	const meta = asRecord(payload.meta)
	const keyEvents = meta?.conversation_key_events
	const fromMeta = Array.isArray(keyEvents)
		? keyEvents.flatMap((item) => {
				if (typeof item === 'string' && item) return [item]
				const record = asRecord(item)
				const encoded = record?.encoded_event
				return typeof encoded === 'string' && encoded ? [encoded] : []
			})
		: []
	const fromData = asRecordList(payload.data).flatMap((item) => {
		const encoded = item.encoded_event
		return typeof encoded === 'string' && encoded ? [encoded] : []
	})
	return [...fromMeta, ...fromData]
}

function shapeConversation(item: XChatConversation): JsonRecord {
	return {
		id: item.id,
		type: item.type,
		participant_ids: item.participant_ids,
		is_muted: item.is_muted,
	}
}

function eventText(event: JsonRecord): string | null {
	if (typeof event.text === 'string') return event.text
	const content = asRecord(event.content)
	if (typeof content?.text === 'string') return content.text
	return null
}

function shapeDecryptedMessage(item: JsonRecord): XChatMessage | null {
	const event = asRecord(item.event) || item
	const type = typeof event.type === 'string' ? event.type : undefined
	if (type && type !== 'Message' && type !== 'message') return null
	return {
		id: typeof event.id === 'string' ? event.id : undefined,
		type: type || 'Message',
		senderId:
			(typeof event.sender_id === 'string' && event.sender_id) ||
			(typeof event.senderId === 'string' && event.senderId) ||
			undefined,
		createdAt:
			(typeof event.created_at === 'string' && event.created_at) ||
			(typeof event.createdAt === 'string' && event.createdAt) ||
			undefined,
		text: eventText(event),
		conversationId:
			(typeof event.conversation_id === 'string' && event.conversation_id) ||
			(typeof event.conversationId === 'string' && event.conversationId) ||
			undefined,
	}
}

async function resolveChatParticipantId(
	params: {
		participantId?: string
		participant_id?: string
		username?: string
	} & XAccountParams,
): Promise<string | null> {
	const participantId = params.participantId || params.participant_id
	if (participantId) return String(participantId)
	if (!params.username) return null
	const user = await getUserByUsername({
		username: params.username,
		...accountSelection(params),
	})
	const id = user?.data?.id
	if (!id) throw new Error(`Could not resolve X user id for username "${params.username}"`)
	return id
}

async function currentUserId(params: XAccountParams): Promise<string> {
	const me = await getMe(accountSelection(params))
	const id = me?.data?.id
	if (!id) throw new Error('Could not determine current X user id')
	return String(id)
}

async function getUserPublicKeys(userId: string, params: XAccountParams): Promise<XChatPublicKeyRecord[]> {
	const payload = await chatGet({
		path: `/users/${encodeURIComponent(userId)}/public_keys`,
		params,
	})
	return publicKeyRecords(payload)
}

async function collectSigningKeys(
	userIds: string[],
	params: XAccountParams,
): Promise<XChatSigningKey[]> {
	const unique = [...new Set(userIds.filter(Boolean))]
	const keys: XChatSigningKey[] = []
	for (const userId of unique) {
		const records = await getUserPublicKeys(userId, params)
		for (const record of records) {
			keys.push(toSigningKey(userId, record))
		}
	}
	return keys
}

async function identityForSidecar(params: XAccountParams): Promise<{
	userId: string
	juiceboxConfig: JsonRecord | string
	signingKeyVersion: string
	ownKeys: XChatSigningKey[]
}> {
	const userId = await currentUserId(params)
	const records = await getUserPublicKeys(userId, params)
	const latest = latestPublicKey(records)
	if (!latest?.juicebox_config) {
		throw new Error(
			'No juicebox_config on this account public key. Unlock the existing X Chat identity — do not register a new keypair.',
		)
	}
	return {
		userId,
		juiceboxConfig: latest.juicebox_config,
		signingKeyVersion: String(latest.public_key_version ?? ''),
		ownKeys: records.map((record) => toSigningKey(userId, record)),
	}
}

/** List encrypted X Chat conversations (metadata only; no decrypt). */
export async function listChatConversations(
	params: XListChatConversationsParams = {},
): Promise<JsonRecord> {
	const payload = await chatGet({
		path: '/chat/conversations',
		query: cleanObject({
			max_results: params.maxResults || params.max_results || 10,
			pagination_token: params.paginationToken || params.pagination_token,
		}) as QueryInput,
		params: accountSelection(params),
	})
	const conversations = asRecordList(payload.data) as XChatConversation[]
	return {
		data: conversations.map((item) => shapeConversation(item)),
		meta: payload.meta,
		errors: payload.errors,
		rateLimit: payload.rateLimit,
	}
}

async function loadConversationEvents(
	conversationId: string,
	params: XGetChatConversationParams,
): Promise<JsonRecord> {
	return chatGet({
		path: `/chat/conversations/${encodeURIComponent(hyphenConversationId(conversationId))}/events`,
		query: cleanObject({
			max_results: params.maxResults || params.max_results || 20,
			pagination_token: params.paginationToken || params.pagination_token,
		}) as QueryInput,
		params: accountSelection(params),
	})
}

/** Read one encrypted X Chat conversation through the Fly XDK sidecar. */
export async function getChatConversation(params: XGetChatConversationParams): Promise<JsonRecord> {
	const conversationId = params.conversationId || params.conversation_id
	const participantId = await resolveChatParticipantId(params)
	if (!conversationId && !participantId) {
		throw new Error('params.conversationId, params.participantId, or params.username is required')
	}

	const identity = await identityForSidecar(params)
	const pathId = hyphenConversationId(String(conversationId || participantId))
	const firstPage = await loadConversationEvents(pathId, params)
	const pages: JsonRecord[] = [firstPage]
	const explicitToken = params.paginationToken || params.pagination_token
	let pageMeta = asRecord(firstPage.meta)
	let pagesLeft = 24
	while (
		!explicitToken &&
		pagesLeft > 0 &&
		pageMeta?.has_more &&
		typeof pageMeta.next_token === 'string' &&
		pageMeta.next_token
	) {
		const nextPage = await loadConversationEvents(pathId, {
			...params,
			paginationToken: pageMeta.next_token,
		})
		pages.push(nextPage)
		pageMeta = asRecord(nextPage.meta)
		pagesLeft -= 1
	}

	const encodedEvents = [...new Set(pages.flatMap((page) => collectEncodedEvents(page)))]
	const conversations = pages.flatMap((page) => asRecordList(page.data))
	const participantIds = new Set<string>([identity.userId])
	if (participantId) participantIds.add(participantId)
	for (const item of conversations) {
		if (typeof item.sender_id === 'string') participantIds.add(item.sender_id)
		const ids = item.participant_ids
		if (Array.isArray(ids)) {
			for (const id of ids) {
				if (typeof id === 'string' && id) participantIds.add(id)
			}
		}
	}

	const signingKeys = [
		...identity.ownKeys,
		...(await collectSigningKeys(
			[...participantIds].filter((id) => id !== identity.userId),
			params,
		)),
	]

	const decrypted = await sidecarRequest(
		'/v1/decrypt-events',
		{
			user_id: identity.userId,
			juicebox_config: identity.juiceboxConfig,
			signing_key_version: identity.signingKeyVersion,
			signing_keys: signingKeys,
			encoded_events: encodedEvents,
			reject_unverified: false,
		},
		params.sidecarUrl,
	)

	const messages = asRecordList(decrypted.messages)
		.map((item) => shapeDecryptedMessage(item))
		.filter((item): item is XChatMessage => item !== null)

	const lastPage = pages[pages.length - 1] ?? firstPage
	const meta = asRecord(lastPage.meta) || {}
	return {
		conversationId: pathId,
		data: messages,
		meta: {
			result_count: messages.length,
			next_token: meta.next_token,
			previous_token: meta.previous_token,
			has_more: meta.has_more,
		},
		decryptErrorCount: decrypted.error_count ?? 0,
		rateLimit: lastPage.rateLimit,
	}
}

/**
 * Encrypt + send an X Chat message through the Fly sidecar.
 * Sending requires confirm: true. X OAuth tokens stay on api.x.com.
 */
/**
 * Encrypt + send an X Chat message through the Fly sidecar.
 * Sending requires confirm: true. X OAuth tokens stay on api.x.com.
 *
 * Pass optional `idempotencyKey` so retries / parallel workers reuse the prior
 * successful send for the same `(account, key)` within 7 days.
 */
export async function sendChat(params: XSendChatParams | Record<string, unknown> = {}): Promise<unknown> {
	const input = parseSendChatInput(params)
	const conversationId = input.conversationId
	const participantId = await resolveChatParticipantId({
		participantId: input.participantId,
		username: input.username,
		account: input.account,
		integration: input.integration,
		sidecarUrl: input.sidecarUrl,
	})
	if (!conversationId && !participantId) {
		throw new Error('send-chat requires conversationId, participantId, or username after alias mapping')
	}
	const pathId = hyphenConversationId(String(conversationId || participantId))
	const previewUrl = `https://api.x.com/2/chat/conversations/${pathId}/messages`
	const targetFields = {
		targetMode: input.targetMode,
		conversationId: conversationId || pathId,
		participantId: input.participantId,
		username: input.username,
	}
	if (input.dryRun || !input.confirm) {
		return withSendTargetResult(
			{
				dryRun: true,
				method: 'POST',
				url: previewUrl,
				authMode: 'oauth',
				requiresConfirm: true,
				payload: {
					conversationId: pathId,
					text: input.text,
					via: 'x-chat-sidecar',
				},
			},
			targetFields,
		)
	}

	const sent = await withSendIdempotency({
		action: 'sendChat',
		account: resolveIntegrationName(input),
		idempotencyKey: input.idempotencyKey,
		run: async () => {
			const identity = await identityForSidecar(input)
			const eventsPayload = await loadConversationEvents(pathId, input).catch(() => ({}) as JsonRecord)
			const encodedEvents = collectEncodedEvents(eventsPayload)
			const participantIds = new Set<string>([identity.userId])
			if (participantId) participantIds.add(participantId)
			const signingKeys = [
				...identity.ownKeys,
				...(await collectSigningKeys(
					[...participantIds].filter((id) => id !== identity.userId),
					input,
				)),
			]

			const encrypted = await sidecarRequest(
				'/v1/encrypt-message',
				{
					user_id: identity.userId,
					juicebox_config: identity.juiceboxConfig,
					signing_key_version: identity.signingKeyVersion,
					signing_keys: signingKeys,
					conversation_id: colonConversationId(pathId),
					text: input.text,
					encoded_events: encodedEvents,
				},
				input.sidecarUrl,
			)
			const sendBody = asRecord(encrypted.payload)
			if (!sendBody?.message_id || !sendBody.encoded_message_create_event) {
				throw new Error('Sidecar encrypt-message did not return a send payload')
			}

			return xRequest({
				path: `/chat/conversations/${encodeURIComponent(pathId)}/messages`,
				method: 'POST',
				body: {
					message_id: sendBody.message_id,
					encoded_message_create_event: sendBody.encoded_message_create_event,
					encoded_message_event_signature: sendBody.encoded_message_event_signature,
				},
				authMode: 'oauth',
				...accountSelection(input),
				confirm: true,
			})
		},
	})
	return withSendTargetResult(sent, targetFields)
}