Skip to content
Draft
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
59 changes: 59 additions & 0 deletions tests/features/processings/log-flood.api.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
import { test, expect } from '@playwright/test'
import { execSync } from 'node:child_process'
import path from 'node:path'
import os from 'node:os'
import fs from 'node:fs'
import { axiosAuth, clean, waitForRunStatus } from '../../support/axios.ts'
import { publishFixturePlugin } from '../../support/registry.ts'

const buildLogFloodTarball = (): string => {
const fixtureDir = path.resolve(import.meta.dirname, '../../fixtures/processing-log-flood')
const outDir = fs.mkdtempSync(path.join(os.tmpdir(), 'log-flood-pack-'))
execSync('npm pack --pack-destination ' + outDir, { cwd: fixtureDir, stdio: 'pipe' })
const tarball = fs.readdirSync(outDir).find(f => f.endsWith('.tgz'))
if (!tarball) throw new Error('npm pack did not produce a tarball')
return path.join(outDir, tarball)
}

// Worker default (worker/config/default.mjs). The fixture writes 20 000 chars
// per entry, so every entry is truncated to maxLogEntryLength.
const maxLogEntryLength = 10000
const mongoDocumentLimit = 16 * 1024 * 1024

test.describe('run log size guard', () => {
test.beforeEach(clean)
test.afterAll(clean)

test('a task flooding its log is failed once the run document is full', async () => {
test.setTimeout(300_000)
const superadmin = await axiosAuth('test_superadmin@test.com')

const plugin = await publishFixturePlugin({
name: '@data-fair-tests/processing-log-flood',
version: '1.0.0',
tarballPath: buildLogFloodTarball()
})
const processing = (await superadmin.post('/api/v1/processings', {
title: 'Log flood test',
plugin: plugin.pluginId,
owner: { type: 'user', id: 'test_superadmin', name: 'Test Super Admin' },
active: true,
config: {}
})).data
const triggered = (await superadmin.post(`/api/v1/processings/${processing._id}/_trigger`)).data
const finalRun = await waitForRunStatus(triggered._id, ['finished', 'error'], 270_000)
expect(finalRun.status).toBe('error')

const log = finalRun.log as Array<{ type: string, msg: string }>
// the loop of the fixture never completes
expect(log.map(l => l.msg)).not.toContain('log-flood fixture finished, the guard did not trigger')
// every flooding entry was truncated
const flooding = log.filter(l => l.msg.startsWith('xxxx'))
expect(flooding.length).toBeGreaterThan(0)
for (const entry of flooding) expect(entry.msg).toBe('x'.repeat(maxLogEntryLength) + '...')
// the document filled up, the log was truncated to make room for the guard message and finish()
expect(log.some(l => l.type === 'error' && /taille maximale/.test(l.msg))).toBe(true)
expect(Buffer.byteLength(JSON.stringify(log))).toBeLessThan(mongoDocumentLimit / 4)
expect(log[log.length - 1].type).toBe('debug') // finish() could still add its entry
})
})
32 changes: 31 additions & 1 deletion tests/features/worker-utils/runs-operations.unit.spec.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { test, expect } from '@playwright/test'
import { shouldDisableForFailures, buildFinishStatusPatch } from '../../../worker/src/utils/runs-operations.ts'
import { shouldDisableForFailures, buildFinishStatusPatch, isDocumentTooLargeError, truncateLogValue } from '../../../worker/src/utils/runs-operations.ts'

test.describe('shouldDisableForFailures', () => {
const maxFailures = 3
Expand Down Expand Up @@ -59,3 +59,33 @@ test.describe('buildFinishStatusPatch', () => {
expect(buildFinishStatusPatch('running', undefined, finishedAt)).toEqual({ status: 'finished', finishedAt })
})
})

test.describe('isDocumentTooLargeError', () => {
test('matches the mongo error of an update overflowing the document limit', () => {
expect(isDocumentTooLargeError(new Error('Plan executor error during findAndModify :: caused by :: Resulting document after update is larger than 16777216'))).toBe(true)
expect(isDocumentTooLargeError(Object.assign(new Error('BSONObjectTooLarge'), { code: 17419 }))).toBe(true)
})

test('ignores other errors', () => {
expect(isDocumentTooLargeError(new Error('Run not found'))).toBe(false)
expect(isDocumentTooLargeError(undefined)).toBe(false)
})
})

