Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion api/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,13 @@ export const start = async () => {
// await upgradeScripts(mongo.db, resolve(import.meta.dirname, '../..'))

await wsServer.start(server, mongo.db, async (channel, sessionState) => {
const [ownerType, ownerId] = channel.split(':')
const [ownerType, ownerId, resource] = channel.split(':')
if (!sessionState.user) return false
if (sessionState.user.adminMode) return true
// same rule as the webhooks routes: an admin of the active account that owns the subscription
if (resource?.startsWith('webhook-subscriptions/')) {
return sessionState.account?.type === ownerType && sessionState.account.id === ownerId && sessionState.accountRole === 'admin'
}
return ownerType === 'user' && ownerId === sessionState.user.id
})
await wsEmitter.init(mongo.db)
Expand Down
5 changes: 3 additions & 2 deletions api/src/webhook-subscriptions/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,9 @@ router.get('', async (req, res, next) => {
query['sender.type'] = req.query.sender.split(':')[0]
query['sender.id'] = req.query.sender.split(':')[1]
}
if (req.query.topic) {
query['topic.key'] = req.query.topic
// a comma-separated list of topic keys, as the subscribe-webhooks embed page receives them
if (req.query.topic && typeof req.query.topic === 'string') {
query['topic.key'] = { $in: req.query.topic.split(',') }
}

const [results, count] = await Promise.all([
Expand Down
3 changes: 3 additions & 0 deletions api/src/webhooks/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { Filter, Sort } from 'mongodb'
import { Router } from 'express'
import { session, mongoPagination, httpError } from '@data-fair/lib-express/index.js'
import mongo from '#mongo'
import { emitWebhook } from './service.ts'

const router = Router()
export default router
Expand Down Expand Up @@ -34,6 +35,7 @@ router.post('/:id/_retry', async (req, res, next) => {
{ returnDocument: 'after' })

if (!webhook) throw httpError(404)
await emitWebhook(webhook)
res.send(webhook)
})

Expand All @@ -47,5 +49,6 @@ router.post('/:id/_cancel', async (req, res, next) => {
{ returnDocument: 'after' })

if (!webhook) throw httpError(404)
await emitWebhook(webhook)
res.send(webhook)
})
14 changes: 12 additions & 2 deletions api/src/webhooks/service.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,18 @@
import type { LocalizedEvent, Webhook, WebhookSubscription } from '#types'

import { nanoid } from 'nanoid'
import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js'
import mongo from '#mongo'
import { coalesceAction, type CoalesceTarget } from './operations.ts'

// the delivery history of a webhook subscription, listened to by its embed page to show progress live
export const webhooksChannel = (owner: Webhook['owner'], subscriptionId: string) =>
`${owner.type}:${owner.id}:webhook-subscriptions/${subscriptionId}/webhooks`

export const emitWebhook = async (webhook: Webhook | null) => {
if (webhook) await wsEmitter.emit(webhooksChannel(webhook.owner, webhook.subscription._id), webhook)
}

export const createWebhook = async (event: LocalizedEvent, webhookSubscription: WebhookSubscription, opts: { coalesce?: boolean } = {}) => {
const notification: Webhook['notification'] = {
title: event.title,
Expand All @@ -25,8 +34,8 @@ export const createWebhook = async (event: LocalizedEvent, webhookSubscription:
? { $set: { notification, status: 'waiting' as const, nbAttempts: 0 }, $unset: { nextAttempt: '' as const, lastAttempt: '' as const } }
: { $set: { notification } }
// matching on the observed status: if the worker grabbed it meanwhile, fall through to an insert
const res = await mongo.webhooks.updateOne({ _id: action._id, status: action.status }, update)
if (res.matchedCount) return
const updated = await mongo.webhooks.findOneAndUpdate({ _id: action._id, status: action.status }, update, { returnDocument: 'after' })
if (updated) return emitWebhook(updated)
}
}

Expand All @@ -43,4 +52,5 @@ export const createWebhook = async (event: LocalizedEvent, webhookSubscription:
nbAttempts: 0
}
await mongo.webhooks.insertOne(webhook)
await emitWebhook(webhook)
}
12 changes: 7 additions & 5 deletions api/src/webhooks/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import axios from '@data-fair/lib-node/axios.js'
import { internalError } from '@data-fair/lib-node/observer.js'
import locks from '@data-fair/lib-node/locks.js'
import { nextAttemptDate } from './operations.ts'
import { emitWebhook } from './service.ts'

const debug = Debug('webhooks-worker')

Expand Down Expand Up @@ -39,18 +40,19 @@ const loop = async () => {
continue
}
debug('work on webhook', webhook)
await emitWebhook(webhook)
const date = new Date().toISOString()
const subscription = await mongo.webhookSubscriptions
.findOne({ _id: webhook.subscription._id, 'owner.type': webhook.owner.type, 'owner.id': webhook.owner.id })
if (!subscription) {
debug('missing subscription for webhook, store as error')
await mongo.webhooks.updateOne({ _id: webhook._id }, {
await emitWebhook(await mongo.webhooks.findOneAndUpdate({ _id: webhook._id }, {
$set: {
status: 'error',
lastAttempt: { date, error: 'missing subscription' }
},
$unset: { nextAttempt: '' }
})
}, { returnDocument: 'after' }))
continue
}
debug('found matching subscription', subscription)
Expand All @@ -62,14 +64,14 @@ const loop = async () => {
debug('send webhook', subscription.url, webhook.notification)
const res = await axios.post(subscription.url, webhook.notification, { headers, timeout: 2000 })
debug('webhook success')
await mongo.webhooks.updateOne({ _id: webhook._id }, {
await emitWebhook(await mongo.webhooks.findOneAndUpdate({ _id: webhook._id }, {
$set: {
status: 'ok',
lastAttempt: { date, status: res.status }
},
$unset: { nextAttempt: '' },
$inc: { nbAttempts: 1 }
})
}, { returnDocument: 'after' }))
} catch (err: any) {
debug('webhook failed', err)
const attempt: Webhook['lastAttempt'] = { date }
Expand All @@ -89,7 +91,7 @@ const loop = async () => {
patch.$set.nextAttempt = nextAttemptDate(webhook.nbAttempts + 1)
debug('webhook failed, progressively backoff', patch.$set.nextAttempt)
}
await mongo.webhooks.updateOne({ _id: webhook._id }, patch)
await emitWebhook(await mongo.webhooks.findOneAndUpdate({ _id: webhook._id }, patch, { returnDocument: 'after' }))
}
}
await locks.release('webhooks-loop')
Expand Down
67 changes: 67 additions & 0 deletions tests/webhooks.api.spec.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,29 @@
import { test, expect } from '@playwright/test'
import { createServer } from 'node:http'
import WebSocket from 'ws'
import { axios, axiosAuth, clean, devBaseURL } from './support/axios.ts'

const axPush = axios({ params: { key: 'SECRET_EVENTS' }, baseURL: devBaseURL })
const axDev = axios({ baseURL: devBaseURL })
const admin1 = await axiosAuth('test1-admin1')
const user1 = await axiosAuth('test-user1')

// open a WS as this user and subscribe to a channel, resolves with the subscription outcome
const wsSubscribe = async (ax: typeof admin1, channel: string) => {
const cookies = ax.cookieJar.getCookiesSync(`http://${process.env.DEV_HOST}:${process.env.NGINX_PORT}`)
const ws = new WebSocket(`ws://localhost:${process.env.DEV_API_PORT}`, { headers: { Cookie: cookies.map(String).join('; ') } })
const messages: any[] = []
const outcome = await new Promise<string>((resolve, reject) => {
ws.on('message', (raw: Buffer) => {
const msg = JSON.parse(raw.toString())
if (msg.type === 'subscribe-confirm' || msg.type === 'error') resolve(msg.type)
if (msg.type === 'message') messages.push(msg.data)
})
ws.on('open', () => ws.send(JSON.stringify({ type: 'subscribe', channel })))
ws.on('error', reject)
})
return { ws, messages, outcome }
}

// helper to post event matching a webhook subscription owned by test1-admin1/test1
const postMatchingEvent = (title: string) => axPush.post('/api/events', [{
Expand Down Expand Up @@ -158,6 +177,54 @@ test.describe('webhooks', () => {
}
})

test('should list webhook subscriptions of several topics', async () => {
for (const key of ['topic1', 'topic2', 'topic3']) {
await admin1.post('/api/webhook-subscriptions', {
title: 'Sub ' + key,
topic: { key },
sender: { type: 'organization', id: 'test1' },
url: 'http://localhost:19881/hook'
})
}
const one = await admin1.get('/api/webhook-subscriptions', { params: { topic: 'topic1' } })
expect(one.data.results.map((s: any) => s.topic.key)).toEqual(['topic1'])
const two = await admin1.get('/api/webhook-subscriptions', { params: { topic: 'topic1,topic3', sort: 'title:1' } })
expect(two.data.results.map((s: any) => s.topic.key)).toEqual(['topic1', 'topic3'])
})

test('should stream the delivery progress of a webhook subscription over WS', async () => {
const hookServer = createServer((req, res) => { req.resume(); req.on('end', () => { res.writeHead(200); res.end() }) })
await new Promise<void>(resolve => hookServer.listen(19882, resolve))
try {
const sub = (await admin1.post('/api/webhook-subscriptions', {
title: 'WS progress',
topic: { key: 'topic1' },
sender: { type: 'organization', id: 'test1' },
url: 'http://localhost:19882/hook'
})).data
const channel = `${sub.owner.type}:${sub.owner.id}:webhook-subscriptions/${sub._id}/webhooks`

// only an admin of the owning account can listen
const other = await wsSubscribe(user1, channel)
other.ws.close()
expect(other.outcome).toBe('error')

const { ws, messages, outcome } = await wsSubscribe(admin1, channel)
try {
expect(outcome).toBe('subscribe-confirm')
await postMatchingEvent('ws progress')
// the worker polls every 4s
await expect.poll(() => messages.map(m => m.status), { timeout: 12000 }).toEqual(['waiting', 'working', 'ok'])
expect(new Set(messages.map(m => m._id)).size).toBe(1)
expect(messages[2].notification.title).toBe('ws progress')
} finally {
ws.close()
}
} finally {
hookServer.close()
}
})

test('should cancel a webhook', async () => {
await admin1.post('/api/webhook-subscriptions', {
title: 'Cancel test',
Expand Down
110 changes: 110 additions & 0 deletions tests/webhooks.e2e.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
import { expect, type Page } from '@playwright/test'
import { test } from './fixtures/login.ts'
import { axiosAuth, clean } from './support/axios.ts'

const pageUrl = (keys: string[], titles: string[]) =>
`/events/embed/subscribe-webhooks?key=${keys.join(',')}&title=${titles.join(',')}&sender=user:test-user1`

const saveButton = (page: Page) => page.getByRole('button', { name: 'Enregistrer' }).first()

test.describe('Webhooks UI', () => {
test.beforeEach(clean)

test('creates a webhook on a chosen topic', async ({ page, goToWithAuth }) => {
await goToWithAuth(pageUrl(['topic1', 'topic2'], ['TopicOne', 'TopicTwo']), 'test-user1')
await page.getByText('Declare a new webhook').click()
await expect(saveButton(page)).toBeDisabled()

await page.getByLabel('Libellé').fill('My webhook')
await page.getByLabel('URL').fill('https://example.com/hook')
await saveButton(page).click()
// several topics are proposed, none is preselected
await expect(page.getByText('Ce paramètre est requis')).toBeVisible()

await page.locator('.v-select').click()
await page.getByRole('option', { name: 'TopicTwo' }).click()
await saveButton(page).click()

await expect(page.getByText('This webhook has not been called yet.')).toBeVisible()
await expect(page.getByRole('button', { name: /My webhook\s*TopicTwo/ })).toBeVisible()
await expect(saveButton(page)).toBeDisabled()
})

test('does not propose a topic selector for a single topic', async ({ page, goToWithAuth }) => {
await goToWithAuth(pageUrl(['topic1'], ['TopicOne']), 'test-user1')
await page.getByText('Declare a new webhook').click()
await expect(page.getByLabel('Libellé')).toBeVisible()
await expect(page.locator('.v-select')).toHaveCount(0)
})

test('edits an existing webhook', async ({ page, goToWithAuth }) => {
const user1 = await axiosAuth('test-user1')
await user1.post('/api/webhook-subscriptions', {
title: 'Existing webhook',
topic: { key: 'topic1', title: 'TopicOne' },
sender: { type: 'user', id: 'test-user1' },
url: 'https://example.com/hook'
})

await goToWithAuth(pageUrl(['topic1', 'topic2'], ['TopicOne', 'TopicTwo']), 'test-user1')
await page.getByRole('button', { name: /Existing webhook/ }).click()
await expect(saveButton(page)).toBeDisabled()

await page.getByLabel('Libellé').first().fill('Renamed webhook')
await expect(saveButton(page)).toBeEnabled()
await page.locator('.v-select').first().click()
await page.getByRole('option', { name: 'TopicTwo' }).click()
const saved = page.waitForResponse(res => res.request().method() === 'POST' && res.url().endsWith('/webhook-subscriptions'))
await saveButton(page).click()
expect((await saved).status()).toBe(200)
// the refreshed subscription (new "updated" date) is not a change to save
await expect(page.getByRole('button', { name: /Renamed webhook\s*TopicTwo/ })).toBeVisible()
await expect(saveButton(page)).toBeDisabled()

const res = await user1.get('/api/webhook-subscriptions')
expect(res.data.results).toHaveLength(1)
expect(res.data.results[0].title).toBe('Renamed webhook')
expect(res.data.results[0].topic.key).toBe('topic2')
})

test('deletes a webhook', async ({ page, goToWithAuth }) => {
const user1 = await axiosAuth('test-user1')
await user1.post('/api/webhook-subscriptions', {
title: 'Doomed webhook',
topic: { key: 'topic1', title: 'TopicOne' },
sender: { type: 'user', id: 'test-user1' },
url: 'https://example.com/hook'
})

await goToWithAuth(pageUrl(['topic1'], ['TopicOne']), 'test-user1')
await page.getByRole('button', { name: /Doomed webhook/ }).click()
await page.getByRole('button', { name: 'Supprimer' }).click()
await page.getByRole('button', { name: 'Yes', exact: true }).click()
await expect(page.getByText('Doomed webhook')).toHaveCount(0)
expect((await user1.get('/api/webhook-subscriptions')).data.count).toBe(0)
})

test('shows the progress of a test call without refreshing', async ({ page, goToWithAuth }) => {
const user1 = await axiosAuth('test-user1')
await user1.post('/api/webhook-subscriptions', {
title: 'Tested webhook',
topic: { key: 'topic1', title: 'TopicOne' },
sender: { type: 'user', id: 'test-user1' },
// nothing listens there, the delivery ends in error
url: 'http://localhost:19898/closed'
})

await goToWithAuth(pageUrl(['topic1'], ['TopicOne']), 'test-user1')
await page.getByRole('button', { name: /Tested webhook/ }).click()
await expect(page.getByText('This webhook has not been called yet.')).toBeVisible()

let listRequests = 0
page.on('request', req => { if (req.url().includes('/api/webhooks?')) listRequests++ })
await page.getByRole('button', { name: 'Test', exact: true }).click()
await expect(page.getByText(/ - waiting| - working/)).toBeVisible()
const listRequestsAfterTest = listRequests
// the worker polls every 4s, its outcome is pushed over WS
await expect(page.getByText(/ - error/)).toBeVisible({ timeout: 15000 })
expect(listRequests).toBe(listRequestsAfterTest)
})
})
Loading
Loading