798 lines
24 KiB
JavaScript
798 lines
24 KiB
JavaScript
import { reactive, readonly } from "vue";
|
|
import { REQUEST_QUEUE_CONFIG } from "@/config.js";
|
|
import { recordReleaseTimelineEvent } from "@/services/releaseTimeline.js";
|
|
|
|
const cloneObject = (value) => ({ ...(value || {}) });
|
|
|
|
const queueConfig = {
|
|
concurrencyByMethod: cloneObject(REQUEST_QUEUE_CONFIG.concurrency),
|
|
spacingMs: REQUEST_QUEUE_CONFIG.spacingMs,
|
|
retryByStatusCode: cloneObject(REQUEST_QUEUE_CONFIG.retryByStatusCode),
|
|
retryDelayBaseMs: REQUEST_QUEUE_CONFIG.retryDelay.baseMs,
|
|
retryDelayMaxMs: REQUEST_QUEUE_CONFIG.retryDelay.maxMs,
|
|
retryDelayJitterMs: REQUEST_QUEUE_CONFIG.retryDelay.jitterMs,
|
|
recentRequestsLimit: REQUEST_QUEUE_CONFIG.inspector.recentRequestsLimit,
|
|
errorHistoryLimit: REQUEST_QUEUE_CONFIG.inspector.errorHistoryLimit,
|
|
missingPermissionsLimit: REQUEST_QUEUE_CONFIG.inspector.missingPermissionsLimit,
|
|
componentPermissionReportWindowMs: REQUEST_QUEUE_CONFIG.inspector.componentPermissionReportWindowMs,
|
|
payloadMaxChars: REQUEST_QUEUE_CONFIG.inspector.payloadMaxChars,
|
|
};
|
|
|
|
const requestQueueStateMutable = reactive({
|
|
batchId: 0,
|
|
pending: 0,
|
|
active: 0,
|
|
batchTotal: 0,
|
|
batchCompleted: 0,
|
|
batchFailed: 0,
|
|
activeRequests: [],
|
|
recentRequests: [],
|
|
errorRequests: [],
|
|
missingPermissions: [],
|
|
networkTotals: {
|
|
outgoingRequests: 0,
|
|
ingoingResponses: 0,
|
|
outgoingBytes: 0,
|
|
ingoingBytes: 0,
|
|
},
|
|
});
|
|
|
|
const requestQueue = [];
|
|
const activeWorkersByKey = {};
|
|
let activeWorkers = 0;
|
|
let requestIdCounter = 0;
|
|
let drainTimer = null;
|
|
let lastRequestStartedAt = 0;
|
|
const recentComponentPermissionReports = {};
|
|
|
|
const isQueueIdle = () => requestQueue.length === 0 && activeWorkers === 0;
|
|
|
|
const startBatchIfNeeded = () => {
|
|
if (!isQueueIdle()) {
|
|
return;
|
|
}
|
|
|
|
requestQueueStateMutable.batchId += 1;
|
|
requestQueueStateMutable.batchTotal = 0;
|
|
requestQueueStateMutable.batchCompleted = 0;
|
|
requestQueueStateMutable.batchFailed = 0;
|
|
requestQueueStateMutable.recentRequests = [];
|
|
};
|
|
|
|
const syncQueueCounters = () => {
|
|
requestQueueStateMutable.pending = requestQueue.length;
|
|
requestQueueStateMutable.active = activeWorkers;
|
|
};
|
|
|
|
const normalizeMethod = (value) => {
|
|
if (typeof value !== "string" || value.trim().length === 0) {
|
|
return "GET";
|
|
}
|
|
return value.trim().toUpperCase();
|
|
};
|
|
|
|
const normalizeQueueGroup = (value) => {
|
|
if (typeof value !== "string") {
|
|
return "";
|
|
}
|
|
|
|
return value.trim().toUpperCase();
|
|
};
|
|
|
|
const normalizeUrl = (value) => {
|
|
if (typeof value !== "string" || value.trim().length === 0) {
|
|
return "(unknown endpoint)";
|
|
}
|
|
return value.trim();
|
|
};
|
|
|
|
const redactHeaders = (headers) => {
|
|
if (!headers || typeof headers !== "object") {
|
|
return headers ?? null;
|
|
}
|
|
|
|
const normalized = {};
|
|
Object.entries(headers).forEach(([key, value]) => {
|
|
if (/authorization|x-api-key|token/i.test(String(key))) {
|
|
normalized[key] = "***";
|
|
return;
|
|
}
|
|
normalized[key] = value;
|
|
});
|
|
return normalized;
|
|
};
|
|
|
|
const toSafeText = (value) => {
|
|
if (value === null || value === undefined) {
|
|
return "";
|
|
}
|
|
|
|
let text = "";
|
|
if (typeof value === "string") {
|
|
text = value;
|
|
} else {
|
|
try {
|
|
text = JSON.stringify(value, null, 2);
|
|
} catch (error) {
|
|
text = String(value);
|
|
}
|
|
}
|
|
|
|
const maxChars = Math.max(200, Number(queueConfig.payloadMaxChars) || 4000);
|
|
if (text.length <= maxChars) {
|
|
return text;
|
|
}
|
|
|
|
return `${text.slice(0, maxChars)}\n... [truncated]`;
|
|
};
|
|
|
|
const measureTextBytes = (value) => {
|
|
const text = toSafeText(value);
|
|
if (!text) {
|
|
return 0;
|
|
}
|
|
|
|
if (typeof TextEncoder === "function") {
|
|
return new TextEncoder().encode(text).length;
|
|
}
|
|
|
|
return text.length;
|
|
};
|
|
|
|
const estimateRequestBytes = (job, method, url) => measureTextBytes({
|
|
method,
|
|
url,
|
|
params: job?.requestData?.params ?? null,
|
|
data: job?.requestData?.data ?? null,
|
|
headers: redactHeaders(job?.requestData?.headers ?? null),
|
|
});
|
|
|
|
const estimateResponseBytes = (response, fallbackMessage = null) => measureTextBytes({
|
|
status: response?.status ?? null,
|
|
statusText: response?.statusText ?? null,
|
|
data: response?.data ?? fallbackMessage,
|
|
headers: redactHeaders(response?.headers ?? null),
|
|
});
|
|
|
|
const addNetworkTotals = ({
|
|
outgoingRequests = 0,
|
|
ingoingResponses = 0,
|
|
outgoingBytes = 0,
|
|
ingoingBytes = 0,
|
|
} = {}) => {
|
|
requestQueueStateMutable.networkTotals = {
|
|
outgoingRequests: Number(requestQueueStateMutable.networkTotals?.outgoingRequests || 0)
|
|
+ Math.max(0, Number(outgoingRequests) || 0),
|
|
ingoingResponses: Number(requestQueueStateMutable.networkTotals?.ingoingResponses || 0)
|
|
+ Math.max(0, Number(ingoingResponses) || 0),
|
|
outgoingBytes: Number(requestQueueStateMutable.networkTotals?.outgoingBytes || 0)
|
|
+ Math.max(0, Number(outgoingBytes) || 0),
|
|
ingoingBytes: Number(requestQueueStateMutable.networkTotals?.ingoingBytes || 0)
|
|
+ Math.max(0, Number(ingoingBytes) || 0),
|
|
};
|
|
};
|
|
|
|
const parseJsonIfString = (value) => {
|
|
if (typeof value !== "string") {
|
|
return value;
|
|
}
|
|
|
|
const trimmedValue = value.trim();
|
|
if (!trimmedValue.startsWith("{") && !trimmedValue.startsWith("[")) {
|
|
return value;
|
|
}
|
|
|
|
try {
|
|
return JSON.parse(trimmedValue);
|
|
} catch (error) {
|
|
return value;
|
|
}
|
|
};
|
|
|
|
const normalizeResponseData = (value) => {
|
|
const parsed = parseJsonIfString(value);
|
|
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) {
|
|
return parsed;
|
|
}
|
|
|
|
const parsedData = parseJsonIfString(parsed.data);
|
|
if (parsedData === parsed.data) {
|
|
return parsed;
|
|
}
|
|
|
|
return {
|
|
...parsed,
|
|
data: parsedData,
|
|
};
|
|
};
|
|
|
|
const extractMissingPermissions = (error) => {
|
|
const responseData = normalizeResponseData(error?.response?.data);
|
|
const messageCandidates = [
|
|
responseData?.data?.message,
|
|
responseData?.message,
|
|
error?.message,
|
|
];
|
|
const hasMissingPermissionMessage = messageCandidates.some((message) =>
|
|
typeof message === "string" && message.toLowerCase().includes("missing permission")
|
|
);
|
|
|
|
if (!hasMissingPermissionMessage) {
|
|
return [];
|
|
}
|
|
|
|
const permissions = responseData?.data?.permissions ?? responseData?.permissions ?? [];
|
|
if (!Array.isArray(permissions)) {
|
|
return [];
|
|
}
|
|
|
|
return permissions
|
|
.map((permission) => String(permission || "").trim())
|
|
.filter((permission) => permission.length > 0);
|
|
};
|
|
|
|
const getMethodConcurrencyLimit = (method) => {
|
|
const normalizedMethod = normalizeMethod(method);
|
|
const configuredLimit = queueConfig.concurrencyByMethod[normalizedMethod]
|
|
?? queueConfig.concurrencyByMethod.DEFAULT
|
|
?? 1;
|
|
return Math.max(1, Number(configuredLimit) || 1);
|
|
};
|
|
|
|
const getJobConcurrencyKey = (job) => {
|
|
const queueGroup = normalizeQueueGroup(job.queueGroup);
|
|
if (queueGroup) {
|
|
return `GROUP:${queueGroup}`;
|
|
}
|
|
|
|
return `METHOD:${normalizeMethod(job.method)}`;
|
|
};
|
|
|
|
const getJobConcurrencyLimit = (job) => {
|
|
const configuredLimit = Number.parseInt(String(job.concurrencyLimit ?? ""), 10);
|
|
if (Number.isInteger(configuredLimit) && configuredLimit > 0) {
|
|
return configuredLimit;
|
|
}
|
|
|
|
return getMethodConcurrencyLimit(job.method);
|
|
};
|
|
|
|
const canRunJob = (job) => {
|
|
const concurrencyKey = getJobConcurrencyKey(job);
|
|
const activeForKey = Number(activeWorkersByKey[concurrencyKey] || 0);
|
|
return activeForKey < getJobConcurrencyLimit(job);
|
|
};
|
|
|
|
const getNextRunnableJobIndex = () => {
|
|
for (let index = 0; index < requestQueue.length; index += 1) {
|
|
if (canRunJob(requestQueue[index])) {
|
|
return index;
|
|
}
|
|
}
|
|
return -1;
|
|
};
|
|
|
|
const clearDrainTimer = () => {
|
|
if (drainTimer !== null) {
|
|
clearTimeout(drainTimer);
|
|
drainTimer = null;
|
|
}
|
|
};
|
|
|
|
const scheduleDrain = (delayMs = 0) => {
|
|
if (delayMs > 0) {
|
|
if (drainTimer !== null) {
|
|
return;
|
|
}
|
|
|
|
drainTimer = setTimeout(() => {
|
|
drainTimer = null;
|
|
void drainQueue();
|
|
}, delayMs);
|
|
return;
|
|
}
|
|
|
|
Promise.resolve().then(() => {
|
|
void drainQueue();
|
|
});
|
|
};
|
|
|
|
const wait = (durationMs) => new Promise((resolve) => {
|
|
setTimeout(resolve, Math.max(0, durationMs || 0));
|
|
});
|
|
|
|
const getErrorStatusCode = (error) => {
|
|
const code = Number.parseInt(error?.response?.status, 10);
|
|
return Number.isFinite(code) ? code : null;
|
|
};
|
|
|
|
const getResponseStatusCode = (response) => {
|
|
const code = Number.parseInt(response?.status, 10);
|
|
return Number.isFinite(code) ? code : null;
|
|
};
|
|
|
|
const getRetriesForStatusCode = (statusCode, retryByStatusCode) => {
|
|
if (statusCode === null) {
|
|
return 0;
|
|
}
|
|
|
|
const retries = Number(retryByStatusCode?.[statusCode] ?? 0);
|
|
return Number.isFinite(retries) && retries > 0 ? Math.floor(retries) : 0;
|
|
};
|
|
|
|
const computeRetryDelayMs = (attemptNumber) => {
|
|
const exponentialDelay = queueConfig.retryDelayBaseMs * (2 ** Math.max(0, attemptNumber));
|
|
const cappedDelay = Math.min(queueConfig.retryDelayMaxMs, exponentialDelay);
|
|
const jitter = queueConfig.retryDelayJitterMs > 0
|
|
? Math.floor(Math.random() * (queueConfig.retryDelayJitterMs + 1))
|
|
: 0;
|
|
return Math.max(0, cappedDelay + jitter);
|
|
};
|
|
|
|
const shouldRetryAttempt = (error, job, attemptNumber) => {
|
|
const statusCode = getErrorStatusCode(error);
|
|
const retryByStatusCode = job.retryByStatusCode ?? queueConfig.retryByStatusCode;
|
|
const retriesForCode = getRetriesForStatusCode(statusCode, retryByStatusCode);
|
|
if (attemptNumber >= retriesForCode) {
|
|
return false;
|
|
}
|
|
|
|
if (typeof job.shouldRetry === "function") {
|
|
return job.shouldRetry(error, { attemptNumber, statusCode }) === true;
|
|
}
|
|
|
|
return true;
|
|
};
|
|
|
|
const executeJobWithRetries = async (job) => {
|
|
let attemptCount = 0;
|
|
|
|
while (true) {
|
|
attemptCount += 1;
|
|
const method = normalizeMethod(job.method);
|
|
const url = normalizeUrl(job.url);
|
|
addNetworkTotals({
|
|
outgoingRequests: 1,
|
|
outgoingBytes: estimateRequestBytes(job, method, url),
|
|
});
|
|
try {
|
|
const response = await job.requestFactory();
|
|
addNetworkTotals({
|
|
ingoingResponses: 1,
|
|
ingoingBytes: estimateResponseBytes(response),
|
|
});
|
|
return { response, attemptCount };
|
|
} catch (error) {
|
|
if (error?.response) {
|
|
addNetworkTotals({
|
|
ingoingResponses: 1,
|
|
ingoingBytes: estimateResponseBytes(error.response, error?.message ?? "Request failed"),
|
|
});
|
|
}
|
|
const attemptNumber = attemptCount - 1;
|
|
if (!shouldRetryAttempt(error, job, attemptNumber)) {
|
|
error.__queueAttemptCount = attemptCount;
|
|
throw error;
|
|
}
|
|
|
|
const retryDelayMs = computeRetryDelayMs(attemptNumber);
|
|
if (retryDelayMs > 0) {
|
|
await wait(retryDelayMs);
|
|
}
|
|
}
|
|
}
|
|
};
|
|
|
|
const upsertActiveRequest = (job, startedAt) => {
|
|
const nextActive = [...requestQueueStateMutable.activeRequests, {
|
|
id: job.id,
|
|
method: normalizeMethod(job.method),
|
|
queueGroup: normalizeQueueGroup(job.queueGroup) || null,
|
|
url: job.url,
|
|
queuedAt: job.enqueuedAt,
|
|
startedAt,
|
|
}];
|
|
requestQueueStateMutable.activeRequests = nextActive;
|
|
};
|
|
|
|
const removeActiveRequest = (requestId) => {
|
|
requestQueueStateMutable.activeRequests = requestQueueStateMutable.activeRequests.filter((item) => item.id !== requestId);
|
|
};
|
|
|
|
const pushRecentRequest = (entry) => {
|
|
const limit = Math.max(1, Number(queueConfig.recentRequestsLimit) || 8);
|
|
const next = [entry, ...requestQueueStateMutable.recentRequests];
|
|
requestQueueStateMutable.recentRequests = next.slice(0, limit);
|
|
};
|
|
|
|
const pushErrorRequest = (entry) => {
|
|
const limit = Math.max(1, Number(queueConfig.errorHistoryLimit) || 10);
|
|
const next = [entry, ...requestQueueStateMutable.errorRequests];
|
|
requestQueueStateMutable.errorRequests = next.slice(0, limit);
|
|
};
|
|
|
|
const pushMissingPermissions = (permissionEntries) => {
|
|
if (!Array.isArray(permissionEntries) || permissionEntries.length === 0) {
|
|
return;
|
|
}
|
|
|
|
const now = Date.now();
|
|
const existing = [...requestQueueStateMutable.missingPermissions];
|
|
permissionEntries.forEach((entry) => {
|
|
const permission = String(entry?.permission || "").trim();
|
|
const method = normalizeMethod(entry?.method);
|
|
const url = normalizeUrl(entry?.url);
|
|
const statusCode = entry?.statusCode ?? null;
|
|
const source = entry?.source ?? null;
|
|
if (!permission) {
|
|
return;
|
|
}
|
|
|
|
const isComponentEntry = method === "COMPONENT";
|
|
const index = existing.findIndex((item) => {
|
|
if (item.permission !== permission) {
|
|
return false;
|
|
}
|
|
if (isComponentEntry) {
|
|
return normalizeMethod(item.method) === "COMPONENT";
|
|
}
|
|
return normalizeMethod(item.method) !== "COMPONENT";
|
|
});
|
|
if (index >= 0) {
|
|
existing[index] = {
|
|
...existing[index],
|
|
method,
|
|
url,
|
|
statusCode,
|
|
source,
|
|
count: Number(existing[index].count || 0) + 1,
|
|
lastSeenAt: now,
|
|
};
|
|
return;
|
|
}
|
|
|
|
existing.push({
|
|
permission,
|
|
method,
|
|
url,
|
|
statusCode,
|
|
source,
|
|
count: 1,
|
|
lastSeenAt: now,
|
|
});
|
|
});
|
|
|
|
existing.sort((a, b) => {
|
|
const aPriority = normalizeMethod(a?.method) === "COMPONENT" ? 1 : 0;
|
|
const bPriority = normalizeMethod(b?.method) === "COMPONENT" ? 1 : 0;
|
|
if (aPriority !== bPriority) {
|
|
return aPriority - bPriority;
|
|
}
|
|
return Number(b.lastSeenAt || 0) - Number(a.lastSeenAt || 0);
|
|
});
|
|
const limit = Math.max(1, Number(queueConfig.missingPermissionsLimit) || 20);
|
|
requestQueueStateMutable.missingPermissions = existing.slice(0, limit);
|
|
};
|
|
|
|
export const reportComponentMissingPermission = (permission, options = {}) => {
|
|
const normalizedPermission = String(permission || "").trim();
|
|
if (!normalizedPermission) {
|
|
return;
|
|
}
|
|
|
|
const source = String(options?.source || "").trim() || "Component";
|
|
const now = Date.now();
|
|
const reportWindowMs = Math.max(0, Number(queueConfig.componentPermissionReportWindowMs) || 0);
|
|
const reportKey = `${normalizedPermission}::${source}`;
|
|
const lastReportedAt = Number(recentComponentPermissionReports[reportKey] || 0);
|
|
if (reportWindowMs > 0 && now - lastReportedAt < reportWindowMs) {
|
|
return;
|
|
}
|
|
|
|
recentComponentPermissionReports[reportKey] = now;
|
|
pushMissingPermissions([{
|
|
permission: normalizedPermission,
|
|
method: "COMPONENT",
|
|
url: source,
|
|
statusCode: null,
|
|
source,
|
|
}]);
|
|
};
|
|
|
|
const runJob = (job) => {
|
|
const method = normalizeMethod(job.method);
|
|
const concurrencyKey = getJobConcurrencyKey(job);
|
|
const queueGroup = normalizeQueueGroup(job.queueGroup) || null;
|
|
const startedAt = Date.now();
|
|
|
|
activeWorkers += 1;
|
|
activeWorkersByKey[concurrencyKey] = Number(activeWorkersByKey[concurrencyKey] || 0) + 1;
|
|
lastRequestStartedAt = startedAt;
|
|
upsertActiveRequest(job, startedAt);
|
|
syncQueueCounters();
|
|
|
|
Promise.resolve()
|
|
.then(() => executeJobWithRetries(job))
|
|
.then(({ response, attemptCount }) => {
|
|
const completedAt = Date.now();
|
|
requestQueueStateMutable.batchCompleted += 1;
|
|
pushRecentRequest({
|
|
id: job.id,
|
|
method,
|
|
queueGroup,
|
|
url: job.url,
|
|
success: true,
|
|
statusCode: getResponseStatusCode(response),
|
|
attemptCount,
|
|
queuedAt: job.enqueuedAt,
|
|
startedAt,
|
|
completedAt,
|
|
queueDurationMs: Math.max(0, startedAt - job.enqueuedAt),
|
|
requestDurationMs: Math.max(0, completedAt - startedAt),
|
|
});
|
|
job.resolve(response);
|
|
})
|
|
.catch((error) => {
|
|
const completedAt = Date.now();
|
|
const statusCode = getErrorStatusCode(error);
|
|
const requestSnapshot = {
|
|
method,
|
|
url: job.url,
|
|
params: job.requestData?.params ?? null,
|
|
data: job.requestData?.data ?? null,
|
|
headers: redactHeaders(job.requestData?.headers ?? null),
|
|
};
|
|
const responseSnapshot = {
|
|
status: statusCode,
|
|
statusText: error?.response?.statusText ?? null,
|
|
code: error?.code ?? null,
|
|
message: error?.message ?? "Request failed",
|
|
data: error?.response?.data ?? null,
|
|
headers: redactHeaders(error?.response?.headers ?? null),
|
|
};
|
|
requestQueueStateMutable.batchFailed += 1;
|
|
pushRecentRequest({
|
|
id: job.id,
|
|
method,
|
|
queueGroup,
|
|
url: job.url,
|
|
success: false,
|
|
statusCode,
|
|
attemptCount: Math.max(1, Number(error?.__queueAttemptCount) || 1),
|
|
queuedAt: job.enqueuedAt,
|
|
startedAt,
|
|
completedAt,
|
|
queueDurationMs: Math.max(0, startedAt - job.enqueuedAt),
|
|
requestDurationMs: Math.max(0, completedAt - startedAt),
|
|
});
|
|
pushErrorRequest({
|
|
id: job.id,
|
|
method,
|
|
queueGroup,
|
|
url: job.url,
|
|
statusCode,
|
|
attemptCount: Math.max(1, Number(error?.__queueAttemptCount) || 1),
|
|
requestDurationMs: Math.max(0, completedAt - startedAt),
|
|
requestText: toSafeText(requestSnapshot),
|
|
responseText: toSafeText(responseSnapshot),
|
|
});
|
|
recordReleaseTimelineEvent("request_failed", {
|
|
request_id: String(job.id),
|
|
method,
|
|
url: job.url,
|
|
statusCode,
|
|
attemptCount: Math.max(1, Number(error?.__queueAttemptCount) || 1),
|
|
requestDurationMs: Math.max(0, completedAt - startedAt),
|
|
request: requestSnapshot,
|
|
response: responseSnapshot,
|
|
}, {
|
|
severity: "error",
|
|
moduleKey: "requestqueue",
|
|
requestId: String(job.id),
|
|
route: job.url,
|
|
});
|
|
|
|
const missingPermissions = extractMissingPermissions(error);
|
|
pushMissingPermissions(
|
|
missingPermissions.map((permission) => ({
|
|
permission,
|
|
method,
|
|
url: job.url,
|
|
statusCode,
|
|
}))
|
|
);
|
|
job.reject(error);
|
|
})
|
|
.finally(() => {
|
|
activeWorkers -= 1;
|
|
activeWorkersByKey[concurrencyKey] = Math.max(0, Number(activeWorkersByKey[concurrencyKey] || 1) - 1);
|
|
if (activeWorkersByKey[concurrencyKey] === 0) {
|
|
delete activeWorkersByKey[concurrencyKey];
|
|
}
|
|
removeActiveRequest(job.id);
|
|
syncQueueCounters();
|
|
scheduleDrain();
|
|
});
|
|
};
|
|
|
|
const drainQueue = async () => {
|
|
while (requestQueue.length > 0) {
|
|
const nextRunnableJobIndex = getNextRunnableJobIndex();
|
|
if (nextRunnableJobIndex < 0) {
|
|
return;
|
|
}
|
|
|
|
const elapsedSinceLastStart = Date.now() - lastRequestStartedAt;
|
|
const waitMs = Math.max(0, queueConfig.spacingMs - elapsedSinceLastStart);
|
|
|
|
if (waitMs > 0) {
|
|
scheduleDrain(waitMs);
|
|
return;
|
|
}
|
|
|
|
const [nextJob] = requestQueue.splice(nextRunnableJobIndex, 1);
|
|
syncQueueCounters();
|
|
runJob(nextJob);
|
|
}
|
|
};
|
|
|
|
export const enqueueRequest = (requestFactory, options = {}) => {
|
|
if (typeof requestFactory !== "function") {
|
|
throw new Error("enqueueRequest requires a function");
|
|
}
|
|
|
|
const method = normalizeMethod(options.method);
|
|
const url = normalizeUrl(options.url);
|
|
startBatchIfNeeded();
|
|
|
|
return new Promise((resolve, reject) => {
|
|
requestQueue.push({
|
|
id: ++requestIdCounter,
|
|
requestFactory,
|
|
resolve,
|
|
reject,
|
|
method,
|
|
url,
|
|
enqueuedAt: Date.now(),
|
|
requestData: options.requestData || null,
|
|
retryByStatusCode: options.retryByStatusCode || null,
|
|
shouldRetry: typeof options.shouldRetry === "function" ? options.shouldRetry : null,
|
|
queueGroup: normalizeQueueGroup(options.queueGroup),
|
|
concurrencyLimit: options.concurrencyLimit || null,
|
|
});
|
|
requestQueueStateMutable.batchTotal += 1;
|
|
syncQueueCounters();
|
|
scheduleDrain();
|
|
});
|
|
};
|
|
|
|
export const requestQueueState = readonly(requestQueueStateMutable);
|
|
|
|
export const clearErrorRequests = () => {
|
|
requestQueueStateMutable.errorRequests = [];
|
|
};
|
|
|
|
export const clearMissingPermissions = () => {
|
|
requestQueueStateMutable.missingPermissions = [];
|
|
Object.keys(recentComponentPermissionReports).forEach((key) => {
|
|
delete recentComponentPermissionReports[key];
|
|
});
|
|
};
|
|
|
|
export const __resetRequestQueueForTests = () => {
|
|
requestQueue.length = 0;
|
|
activeWorkers = 0;
|
|
requestIdCounter = 0;
|
|
Object.keys(activeWorkersByKey).forEach((key) => delete activeWorkersByKey[key]);
|
|
lastRequestStartedAt = 0;
|
|
clearDrainTimer();
|
|
|
|
requestQueueStateMutable.batchId = 0;
|
|
requestQueueStateMutable.pending = 0;
|
|
requestQueueStateMutable.active = 0;
|
|
requestQueueStateMutable.batchTotal = 0;
|
|
requestQueueStateMutable.batchCompleted = 0;
|
|
requestQueueStateMutable.batchFailed = 0;
|
|
requestQueueStateMutable.activeRequests = [];
|
|
requestQueueStateMutable.recentRequests = [];
|
|
requestQueueStateMutable.errorRequests = [];
|
|
requestQueueStateMutable.missingPermissions = [];
|
|
requestQueueStateMutable.networkTotals = {
|
|
outgoingRequests: 0,
|
|
ingoingResponses: 0,
|
|
outgoingBytes: 0,
|
|
ingoingBytes: 0,
|
|
};
|
|
|
|
queueConfig.concurrencyByMethod = cloneObject(REQUEST_QUEUE_CONFIG.concurrency);
|
|
queueConfig.spacingMs = 0;
|
|
queueConfig.retryByStatusCode = cloneObject(REQUEST_QUEUE_CONFIG.retryByStatusCode);
|
|
queueConfig.retryDelayBaseMs = 0;
|
|
queueConfig.retryDelayMaxMs = 0;
|
|
queueConfig.retryDelayJitterMs = 0;
|
|
queueConfig.recentRequestsLimit = REQUEST_QUEUE_CONFIG.inspector.recentRequestsLimit;
|
|
queueConfig.errorHistoryLimit = REQUEST_QUEUE_CONFIG.inspector.errorHistoryLimit;
|
|
queueConfig.missingPermissionsLimit = REQUEST_QUEUE_CONFIG.inspector.missingPermissionsLimit;
|
|
queueConfig.componentPermissionReportWindowMs = REQUEST_QUEUE_CONFIG.inspector.componentPermissionReportWindowMs;
|
|
queueConfig.payloadMaxChars = REQUEST_QUEUE_CONFIG.inspector.payloadMaxChars;
|
|
Object.keys(recentComponentPermissionReports).forEach((key) => {
|
|
delete recentComponentPermissionReports[key];
|
|
});
|
|
};
|
|
|
|
export const __configureRequestQueueForTests = ({
|
|
concurrencyByMethod,
|
|
maxConcurrentGet,
|
|
maxConcurrentOther,
|
|
spacingMs,
|
|
retryByStatusCode,
|
|
retryDelayBaseMs,
|
|
retryDelayMaxMs,
|
|
retryDelayJitterMs,
|
|
recentRequestsLimit,
|
|
errorHistoryLimit,
|
|
missingPermissionsLimit,
|
|
componentPermissionReportWindowMs,
|
|
payloadMaxChars,
|
|
} = {}) => {
|
|
if (concurrencyByMethod && typeof concurrencyByMethod === "object") {
|
|
queueConfig.concurrencyByMethod = {
|
|
...queueConfig.concurrencyByMethod,
|
|
...concurrencyByMethod,
|
|
};
|
|
}
|
|
|
|
if (typeof maxConcurrentOther === "number" && maxConcurrentOther > 0) {
|
|
const normalized = Math.max(1, Math.floor(maxConcurrentOther));
|
|
queueConfig.concurrencyByMethod.DEFAULT = normalized;
|
|
queueConfig.concurrencyByMethod.POST = normalized;
|
|
queueConfig.concurrencyByMethod.PATCH = normalized;
|
|
queueConfig.concurrencyByMethod.PUT = normalized;
|
|
queueConfig.concurrencyByMethod.DELETE = normalized;
|
|
}
|
|
|
|
if (typeof maxConcurrentGet === "number" && maxConcurrentGet > 0) {
|
|
queueConfig.concurrencyByMethod.GET = Math.max(1, Math.floor(maxConcurrentGet));
|
|
}
|
|
|
|
if (typeof spacingMs === "number" && spacingMs >= 0) {
|
|
queueConfig.spacingMs = Math.floor(spacingMs);
|
|
}
|
|
|
|
if (retryByStatusCode && typeof retryByStatusCode === "object") {
|
|
queueConfig.retryByStatusCode = cloneObject(retryByStatusCode);
|
|
}
|
|
|
|
if (typeof retryDelayBaseMs === "number" && retryDelayBaseMs >= 0) {
|
|
queueConfig.retryDelayBaseMs = Math.floor(retryDelayBaseMs);
|
|
}
|
|
|
|
if (typeof retryDelayMaxMs === "number" && retryDelayMaxMs >= 0) {
|
|
queueConfig.retryDelayMaxMs = Math.floor(retryDelayMaxMs);
|
|
}
|
|
|
|
if (typeof retryDelayJitterMs === "number" && retryDelayJitterMs >= 0) {
|
|
queueConfig.retryDelayJitterMs = Math.floor(retryDelayJitterMs);
|
|
}
|
|
|
|
if (typeof recentRequestsLimit === "number" && recentRequestsLimit > 0) {
|
|
queueConfig.recentRequestsLimit = Math.max(1, Math.floor(recentRequestsLimit));
|
|
}
|
|
|
|
if (typeof errorHistoryLimit === "number" && errorHistoryLimit > 0) {
|
|
queueConfig.errorHistoryLimit = Math.max(1, Math.floor(errorHistoryLimit));
|
|
}
|
|
|
|
if (typeof missingPermissionsLimit === "number" && missingPermissionsLimit > 0) {
|
|
queueConfig.missingPermissionsLimit = Math.max(1, Math.floor(missingPermissionsLimit));
|
|
}
|
|
|
|
if (typeof componentPermissionReportWindowMs === "number" && componentPermissionReportWindowMs >= 0) {
|
|
queueConfig.componentPermissionReportWindowMs = Math.floor(componentPermissionReportWindowMs);
|
|
}
|
|
|
|
if (typeof payloadMaxChars === "number" && payloadMaxChars >= 200) {
|
|
queueConfig.payloadMaxChars = Math.floor(payloadMaxChars);
|
|
}
|
|
};
|