← 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 · TypeScriptimport { 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
}