1
0

github-webhook-real.e2e.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442
  1. /** Real CLI and DeepSeek evidence for a GitHub webhook-created Session. */
  2. import type { ChildProcess } from 'node:child_process'
  3. import { spawn } from 'node:child_process'
  4. import { createHmac, randomUUID } from 'node:crypto'
  5. import { existsSync } from 'node:fs'
  6. import { mkdir, mkdtemp, realpath, rm } from 'node:fs/promises'
  7. import { createServer } from 'node:net'
  8. import type { AddressInfo } from 'node:net'
  9. import { tmpdir } from 'node:os'
  10. import { join } from 'node:path'
  11. import { setTimeout as delay } from 'node:timers/promises'
  12. import { fileURLToPath } from 'node:url'
  13. import { describe, expect, it } from 'vitest'
  14. const REPO_ROOT = fileURLToPath(new URL('../../..', import.meta.url))
  15. const BUILT_BIN = join(REPO_ROOT, 'apps/cli/lib/bin.js')
  16. const OVERLAY = fileURLToPath(new URL(
  17. '../../../examples/web-github-review/tests/fixtures/real-cli/cordis.yml',
  18. import.meta.url,
  19. ))
  20. const SECRET = 'github-webhook-real-e2e-secret'
  21. const DELIVERY = 'github-webhook-real-e2e-delivery'
  22. const MARKER = 'DSH_GITHUB_WEBHOOK_REAL_E2E_OK'
  23. const TITLE = 'GitHub webhook real e2e'
  24. interface SessionList {
  25. items: Array<{
  26. sessionId: string
  27. cwd?: string
  28. agentPreset?: string
  29. blank: boolean
  30. }>
  31. }
  32. interface WorkspaceBaseline {
  33. items: Array<{
  34. path: string
  35. sessionIds: string[]
  36. }>
  37. }
  38. interface HistoryPage {
  39. events: Array<{
  40. event: {
  41. type: string
  42. data: unknown
  43. }
  44. }>
  45. hasMore: boolean
  46. }
  47. interface ProcessObservation {
  48. readonly ready: Promise<string>
  49. readonly text: () => string
  50. }
  51. function isRecord(value: unknown): value is Record<string, unknown> {
  52. return typeof value === 'object' && value !== null
  53. }
  54. /** Capture bounded process output and resolve the public Web URL after settled boot. */
  55. function observeProcess(child: ChildProcess): ProcessObservation {
  56. let output = ''
  57. let settled = false
  58. let resolveReady!: (url: string) => void
  59. let rejectReady!: (error: Error) => void
  60. const ready = new Promise<string>((resolve, reject) => {
  61. resolveReady = resolve
  62. rejectReady = reject
  63. })
  64. const timer = setTimeout(() => {
  65. if (!settled) rejectReady(new Error(`dsh web did not become ready within 90s:\n${output}`))
  66. }, 90_000)
  67. timer.unref()
  68. const append = (chunk: Buffer | string): void => {
  69. output = `${output}${String(chunk)}`.slice(-100_000)
  70. const match = /dsh web: (http:\/\/[^\s]+)/u.exec(output)
  71. if (settled || match?.[1] === undefined) return
  72. settled = true
  73. clearTimeout(timer)
  74. resolveReady(match[1].replace('0.0.0.0', '127.0.0.1'))
  75. }
  76. child.stdout?.on('data', append)
  77. child.stderr?.on('data', append)
  78. child.once('error', (error) => {
  79. if (!settled) rejectReady(error)
  80. })
  81. child.once('exit', (code) => {
  82. if (!settled) rejectReady(new Error(`dsh web exited before readiness (code ${String(code)}):\n${output}`))
  83. })
  84. return { ready, text: () => output }
  85. }
  86. /** Reserve and release one loopback port for the isolated webhook listener. */
  87. async function freePort(): Promise<number> {
  88. const server = createServer()
  89. await new Promise<void>((resolve, reject) => {
  90. server.once('error', reject)
  91. server.listen(0, '127.0.0.1', resolve)
  92. })
  93. const port = (server.address() as AddressInfo).port
  94. await new Promise<void>((resolve, reject) => {
  95. server.close((error) => {
  96. if (error === undefined) resolve()
  97. else reject(error)
  98. })
  99. })
  100. return port
  101. }
  102. /** Invoke one public Remote method over its HTTP carrier. */
  103. async function remoteRpc<T>(baseUrl: string, endpoint: string, args: object): Promise<T> {
  104. const response = await fetch(`${baseUrl}/api/${endpoint}`, {
  105. method: 'POST',
  106. headers: { 'content-type': 'application/json' },
  107. body: JSON.stringify({
  108. type: 'client-request',
  109. rpcId: `github-webhook-real-${endpoint}-${randomUUID()}`,
  110. method: endpoint,
  111. payload: { args },
  112. }),
  113. })
  114. if (!response.ok) {
  115. throw new Error(`${endpoint} returned HTTP ${String(response.status)}: ${await response.text()}`)
  116. }
  117. const envelope = await response.json() as {
  118. result: { ok: true; value: T } | { ok: false; error: { code: string; message: string } }
  119. }
  120. if (!envelope.result.ok) {
  121. throw new Error(`${endpoint} failed: ${envelope.result.error.code}: ${envelope.result.error.message}`)
  122. }
  123. return envelope.result.value
  124. }
  125. /** Read one opening item from a public Remote stream. */
  126. async function openingStreamItem(
  127. baseUrl: string,
  128. endpoint: string,
  129. args: object,
  130. accepts: (value: unknown) => boolean,
  131. ): Promise<Record<string, unknown>> {
  132. const socket = new WebSocket(`${baseUrl.replace(/^http/u, 'ws')}/api/remote.mux`)
  133. const streamId = `github-webhook-real-${endpoint}-${randomUUID()}`
  134. try {
  135. await new Promise<void>((resolve, reject) => {
  136. const cleanup = (): void => {
  137. socket.removeEventListener('open', opened)
  138. socket.removeEventListener('error', failed)
  139. socket.removeEventListener('close', closed)
  140. }
  141. const opened = (): void => {
  142. cleanup()
  143. resolve()
  144. }
  145. const failed = (): void => {
  146. cleanup()
  147. reject(new Error(`${endpoint} carrier failed before opening`))
  148. }
  149. const closed = (): void => {
  150. cleanup()
  151. reject(new Error(`${endpoint} carrier closed before opening`))
  152. }
  153. socket.addEventListener('open', opened)
  154. socket.addEventListener('error', failed)
  155. socket.addEventListener('close', closed)
  156. })
  157. return await new Promise<Record<string, unknown>>((resolve, reject) => {
  158. const timer = setTimeout(() => { finish(new Error(`${endpoint} did not publish its opening item`)) }, 10_000)
  159. const cleanup = (): void => {
  160. clearTimeout(timer)
  161. socket.removeEventListener('message', message)
  162. socket.removeEventListener('error', failed)
  163. socket.removeEventListener('close', closed)
  164. }
  165. const finish = (error: Error | undefined, value?: Record<string, unknown>): void => {
  166. cleanup()
  167. if (error !== undefined) reject(error)
  168. else if (value === undefined) reject(new Error(`${endpoint} opening item was absent`))
  169. else resolve(value)
  170. }
  171. const message = (event: MessageEvent<unknown>): void => {
  172. try {
  173. if (typeof event.data !== 'string') throw new Error(`${endpoint} published a non-text frame`)
  174. const frame: unknown = JSON.parse(event.data)
  175. if (!isRecord(frame) || frame.streamId !== streamId) return
  176. if (frame.type === 'error') {
  177. finish(new Error(`${endpoint} failed: ${JSON.stringify(frame.error)}`))
  178. return
  179. }
  180. if (frame.type === 'end') {
  181. finish(new Error(`${endpoint} ended before its opening item`))
  182. return
  183. }
  184. if (frame.type === 'item' && isRecord(frame.value) && accepts(frame.value)) {
  185. finish(undefined, frame.value)
  186. }
  187. } catch (error) {
  188. finish(error instanceof Error ? error : new Error(String(error)))
  189. }
  190. }
  191. const failed = (): void => { finish(new Error(`${endpoint} carrier failed before its opening item`)) }
  192. const closed = (): void => { finish(new Error(`${endpoint} carrier closed before its opening item`)) }
  193. socket.addEventListener('message', message)
  194. socket.addEventListener('error', failed)
  195. socket.addEventListener('close', closed)
  196. socket.send(JSON.stringify({ type: 'open', streamId, endpoint, payload: { args } }))
  197. })
  198. } finally {
  199. socket.close()
  200. }
  201. }
  202. /** Read the current Workspace baseline from a fresh follow generation. */
  203. async function workspaceBaseline(baseUrl: string): Promise<WorkspaceBaseline> {
  204. const frame = await openingStreamItem(
  205. baseUrl,
  206. 'workspace/follow',
  207. {},
  208. value => isRecord(value) && value.type === 'baseline' && isRecord(value.value),
  209. )
  210. return frame.value as WorkspaceBaseline
  211. }
  212. /** Read the explicit page cut from a fresh Session follow generation. */
  213. async function sessionCursor(baseUrl: string, sessionId: string): Promise<number> {
  214. const frame = await openingStreamItem(
  215. baseUrl,
  216. 'session/follow',
  217. { request: { address: { kind: 'session', sessionId } } },
  218. value => isRecord(value) && value.type === 'opened' && Number.isSafeInteger(value.cursor),
  219. )
  220. return frame.cursor as number
  221. }
  222. /** Read Session history at the cursor explicitly opened for this page. */
  223. async function history(baseUrl: string, sessionId: string): Promise<HistoryPage> {
  224. const throughSeq = await sessionCursor(baseUrl, sessionId)
  225. return remoteRpc<HistoryPage>(baseUrl, 'session/page', {
  226. request: { address: { kind: 'session', sessionId }, throughSeq, maxMessages: 100 },
  227. })
  228. }
  229. /** Poll a public observation until it satisfies the test's behavior predicate. */
  230. async function eventually<T>(
  231. child: ChildProcess,
  232. processOutput: () => string,
  233. label: string,
  234. probe: () => Promise<T>,
  235. accepts: (value: T) => boolean,
  236. timeoutMs: number,
  237. ): Promise<T> {
  238. const deadline = Date.now() + timeoutMs
  239. let lastValue: T | undefined
  240. let lastError: unknown
  241. while (Date.now() < deadline) {
  242. if (child.exitCode !== null) {
  243. throw new Error(`dsh web exited while waiting for ${label} (code ${String(child.exitCode)}):\n${processOutput()}`)
  244. }
  245. try {
  246. lastValue = await probe()
  247. if (accepts(lastValue)) return lastValue
  248. } catch (error) {
  249. lastError = error
  250. }
  251. await delay(300)
  252. }
  253. throw new Error(
  254. `timed out waiting for ${label}; last value=${JSON.stringify(lastValue)}; `
  255. + `last error=${String(lastError)}; process output:\n${processOutput()}`,
  256. )
  257. }
  258. /** Return every text block from durable assistant messages. */
  259. function assistantText(page: HistoryPage): string {
  260. const text: string[] = []
  261. for (const { event } of page.events) {
  262. if (event.type !== 'assistant/message' || !isRecord(event.data) || !isRecord(event.data.message)) continue
  263. const content = event.data.message.content
  264. if (!Array.isArray(content)) continue
  265. for (const block of content) {
  266. if (isRecord(block) && block.type === 'text' && typeof block.text === 'string') text.push(block.text)
  267. }
  268. }
  269. return text.join('\n')
  270. }
  271. /** Stop the spawned CLI through its normal signal path, escalating only on a stuck teardown. */
  272. async function stop(child: ChildProcess): Promise<void> {
  273. if (child.exitCode !== null) return
  274. let resolveClosed!: () => void
  275. const closed = new Promise<void>((resolve) => { resolveClosed = resolve })
  276. child.once('close', resolveClosed)
  277. child.kill('SIGTERM')
  278. if (await Promise.race([closed.then(() => true), delay(10_000, false, { ref: false })])) return
  279. if (child.exitCode === null) child.kill('SIGKILL')
  280. await Promise.race([closed, delay(5_000, undefined, { ref: false })])
  281. }
  282. /** Send the sole synthetic external interaction: one signed GitHub delivery. */
  283. async function sendGitHubDelivery(origin: string): Promise<Response> {
  284. const body = JSON.stringify({
  285. action: 'ready_for_review',
  286. number: 4242,
  287. repository: { full_name: 'deepseek-harness/deepseek-harness' },
  288. pull_request: {
  289. title: 'Real CLI webhook e2e',
  290. html_url: 'https://github.com/deepseek-harness/deepseek-harness/pull/4242',
  291. draft: false,
  292. user: { login: 'octocat' },
  293. base: { ref: 'master', sha: 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa' },
  294. head: { ref: 'webhook-e2e', sha: 'bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb' },
  295. },
  296. })
  297. const signature = `sha256=${createHmac('sha256', SECRET).update(body).digest('hex')}`
  298. return await fetch(`${origin}/github`, {
  299. method: 'POST',
  300. headers: {
  301. 'content-type': 'application/json',
  302. 'x-github-delivery': DELIVERY,
  303. 'x-github-event': 'pull_request',
  304. 'x-hub-signature-256': signature,
  305. },
  306. body,
  307. })
  308. }
  309. describe.skipIf(!process.env.DEEPSEEK_API_KEY)('GitHub webhook through the real dsh CLI and model', () => {
  310. it('creates, attaches, prompts, and completes a Workspace Session', async () => {
  311. expect(existsSync(BUILT_BIN), `missing built CLI ${BUILT_BIN}; run pnpm run build:official`).toBe(true)
  312. const root = await mkdtemp(join(tmpdir(), 'dsh-github-webhook-real-'))
  313. const workspacePath = join(root, 'workspace')
  314. await mkdir(workspacePath)
  315. const canonicalWorkspacePath = await realpath(workspacePath)
  316. const webhookPort = await freePort()
  317. const child = spawn(process.execPath, [
  318. BUILT_BIN,
  319. 'web',
  320. '--patch', OVERLAY,
  321. '--no-open',
  322. '--host', '127.0.0.1',
  323. '--port', '0',
  324. ], {
  325. cwd: root,
  326. env: {
  327. ...process.env,
  328. DSH_AGENTS_HOME: join(root, '.agents'),
  329. DSH_GITHUB_E2E_MARKER: MARKER,
  330. DSH_GITHUB_E2E_WORKSPACE: workspacePath,
  331. DSH_GITHUB_WEBHOOK_PORT: String(webhookPort),
  332. DSH_GITHUB_WEBHOOK_SECRET: SECRET,
  333. DSH_HOME: join(root, '.dsh'),
  334. DSH_TELEMETRY_DISABLED: '1',
  335. },
  336. stdio: ['ignore', 'pipe', 'pipe'],
  337. })
  338. const observation = observeProcess(child)
  339. try {
  340. const baseUrl = await observation.ready
  341. const webhookOrigin = `http://127.0.0.1:${String(webhookPort)}`
  342. expect((await fetch(`${webhookOrigin}/api`)).status).toBe(404)
  343. expect((await sendGitHubDelivery(baseUrl)).status).not.toBe(202)
  344. expect((await sendGitHubDelivery(webhookOrigin)).status).toBe(202)
  345. const workspaces = await eventually(
  346. child,
  347. observation.text,
  348. 'one Workspace-attached Session',
  349. async () => await workspaceBaseline(baseUrl),
  350. value => value.items.some(workspace =>
  351. workspace.path === canonicalWorkspacePath && workspace.sessionIds.length === 1),
  352. 30_000,
  353. )
  354. const workspace = workspaces.items.find(item => item.path === canonicalWorkspacePath)
  355. const sessionId = workspace?.sessionIds[0]
  356. if (sessionId === undefined) throw new Error('workspace/follow did not expose the webhook Session')
  357. const sessions = await remoteRpc<SessionList>(baseUrl, 'session/list', { _request: {} })
  358. expect(sessions.items.find(session => session.sessionId === sessionId)).toMatchObject({
  359. agentPreset: 'minimal',
  360. blank: false,
  361. cwd: canonicalWorkspacePath,
  362. })
  363. const admitted = await eventually(
  364. child,
  365. observation.text,
  366. 'webhook provenance, title, and permission events',
  367. async () => await history(baseUrl, sessionId),
  368. (page) => {
  369. const events = page.events.map(item => item.event)
  370. const title = events.find(event => event.type === 'session/title')
  371. const permission = events.find(event =>
  372. event.type === 'permission/preset'
  373. && isRecord(event.data)
  374. && event.data.preset === 'read-only')
  375. const message = events.find(event =>
  376. event.type === 'user/message'
  377. && isRecord(event.data)
  378. && isRecord(event.data.source)
  379. && event.data.source.kind === 'webhook')
  380. return isRecord(title?.data) && title.data.title === TITLE
  381. && permission !== undefined
  382. && isRecord(message?.data) && isRecord(message.data.source)
  383. && message.data.source.provider === 'github'
  384. && message.data.source.deliveryId === DELIVERY
  385. },
  386. 30_000,
  387. )
  388. const webhookMessage = admitted.events.map(item => item.event)
  389. .find(event => event.type === 'user/message'
  390. && isRecord(event.data)
  391. && isRecord(event.data.source)
  392. && event.data.source.kind === 'webhook')
  393. expect(webhookMessage?.data).toMatchObject({
  394. content: [{ type: 'text', text: `Reply with exactly ${MARKER} and no other text. Do not call tools.` }],
  395. source: {
  396. kind: 'webhook',
  397. provider: 'github',
  398. deliveryId: DELIVERY,
  399. ruleId: 'github-real-e2e',
  400. source: 'github-real-e2e',
  401. },
  402. })
  403. const completed = await eventually(
  404. child,
  405. observation.text,
  406. 'a real DeepSeek assistant response',
  407. async () => await history(baseUrl, sessionId),
  408. page => assistantText(page).includes(MARKER),
  409. 150_000,
  410. )
  411. expect(assistantText(completed)).toContain(MARKER)
  412. } finally {
  413. await stop(child)
  414. await rm(root, { recursive: true, force: true })
  415. }
  416. }, 330_000)
  417. })