/* This Source Code Form is subject to the terms of the Mozilla Public * License, v. 2.0. If a copy of the MPL was not distributed with this * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ import crypto from 'node:crypto'; import fs from 'node:fs'; import type http from 'node:http'; import path from 'node:path'; import qs from 'node:querystring'; import { parse as urlparse } from 'node:url'; import { callbackify, promisify } from 'node:util'; import HttpAgent from 'agentkeepalive'; import async from 'async'; import createDebug from 'debug'; import decompressResponse from 'decompress-response'; import { compileExpression as filtrex } from 'filtrex'; import FormData from 'form-data'; import { HttpProxyAgent, HttpsProxyAgent } from 'hpagent'; import _ from 'lodash'; import * as tough from 'tough-cookie'; import { engine_util as engineUtil } from '../commons/index.ts'; const debug = createDebug('http'); const debugRequests = createDebug('http:request'); const debugResponse = createDebug('http:response'); const USER_AGENT = 'Artillery (https://artillery.io)'; const ensurePropertyIsAList = engineUtil.ensurePropertyIsAList; const template = engineUtil.template; const { HttpsAgent } = HttpAgent; // NOTE on typing: request specs, request params and VU variables are // dynamic by design (user scripts + processor hooks may attach // arbitrary properties), so they are typed as Record // until the canonical script/flow input types exist (modernization // plan, phase 2). Structural contracts that are stable - engine // surface, VU context internals, response handling - are typed. type ProcessorFunction = (...args: any[]) => any; // Structural view of got's request function: its own overloaded types // reject dynamically-assembled option objects. interface CancelableRequestLike extends Promise { on(event: string, listener: (...args: any[]) => void): CancelableRequestLike; } type GotRequester = (options: Record) => CancelableRequestLike; type StepCallback = ( err: Error | null | undefined, context?: VUContext ) => void; type StepFunction = (context: VUContext, callback: StepCallback) => void; interface EventEmitterLike { emit(event: string, ...args: any[]): unknown; } // What callers hand to a compiled scenario: VU variables plus // whatever the runner attached. Engine internals (_jar, agents, // counters) are attached by setInitialContext(). export interface InitialVUContext { vars: Record; [key: string]: any; } export interface VUContext { vars: Record; _successCount: number; _enableCookieJar: boolean; _jar: tough.CookieJar; _defaultCookie?: Record; _defaultStrictCapture?: boolean; _httpAgent: http.Agent; _httpsAgent: http.Agent; // Engines, plugins and user hooks attach ad-hoc state: [key: string]: any; } // Raw response as surfaced by got's 'response' event: an // http.IncomingMessage with got's timing info attached. interface WireResponse { statusCode: number; headers: Record; timings: { phases: Record }; body?: unknown; on(event: string, listener: (...args: any[]) => void): unknown; } interface AgentOptions { keepAlive?: boolean; keepAliveMsec?: number; maxSockets?: number; maxFreeSockets?: number; timeout?: number; proxy?: string; } interface HttpDefaultsConfig { headers?: Record; cookie?: Record; strictCapture?: boolean | string; think?: Record; [key: string]: any; } interface NormalizedHttpEngineConfig { target?: string; timeout?: number; tls?: Record; processor?: Record; defaults: HttpDefaultsConfig; http: { defaults: HttpDefaultsConfig; cookieJarOptions: Record; pool?: number | string; timeout?: number; extendedMetrics?: boolean; distributedTracing?: | boolean | { enabled?: boolean; sampled?: boolean; traceIdPrefix?: string }; [key: string]: any; }; [key: string]: any; } const GOT_OPTION_NAMES = [ 'url', 'searchParams', 'method', 'headers', 'body', 'json', 'form', 'allowGetBody', 'timeout', 'retry', 'encoding', 'cookieJar', 'followRedirect', 'maxRedirects', 'decompress', 'http2', 'agent', 'username', 'password', 'https', 'throwHttpErrors' ]; // NOTE: pre-existing quirk preserved: this object is mutated by // Object.assign() in the constructor and in setInitialContext(), so // options accumulate on the shared constant across instances. const DEFAULT_AGENT_OPTIONS: AgentOptions = { keepAlive: true, keepAliveMsec: 1000 }; // agentkeepalive defaults `timeout` (active socket inactivity) to a hardcoded // floor of 8000ms (Math.max(freeSocketTimeout * 2, 8000)). When a request's // origin doesn't send response bytes within 8s, the socket is destroyed with // `error.code = 'ERR_SOCKET_TIMEOUT'` regardless of the user's configured // `http.timeout`. Mirror `http.timeout` (or top-level `timeout`) into the agent // so the YAML config becomes the actual binding constraint. function deriveAgentTimeoutMs( scriptConfig: Record ): number | undefined { const timeoutSec = scriptConfig?.timeout || scriptConfig?.http?.timeout; if (typeof timeoutSec !== 'number') return undefined; // Never go below the existing 8s floor — preserves current behaviour for // users who explicitly configure a smaller http.timeout. return Math.max(8000, timeoutSec * 1000); } function createAgents( proxies: { http?: string; https?: string }, opts: AgentOptions ): { httpAgent: http.Agent; httpsAgent: http.Agent } { const agentOpts: AgentOptions = Object.assign( {}, DEFAULT_AGENT_OPTIONS, opts ); // HTTP proxy endpoint will be used for all requests, unless a separate // HTTPS proxy URL is also set, which will be used for HTTPS requests: if (proxies.http) { agentOpts.proxy = proxies.http; const httpAgent = new HttpProxyAgent( agentOpts as ConstructorParameters[0] ); if (proxies.https) { agentOpts.proxy = proxies.https; } const httpsAgent = new HttpsProxyAgent( agentOpts as ConstructorParameters[0] ); return { httpAgent, httpsAgent }; } // If only HTTPS proxy is provided, it will be used for HTTPS requests, // but not for HTTP requests: if (proxies.https) { return { httpAgent: new HttpAgent(agentOpts), httpsAgent: new HttpsProxyAgent( Object.assign( { proxy: proxies.https }, agentOpts ) as ConstructorParameters[0] ) }; } // By default nothing is proxied: return { httpAgent: new HttpAgent(agentOpts), httpsAgent: new HttpsAgent(agentOpts) }; } export default class HttpEngine { declare config: NormalizedHttpEngineConfig; declare maxSockets: number; declare _httpAgent: http.Agent; declare _httpsAgent: http.Agent; declare extendedHTTPMetrics?: boolean; declare request: GotRequester; // Set by the runner after loading: declare __name?: string; constructor(script: { config: Record }) { this.config = script.config as NormalizedHttpEngineConfig; if (typeof this.config.defaults === 'undefined') { this.config.defaults = {}; } if (typeof this.config.http === 'undefined') { this.config.http = {} as NormalizedHttpEngineConfig['http']; } if (typeof this.config.http.defaults === 'undefined') { this.config.http.defaults = {}; } if (typeof this.config.http.cookieJarOptions === 'undefined') { this.config.http.cookieJarOptions = {}; } // If config.http.pool is set, create & reuse agents for all requests (with // max sockets set). That's what we're done here. // If config.http.pool is not set, we create new agents for each virtual user. // That's done when the VU is initialized. this.maxSockets = Infinity; if (script.config.http?.pool) { this.maxSockets = Number(script.config.http.pool); } const agentOpts: AgentOptions = Object.assign(DEFAULT_AGENT_OPTIONS, { maxSockets: this.maxSockets, maxFreeSockets: this.maxSockets }); const agentTimeoutMs = deriveAgentTimeoutMs(script.config); if (agentTimeoutMs !== undefined) { agentOpts.timeout = agentTimeoutMs; } const agents = createAgents( { http: process.env.HTTP_PROXY, https: process.env.HTTPS_PROXY }, agentOpts ); this._httpAgent = agents.httpAgent; this._httpsAgent = agents.httpsAgent; if ( (script.config.http && script.config.http.extendedMetrics === true) || global.artillery?.runtimeOptions.extendedHTTPMetrics ) { this.extendedHTTPMetrics = true; } } async init(): Promise { this.request = (await import('got')).default as unknown as GotRequester; } _isDistributedTracingEnabled(config: NormalizedHttpEngineConfig): boolean { const dtConfig = config.http?.distributedTracing; if (!dtConfig) { return false; } // Handle both boolean and object forms if (typeof dtConfig === 'boolean') { return dtConfig; } if (typeof dtConfig === 'object' && dtConfig.enabled !== undefined) { return dtConfig.enabled; } // Default to true if distributedTracing is set but enabled is not specified return true; } _generateTraceparent(config: NormalizedHttpEngineConfig): string { // W3C Trace Context format: version-trace-id-parent-id-trace-flags const version = '00'; // Get configuration const dtConfig = config.http?.distributedTracing; let sampled = true; // Default to sampled let traceIdPrefix = 'a9'; // Default prefix if (typeof dtConfig === 'object') { if (dtConfig.sampled !== undefined) { sampled = dtConfig.sampled; } if (dtConfig.traceIdPrefix !== undefined) { traceIdPrefix = dtConfig.traceIdPrefix; } } // Validate and normalize prefix (must be valid hex, max 8 chars) traceIdPrefix = traceIdPrefix .toLowerCase() .replace(/[^0-9a-f]/g, '') .slice(0, 8); if (traceIdPrefix.length === 0) { traceIdPrefix = 'a9'; // Fallback to default if invalid } // Generate trace-id with prefix (32 hex chars total) const remainingBytes = Math.ceil((32 - traceIdPrefix.length) / 2); const randomPart = crypto.randomBytes(remainingBytes).toString('hex'); const traceId = (traceIdPrefix + randomPart).slice(0, 32); // Generate 8-byte parent-id (16 hex chars) const parentId = crypto.randomBytes(8).toString('hex'); const traceFlags = sampled ? '01' : '00'; return `${version}-${traceId}-${parentId}-${traceFlags}`; } createScenario(scenarioSpec: Record, ee: EventEmitterLike) { ensurePropertyIsAList(scenarioSpec, 'beforeRequest'); ensurePropertyIsAList(scenarioSpec, 'afterResponse'); ensurePropertyIsAList(scenarioSpec, 'beforeScenario'); ensurePropertyIsAList(scenarioSpec, 'afterScenario'); ensurePropertyIsAList(scenarioSpec, 'onError'); // Add scenario-level hooks if needed: // For now, just turn them into function steps and insert them // directly into the flow array. // TODO: Scenario-level hooks will probably want access to the // entire scenario spec rather than just the userContext. const beforeScenarioFns = _.map( scenarioSpec.beforeScenario, (hookFunctionName) => ({ function: hookFunctionName }) ); const afterScenarioFns = _.map( scenarioSpec.afterScenario, (hookFunctionName) => ({ function: hookFunctionName }) ); const newFlow = beforeScenarioFns.concat( scenarioSpec.flow.concat(afterScenarioFns) ); scenarioSpec.flow = newFlow; const tasks = _.map(scenarioSpec.flow, (rs) => this.step(rs, ee, { beforeRequest: scenarioSpec.beforeRequest, afterResponse: scenarioSpec.afterResponse, onError: scenarioSpec.onError }) ); return this.compile(tasks, scenarioSpec, ee); } step( requestSpec: Record, ee: EventEmitterLike, opts?: { beforeRequest?: string[]; afterResponse?: string[]; onError?: string[]; } ): StepFunction { opts = opts || {}; const self = this; const config = this.config; if (requestSpec.loop) { const steps = _.map(requestSpec.loop, (rs) => self.step(rs, ee, opts)); return engineUtil.createLoopWithCount(requestSpec.count || -1, steps, { loopValue: requestSpec.loopValue || '$loopCount', loopElement: requestSpec.loopElement || '$loopElement', overValues: requestSpec.over, whileTrue: self.config.processor ? self.config.processor[requestSpec.whileTrue] : undefined }); } if (requestSpec.parallel) { const steps = _.map(requestSpec.parallel, (rs) => self.step(rs, ee, opts) ); return engineUtil.createParallel(steps, { limitValue: requestSpec.limit }); } if (typeof requestSpec.think !== 'undefined') { return engineUtil.createThink( requestSpec, self.config.http.defaults.think || self.config.defaults.think ); } if (typeof requestSpec.log !== 'undefined') { return (context, callback) => { console.log(template(requestSpec.log, context)); return process.nextTick(() => { callback(null, context); }); }; } if (requestSpec.function) { return (context, callback) => { const processFunc = self.config.processor?.[requestSpec.function]; if (processFunc) { let f: ProcessorFunction; if (processFunc.constructor.name === 'Function') { f = processFunc; } else { f = callbackify( processFunc as (...args: any[]) => Promise ); } return f(context, ee, (hookErr: Error | null | undefined) => callback(hookErr, context) ); } else { debug(`Function "${requestSpec.function}" not defined`); debug('processor: %o', self.config.processor); ee.emit('error', `Undefined function "${requestSpec.function}"`); return process.nextTick(() => { callback(null, context); }); } }; } const f: StepFunction = (context, callback) => { const method = _.keys(requestSpec)[0].toUpperCase(); const params: Record = requestSpec[method.toLowerCase()]; const onErrorHandlers = opts.onError; // only scenario-lever onError handlers are supported // A special case for when "url" attribute is missing. We need to check for // it manually as request.js won't emit an 'error' event when the argument // is missing. // This will be obsoleted by better script validation. if (!params.url && !params.uri) { const err = new Error('an URL must be specified'); return callback(err, context); } const tls = config.tls || {}; const timeout = config.timeout || _.get(config, 'http.timeout') || 10; if (!engineUtil.isProbableEnough(params)) { return process.nextTick(() => { callback(null, context); }); } if (!_.isUndefined(params.ifTrue)) { let result; try { const cond = _.has(config.processor, params.ifTrue) ? (config.processor as Record)[ params.ifTrue ] : filtrex(params.ifTrue); result = cond(context.vars); } catch (err) { debug('ifTrue error:', err); result = 1; // if the expression is incorrect, just proceed } if (!result) { return process.nextTick(() => { callback(null, context); }); } } // Run beforeRequest processors (scenario-level ones too) const requestParams: Record = _.extend(_.clone(params), { url: maybePrependBase(params.url || params.uri, config), // *NOT* templating here method: method, timeout: timeout, uuid: crypto.randomUUID() }); if (context._enableCookieJar) { requestParams.cookieJar = context._jar; } if (tls) { requestParams.https = requestParams.https || {}; requestParams.https = _.extend(requestParams.https, tls); } const functionNames = _.concat( opts.beforeRequest || [], params.beforeRequest || [] ); async.eachSeries( functionNames, function iteratee(functionName: string, next) { const fn = template(functionName, context); let processFunc = config.processor?.[fn]; if (!processFunc) { processFunc = (_r: unknown, _c: unknown, _e: unknown, cb: any) => cb(null); console.log(`WARNING: custom function ${fn} could not be found`); // TODO: a 'warning' event } if (processFunc.constructor.name === 'Function') { processFunc(requestParams, context, ee, (err: unknown) => { if (err) { return next(err as Error); } return next(null); }); } else { processFunc(requestParams, context, ee).then(next).catch(next); } }, function done(err) { if (err) { debug(err); return callback(err as Error, context); } // Order of precedence: json set in a function, json set in the script, body set in a function, body set in the script. if (requestParams.json) { requestParams.json = template(requestParams.json, context); delete requestParams.body; } else if (requestParams.body) { requestParams.body = template(requestParams.body, context); // TODO: Warn if body is not a string or a buffer } // add loop, name & uri elements to be interpolated if (context.vars.$loopElement) { context.vars.$loopElement = template( context.vars.$loopElement, context ); } if (requestParams.name) { requestParams.name = template(requestParams.name, context); } if (requestParams.uri) { requestParams.uri = template(requestParams.uri, context); } if (requestParams.url) { requestParams.url = template(requestParams.url, context); } // Follow all redirects by default unless specified otherwise if (typeof requestParams.followRedirect === 'undefined') { requestParams.followRedirect = true; } // TODO: Use traverse on the entire flow instead // Request.js -> Got.js translation if (params.qs) { requestParams.searchParams = qs.stringify( template(params.qs, context) ); } if (typeof params.gzip === 'boolean') { requestParams.decompress = params.gzip; } else { requestParams.decompress = true; } if (params.form) { requestParams.form = _.reduce( requestParams.form, (acc: Record, v, k) => { acc[k] = template(v, context); return acc; }, {} ); } if (params.formData) { let fileUpload: string | undefined; const f = new FormData(); requestParams.body = _.reduce( requestParams.formData, (acc: FormData, v, k) => { let V = template(v, context); let options; if (V && _.isPlainObject(V)) { if (V.contentType) { options = { contentType: V.contentType }; } if (V.fromFile) { const absPath = path.resolve( path.dirname(context.vars.$scenarioFile), V.fromFile ); fileUpload = absPath; V = fs.createReadStream(absPath); } else if (V.value) { V = V.value; } } acc.append(k, V, options); return acc; }, f ); if (params.setContentLengthHeader && fileUpload) { try { requestParams.headers = requestParams.headers || {}; requestParams.headers['content-length'] = fs.statSync(fileUpload).size; } catch (err) { debug(`stat() on ${fileUpload} failed with ${err}`); } } } // Assign default headers then overwrite as needed const defaultHeaders = lowcaseKeys( config.http.defaults.headers || config.defaults.headers || { 'user-agent': USER_AGENT } ); const combinedHeaders = _.extend( defaultHeaders, lowcaseKeys(params.headers), lowcaseKeys(requestParams.headers) ); const templatedHeaders = _.mapValues(combinedHeaders, (v, _k, _obj) => template(v, context) ); requestParams.headers = templatedHeaders; // We compute the url here so that the cookies are set properly afterwards const url = maybePrependBase( template(requestParams.uri || requestParams.url, context), config ); if (requestParams.uri) { // If a hook function sets requestParams.uri to something, request.js // will pick that over .url, so we need to delete it. delete requestParams.uri; } requestParams.url = url; if ( typeof requestParams.cookie === 'object' || typeof context._defaultCookie === 'object' ) { requestParams.cookieJar = context._jar; const cookie: Record = Object.assign( {}, context._defaultCookie, requestParams.cookie ); Object.keys(cookie).forEach((k) => { context._jar.setCookieSync( `${k}=${template(cookie[k], context)}`, requestParams.url ); }); } if (typeof requestParams.auth === 'object') { requestParams.username = template(requestParams.auth.user, context); requestParams.password = template(requestParams.auth.pass, context); delete requestParams.auth; } // TODO: Bypass proxy if "proxy: false" is set requestParams.agent = { http: context._httpAgent, https: context._httpsAgent }; requestParams.throwHttpErrors = false; if (!requestParams.url || !requestParams.url.startsWith('http')) { const err = new Error(`Invalid URL - ${requestParams.url}`); return callback(err, context); } function responseProcessor( isLast: boolean, res: WireResponse, body: string, done: StepCallback ) { if (process.env.DEBUG) { let requestInfo: Record = { url: requestParams.url, method: requestParams.method, headers: requestParams.headers }; // Internal handle of the wrapped cookie jar - debug only: const jarInternal = ( context._jar as unknown as { _jar?: { getCookieStringSync?: (url: string) => string }; } )._jar; if ( jarInternal && typeof jarInternal.getCookieStringSync === 'function' ) { requestInfo = Object.assign(requestInfo, { cookie: jarInternal.getCookieStringSync(requestParams.url) }); } if ( requestParams.json && typeof requestParams.json !== 'boolean' ) { requestInfo.json = requestParams.json; } // If "json" is set to an object, it will be serialised and sent as body and the value of the "body" attribute will be ignored. if ( requestParams.body && typeof requestParams.json !== 'object' ) { if (process.env.DEBUG.indexOf('http:full_body') > -1) { // Show the entire body requestInfo.body = requestParams.body; } else { // Only show the beginning of long bodies if (typeof requestParams.body === 'string') { requestInfo.body = requestParams.body.substring(0, 512); if (requestParams.body.length > 512) { requestInfo.body += ' ...'; } } else if (typeof requestParams.body === 'object') { requestInfo.body = `< ${requestParams.body.constructor.name} >`; } else { requestInfo.body = String(requestInfo.body); } } } if (requestParams.qs) { requestInfo.qs = qs.encode( Object.assign( qs.parse(urlparse(requestParams.url).query ?? ''), template(requestParams.qs, context) ) ); } debug('request: %s', JSON.stringify(requestInfo, null, 2)); } debugResponse(JSON.stringify(res.headers, null, 2)); debugResponse(JSON.stringify(body, null, 2)); // capture/match/response hooks run only for last request in a task if (!isLast) { return done(null, context); } const resForCapture = { headers: res.headers, body: body }; engineUtil.captureOrMatch( params, resForCapture, context, function captured(err: Error | null, result: any) { if (err) { // Run onError hooks and end the scenario: runOnErrorHooks( onErrorHandlers, config.processor, err, requestParams, context, ee, (_asyncErr: unknown) => done(err, context) ); } let haveFailedMatches = false; let haveFailedCaptures = false; if (result !== null) { ee.emit('trace:http:capture', result, requestParams.uuid); if ( Object.keys(result.matches).length > 0 || Object.keys(result.captures).length > 0 ) { debug('captures and matches:'); debug(result.matches); debug(result.captures); } // match and capture are strict by default: haveFailedMatches = _.some( result.matches, (v: any, _k) => !v.success && v.strict !== false ); haveFailedCaptures = _.some( result.captures, (v: any, _k) => v.failed ); if (haveFailedMatches || haveFailedCaptures) { // TODO: Emit the details of each failed capture/match } else { _.each(result.matches, (v: any, _k) => { ee.emit('match', v.success, { expected: v.expected, got: v.got, expression: v.expression, strict: v.strict }); }); _.each(result.captures, (v: any, k) => { _.set(context.vars, k, v.value); }); } } // Now run afterResponse processors const functionNames = _.concat( opts?.afterResponse || [], params.afterResponse || [] ); async.eachSeries( functionNames, function iteratee(functionName: string, next) { const fn = template(functionName, context); let processFunc = config.processor?.[fn]; if (!processFunc) { // TODO: DRY - #223 processFunc = ( _r: unknown, _res: unknown, _c: unknown, _e: unknown, cb: any ) => cb(null); console.log( `WARNING: custom function ${fn} could not be found` ); // TODO: a 'warning' event } // Got does not have res.body which Request.js used to have, so we attach it here: res.body = body; if (processFunc.constructor.name === 'Function') { processFunc( requestParams, res, context, ee, (err: unknown) => { if (err) { return next(err as Error); } return next(null); } ); } else { processFunc(requestParams, res, context, ee) .then(next) .catch(next); } }, (err) => { if (err) { debug(err); return done(err as Error, context); } if (haveFailedMatches || haveFailedCaptures) { // FIXME: This means only one error in the report even if multiple captures failed for the same request. return done( new Error('Failed capture or match'), context ); } return done(null, context); } ); } ); } let needToProcessResponse = false; if ( typeof requestParams.capture === 'object' || typeof requestParams.match === 'object' || requestParams.afterResponse || (typeof opts?.afterResponse === 'object' && opts.afterResponse.length > 0) || process.env.DEBUG ) { needToProcessResponse = true; } if (!requestParams.url) { const err = new Error('an URL must be specified'); // Run onError hooks and end the scenario runOnErrorHooks( onErrorHandlers, config.processor, err, requestParams, context, ee, (_asyncErr: unknown) => callback(err, context) ); } requestParams.retry = { limit: 0 }; // disable retries - ignored when using streams // Convert scalar seconds to Got v14 timeout object right before request const gotOptions: Record = _.pick( requestParams, GOT_OPTION_NAMES ); gotOptions.timeout = { response: requestParams.timeout * 1000 }; // Add W3C Trace Context headers if distributed tracing is enabled if (self._isDistributedTracingEnabled(config)) { const traceparent = self._generateTraceparent(config); gotOptions.headers = gotOptions.headers || {}; gotOptions.headers.traceparent = traceparent; } let totalDownloaded: number | undefined = 0; self .request(gotOptions) .on('request', (req: http.ClientRequest) => { ee.emit('trace:http:request', requestParams, requestParams.uuid); debugRequests('request start: %s', req.path); ee.emit('counter', 'http.requests', 1); ee.emit('rate', 'http.request_rate'); req.on('response', (res) => { res.on('end', () => { ee.emit('counter', 'http.downloaded_bytes', totalDownloaded); }); ee.emit('trace:http:response', res, requestParams.uuid); self._handleResponse( requestParams, res as unknown as WireResponse, ee, context, needToProcessResponse ? responseProcessor : null, callback ); }); }) .on('downloadProgress', (progress: { total?: number }) => { totalDownloaded = progress.total; }) .on('error', (err: NodeJS.ErrnoException) => { ee.emit('trace:http:error', err, requestParams.uuid); if (err.name === 'HTTPError') { return; } // this is an ENOTFOUND, ECONNRESET etc debug(err); // Run onError hooks and end the scenario: runOnErrorHooks( onErrorHandlers, config.processor, err, requestParams, context, ee, (_asyncErr: unknown) => callback(err, context) ); }) .catch((gotErr: Error) => { // TODO: Handle the error properly with run hooks debug(gotErr); runOnErrorHooks( onErrorHandlers, config.processor, gotErr, requestParams, context, ee, (_asyncErr: unknown) => callback(gotErr, context) ); }); } ); // eachSeries }; return f; } _handleResponse( requestParams: Record, res: WireResponse, ee: EventEmitterLike, context: VUContext, responseProcessor: | (( isLast: boolean, res: WireResponse, body: string, done: StepCallback ) => void) | null, callback: StepCallback ): void { const url = requestParams.url; if (requestParams.decompress) { res = decompressResponse( res as unknown as Parameters[0] ) as unknown as WireResponse; } const code = res.statusCode; if (!context._enableCookieJar) { const rawCookies = res.headers['set-cookie']; if (rawCookies) { context._enableCookieJar = true; (rawCookies as string[]).forEach((cookieString) => { try { context._jar.setCookieSync(cookieString, url); } catch (err) { debug( `Could not parse cookieString "${cookieString}" from response header, skipping it` ); debug(err); ee.emit('error', 'cookie_parse_error_invalid_cookie'); } }); } } ee.emit('counter', `http.codes.${code}`, 1); ee.emit('counter', 'http.responses', 1); // ee.emit('rate', 'http.response_rate'); ee.emit('histogram', 'http.response_time', res.timings.phases.firstByte); const statusCode = res.statusCode; if (statusCode >= 200 && statusCode < 300) { ee.emit( 'histogram', 'http.response_time.2xx', res.timings.phases.firstByte ); } else if (statusCode >= 300 && statusCode < 400) { ee.emit( 'histogram', 'http.response_time.3xx', res.timings.phases.firstByte ); } else if (statusCode >= 400 && statusCode < 500) { ee.emit( 'histogram', 'http.response_time.4xx', res.timings.phases.firstByte ); } else if (statusCode >= 500 && statusCode < 600) { ee.emit( 'histogram', 'http.response_time.5xx', res.timings.phases.firstByte ); } if (this.extendedHTTPMetrics) { ee.emit('histogram', 'http.dns', res.timings.phases.dns); ee.emit('histogram', 'http.tcp', res.timings.phases.tcp); ee.emit('histogram', 'http.tls', res.timings.phases.tls); } let body = ''; if (responseProcessor) { res.on('data', (d: Buffer | string) => { body += d; }); } else { res.on('data', () => {}); } res.on('end', () => { if (this.extendedHTTPMetrics) { ee.emit('histogram', 'http.total', res.timings.phases.total); } context._successCount++; // config.defaults won't be taken into account for this const isLastRequest = lastRequest(res, requestParams); if (responseProcessor) { responseProcessor(isLastRequest, res, body, (processResponseErr) => { // capture/match returned an error object, or a hook function returned // with an error if (processResponseErr) { return callback(processResponseErr, context); } if (isLastRequest) { return callback(null, context); } }); } else { if (isLastRequest) { return callback(null, context); } } }); } setInitialContext(initialContext: InitialVUContext): VUContext { initialContext._successCount = 0; initialContext._defaultStrictCapture = true; if ( this.config.http?.defaults?.strictCapture === false || this.config.defaults?.strictCapture === false ) { initialContext._defaultStrictCapture = false; } initialContext._jar = new tough.CookieJar( null, this.config.http.cookieJarOptions ); initialContext._enableCookieJar = false; // If a default cookie is set, we will use the jar straightaway: if ( typeof this.config.http.defaults.cookie === 'object' || typeof this.config.defaults.cookie === 'object' ) { initialContext._defaultCookie = this.config.http.defaults.cookie || this.config.defaults.cookie; initialContext._enableCookieJar = true; } if (this.config.http && typeof this.config.http.pool !== 'undefined') { // Reuse common agents (created in the engine instance constructor) initialContext._httpAgent = this._httpAgent; initialContext._httpsAgent = this._httpsAgent; } else { // Create agents just for this VU const agentOpts: AgentOptions = Object.assign(DEFAULT_AGENT_OPTIONS, { maxSockets: 1, maxFreeSockets: 1 }); const agentTimeoutMs = deriveAgentTimeoutMs(this.config); if (agentTimeoutMs !== undefined) { agentOpts.timeout = agentTimeoutMs; } const agents = createAgents( { http: process.env.HTTP_PROXY, https: process.env.HTTPS_PROXY }, agentOpts ); initialContext._httpAgent = agents.httpAgent; initialContext._httpsAgent = agents.httpsAgent; } return initialContext as VUContext; } compile( tasks: StepFunction[], _scenarioSpec: Record, ee: EventEmitterLike ) { const self = this; return async function scenario( initialContext: InitialVUContext, callback?: StepCallback ): Promise { let context = self.setInitialContext(initialContext); ee.emit('started'); for (const task of tasks) { try { context = (await promisify(task)(context)) as VUContext; } catch (taskErr) { const err = taskErr as NodeJS.ErrnoException; ee.emit('error', err.code || err.message); if (callback) { return callback(err, context) as undefined; } throw taskErr; } } if (callback) { return callback(null, context) as undefined; } return context; }; } } function lastRequest( res: WireResponse, requestParams: Record ): boolean { // We're done when: // - 3xx response and not following redirects // - not a 3xx response return ( (res.statusCode >= 300 && res.statusCode < 400 && !requestParams.followRedirect) || res.statusCode < 300 || res.statusCode >= 400 ); } function maybePrependBase(uri: string, config: { target?: string }): string { if (_.startsWith(uri, '/')) { return config.target + uri; } else { return uri; } } /* * Given a dictionary, return a dictionary with all keys lowercased. */ function lowcaseKeys(h: unknown): Record { return _.transform( (h || {}) as Record, (result: Record, v, k) => { result[String(k).toLowerCase()] = v; } ); } function runOnErrorHooks( functionNames: string[] | undefined, functions: Record | undefined, err: Error, requestParams: Record, context: VUContext, ee: EventEmitterLike, callback: (err?: Error | null) => void ): void { async.eachSeries( functionNames, function iteratee(functionName: string, next) { const processFunc = functions?.[functionName] as ProcessorFunction; processFunc(err, requestParams, context, ee, (asyncErr: unknown) => { if (asyncErr) { return next(asyncErr as Error); } return next(null); }); }, function done(asyncErr) { return callback(asyncErr as Error | null | undefined); } ); }