Skip to content
This repository was archived by the owner on Nov 4, 2021. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
e85dddd
extract redlock from schedule
mariusandra Apr 20, 2021
7333c0c
implement generic retrying
mariusandra Apr 20, 2021
9d75e8b
capture console.log in tests via a temp file
mariusandra Apr 20, 2021
6b3a8c0
add graphile queue
mariusandra Apr 20, 2021
c75b7a6
make it prettier and safe
mariusandra Apr 20, 2021
5eab403
style fixes
mariusandra Apr 20, 2021
185a187
fix some tests
mariusandra Apr 20, 2021
39c28af
release if there
mariusandra Apr 20, 2021
9665f44
split postgres tests
mariusandra Apr 21, 2021
6e3e034
don't make a graphile worker in all tests
mariusandra Apr 21, 2021
3e34bb2
revert "split postgres tests"
mariusandra Apr 21, 2021
806fc00
Merge branch 'master' into retry-queue
mariusandra Apr 28, 2021
c4845e6
skip retries if pluginConfig not found
mariusandra Apr 28, 2021
6fe7eba
reset graphile schema before test
mariusandra Apr 29, 2021
16e14b6
fix failing tests by clearing the retry consumer redlock
mariusandra Apr 29, 2021
c20c67c
bust github actions cache
mariusandra Apr 29, 2021
a2f32e6
slight cleanup
mariusandra Apr 29, 2021
b1ee2c9
fix github/eslint complaining about an `any`
mariusandra Apr 29, 2021
10eb8c4
separate url for graphile retry queue, otherwise use existing postgre…
mariusandra Apr 29, 2021
10a4f1c
Merge branch 'master' into retry-queue
mariusandra Apr 29, 2021
a2e05b9
convert startRedlock params to options object
mariusandra Apr 29, 2021
f88ac09
Merge branch 'retry-queue' of github.com:PostHog/posthog-plugin-serve…
mariusandra Apr 29, 2021
a026680
move type around
mariusandra Apr 29, 2021
8d002d9
use an enum
mariusandra Apr 29, 2021
3990953
update typo in comment
mariusandra Apr 30, 2021
b70e7c1
Merge branch 'master' into retry-queue
Twixes Apr 30, 2021
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 .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,4 @@ dist/
yalc.lock
.yalc/
src/idl/protos.*
tmp
2 changes: 1 addition & 1 deletion benchmarks/clickhouse/e2e.kafka.benchmark.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ describe('e2e kafka & clickhouse benchmark', () => {

// hope that 5sec is enough to load kafka with all the events (posthog.capture can't be awaited)
await delay(5000)
queue.resume()
await queue.resume()

console.log('Starting timer')
const startTime = performance.now()
Expand Down
2 changes: 1 addition & 1 deletion benchmarks/clickhouse/e2e.timeout.benchmark.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ describe('e2e kafka processing timeout benchmark', () => {

// hope that 5sec is enough to load kafka with all the events (posthog.capture can't be awaited)
await delay(5000)
queue.resume()
await queue.resume()

console.log('Starting timer')
const startTime = performance.now()
Expand Down
2 changes: 1 addition & 1 deletion benchmarks/postgres/e2e.celery.benchmark.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ describe('e2e celery & postgres benchmark', () => {
}
await delay(3000)
expect(await redis.llen(server.PLUGINS_CELERY_QUEUE)).toEqual(count)
queue.resume()
await queue.resume()

console.log('Starting timer')
const startTime = performance.now()
Expand Down
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@
"fast-deep-equal": "^3.1.3",
"fastify": "^3.8.0",
"generic-pool": "^3.7.1",
"graphile-worker": "^0.11.0",
"hot-shots": "^8.2.1",
"ioredis": "^4.19.2",
"kafkajs": "^1.15.0",
Expand Down
10 changes: 8 additions & 2 deletions src/main/pluginsServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ import { defaultConfig } from '../shared/config'
import { createServer } from '../shared/server'
import { status } from '../shared/status'
import { createRedis, delay, getPiscinaStats } from '../shared/utils'
import { PluginsServer, PluginsServerConfig, Queue, ScheduleControl } from '../types'
import { PluginsServer, PluginsServerConfig, Queue, RetryQueueConsumerControl, ScheduleControl } from '../types'
import { createMmdbServer, performMmdbStalenessCheck, prepareMmdb } from './mmdb'
import { startQueue } from './queue'
import { startRetryQueueConsumer } from './services/retry-queue-consumer'
import { startSchedule } from './services/schedule'
import { startFastifyInstance, stopFastifyInstance } from './web/server'

Expand Down Expand Up @@ -46,6 +47,7 @@ export async function startPluginsServer(
let statsJob: schedule.Job | undefined
let piscina: Piscina | undefined
let queue: Queue | undefined
let retryQueueConsumer: RetryQueueConsumerControl | undefined
let closeServer: () => Promise<void> | undefined
let scheduleControl: ScheduleControl | undefined
let mmdbServer: net.Server | undefined
Expand Down Expand Up @@ -73,6 +75,7 @@ export async function startPluginsServer(
await pubSub?.quit()
pingJob && schedule.cancelJob(pingJob)
statsJob && schedule.cancelJob(statsJob)
await retryQueueConsumer?.stop()
await scheduleControl?.stopSchedule()
await new Promise<void>((resolve, reject) =>
!mmdbServer
Expand Down Expand Up @@ -129,9 +132,12 @@ export async function startPluginsServer(
}

scheduleControl = await startSchedule(server, piscina)
retryQueueConsumer = await startRetryQueueConsumer(server, piscina)

queue = await startQueue(server, piscina)
piscina.on('drain', () => {
queue?.resume()
void queue?.resume()
void retryQueueConsumer?.resume()
})

// use one extra connection for redis pubsub
Expand Down
14 changes: 9 additions & 5 deletions src/main/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,13 @@ export type WorkerMethods = {
ingestEvent: (event: PluginEvent) => Promise<IngestEventResponse>
}

function pauseQueueIfWorkerFull(queue: Queue | undefined, server: PluginsServer, piscina?: Piscina) {
if (queue && (piscina?.queueSize || 0) > (server.WORKER_CONCURRENCY || 4) * (server.WORKER_CONCURRENCY || 4)) {
void queue.pause()
export function pauseQueueIfWorkerFull(
pause: undefined | (() => void | Promise<void>),
server: PluginsServer,
piscina?: Piscina
): void {
if (pause && (piscina?.queueSize || 0) > (server.WORKER_CONCURRENCY || 4) * (server.WORKER_CONCURRENCY || 4)) {
void pause()
}
}

Expand Down Expand Up @@ -81,10 +85,10 @@ function startQueueRedis(server: PluginsServer, piscina: Piscina | undefined, wo
...data,
} as PluginEvent)
try {
pauseQueueIfWorkerFull(celeryQueue, server, piscina)
pauseQueueIfWorkerFull(() => celeryQueue.pause(), server, piscina)
const processedEvent = await workerMethods.processEvent(event)
if (processedEvent) {
pauseQueueIfWorkerFull(celeryQueue, server, piscina)
pauseQueueIfWorkerFull(() => celeryQueue.pause(), server, piscina)
await workerMethods.ingestEvent(processedEvent)
}
} catch (e) {
Expand Down
75 changes: 75 additions & 0 deletions src/main/retry/fs-queue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
import { EnqueuedRetry, OnRetryCallback, RetryQueue } from '../../types'
import Timeout = NodeJS.Timeout
import * as fs from 'fs'
import * as path from 'path'

export class FsQueue implements RetryQueue {
paused: boolean
started: boolean
interval: Timeout | null
filename: string

constructor(filename?: string) {
if (process.env.NODE_ENV !== 'test') {
throw new Error('Can not use FsQueue outside tests')
}
this.paused = false
this.started = false
this.interval = null
this.filename = filename || path.join(process.cwd(), 'tmp', 'fs-queue.txt')

fs.mkdirSync(path.dirname(this.filename), { recursive: true })
fs.writeFileSync(this.filename, '')
}

enqueue(retry: EnqueuedRetry): Promise<void> | void {
fs.appendFileSync(this.filename, `${JSON.stringify(retry)}\n`)
}

quit(): void {
// nothing to do
}

startConsumer(onRetry: OnRetryCallback): void {
fs.writeFileSync(this.filename, '')
this.started = true
this.interval = setInterval(() => {
if (this.paused) {
return
}
const timestamp = new Date().valueOf()
const queue = fs
.readFileSync(this.filename)
.toString()
.split('\n')
.filter((a) => a)
.map((s) => JSON.parse(s) as EnqueuedRetry)

const newQueue = queue.filter((element) => element.timestamp < timestamp)
if (newQueue.length > 0) {
const oldQueue = queue.filter((element) => element.timestamp >= timestamp)
fs.writeFileSync(this.filename, `${oldQueue.map((q) => JSON.stringify(q)).join('\n')}\n`)

void onRetry(newQueue)
}
}, 1000)
}

stopConsumer(): void {
this.started = false
this.interval && clearInterval(this.interval)
fs.unlinkSync(this.filename)
}

pauseConsumer(): void {
this.paused = true
}

isConsumerPaused(): boolean {
return this.paused
}

resumeConsumer(): void {
this.paused = false
}
}
94 changes: 94 additions & 0 deletions src/main/retry/graphile-queue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
import { makeWorkerUtils, run, Runner, WorkerUtils, WorkerUtilsOptions } from 'graphile-worker'

import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types'

export class GraphileQueue implements RetryQueue {
pluginsServer: PluginsServer
started: boolean
paused: boolean
onRetry: OnRetryCallback | null
runner: Runner | null
workerUtils: WorkerUtils | null

constructor(pluginsServer: PluginsServer) {
this.pluginsServer = pluginsServer
this.started = false
this.paused = false
this.onRetry = null
this.runner = null
this.workerUtils = null
}

async enqueue(retry: EnqueuedRetry): Promise<void> {
if (!this.workerUtils) {
this.workerUtils = await makeWorkerUtils(
this.pluginsServer.RETRY_QUEUE_GRAPHILE_URL
? {
connectionString: this.pluginsServer.RETRY_QUEUE_GRAPHILE_URL,
}
: ({
pgPool: this.pluginsServer.postgres,
} as WorkerUtilsOptions)
)
await this.workerUtils.migrate()
}
await this.workerUtils.addJob('retryTask', retry, { runAt: new Date(retry.timestamp), maxAttempts: 1 })
}

async quit(): Promise<void> {
const oldWorkerUtils = this.workerUtils
this.workerUtils = null
await oldWorkerUtils?.release()
}

async startConsumer(onRetry: OnRetryCallback): Promise<void> {
this.started = true
this.onRetry = onRetry
await this.syncState()
}

async stopConsumer(): Promise<void> {
this.started = false
await this.syncState()
}

async pauseConsumer(): Promise<void> {
this.paused = true
await this.syncState()
}

isConsumerPaused(): boolean {
return this.paused
}

async resumeConsumer(): Promise<void> {
this.paused = false
await this.syncState()
}

async syncState(): Promise<void> {
if (this.started && !this.paused) {
if (!this.runner) {
this.runner = await run({
connectionString: this.pluginsServer.DATABASE_URL,
concurrency: 1,
// Install signal handlers for graceful shutdown on SIGINT, SIGTERM, etc
noHandleSignals: false,
pollInterval: 100,
// you can set the taskList or taskDirectory but not both
taskList: {
retryTask: (payload) => {
void this.onRetry?.([payload as EnqueuedRetry])
},
},
})
}
} else {
if (this.runner) {
const oldRunner = this.runner
this.runner = null
await oldRunner?.stop()
}
}
}
}
80 changes: 80 additions & 0 deletions src/main/retry/retry-queue-manager.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
import * as Sentry from '@sentry/node'

import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types'
import { FsQueue } from './fs-queue'
import { GraphileQueue } from './graphile-queue'

enum QueueType {
FS = 'fs',
Graphile = 'graphile',
}

const queues: Record<QueueType, (server: PluginsServer) => RetryQueue> = {
fs: () => new FsQueue(),
graphile: (pluginsServer: PluginsServer) => new GraphileQueue(pluginsServer),
}

export class RetryQueueManager implements RetryQueue {
pluginsServer: PluginsServer
retryQueues: RetryQueue[]

constructor(pluginsServer: PluginsServer) {
this.pluginsServer = pluginsServer

this.retryQueues = pluginsServer.RETRY_QUEUES.split(',')
.map((q) => q.trim() as QueueType)
.filter((q) => !!q)
.map(
(queue): RetryQueue => {
if (queues[queue]) {
return queues[queue](pluginsServer)
} else {
throw new Error(`Unknown retry queue "${queue}"`)
}
}
)
}

async enqueue(retry: EnqueuedRetry): Promise<void> {
for (const retryQueue of this.retryQueues) {
try {
await retryQueue.enqueue(retry)
return
} catch (error) {
// if one fails, take the next queue
Sentry.captureException(error, {
extra: {
retry: JSON.stringify(retry),
queue: retryQueue.toString(),
queues: this.retryQueues.map((q) => q.toString()),
},
})
}
}
throw new Error('No RetryQueue available')
}

async quit(): Promise<void> {
await Promise.all(this.retryQueues.map((r) => r.quit()))
}

async startConsumer(onRetry: OnRetryCallback): Promise<void> {
await Promise.all(this.retryQueues.map((r) => r.startConsumer(onRetry)))
}

async stopConsumer(): Promise<void> {
await Promise.all(this.retryQueues.map((r) => r.stopConsumer()))
}

async pauseConsumer(): Promise<void> {
await Promise.all(this.retryQueues.map((r) => r.pauseConsumer()))
}

isConsumerPaused(): boolean {
return !!this.retryQueues.find((r) => r.isConsumerPaused())
}

async resumeConsumer(): Promise<void> {
await Promise.all(this.retryQueues.map((r) => r.resumeConsumer()))
}
}
Loading