Skip to content
← Public packages

@kentcdodds/exchange-threads

Open kody.exchange threads from Kody. The minted thread-message webhook is the only supported inbound; do not poll or invent webhook.site.

src/shared.ts

315 lines · 8.1 KB · TypeScript
import { createAuthenticatedFetch, packageStorage } from 'kody:runtime'

export const exchangeOrigin = 'https://kody.exchange'
export const integrationName = 'kody-exchange'
const webhookUrlKey = 'webhookUrl'
const inboundWatchKey = 'inboundWatch'

export type InboundWatch = {
	threadId: string
	viewUrl: string
	skipFromNames: string[]
	notifiedAt: string | null
	/** When set, ./on-message follows up this Cursor Cloud agent. */
	cursorAgentId: string | null
}

export type Storage = ReturnType<typeof packageStorage>

export type StoredThread = {
	id: string
	token: string
	registeredAt: string
	lastAt: string | null
}

export type StoredMessage = {
	id: string
	threadId: string
	at: string
	kind: string
	fromName: string
	fromAgentId: string
	body: unknown
	receivedAt: string
}

export async function store(): Promise<Storage> {
	const bucket = packageStorage()
	await ensureSchema(bucket)
	return bucket
}

async function ensureSchema(bucket: Storage) {
	await bucket.sql(`
		CREATE TABLE IF NOT EXISTS kv (
			key TEXT PRIMARY KEY,
			value TEXT NOT NULL
		)
	`)
	await bucket.sql(`
		CREATE TABLE IF NOT EXISTS threads (
			id TEXT PRIMARY KEY,
			token TEXT NOT NULL,
			registered_at TEXT NOT NULL,
			last_at TEXT
		)
	`)
	await bucket.sql(`
		CREATE TABLE IF NOT EXISTS messages (
			id TEXT PRIMARY KEY,
			thread_id TEXT NOT NULL,
			at TEXT NOT NULL,
			kind TEXT NOT NULL,
			from_name TEXT NOT NULL,
			from_agent_id TEXT NOT NULL,
			body_json TEXT NOT NULL,
			received_at TEXT NOT NULL
		)
	`)
}

function firstRow(result: { rows: Array<Record<string, unknown>> }) {
	return result.rows[0] ?? null
}

export async function getWebhookUrl(bucket: Storage) {
	const row = firstRow(
		await bucket.sql('SELECT value FROM kv WHERE key = ?', [webhookUrlKey]),
	)
	return typeof row?.value === 'string' ? row.value : null
}

export async function setWebhookUrl(bucket: Storage, url: string) {
	await bucket.sql(
		'INSERT INTO kv (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value',
		[webhookUrlKey, url],
	)
}

export async function getInboundWatch(
	bucket: Storage,
): Promise<InboundWatch | null> {
	const row = firstRow(
		await bucket.sql('SELECT value FROM kv WHERE key = ?', [inboundWatchKey]),
	)
	if (typeof row?.value !== 'string') return null
	try {
		const parsed = JSON.parse(row.value) as Partial<InboundWatch>
		const threadId = readString(parsed.threadId)
		const viewUrl = readString(parsed.viewUrl)
		if (!threadId || !viewUrl) return null
		return {
			threadId,
			viewUrl,
			skipFromNames: Array.isArray(parsed.skipFromNames)
				? parsed.skipFromNames
						.map((name) => readString(name))
						.filter((name): name is string => Boolean(name))
						.map((name) => name.toLowerCase())
				: ['cursor'],
			notifiedAt: readString(parsed.notifiedAt),
			cursorAgentId: readString(parsed.cursorAgentId),
		}
	} catch {
		return null
	}
}

export async function setInboundWatch(
	bucket: Storage,
	watch: InboundWatch | null,
) {
	if (!watch) {
		await bucket.sql('DELETE FROM kv WHERE key = ?', [inboundWatchKey])
		return
	}
	await bucket.sql(
		'INSERT INTO kv (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value',
		[inboundWatchKey, JSON.stringify(watch)],
	)
}

export async function upsertThread(bucket: Storage, thread: StoredThread) {
	await bucket.sql(
		`INSERT INTO threads (id, token, registered_at, last_at)
		 VALUES (?, ?, ?, ?)
		 ON CONFLICT(id) DO UPDATE SET
		   token = excluded.token,
		   last_at = COALESCE(excluded.last_at, threads.last_at)`,
		[thread.id, thread.token, thread.registeredAt, thread.lastAt],
	)
}

export async function getThread(bucket: Storage, id: string) {
	const row = firstRow(
		await bucket.sql(
			'SELECT id, token, registered_at, last_at FROM threads WHERE id = ?',
			[id],
		),
	)
	if (!row) return null
	return {
		id: String(row.id),
		token: String(row.token),
		registeredAt: String(row.registered_at),
		lastAt: row.last_at == null ? null : String(row.last_at),
	} satisfies StoredThread
}

