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
1 change: 1 addition & 0 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@
"platformatic": "^3.52.0",
"postgres-migrations": "^5.3.0",
"pprof-format": "^2.2.1",
"undici": "^7.24.0",
"wattpm": "^3.52.0",
"xml2js": "^0.6.2"
},
Expand Down
100 changes: 75 additions & 25 deletions src/storage/events/lifecycle/webhook.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
import { getTenantConfig } from '@internal/database'
import { logger, logSchema } from '@internal/monitoring'
import HttpAgent, { HttpsAgent } from 'agentkeepalive'
import axios from 'axios'
import { Job, SendOptions, WorkOptions } from 'pg-boss'
import { Agent } from 'undici'
import { getConfig } from '../../../config'
import { BaseEvent } from '../base-event'

Expand All @@ -14,6 +13,11 @@ const {
webhookMaxConnections,
webhookQueueMaxFreeSockets,
} = getConfig()
const WEBHOOK_TIMEOUT_MS = 4000
const webhookKeepAliveTimeoutMs = Math.min(
Math.max(webhookQueueMaxFreeSockets, 1) * 1000,
WEBHOOK_TIMEOUT_MS
)
Comment thread
ferhatelmas marked this conversation as resolved.

interface WebhookEvent {
event: {
Expand All @@ -30,31 +34,77 @@ interface WebhookEvent {
}
}

const httpAgent = webhookURL?.startsWith('https://')
? {
httpsAgent: new HttpsAgent({
maxSockets: webhookMaxConnections,
maxFreeSockets: webhookQueueMaxFreeSockets,
}),
}
: {
httpAgent: new HttpAgent({
maxSockets: webhookMaxConnections,
maxFreeSockets: webhookQueueMaxFreeSockets,
}),
interface WebhookRequest {
type: 'Webhook'
event: WebhookEvent['event']
sentAt: Date
Comment thread
ferhatelmas marked this conversation as resolved.
tenant: WebhookEvent['tenant']
}

interface WebhookClient {
post(url: string, payload: WebhookRequest): Promise<void>
}

const dispatcher = new Agent({
connections: webhookMaxConnections,
// `undici` cannot cap idle socket count like `agentkeepalive.maxFreeSockets`,
// so use the old knob to make idle sockets expire sooner when a small free pool is desired.
keepAliveTimeout: webhookKeepAliveTimeoutMs,
keepAliveMaxTimeout: webhookKeepAliveTimeoutMs,
})

const defaultHeaders = new Headers({
'content-type': 'application/json',
...(webhookApiKey ? { authorization: `Bearer ${webhookApiKey}` } : {}),
})

async function assertOkResponse(response: Response) {
if (response.ok) {
return
}

throw new Error(`Request failed with status code ${response.status}`)
}

function normalizeWebhookError(error: unknown) {
if (error instanceof DOMException && error.name === 'TimeoutError') {
return new Error(`timeout of ${WEBHOOK_TIMEOUT_MS}ms exceeded`)
}

if (error instanceof Error) {
return error
}

return new Error(String(error))
}

const client: WebhookClient = {
async post(url, payload) {
const requestInit: RequestInit & { dispatcher: Agent } = {
method: 'POST',
body: JSON.stringify(payload),
headers: defaultHeaders,
dispatcher,
signal: AbortSignal.timeout(WEBHOOK_TIMEOUT_MS),
Comment thread
ferhatelmas marked this conversation as resolved.
}

const client = axios.create({
...httpAgent,
timeout: 4000,
headers: {
...(webhookApiKey ? { authorization: `Bearer ${webhookApiKey}` } : {}),
const response = await fetch(url, requestInit)

try {
await assertOkResponse(response)
} finally {
await response.body?.cancel().catch(() => {})
}
},
})
}

export class Webhook extends BaseEvent<WebhookEvent> {
static queueName = 'webhooks'

protected static getClient() {
return client
}

static getWorkerOptions(): WorkOptions {
return {
pollingIntervalSeconds: webhookQueuePullInterval
Expand Down Expand Up @@ -110,16 +160,18 @@ export class Webhook extends BaseEvent<WebhookEvent> {
})

try {
await client.post(webhookURL, {
await this.getClient().post(webhookURL, {
type: 'Webhook',
event: job.data.event,
sentAt: new Date(),
tenant: job.data.tenant,
})
} catch (e) {
const error = normalizeWebhookError(e)

logger.error(
{
error: (e as Error)?.message,
error: error.message,
jodId: job.id,
type: 'event',
event: job.data.event.type,
Expand All @@ -134,9 +186,7 @@ export class Webhook extends BaseEvent<WebhookEvent> {
`[Lifecycle]: ${job.data.event.type} ${path} - FAILED`
)
throw new Error(
`Failed to send webhook for event ${job.data.event.type} to ${webhookURL}: ${
(e as Error).message
}`
`Failed to send webhook for event ${job.data.event.type} to ${webhookURL}: ${error.message}`
)
}

Expand Down
4 changes: 3 additions & 1 deletion src/storage/renderer/image.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,14 +72,16 @@ interface TransformLimits {
maxResolution?: number
}

type ImageRendererClient = Pick<Axios, 'get'>

/**
* ImageRenderer
* renders an image by applying transformations
*
* Interacts with an imgproxy backend for the actual transformation
*/
export class ImageRenderer extends Renderer {
private readonly client: Axios
private readonly client: ImageRendererClient
private transformOptions?: TransformOptions
private limits?: TransformLimits

Expand Down
48 changes: 10 additions & 38 deletions src/test/cdn.test.ts
Original file line number Diff line number Diff line change
@@ -1,41 +1,16 @@
import { CdnCacheManager } from '@storage/cdn/cdn-cache-manager'
import { FastifyInstance } from 'fastify'
import { SignJWT } from 'jose'
import { Readable } from 'stream'
import { getConfig, mergeConfig } from '../config'
import { useStorage } from './utils/storage'

getConfig()
mergeConfig({
cdnPurgeEndpointURL: 'http://localhost/stub/cache',
cdnPurgeEndpointKey: 'test-key',
})

vi.mock('axios', () => {
const instance = {
post: vi.fn(),
interceptors: {
request: {
use: vi.fn(),
},
response: {
use: vi.fn(),
},
},
}

const axiosMock = {
create: vi.fn().mockReturnValue(instance),
...instance,
}

return {
default: axiosMock,
...axiosMock,
}
})

import axios from 'axios'
import { FastifyInstance } from 'fastify'
import { SignJWT } from 'jose'
import { Readable } from 'stream'
import { useStorage } from './utils/storage'

const { serviceKeyAsync, anonKeyAsync, tenantId, jwtSecret } = getConfig()

describe('CDN Cache Manager', () => {
Expand All @@ -58,6 +33,7 @@ describe('CDN Cache Manager', () => {

afterEach(async () => {
await appInstance.close()
vi.restoreAllMocks()
vi.clearAllMocks()
})

Expand Down Expand Up @@ -108,9 +84,7 @@ describe('CDN Cache Manager', () => {
},
})

const spy = vi
.spyOn(axios, 'post')
.mockReturnValue(Promise.resolve({ data: { message: 'success' } }))
const purgeSpy = vi.spyOn(CdnCacheManager.prototype, 'purge').mockResolvedValue(undefined)

Comment thread
ferhatelmas marked this conversation as resolved.
const response = await appInstance.inject({
method: 'DELETE',
Expand All @@ -124,11 +98,9 @@ describe('CDN Cache Manager', () => {

const body = await response.json()
expect(body).toEqual({ message: 'success' })
expect(spy).toHaveBeenCalledWith('/purge', {
tenant: {
ref: tenantId,
},
bucketId: bucketName,
expect(purgeSpy).toHaveBeenCalledWith({
tenant: tenantId,
bucket: bucketName,
objectName,
})
})
Expand Down
Loading