test.describe('truncateLogValue', () => {
test('leaves short values untouched', () => {
expect(truncateLogValue('hello', 100)).toBe('hello')
expect(truncateLogValue({ a: 1 }, 100)).toEqual({ a: 1 })
expect(truncateLogValue(undefined, 100)).toBeUndefined()
expect(truncateLogValue('', 100)).toBe('')
})

test('truncates a long string', () => {
expect(truncateLogValue('a'.repeat(50), 10)).toBe('a'.repeat(10) + '...')
})

test('serializes then truncates a long non-string value', () => {
const value = { items: new Array(100).fill('x') }
expect(truncateLogValue(value, 10)).toBe(JSON.stringify(value).slice(0, 10) + '...')
})
})
14 changes: 14 additions & 0 deletions tests/fixtures/processing-log-flood/index.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
// Test fixture: push far more log data than fits in the run document (mongo's
// 16MB limit). Each message is longer than worker.task.maxLogEntryLength so
// the per-entry truncation is exercised too, and so the document fills up in
// a couple thousand entries instead of hundreds of thousands. The worker is
// expected to fail the run before the loop ends.
export const run = async (context) => {
const { log } = context
await log.step('starting log-flood fixture')
const message = 'x'.repeat(20_000)
for (let i = 0; i < 5000; i++) {
await log.info(message)
}
await log.info('log-flood fixture finished, the guard did not trigger')
}
8 changes: 8 additions & 0 deletions tests/fixtures/processing-log-flood/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"name": "@data-fair-tests/processing-log-flood",
"version": "1.0.0",
"description": "Test plugin that floods the run log to trigger the worker log size guard.",
"main": "index.js",
"type": "module",
"files": ["index.js", "processing-config-schema.json"]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
{
"type": "object",
"additionalProperties": true,
"properties": {}
}
3 changes: 2 additions & 1 deletion worker/config/custom-environment-variables.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,8 @@ export default {
externalSamplerEnabled: {
__name: 'WORKER_TASK_EXTERNAL_SAMPLER_ENABLED',
__format: 'json'
}
},
maxLogEntryLength: 'WORKER_TASK_MAX_LOG_ENTRY_LENGTH'
}
}
}
5 changes: 4 additions & 1 deletion worker/config/default.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,10 @@ export default {
// for each task child. Becomes the authoritative writer for the per-slot
// RSS gauge (the in-process df-mem RSS write is suppressed). Auto-disabled
// at boot on non-Linux platforms.
externalSamplerEnabled: true
externalSamplerEnabled: true,
// Max length (chars) of the msg / extra of a single run log entry,
// longer values are truncated.
maxLogEntryLength: 10000
}
},
upgradeRoot: '/app/'
Expand Down
5 changes: 3 additions & 2 deletions worker/config/type/schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -146,12 +146,13 @@
"gracePeriod": { "type": "number" },
"task": {
"type": "object",
"required": ["maxHeapMB", "memorySampleIntervalMs", "memoryHeadroomWarnPct", "externalSamplerEnabled"],
"required": ["maxHeapMB", "memorySampleIntervalMs", "memoryHeadroomWarnPct", "externalSamplerEnabled", "maxLogEntryLength"],
"properties": {
"maxHeapMB": { "type": "number", "minimum": 64, "description": "Per-task V8 old-space heap cap in MB, passed as --max-old-space-size to the child. Defaults to V8's own heap_size_limit for the host." },
"memorySampleIntervalMs": { "type": "number", "minimum": 0 },
"memoryHeadroomWarnPct": { "type": "number", "minimum": 0, "maximum": 100 },
"externalSamplerEnabled": { "type": "boolean", "description": "Enable parent-side /proc-based RSS/CPU sampler for each task child. Auto-disabled at boot on non-Linux platforms." }
"externalSamplerEnabled": { "type": "boolean", "description": "Enable parent-side /proc-based RSS/CPU sampler for each task child. Auto-disabled at boot on non-Linux platforms." },
"maxLogEntryLength": { "type": "number", "minimum": 100, "description": "Max length in chars of the msg/extra of a single run log entry, longer values are truncated." }
}
}
}
Expand Down
31 changes: 24 additions & 7 deletions worker/src/task/task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@ import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js'
import { ensureArtefact } from '@data-fair/lib-node-registry'
import { decipher } from '@data-fair/processings-shared/cipher.ts'
import { importPluginModule } from '@data-fair/processings-shared/plugin-load.ts'
import { running } from '../utils/runs.ts'
import { running, truncateLog } from '../utils/runs.ts'
import { isDocumentTooLargeError, truncateLogValue } from '../utils/runs-operations.ts'
import config, { registryCacheDir } from '#config'
import mongo from '#mongo'
import { getAxiosInstance, getHttpErrorMessage, prepareAxiosError } from './axios.ts'
Expand Down Expand Up @@ -44,9 +45,29 @@ const wsInstance = (log: LogFunctions, owner: Account): DataFairWsClient => {
* Prepare log functions.
*/
const prepareLog = (processing: Processing, run: Run): LogFunctions => {
const { maxLogEntryLength } = config.worker.task
let overflowed = false
const pushLog = async (log: any) => {
// the task is failing because of its log, drop what it writes on the way out
if (overflowed) return
log.msg = truncateLogValue(log.msg, maxLogEntryLength)
log.extra = truncateLogValue(log.extra, maxLogEntryLength)
log.date = new Date().toISOString()
await mongo.runs.updateOne({ _id: run._id }, { $push: { log } })
try {
await mongo.runs.updateOne({ _id: run._id }, { $push: { log } })
} catch (err) {
if (!isDocumentTooLargeError(err)) throw err
// The whole log lives in the run document and it reached mongo's 16MB
// limit. Make room so the run can still be finished, leave an explicit
// error and fail the task.
overflowed = true
const msg = 'Le journal d\'exécution a atteint la taille maximale autorisée, l\'exécution est interrompue. Le traitement produit trop de messages, réduisez sa verbosité.'
await truncateLog(run._id)
const errorLog: any = { type: 'error', msg, extra: '', date: log.date }
await mongo.runs.updateOne({ _id: run._id }, { $push: { log: errorLog } })
await wsEmitter.emit(`processings/${processing._id}/run-log`, { _id: run._id, log: errorLog })
throw new Error(msg)
}
await wsEmitter.emit(`processings/${processing._id}/run-log`, { _id: run._id, log })
}

Expand Down Expand Up @@ -185,13 +206,9 @@ export const run = async (mailTransport: any) => {
const httpMessage = getHttpErrorMessage(err)

if (httpMessage) {
let errStr = util.inspect(err, { depth: 5 })
if (errStr.length > 10000) {
errStr = errStr.slice(0, 10000) + '...'
}
console.error(httpMessage)
await log.error(httpMessage)
await log.debug('axios error', errStr)
await log.debug('axios error', util.inspect(err, { depth: 5 }))
} else {
console.error(err.message || err)
await log.error(err.message)
Expand Down
17 changes: 17 additions & 0 deletions worker/src/utils/runs-operations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,3 +36,20 @@ export const buildFinishStatusPatch = (
if (errorMessage) return { status: 'error', finishedAt }
return { status: 'finished', finishedAt }
}

/**
* Whether a mongo error means the updated document would exceed the 16MB limit
* ("Resulting document after update is larger than 16777216").
*/
export const isDocumentTooLargeError = (err: any): boolean =>
err?.code === 17419 || /larger than \d+/.test(err?.message ?? '')

/**
* Cap a log msg/extra at `maxLength` chars. A non-string value that is too
* long once serialized is replaced by its truncated serialization.
*/
export const truncateLogValue = (value: any, maxLength: number): any => {
const str = typeof value === 'string' ? value : JSON.stringify(value)
if (!str || str.length <= maxLength) return value
return str.slice(0, maxLength) + '...'
}
35 changes: 29 additions & 6 deletions worker/src/utils/runs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import eventsQueue from '@data-fair/lib-node/events-queue.js'
import config from '#config'
import mongo from '#mongo'
import { internalError } from '@data-fair/lib-node/observer.js'
import { shouldDisableForFailures } from './runs-operations.ts'
import { isDocumentTooLargeError, shouldDisableForFailures } from './runs-operations.ts'

const sendProcessingEvent = (
run: Run,
Expand Down Expand Up @@ -54,6 +54,15 @@ export const running = async (run: Run) => {
.updateOne({ _id: run.processing._id }, { $set: { lastRun }, $unset: { nextRun: '' } })
}

/**
* Keep only the tail of a run log. Used when the run document reached mongo's
* 16MB limit and can not be updated anymore.
*/
export const truncateLog = async (runId: string, keep = 100) => {
const truncate: Record<string, any> = { $push: { log: { $each: [], $slice: -keep } } }
await mongo.runs.updateOne({ _id: runId }, truncate)
}

/**
* Update the database when a run is finished (edit status, log, duration, etc.)
*/
Expand All @@ -75,11 +84,25 @@ export const finish = async (run: Run, errorMessage: string | undefined = undefi
if (metricsMessage) logs.push({ type: 'debug', msg: metricsMessage, date: new Date(now + 1).toISOString() })
query.$push = { log: { $each: logs } }
}
let lastRun = (await mongo.runs.findOneAndUpdate(
{ _id: run._id },
query,
{ returnDocument: 'after', projection: { processing: 0, owner: 0 } }
))
let lastRun
try {
lastRun = await mongo.runs.findOneAndUpdate(
{ _id: run._id },
query,
{ returnDocument: 'after', projection: { processing: 0, owner: 0 } }
)
} catch (err) {
if (!isDocumentTooLargeError(err)) throw err
// A run whose log filled the document would otherwise be impossible to
// finish or kill, and the kill loop would retry it forever.
console.warn('run document too large to be finished, truncating its log', run._id)
await truncateLog(run._id)
lastRun = await mongo.runs.findOneAndUpdate(
{ _id: run._id },
query,
{ returnDocument: 'after', projection: { processing: 0, owner: 0 } }
)
}
if (!lastRun) return internalError('processing-worker', 'Last run not found after finish update')
if (!lastRun.startedAt) {
lastRun = (await mongo.runs.findOneAndUpdate(
Expand Down
Loading