github-webhook-real.e2e.ts 18 KB

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