@@ -92,7 +92,7 @@
|
||||
"octokit": "^2.0.0",
|
||||
"octokit-plugin-create-pull-request": "^3.10.0",
|
||||
"p-limit": "^3.1.0",
|
||||
"p-queue": "^6.6.2",
|
||||
"p-queue": "^7.4.1",
|
||||
"prom-client": "^14.0.1",
|
||||
"uuid": "^8.2.0",
|
||||
"winston": "^3.2.1",
|
||||
|
||||
@@ -14,6 +14,7 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
import fetch from 'node-fetch';
|
||||
import {
|
||||
createBackendPlugin,
|
||||
coreServices,
|
||||
@@ -79,6 +80,7 @@ export const scaffolderPlugin = createBackendPlugin({
|
||||
reader: coreServices.urlReader,
|
||||
permissions: coreServices.permissions,
|
||||
database: coreServices.database,
|
||||
discovery: coreServices.discovery,
|
||||
httpRouter: coreServices.httpRouter,
|
||||
catalogClient: catalogServiceRef,
|
||||
},
|
||||
@@ -88,6 +90,7 @@ export const scaffolderPlugin = createBackendPlugin({
|
||||
lifecycle,
|
||||
reader,
|
||||
database,
|
||||
discovery,
|
||||
httpRouter,
|
||||
catalogClient,
|
||||
permissions,
|
||||
@@ -107,23 +110,11 @@ export const scaffolderPlugin = createBackendPlugin({
|
||||
];
|
||||
|
||||
lifecycle.addShutdownHook(async () => {
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({
|
||||
database,
|
||||
const baseUrl = await discovery.getBaseUrl('scaffolder');
|
||||
const url = `${baseUrl}/v2/tasks/cancel`;
|
||||
await fetch(url, {
|
||||
method: 'POST',
|
||||
});
|
||||
const { tasks: processingTasks } = await databaseTaskStore.list({
|
||||
status: 'processing',
|
||||
});
|
||||
|
||||
if (processingTasks.length > 0) {
|
||||
await Promise.all(
|
||||
processingTasks.map(task =>
|
||||
databaseTaskStore.shutdownTask({ taskId: task.id }),
|
||||
),
|
||||
);
|
||||
logger.info(
|
||||
`Successfully shut ${processingTasks.length} processing tasks down.`,
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
const actionIds = actions.map(action => action.id).join(', ');
|
||||
|
||||
@@ -78,8 +78,10 @@ export type CreateWorkerOptions = {
|
||||
*/
|
||||
export class TaskWorker {
|
||||
private taskQueue: PQueue;
|
||||
private taskQueueAbortController: AbortController;
|
||||
|
||||
private constructor(private readonly options: TaskWorkerOptions) {
|
||||
this.taskQueueAbortController = new AbortController();
|
||||
this.taskQueue = new PQueue({
|
||||
concurrency: options.concurrentTasksLimit,
|
||||
});
|
||||
@@ -121,11 +123,31 @@ export class TaskWorker {
|
||||
for (;;) {
|
||||
await this.onReadyToClaimTask();
|
||||
const task = await this.options.taskBroker.claim();
|
||||
this.taskQueue.add(() => this.runOneTask(task));
|
||||
const taskId = await task.getWorkspaceName();
|
||||
|
||||
try {
|
||||
this.taskQueue.add(
|
||||
({ signal }) => {
|
||||
this.runOneTask(task);
|
||||
signal?.addEventListener('abort', () => {
|
||||
this.options.taskBroker.cancel?.(taskId);
|
||||
});
|
||||
},
|
||||
{ signal: this.taskQueueAbortController.signal },
|
||||
);
|
||||
} catch (error) {
|
||||
if (!(error instanceof AbortError)) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
})();
|
||||
}
|
||||
|
||||
cancelAllRunningTasks() {
|
||||
this.taskQueueAbortController.abort();
|
||||
}
|
||||
|
||||
protected onReadyToClaimTask(): Promise<void> {
|
||||
if (this.taskQueue.pending < this.options.concurrentTasksLimit) {
|
||||
return Promise.resolve();
|
||||
|
||||
@@ -301,7 +301,7 @@ export async function createRouter(
|
||||
|
||||
const actionRegistry = new TemplateActionRegistry();
|
||||
|
||||
const workers = [];
|
||||
const workers: TaskWorker[] = [];
|
||||
if (concurrentTasksLimit !== 0) {
|
||||
for (let i = 0; i < (taskWorkers || 1); i++) {
|
||||
const worker = await TaskWorker.create({
|
||||
@@ -516,6 +516,10 @@ export async function createRouter(
|
||||
delete task.secrets;
|
||||
res.status(200).json(task);
|
||||
})
|
||||
.post('/v2/tasks/cancel', async (req, res) => {
|
||||
workers.forEach(worker => worker.cancelAllRunningTasks());
|
||||
res.status(200).json({ status: 'cancelled' });
|
||||
})
|
||||
.post('/v2/tasks/:taskId/cancel', async (req, res) => {
|
||||
const { taskId } = req.params;
|
||||
await taskBroker.cancel?.(taskId);
|
||||
|
||||
@@ -8876,7 +8876,7 @@ __metadata:
|
||||
octokit: ^2.0.0
|
||||
octokit-plugin-create-pull-request: ^3.10.0
|
||||
p-limit: ^3.1.0
|
||||
p-queue: ^6.6.2
|
||||
p-queue: ^7.4.1
|
||||
prom-client: ^14.0.1
|
||||
supertest: ^6.1.3
|
||||
uuid: ^8.2.0
|
||||
@@ -37014,6 +37014,16 @@ __metadata:
|
||||
languageName: node
|
||||
linkType: hard
|
||||
|
||||
"p-queue@npm:^7.4.1":
|
||||
version: 7.4.1
|
||||
resolution: "p-queue@npm:7.4.1"
|
||||
dependencies:
|
||||
eventemitter3: ^5.0.1
|
||||
p-timeout: ^5.0.2
|
||||
checksum: 1c6888aa994d399262a9fbdd49c7066f8359732397f7a42ecf03f22875a1d65899797b46413f97e44acc18dddafbcc101eb135c284714c931dbbc83c3967f450
|
||||
languageName: node
|
||||
linkType: hard
|
||||
|
||||
"p-retry@npm:^4.5.0":
|
||||
version: 4.5.0
|
||||
resolution: "p-retry@npm:4.5.0"
|
||||
@@ -37033,6 +37043,13 @@ __metadata:
|
||||
languageName: node
|
||||
linkType: hard
|
||||
|
||||
"p-timeout@npm:^5.0.2":
|
||||
version: 5.1.0
|
||||
resolution: "p-timeout@npm:5.1.0"
|
||||
checksum: f5cd4e17301ff1ff1d8dbf2817df0ad88c6bba99349fc24d8d181827176ad4f8aca649190b8a5b1a428dfd6ddc091af4606835d3e0cb0656e04045da5c9e270c
|
||||
languageName: node
|
||||
linkType: hard
|
||||
|
||||
"p-transform@npm:^1.3.0":
|
||||
version: 1.3.0
|
||||
resolution: "p-transform@npm:1.3.0"
|
||||
|
||||
Reference in New Issue
Block a user