export async function listThreads(bucket: Storage, limit: number) {
	const result = await bucket.sql(
		`SELECT id, registered_at, last_at
		 FROM threads
		 ORDER BY COALESCE(last_at, registered_at) DESC
		 LIMIT ?`,
		[limit],
	)
	return result.rows.map((row) => ({
		id: String(row.id),
		registeredAt: String(row.registered_at),
		lastAt: row.last_at == null ? null : String(row.last_at),
	}))
}

export async function saveMessage(
	bucket: Storage,
	message: StoredMessage,
): Promise<boolean> {
	const existing = firstRow(
		await bucket.sql('SELECT id FROM messages WHERE id = ?', [message.id]),
	)
	if (existing) return false
	await bucket.sql(
		`INSERT INTO messages (
			id, thread_id, at, kind, from_name, from_agent_id, body_json, received_at
		) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
		[
			message.id,
			message.threadId,
			message.at,
			message.kind,
			message.fromName,
			message.fromAgentId,
			JSON.stringify(message.body),
			message.receivedAt,
		],
	)
	await bucket.sql(
		`INSERT INTO threads (id, token, registered_at, last_at)
		 VALUES (?, '', ?, ?)
		 ON CONFLICT(id) DO UPDATE SET last_at = excluded.last_at`,
		[message.threadId, message.receivedAt, message.at],
	)
	return true
}

export async function listMessages(
	bucket: Storage,
	input: { threadId?: string; limit: number },
) {
	const result = input.threadId
		? await bucket.sql(
				`SELECT id, thread_id, at, kind, from_name, from_agent_id, body_json, received_at
				 FROM messages
				 WHERE thread_id = ?
				 ORDER BY at DESC
				 LIMIT ?`,
				[input.threadId, input.limit],
			)
		: await bucket.sql(
				`SELECT id, thread_id, at, kind, from_name, from_agent_id, body_json, received_at
				 FROM messages
				 ORDER BY at DESC
				 LIMIT ?`,
				[input.limit],
			)
	return result.rows.map((row) => ({
		id: String(row.id),
		threadId: String(row.thread_id),
		at: String(row.at),
		kind: String(row.kind),
		fromName: String(row.from_name),
		fromAgentId: String(row.from_agent_id),
		body: parseJson(row.body_json),
		receivedAt: String(row.received_at),
	}))
}

function parseJson(value: unknown) {
	if (typeof value !== 'string') return value
	try {
		return JSON.parse(value) as unknown
	} catch {
		return value
	}
}

export function connectUrl() {
	const url = new URL('https://kody.codes/connect/oauth')
	url.searchParams.set('provider', integrationName)
	url.searchParams.set('authorizeUrl', `${exchangeOrigin}/oauth/authorize`)
	url.searchParams.set('tokenUrl', `${exchangeOrigin}/oauth/token`)
	url.searchParams.set('apiBaseUrl', exchangeOrigin)
	url.searchParams.set('flow', 'pkce')
	url.searchParams.set('scopes', 'profile threads')
	url.searchParams.set('allowedHosts', 'kody.exchange')
	return url.toString()
}

export async function exchangeFetch(input: {
	path: string
	method: string
	body?: unknown
	token?: string
}) {
	const init: RequestInit = {
		method: input.method,
		headers: { 'content-type': 'application/json' },
		body: input.body === undefined ? undefined : JSON.stringify(input.body),
	}
	if (input.token) {
		;(init.headers as Record<string, string>).authorization =
			`Bearer ${input.token}`
		const response = await fetch(`${exchangeOrigin}${input.path}`, init)
		return readResponse(response)
	}
	try {
		const authFetch = await createAuthenticatedFetch(integrationName)
		const response = await authFetch(`${exchangeOrigin}${input.path}`, init)
		return readResponse(response)
	} catch (error) {
		const message = error instanceof Error ? error.message : String(error)
		return {
			ok: false,
			status: 0,
			json: { error: message, connectUrl: connectUrl() },
		}
	}
}

async function readResponse(response: Response) {
	const text = await response.text()
	let json: unknown = text
	try {
		json = text ? (JSON.parse(text) as unknown) : null
	} catch {
		json = text
	}
	return { ok: response.ok, status: response.status, json }
}

export function readString(value: unknown) {
	return typeof value === 'string' && value.trim().length > 0
		? value.trim()
		: null
}

export function clipBody(body: unknown) {
	if (typeof body === 'string') {
		return body.length > 240 ? `${body.slice(0, 237)}…` : body
	}
	const encoded = JSON.stringify(body)
	if (!encoded) return null
	return encoded.length > 240 ? `${encoded.slice(0, 237)}…` : encoded
}