mirror of
https://github.com/morten-olsen/mini-loader.git
synced 2026-02-08 01:36:26 +01:00
feat: add http gateway (#3)
This commit is contained in:
@@ -17,12 +17,12 @@ const bundle = async ({ entry, autoInstall }: BundleOptions) => {
|
|||||||
const entryFile = resolve(entry);
|
const entryFile = resolve(entry);
|
||||||
const codeBundler = await rollup({
|
const codeBundler = await rollup({
|
||||||
plugins: [
|
plugins: [
|
||||||
|
fix(json)(),
|
||||||
fix(sucrase)({
|
fix(sucrase)({
|
||||||
transforms: ['typescript', 'jsx'],
|
transforms: ['typescript', 'jsx'],
|
||||||
}),
|
}),
|
||||||
...[autoInstall ? fix(auto) : []],
|
...[autoInstall ? fix(auto) : []],
|
||||||
nodeResolve({ extensions: ['.js', '.jsx', '.ts', '.tsx'] }),
|
nodeResolve({ preferBuiltins: true, extensions: ['.js', '.jsx', '.ts', '.tsx'] }),
|
||||||
fix(json)(),
|
|
||||||
fix(commonjs)({ include: /node_modules/ }),
|
fix(commonjs)({ include: /node_modules/ }),
|
||||||
],
|
],
|
||||||
input: entryFile,
|
input: entryFile,
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ push
|
|||||||
const code = await step('Bundling', async () => {
|
const code = await step('Bundling', async () => {
|
||||||
return await bundle({ entry: location, autoInstall: opts.autoInstall });
|
return await bundle({ entry: location, autoInstall: opts.autoInstall });
|
||||||
});
|
});
|
||||||
const id = await step('Creating load', async () => {
|
const id = await step(`Creating load ${(code.length / 1024).toFixed(0)}`, async () => {
|
||||||
return await client.loads.set.mutate({
|
return await client.loads.set.mutate({
|
||||||
id: opts.id,
|
id: opts.id,
|
||||||
name: opts.name,
|
name: opts.name,
|
||||||
@@ -34,9 +34,10 @@ push
|
|||||||
});
|
});
|
||||||
console.log('created load with id', id);
|
console.log('created load with id', id);
|
||||||
if (opts.run) {
|
if (opts.run) {
|
||||||
await step('Creating run', async () => {
|
const runId = await step('Creating run', async () => {
|
||||||
await client.runs.create.mutate({ loadId: id });
|
return await client.runs.create.mutate({ loadId: id });
|
||||||
});
|
});
|
||||||
|
console.log('created run with id', runId);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
59
packages/cli/src/commands/runs/runs.remove.ts
Normal file
59
packages/cli/src/commands/runs/runs.remove.ts
Normal file
@@ -0,0 +1,59 @@
|
|||||||
|
import { Command } from 'commander';
|
||||||
|
import { createClient } from '../../client/client.js';
|
||||||
|
import { step } from '../../utils/step.js';
|
||||||
|
import { Context } from '../../context/context.js';
|
||||||
|
import inquirer from 'inquirer';
|
||||||
|
import { Config } from '../../config/config.js';
|
||||||
|
|
||||||
|
const remove = new Command('remove');
|
||||||
|
|
||||||
|
const toInt = (value?: string) => {
|
||||||
|
if (!value) {
|
||||||
|
return undefined;
|
||||||
|
}
|
||||||
|
return parseInt(value, 10);
|
||||||
|
};
|
||||||
|
|
||||||
|
remove
|
||||||
|
.alias('ls')
|
||||||
|
.description('List logs')
|
||||||
|
.option('-l, --load-id <loadId>', 'Load ID')
|
||||||
|
.option('-o, --offset <offset>', 'Offset')
|
||||||
|
.option('-a, --limit <limit>', 'Limit', '1000')
|
||||||
|
.action(async () => {
|
||||||
|
const { loadId, offset, limit } = remove.opts();
|
||||||
|
const config = new Config();
|
||||||
|
const context = new Context(config.context);
|
||||||
|
const client = await step('Connecting to server', async () => {
|
||||||
|
return createClient(context);
|
||||||
|
});
|
||||||
|
const response = await step('Preparing to delete', async () => {
|
||||||
|
return await client.runs.prepareRemove.query({
|
||||||
|
loadId,
|
||||||
|
offset: toInt(offset),
|
||||||
|
limit: toInt(limit),
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!response.ids.length) {
|
||||||
|
console.log('No logs to delete');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const { confirm } = await inquirer.prompt([
|
||||||
|
{
|
||||||
|
type: 'confirm',
|
||||||
|
name: 'confirm',
|
||||||
|
message: `Are you sure you want to delete ${response.ids.length} logs?`,
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
|
||||||
|
if (!confirm) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
await step('Deleting artifacts', async () => {
|
||||||
|
await client.runs.remove.mutate(response);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
export { remove };
|
||||||
23
packages/cli/src/commands/runs/runs.terminate.ts
Normal file
23
packages/cli/src/commands/runs/runs.terminate.ts
Normal file
@@ -0,0 +1,23 @@
|
|||||||
|
import { Command } from 'commander';
|
||||||
|
import { createClient } from '../../client/client.js';
|
||||||
|
import { step } from '../../utils/step.js';
|
||||||
|
import { Context } from '../../context/context.js';
|
||||||
|
import { Config } from '../../config/config.js';
|
||||||
|
|
||||||
|
const terminate = new Command('terminate');
|
||||||
|
|
||||||
|
terminate
|
||||||
|
.description('Terminate an in progress run')
|
||||||
|
.argument('run-id', 'Run ID')
|
||||||
|
.action(async (runId) => {
|
||||||
|
const config = new Config();
|
||||||
|
const context = new Context(config.context);
|
||||||
|
const client = await step('Connecting to server', async () => {
|
||||||
|
return createClient(context);
|
||||||
|
});
|
||||||
|
await step('Terminating run', async () => {
|
||||||
|
await client.runs.terminate.mutate(runId);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
export { terminate };
|
||||||
@@ -1,8 +1,14 @@
|
|||||||
import { Command } from 'commander';
|
import { Command } from 'commander';
|
||||||
import { create } from './runs.create.js';
|
import { create } from './runs.create.js';
|
||||||
import { list } from './runs.list.js';
|
import { list } from './runs.list.js';
|
||||||
|
import { remove } from './runs.remove.js';
|
||||||
|
import { terminate } from './runs.terminate.js';
|
||||||
|
|
||||||
const runs = new Command('runs');
|
const runs = new Command('runs');
|
||||||
runs.description('Manage runs').addCommand(create).addCommand(list);
|
runs.description('Manage runs');
|
||||||
|
runs.addCommand(create);
|
||||||
|
runs.addCommand(list);
|
||||||
|
runs.addCommand(remove);
|
||||||
|
runs.addCommand(terminate);
|
||||||
|
|
||||||
export { runs };
|
export { runs };
|
||||||
|
|||||||
@@ -18,9 +18,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@morten-olsen/mini-loader-configs": "workspace:^",
|
|
||||||
"@morten-olsen/mini-loader-cli": "workspace:^",
|
|
||||||
"@morten-olsen/mini-loader": "workspace:^",
|
"@morten-olsen/mini-loader": "workspace:^",
|
||||||
|
"@morten-olsen/mini-loader-cli": "workspace:^",
|
||||||
|
"@morten-olsen/mini-loader-configs": "workspace:^",
|
||||||
"@types/node": "^20.10.8",
|
"@types/node": "^20.10.8",
|
||||||
"typescript": "^5.3.3"
|
"typescript": "^5.3.3"
|
||||||
},
|
},
|
||||||
@@ -28,5 +28,8 @@
|
|||||||
"repository": {
|
"repository": {
|
||||||
"type": "git",
|
"type": "git",
|
||||||
"url": "https://github.com/morten-olsen/mini-loader"
|
"url": "https://github.com/morten-olsen/mini-loader"
|
||||||
|
},
|
||||||
|
"dependencies": {
|
||||||
|
"fastify": "^4.25.2"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
12
packages/examples/src/http.ts
Normal file
12
packages/examples/src/http.ts
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
import { http } from '@morten-olsen/mini-loader';
|
||||||
|
import fastify from 'fastify';
|
||||||
|
|
||||||
|
const server = fastify();
|
||||||
|
|
||||||
|
server.all('*', async (req) => {
|
||||||
|
return req.url;
|
||||||
|
});
|
||||||
|
|
||||||
|
server.listen({
|
||||||
|
path: http.getPath(),
|
||||||
|
});
|
||||||
@@ -3,6 +3,7 @@ import { artifacts, logger } from '@morten-olsen/mini-loader';
|
|||||||
const run = async () => {
|
const run = async () => {
|
||||||
await logger.info('Hello world');
|
await logger.info('Hello world');
|
||||||
await artifacts.create('foo', 'bar');
|
await artifacts.create('foo', 'bar');
|
||||||
|
process.exit(0);
|
||||||
};
|
};
|
||||||
|
|
||||||
run();
|
run();
|
||||||
|
|||||||
7
packages/mini-loader/src/http/http.ts
Normal file
7
packages/mini-loader/src/http/http.ts
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
const getPath = () => process.env.HTTP_GATEWAY_PATH!;
|
||||||
|
|
||||||
|
const http = {
|
||||||
|
getPath,
|
||||||
|
};
|
||||||
|
|
||||||
|
export { http };
|
||||||
@@ -8,3 +8,4 @@ export { logger } from './logger/logger.js';
|
|||||||
export { artifacts } from './artifacts/artifacts.js';
|
export { artifacts } from './artifacts/artifacts.js';
|
||||||
export { input } from './input/input.js';
|
export { input } from './input/input.js';
|
||||||
export { secrets } from './secrets/secrets.js';
|
export { secrets } from './secrets/secrets.js';
|
||||||
|
export { http } from './http/http.js';
|
||||||
|
|||||||
@@ -17,12 +17,12 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@morten-olsen/mini-loader": "workspace:^",
|
|
||||||
"@morten-olsen/mini-loader-configs": "workspace:^",
|
"@morten-olsen/mini-loader-configs": "workspace:^",
|
||||||
"@types/node": "^20.10.8",
|
"@types/node": "^20.10.8",
|
||||||
"typescript": "^5.3.3"
|
"typescript": "^5.3.3"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
|
"@morten-olsen/mini-loader": "workspace:^",
|
||||||
"eventemitter3": "^5.0.1",
|
"eventemitter3": "^5.0.1",
|
||||||
"nanoid": "^5.0.4"
|
"nanoid": "^5.0.4"
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -1,17 +1,5 @@
|
|||||||
import { Worker } from 'worker_threads';
|
import { Worker } from 'worker_threads';
|
||||||
import os from 'os';
|
import { setup } from './setup/setup.js';
|
||||||
import { EventEmitter } from 'eventemitter3';
|
|
||||||
import { Event } from '@morten-olsen/mini-loader';
|
|
||||||
import { join } from 'path';
|
|
||||||
import { createServer } from 'http';
|
|
||||||
import { nanoid } from 'nanoid';
|
|
||||||
import { chmod, mkdir, rm, writeFile } from 'fs/promises';
|
|
||||||
|
|
||||||
type RunEvents = {
|
|
||||||
message: (event: Event) => void;
|
|
||||||
error: (error: Error) => void;
|
|
||||||
exit: () => void;
|
|
||||||
};
|
|
||||||
|
|
||||||
type RunOptions = {
|
type RunOptions = {
|
||||||
script: string;
|
script: string;
|
||||||
@@ -20,44 +8,17 @@ type RunOptions = {
|
|||||||
};
|
};
|
||||||
|
|
||||||
const run = async ({ script, input, secrets }: RunOptions) => {
|
const run = async ({ script, input, secrets }: RunOptions) => {
|
||||||
const dataDir = join(os.tmpdir(), 'mini-loader', nanoid());
|
const info = await setup({ script, input, secrets });
|
||||||
await mkdir(dataDir, { recursive: true });
|
|
||||||
await chmod(dataDir, 0o700);
|
|
||||||
const hostSocket = join(dataDir, 'host');
|
|
||||||
const server = createServer();
|
|
||||||
const inputLocation = join(dataDir, 'input');
|
|
||||||
|
|
||||||
if (input) {
|
const worker = new Worker(info.scriptLocation, {
|
||||||
await writeFile(inputLocation, input);
|
|
||||||
}
|
|
||||||
|
|
||||||
const emitter = new EventEmitter<RunEvents>();
|
|
||||||
|
|
||||||
server.on('connection', (socket) => {
|
|
||||||
socket.on('data', (data) => {
|
|
||||||
const message = JSON.parse(data.toString());
|
|
||||||
emitter.emit('message', message);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
server.listen(hostSocket);
|
|
||||||
|
|
||||||
const worker = new Worker(script, {
|
|
||||||
eval: true,
|
|
||||||
stdin: false,
|
stdin: false,
|
||||||
stdout: false,
|
stdout: false,
|
||||||
stderr: false,
|
stderr: false,
|
||||||
env: {
|
env: info.env,
|
||||||
HOST_SOCKET: hostSocket,
|
|
||||||
SECRETS: JSON.stringify(secrets),
|
|
||||||
INPUT_PATH: inputLocation,
|
|
||||||
},
|
|
||||||
workerData: {
|
|
||||||
input,
|
|
||||||
},
|
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.stdout?.on('data', (data) => {
|
worker.stdout?.on('data', (data) => {
|
||||||
emitter.emit('message', {
|
info.emitter.emit('message', {
|
||||||
type: 'log',
|
type: 'log',
|
||||||
payload: {
|
payload: {
|
||||||
severity: 'info',
|
severity: 'info',
|
||||||
@@ -67,7 +28,7 @@ const run = async ({ script, input, secrets }: RunOptions) => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
worker.stderr?.on('data', (data) => {
|
worker.stderr?.on('data', (data) => {
|
||||||
emitter.emit('message', {
|
info.emitter.emit('message', {
|
||||||
type: 'log',
|
type: 'log',
|
||||||
payload: {
|
payload: {
|
||||||
severity: 'error',
|
severity: 'error',
|
||||||
@@ -78,20 +39,24 @@ const run = async ({ script, input, secrets }: RunOptions) => {
|
|||||||
|
|
||||||
const promise = new Promise<void>((resolve, reject) => {
|
const promise = new Promise<void>((resolve, reject) => {
|
||||||
worker.on('exit', async () => {
|
worker.on('exit', async () => {
|
||||||
server.close();
|
await info.teardown();
|
||||||
await rm(dataDir, { recursive: true, force: true });
|
|
||||||
resolve();
|
resolve();
|
||||||
});
|
});
|
||||||
worker.on('error', async (error) => {
|
worker.on('error', async (error) => {
|
||||||
server.close();
|
|
||||||
reject(error);
|
reject(error);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
return {
|
return {
|
||||||
emitter,
|
...info,
|
||||||
|
teardown: async () => {
|
||||||
|
worker.terminate();
|
||||||
|
},
|
||||||
promise,
|
promise,
|
||||||
};
|
};
|
||||||
};
|
};
|
||||||
|
|
||||||
|
type RunInfo = Awaited<ReturnType<typeof run>>;
|
||||||
|
|
||||||
|
export type { RunInfo };
|
||||||
export { run };
|
export { run };
|
||||||
|
|||||||
71
packages/runner/src/setup/setup.ts
Normal file
71
packages/runner/src/setup/setup.ts
Normal file
@@ -0,0 +1,71 @@
|
|||||||
|
import { join } from 'path';
|
||||||
|
import os from 'os';
|
||||||
|
import { nanoid } from 'nanoid';
|
||||||
|
import { chmod, mkdir, rm, writeFile } from 'fs/promises';
|
||||||
|
import { createServer } from 'net';
|
||||||
|
import { EventEmitter } from 'eventemitter3';
|
||||||
|
|
||||||
|
type SetupOptions = {
|
||||||
|
input?: Buffer | string;
|
||||||
|
script: string;
|
||||||
|
secrets?: Record<string, string>;
|
||||||
|
};
|
||||||
|
|
||||||
|
type RunEvents = {
|
||||||
|
message: (event: any) => void;
|
||||||
|
error: (error: Error) => void;
|
||||||
|
exit: () => void;
|
||||||
|
};
|
||||||
|
|
||||||
|
const setup = async (options: SetupOptions) => {
|
||||||
|
const { input, script, secrets } = options;
|
||||||
|
const emitter = new EventEmitter<RunEvents>();
|
||||||
|
const dataDir = join(os.tmpdir(), 'mini-loader', nanoid());
|
||||||
|
|
||||||
|
await mkdir(dataDir, { recursive: true });
|
||||||
|
await chmod(dataDir, 0o700);
|
||||||
|
const hostSocket = join(dataDir, 'host');
|
||||||
|
const httpGatewaySocket = join(dataDir, 'socket');
|
||||||
|
const server = createServer();
|
||||||
|
const inputLocation = join(dataDir, 'input');
|
||||||
|
const scriptLocation = join(dataDir, 'script.js');
|
||||||
|
|
||||||
|
if (input) {
|
||||||
|
await writeFile(inputLocation, input);
|
||||||
|
}
|
||||||
|
await writeFile(scriptLocation, script);
|
||||||
|
const env = {
|
||||||
|
HOST_SOCKET: hostSocket,
|
||||||
|
SECRETS: JSON.stringify(secrets || {}),
|
||||||
|
INPUT_PATH: inputLocation,
|
||||||
|
HTTP_GATEWAY_PATH: httpGatewaySocket,
|
||||||
|
};
|
||||||
|
|
||||||
|
const teardown = async () => {
|
||||||
|
server.close();
|
||||||
|
await rm(dataDir, { recursive: true, force: true });
|
||||||
|
};
|
||||||
|
|
||||||
|
server.on('connection', (socket) => {
|
||||||
|
socket.on('data', (data) => {
|
||||||
|
const message = JSON.parse(data.toString());
|
||||||
|
emitter.emit('message', message);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
server.listen(hostSocket);
|
||||||
|
|
||||||
|
return {
|
||||||
|
env,
|
||||||
|
emitter,
|
||||||
|
teardown,
|
||||||
|
httpGatewaySocket,
|
||||||
|
scriptLocation,
|
||||||
|
hostSocket,
|
||||||
|
};
|
||||||
|
};
|
||||||
|
|
||||||
|
type Setup = Awaited<ReturnType<typeof setup>>;
|
||||||
|
|
||||||
|
export type { Setup };
|
||||||
|
export { setup };
|
||||||
@@ -27,6 +27,7 @@
|
|||||||
"typescript": "^5.3.3"
|
"typescript": "^5.3.3"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
|
"@fastify/reply-from": "^9.7.0",
|
||||||
"@trpc/client": "^10.45.0",
|
"@trpc/client": "^10.45.0",
|
||||||
"@trpc/server": "^10.45.0",
|
"@trpc/server": "^10.45.0",
|
||||||
"commander": "^11.1.0",
|
"commander": "^11.1.0",
|
||||||
|
|||||||
34
packages/server/src/gateway/gateway.ts
Normal file
34
packages/server/src/gateway/gateway.ts
Normal file
@@ -0,0 +1,34 @@
|
|||||||
|
import { FastifyPluginAsync } from 'fastify';
|
||||||
|
import FastifyReplyFrom from '@fastify/reply-from';
|
||||||
|
import { escape } from 'querystring';
|
||||||
|
import { Runtime } from '../runtime/runtime.js';
|
||||||
|
|
||||||
|
type Options = {
|
||||||
|
runtime: Runtime;
|
||||||
|
};
|
||||||
|
|
||||||
|
const gateway: FastifyPluginAsync<Options> = async (fastify, { runtime }) => {
|
||||||
|
await fastify.register(FastifyReplyFrom, {
|
||||||
|
http: {},
|
||||||
|
});
|
||||||
|
|
||||||
|
fastify.all('/gateway/*', (req, res) => {
|
||||||
|
const [runId, ...pathSegments] = (req.params as any)['*'].split('/').filter(Boolean);
|
||||||
|
const run = runtime.runner.getInstance(runId);
|
||||||
|
if (!run) {
|
||||||
|
res.statusCode = 404;
|
||||||
|
res.send({ error: 'Run not found' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const socketPath = run.run?.httpGatewaySocket;
|
||||||
|
if (!socketPath) {
|
||||||
|
res.statusCode = 404;
|
||||||
|
res.send({ error: 'No socket path to run' });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const path = pathSegments.join('/');
|
||||||
|
res.from(`unix+http://${escape(socketPath)}/${path}`);
|
||||||
|
});
|
||||||
|
};
|
||||||
|
|
||||||
|
export { gateway };
|
||||||
@@ -3,6 +3,7 @@ import { EventEmitter } from 'eventemitter3';
|
|||||||
import { Database } from '../../database/database.js';
|
import { Database } from '../../database/database.js';
|
||||||
import { CreateRunOptions, FindRunsOptions, UpdateRunOptions } from './runs.schemas.js';
|
import { CreateRunOptions, FindRunsOptions, UpdateRunOptions } from './runs.schemas.js';
|
||||||
import { LoadRepo } from '../loads/loads.js';
|
import { LoadRepo } from '../loads/loads.js';
|
||||||
|
import { createHash } from 'crypto';
|
||||||
|
|
||||||
type RunRepoEvents = {
|
type RunRepoEvents = {
|
||||||
created: (args: { id: string; loadId: string }) => void;
|
created: (args: { id: string; loadId: string }) => void;
|
||||||
@@ -18,13 +19,22 @@ type RunRepoOptions = {
|
|||||||
|
|
||||||
class RunRepo extends EventEmitter<RunRepoEvents> {
|
class RunRepo extends EventEmitter<RunRepoEvents> {
|
||||||
#options: RunRepoOptions;
|
#options: RunRepoOptions;
|
||||||
|
#isReady: Promise<void>;
|
||||||
|
|
||||||
constructor(options: RunRepoOptions) {
|
constructor(options: RunRepoOptions) {
|
||||||
super();
|
super();
|
||||||
this.#options = options;
|
this.#options = options;
|
||||||
|
this.#isReady = this.#setup();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#setup = async () => {
|
||||||
|
const { database } = this.#options;
|
||||||
|
const db = await database.instance;
|
||||||
|
await db('runs').update({ status: 'failed', error: 'server was shut down' }).where({ status: 'running' });
|
||||||
|
};
|
||||||
|
|
||||||
public getById = async (id: string) => {
|
public getById = async (id: string) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
|
|
||||||
@@ -36,6 +46,7 @@ class RunRepo extends EventEmitter<RunRepoEvents> {
|
|||||||
};
|
};
|
||||||
|
|
||||||
public getByLoadId = async (loadId: string) => {
|
public getByLoadId = async (loadId: string) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
|
|
||||||
@@ -44,6 +55,7 @@ class RunRepo extends EventEmitter<RunRepoEvents> {
|
|||||||
};
|
};
|
||||||
|
|
||||||
public find = async (options: FindRunsOptions) => {
|
public find = async (options: FindRunsOptions) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
const query = db('runs').select(['id', 'status', 'startedAt', 'status', 'error', 'endedAt']);
|
const query = db('runs').select(['id', 'status', 'startedAt', 'status', 'error', 'endedAt']);
|
||||||
@@ -62,19 +74,41 @@ class RunRepo extends EventEmitter<RunRepoEvents> {
|
|||||||
return runs;
|
return runs;
|
||||||
};
|
};
|
||||||
|
|
||||||
public remove = async (options: FindRunsOptions) => {
|
public prepareRemove = async (options: FindRunsOptions) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
const query = db('runs');
|
const query = db('runs').select('id');
|
||||||
|
|
||||||
if (options.loadId) {
|
if (options.loadId) {
|
||||||
query.where({ loadId: options.loadId });
|
query.where({ loadId: options.loadId });
|
||||||
}
|
}
|
||||||
|
|
||||||
await query.del();
|
const result = await query;
|
||||||
|
const ids = result.map((row) => row.id);
|
||||||
|
const token = ids.map((id) => Buffer.from(id).toString('base64')).join('|');
|
||||||
|
const hash = createHash('sha256').update(token).digest('hex');
|
||||||
|
return {
|
||||||
|
ids,
|
||||||
|
hash,
|
||||||
|
};
|
||||||
|
};
|
||||||
|
|
||||||
|
public remove = async (hash: string, ids: string[]) => {
|
||||||
|
const { database } = this.#options;
|
||||||
|
const db = await database.instance;
|
||||||
|
const token = ids.map((id) => Buffer.from(id).toString('base64')).join('|');
|
||||||
|
const actualHash = createHash('sha256').update(token).digest('hex');
|
||||||
|
|
||||||
|
if (hash !== actualHash) {
|
||||||
|
throw new Error('Invalid hash');
|
||||||
|
}
|
||||||
|
|
||||||
|
await db('runs').whereIn('id', ids).delete();
|
||||||
};
|
};
|
||||||
|
|
||||||
public started = async (id: string) => {
|
public started = async (id: string) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
const current = await this.getById(id);
|
const current = await this.getById(id);
|
||||||
@@ -92,6 +126,7 @@ class RunRepo extends EventEmitter<RunRepoEvents> {
|
|||||||
};
|
};
|
||||||
|
|
||||||
public finished = async (id: string, options: UpdateRunOptions) => {
|
public finished = async (id: string, options: UpdateRunOptions) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database } = this.#options;
|
const { database } = this.#options;
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
const { loadId } = await this.getById(id);
|
const { loadId } = await this.getById(id);
|
||||||
@@ -114,6 +149,7 @@ class RunRepo extends EventEmitter<RunRepoEvents> {
|
|||||||
};
|
};
|
||||||
|
|
||||||
public create = async (options: CreateRunOptions) => {
|
public create = async (options: CreateRunOptions) => {
|
||||||
|
await this.#isReady;
|
||||||
const { database, loads } = this.#options;
|
const { database, loads } = this.#options;
|
||||||
const id = nanoid();
|
const id = nanoid();
|
||||||
const db = await database.instance;
|
const db = await database.instance;
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
|
import { z } from 'zod';
|
||||||
import { createRunSchema, findRunsSchema } from '../repos/repos.js';
|
import { createRunSchema, findRunsSchema } from '../repos/repos.js';
|
||||||
import { publicProcedure, router } from './router.utils.js';
|
import { publicProcedure, router } from './router.utils.js';
|
||||||
|
|
||||||
@@ -17,17 +18,50 @@ const find = publicProcedure.input(findRunsSchema).query(async ({ input, ctx })
|
|||||||
return results;
|
return results;
|
||||||
});
|
});
|
||||||
|
|
||||||
const remove = publicProcedure.input(findRunsSchema).mutation(async ({ input, ctx }) => {
|
const prepareRemove = publicProcedure.input(findRunsSchema).query(async ({ input, ctx }) => {
|
||||||
const { runtime } = ctx;
|
const { runtime } = ctx;
|
||||||
const { repos } = runtime;
|
const { repos } = runtime;
|
||||||
const { runs } = repos;
|
const { runs } = repos;
|
||||||
await runs.remove(input);
|
return await runs.prepareRemove(input);
|
||||||
|
});
|
||||||
|
|
||||||
|
const remove = publicProcedure
|
||||||
|
|
||||||
|
.input(
|
||||||
|
z.object({
|
||||||
|
hash: z.string(),
|
||||||
|
ids: z.array(z.string()),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
.mutation(async ({ input, ctx }) => {
|
||||||
|
const { runtime } = ctx;
|
||||||
|
const { repos } = runtime;
|
||||||
|
const { runs } = repos;
|
||||||
|
for (const id of input.ids) {
|
||||||
|
const instance = runtime.runner.getInstance(id);
|
||||||
|
if (instance) {
|
||||||
|
await instance.run?.teardown();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
await runs.remove(input.hash, input.ids);
|
||||||
|
});
|
||||||
|
|
||||||
|
const terminate = publicProcedure.input(z.string()).mutation(async ({ input, ctx }) => {
|
||||||
|
const { runtime } = ctx;
|
||||||
|
const { runner } = runtime;
|
||||||
|
const instance = runner.getInstance(input);
|
||||||
|
if (!instance || !instance.run) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
await instance.run.teardown();
|
||||||
});
|
});
|
||||||
|
|
||||||
const runsRouter = router({
|
const runsRouter = router({
|
||||||
create,
|
create,
|
||||||
find,
|
find,
|
||||||
remove,
|
remove,
|
||||||
|
prepareRemove,
|
||||||
|
terminate,
|
||||||
});
|
});
|
||||||
|
|
||||||
export { runsRouter };
|
export { runsRouter };
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
import { EventEmitter } from 'eventemitter3';
|
import { EventEmitter } from 'eventemitter3';
|
||||||
import { run } from '@morten-olsen/mini-loader-runner';
|
import { RunInfo, run } from '@morten-olsen/mini-loader-runner';
|
||||||
import { Repos } from '../repos/repos.js';
|
import { Repos } from '../repos/repos.js';
|
||||||
import { LoggerEvent } from '../../../mini-loader/dist/esm/logger/logger.js';
|
import { LoggerEvent } from '../../../mini-loader/dist/esm/logger/logger.js';
|
||||||
import { ArtifactCreateEvent } from '../../../mini-loader/dist/esm/artifacts/artifacts.js';
|
import { ArtifactCreateEvent } from '../../../mini-loader/dist/esm/artifacts/artifacts.js';
|
||||||
@@ -20,12 +20,17 @@ type RunnerInstanceOptions = {
|
|||||||
|
|
||||||
class RunnerInstance extends EventEmitter<RunnerInstanceEvents> {
|
class RunnerInstance extends EventEmitter<RunnerInstanceEvents> {
|
||||||
#options: RunnerInstanceOptions;
|
#options: RunnerInstanceOptions;
|
||||||
|
#run?: RunInfo;
|
||||||
|
|
||||||
constructor(options: RunnerInstanceOptions) {
|
constructor(options: RunnerInstanceOptions) {
|
||||||
super();
|
super();
|
||||||
this.#options = options;
|
this.#options = options;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public get run() {
|
||||||
|
return this.#run;
|
||||||
|
}
|
||||||
|
|
||||||
#addLog = async (event: LoggerEvent['payload']) => {
|
#addLog = async (event: LoggerEvent['payload']) => {
|
||||||
const { repos, id, loadId } = this.#options;
|
const { repos, id, loadId } = this.#options;
|
||||||
const { logs } = repos;
|
const { logs } = repos;
|
||||||
@@ -58,11 +63,13 @@ class RunnerInstance extends EventEmitter<RunnerInstanceEvents> {
|
|||||||
const script = await readFile(scriptLocation, 'utf-8');
|
const script = await readFile(scriptLocation, 'utf-8');
|
||||||
const allSecrets = await secrets.getAll();
|
const allSecrets = await secrets.getAll();
|
||||||
await runs.started(id);
|
await runs.started(id);
|
||||||
const { promise, emitter } = await run({
|
const current = await run({
|
||||||
script,
|
script,
|
||||||
secrets: allSecrets,
|
secrets: allSecrets,
|
||||||
input,
|
input,
|
||||||
});
|
});
|
||||||
|
this.#run = current;
|
||||||
|
const { promise, emitter } = current;
|
||||||
emitter.on('message', (message) => {
|
emitter.on('message', (message) => {
|
||||||
switch (message.type) {
|
switch (message.type) {
|
||||||
case 'log': {
|
case 'log': {
|
||||||
@@ -84,9 +91,11 @@ class RunnerInstance extends EventEmitter<RunnerInstanceEvents> {
|
|||||||
}
|
}
|
||||||
await runs.finished(id, { status: 'failed', error: errorMessage });
|
await runs.finished(id, { status: 'failed', error: errorMessage });
|
||||||
} finally {
|
} finally {
|
||||||
|
this.#run = undefined;
|
||||||
this.emit('completed', { id });
|
this.emit('completed', { id });
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type { RunInfo };
|
||||||
export { RunnerInstance };
|
export { RunnerInstance };
|
||||||
|
|||||||
@@ -36,6 +36,10 @@ class Runner {
|
|||||||
this.#instances.set(args.id, instance);
|
this.#instances.set(args.id, instance);
|
||||||
await instance.start();
|
await instance.start();
|
||||||
};
|
};
|
||||||
|
|
||||||
|
public getInstance = (id: string) => {
|
||||||
|
return this.#instances.get(id);
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
export { Runner };
|
export { Runner };
|
||||||
|
|||||||
@@ -3,9 +3,16 @@ import fastify from 'fastify';
|
|||||||
import { RootRouter, rootRouter } from '../router/router.js';
|
import { RootRouter, rootRouter } from '../router/router.js';
|
||||||
import { createContext } from '../router/router.utils.js';
|
import { createContext } from '../router/router.utils.js';
|
||||||
import { Runtime } from '../runtime/runtime.js';
|
import { Runtime } from '../runtime/runtime.js';
|
||||||
|
import { gateway } from '../gateway/gateway.js';
|
||||||
|
|
||||||
const createServer = async (runtime: Runtime) => {
|
const createServer = async (runtime: Runtime) => {
|
||||||
const server = fastify({});
|
const server = fastify({
|
||||||
|
maxParamLength: 10000,
|
||||||
|
bodyLimit: 30 * 1024 * 1024,
|
||||||
|
logger: {
|
||||||
|
level: 'warn',
|
||||||
|
},
|
||||||
|
});
|
||||||
server.get('/', async () => {
|
server.get('/', async () => {
|
||||||
return { hello: 'world' };
|
return { hello: 'world' };
|
||||||
});
|
});
|
||||||
@@ -33,6 +40,14 @@ const createServer = async (runtime: Runtime) => {
|
|||||||
},
|
},
|
||||||
} satisfies FastifyTRPCPluginOptions<RootRouter>['trpcOptions'],
|
} satisfies FastifyTRPCPluginOptions<RootRouter>['trpcOptions'],
|
||||||
});
|
});
|
||||||
|
|
||||||
|
server.register(gateway, {
|
||||||
|
runtime,
|
||||||
|
});
|
||||||
|
|
||||||
|
server.addHook('onError', async (request, reply, error) => {
|
||||||
|
console.error(error);
|
||||||
|
});
|
||||||
await server.ready();
|
await server.ready();
|
||||||
|
|
||||||
return server;
|
return server;
|
||||||
|
|||||||
47
pnpm-lock.yaml
generated
47
pnpm-lock.yaml
generated
@@ -101,6 +101,10 @@ importers:
|
|||||||
packages/configs: {}
|
packages/configs: {}
|
||||||
|
|
||||||
packages/examples:
|
packages/examples:
|
||||||
|
dependencies:
|
||||||
|
fastify:
|
||||||
|
specifier: ^4.25.2
|
||||||
|
version: 4.25.2
|
||||||
devDependencies:
|
devDependencies:
|
||||||
'@morten-olsen/mini-loader':
|
'@morten-olsen/mini-loader':
|
||||||
specifier: workspace:^
|
specifier: workspace:^
|
||||||
@@ -132,6 +136,9 @@ importers:
|
|||||||
|
|
||||||
packages/runner:
|
packages/runner:
|
||||||
dependencies:
|
dependencies:
|
||||||
|
'@morten-olsen/mini-loader':
|
||||||
|
specifier: workspace:^
|
||||||
|
version: link:../mini-loader
|
||||||
eventemitter3:
|
eventemitter3:
|
||||||
specifier: ^5.0.1
|
specifier: ^5.0.1
|
||||||
version: 5.0.1
|
version: 5.0.1
|
||||||
@@ -139,9 +146,6 @@ importers:
|
|||||||
specifier: ^5.0.4
|
specifier: ^5.0.4
|
||||||
version: 5.0.4
|
version: 5.0.4
|
||||||
devDependencies:
|
devDependencies:
|
||||||
'@morten-olsen/mini-loader':
|
|
||||||
specifier: workspace:^
|
|
||||||
version: link:../mini-loader
|
|
||||||
'@morten-olsen/mini-loader-configs':
|
'@morten-olsen/mini-loader-configs':
|
||||||
specifier: workspace:^
|
specifier: workspace:^
|
||||||
version: link:../configs
|
version: link:../configs
|
||||||
@@ -154,6 +158,9 @@ importers:
|
|||||||
|
|
||||||
packages/server:
|
packages/server:
|
||||||
dependencies:
|
dependencies:
|
||||||
|
'@fastify/reply-from':
|
||||||
|
specifier: ^9.7.0
|
||||||
|
version: 9.7.0
|
||||||
'@trpc/client':
|
'@trpc/client':
|
||||||
specifier: ^10.45.0
|
specifier: ^10.45.0
|
||||||
version: 10.45.0(@trpc/server@10.45.0)
|
version: 10.45.0(@trpc/server@10.45.0)
|
||||||
@@ -476,6 +483,11 @@ packages:
|
|||||||
fast-uri: 2.3.0
|
fast-uri: 2.3.0
|
||||||
dev: false
|
dev: false
|
||||||
|
|
||||||
|
/@fastify/busboy@2.1.0:
|
||||||
|
resolution: {integrity: sha512-+KpH+QxZU7O4675t3mnkQKcZZg56u+K/Ct2K+N2AZYNVK8kyeo/bI18tI8aPm3tvNNRyTWfj6s5tnGNlcbQRsA==}
|
||||||
|
engines: {node: '>=14'}
|
||||||
|
dev: false
|
||||||
|
|
||||||
/@fastify/deepmerge@1.3.0:
|
/@fastify/deepmerge@1.3.0:
|
||||||
resolution: {integrity: sha512-J8TOSBq3SoZbDhM9+R/u77hP93gz/rajSA+K2kGyijPpORPWUXHUpTaleoj+92As0S9uPRP7Oi8IqMf0u+ro6A==}
|
resolution: {integrity: sha512-J8TOSBq3SoZbDhM9+R/u77hP93gz/rajSA+K2kGyijPpORPWUXHUpTaleoj+92As0S9uPRP7Oi8IqMf0u+ro6A==}
|
||||||
dev: false
|
dev: false
|
||||||
@@ -490,6 +502,19 @@ packages:
|
|||||||
fast-json-stringify: 5.10.0
|
fast-json-stringify: 5.10.0
|
||||||
dev: false
|
dev: false
|
||||||
|
|
||||||
|
/@fastify/reply-from@9.7.0:
|
||||||
|
resolution: {integrity: sha512-/F1QBl3FGlTqStjmiuoLRDchVxP967TZh6FZPwQteWhdLsDec8mqSACE+cRzw6qHUj3v9hfdd7JNgmb++fyFhQ==}
|
||||||
|
dependencies:
|
||||||
|
'@fastify/error': 3.4.1
|
||||||
|
end-of-stream: 1.4.4
|
||||||
|
fast-content-type-parse: 1.1.0
|
||||||
|
fast-querystring: 1.1.2
|
||||||
|
fastify-plugin: 4.5.1
|
||||||
|
pump: 3.0.0
|
||||||
|
tiny-lru: 11.2.5
|
||||||
|
undici: 5.28.2
|
||||||
|
dev: false
|
||||||
|
|
||||||
/@gar/promisify@1.1.3:
|
/@gar/promisify@1.1.3:
|
||||||
resolution: {integrity: sha512-k2Ty1JcVojjJFwrg/ThKi2ujJ7XNLYaFGNB/bWT9wGR+oSMJHMa5w+CUq6p/pVrKeNNgA7pCqEcjSnHVoqJQFw==}
|
resolution: {integrity: sha512-k2Ty1JcVojjJFwrg/ThKi2ujJ7XNLYaFGNB/bWT9wGR+oSMJHMa5w+CUq6p/pVrKeNNgA7pCqEcjSnHVoqJQFw==}
|
||||||
requiresBuild: true
|
requiresBuild: true
|
||||||
@@ -2683,6 +2708,10 @@ packages:
|
|||||||
resolution: {integrity: sha512-eel5UKGn369gGEWOqBShmFJWfq/xSJvsgDzgLYC845GneayWvXBf0lJCBn5qTABfewy1ZDPoaR5OZCP+kssfuw==}
|
resolution: {integrity: sha512-eel5UKGn369gGEWOqBShmFJWfq/xSJvsgDzgLYC845GneayWvXBf0lJCBn5qTABfewy1ZDPoaR5OZCP+kssfuw==}
|
||||||
dev: false
|
dev: false
|
||||||
|
|
||||||
|
/fastify-plugin@4.5.1:
|
||||||
|
resolution: {integrity: sha512-stRHYGeuqpEZTL1Ef0Ovr2ltazUT9g844X5z/zEBFLG8RYlpDiOCIG+ATvYEp+/zmc7sN29mcIMp8gvYplYPIQ==}
|
||||||
|
dev: false
|
||||||
|
|
||||||
/fastify@4.25.2:
|
/fastify@4.25.2:
|
||||||
resolution: {integrity: sha512-SywRouGleDHvRh054onj+lEZnbC1sBCLkR0UY3oyJwjD4BdZJUrxBqfkfCaqn74pVCwBaRHGuL3nEWeHbHzAfw==}
|
resolution: {integrity: sha512-SywRouGleDHvRh054onj+lEZnbC1sBCLkR0UY3oyJwjD4BdZJUrxBqfkfCaqn74pVCwBaRHGuL3nEWeHbHzAfw==}
|
||||||
dependencies:
|
dependencies:
|
||||||
@@ -5151,6 +5180,11 @@ packages:
|
|||||||
engines: {node: '>=8'}
|
engines: {node: '>=8'}
|
||||||
dev: false
|
dev: false
|
||||||
|
|
||||||
|
/tiny-lru@11.2.5:
|
||||||
|
resolution: {integrity: sha512-JpqM0K33lG6iQGKiigcwuURAKZlq6rHXfrgeL4/I8/REoyJTGU+tEMszvT/oTRVHG2OiylhGDjqPp1jWMlr3bw==}
|
||||||
|
engines: {node: '>=12'}
|
||||||
|
dev: false
|
||||||
|
|
||||||
/tmp@0.0.33:
|
/tmp@0.0.33:
|
||||||
resolution: {integrity: sha512-jRCJlojKnZ3addtTOjdIqoRuPEKBvNXcGYqzO6zWZX8KfKEpnGY5jfggJQ3EjKuu8D4bJRr0y+cYJFmYbImXGw==}
|
resolution: {integrity: sha512-jRCJlojKnZ3addtTOjdIqoRuPEKBvNXcGYqzO6zWZX8KfKEpnGY5jfggJQ3EjKuu8D4bJRr0y+cYJFmYbImXGw==}
|
||||||
engines: {node: '>=0.6.0'}
|
engines: {node: '>=0.6.0'}
|
||||||
@@ -5368,6 +5402,13 @@ packages:
|
|||||||
/undici-types@5.26.5:
|
/undici-types@5.26.5:
|
||||||
resolution: {integrity: sha512-JlCMO+ehdEIKqlFxk6IfVoAUVmgz7cU7zD/h9XZ0qzeosSHmUJVOzSQvvYSYWXkFXC+IfLKSIffhv0sVZup6pA==}
|
resolution: {integrity: sha512-JlCMO+ehdEIKqlFxk6IfVoAUVmgz7cU7zD/h9XZ0qzeosSHmUJVOzSQvvYSYWXkFXC+IfLKSIffhv0sVZup6pA==}
|
||||||
|
|
||||||
|
/undici@5.28.2:
|
||||||
|
resolution: {integrity: sha512-wh1pHJHnUeQV5Xa8/kyQhO7WFa8M34l026L5P/+2TYiakvGy5Rdc8jWZVyG7ieht/0WgJLEd3kcU5gKx+6GC8w==}
|
||||||
|
engines: {node: '>=14.0'}
|
||||||
|
dependencies:
|
||||||
|
'@fastify/busboy': 2.1.0
|
||||||
|
dev: false
|
||||||
|
|
||||||
/unique-filename@1.1.1:
|
/unique-filename@1.1.1:
|
||||||
resolution: {integrity: sha512-Vmp0jIp2ln35UTXuryvjzkjGdRyf9b2lTXuSYUiPmzRcl3FDtYqAwOnTJkAngD9SWhnoJzDbTKwaOrZ+STtxNQ==}
|
resolution: {integrity: sha512-Vmp0jIp2ln35UTXuryvjzkjGdRyf9b2lTXuSYUiPmzRcl3FDtYqAwOnTJkAngD9SWhnoJzDbTKwaOrZ+STtxNQ==}
|
||||||
requiresBuild: true
|
requiresBuild: true
|
||||||
|
|||||||
@@ -1,15 +1,15 @@
|
|||||||
{
|
{
|
||||||
"include": [],
|
"include": [],
|
||||||
"references": [
|
"references": [
|
||||||
|
{
|
||||||
|
"path": "./packages/mini-loader/tsconfig.json"
|
||||||
|
},
|
||||||
{
|
{
|
||||||
"path": "./packages/runner/tsconfig.json"
|
"path": "./packages/runner/tsconfig.json"
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"path": "./packages/server/tsconfig.json"
|
"path": "./packages/server/tsconfig.json"
|
||||||
},
|
},
|
||||||
{
|
|
||||||
"path": "./packages/mini-loader/tsconfig.json"
|
|
||||||
},
|
|
||||||
{
|
{
|
||||||
"path": "./packages/cli/tsconfig.json"
|
"path": "./packages/cli/tsconfig.json"
|
||||||
},
|
},
|
||||||
|
|||||||
Reference in New Issue
Block a user