Skip to content

Commit b923c17

Browse files
authored
Merge pull request #4 from sroussey/job-queue-task-worker
Job queue task worker
2 parents 50b7b57 + 3129e4a commit b923c17

55 files changed

Lines changed: 2517 additions & 316 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

bun.lockb

3.8 KB
Binary file not shown.

package.json

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,15 +17,18 @@
1717
},
1818
"dependencies": {
1919
"@sroussey/typescript-graph": "^0.3.12",
20+
"@types/better-sqlite3": "^7.6.9",
2021
"@types/pg": "^8.11.2",
2122
"@xyflow/react": "12.0.0-next.11",
23+
"better-sqlite3": "^9.4.3",
2224
"chalk": "^5.3.0",
2325
"commander": "^11.1.0",
2426
"eventemitter3": "^5.0.1",
2527
"listr2": "^8.0.2",
2628
"nanoid": "^5.0.6",
2729
"pg": "^8.11.3",
2830
"postcss": "^8.4.35",
31+
"postgres": "^3.4.3",
2932
"react-hotkeys-hook": "^4.5.0",
3033
"react-icons": "^5.0.1",
3134
"rxjs": "^7.8.1",
@@ -40,11 +43,11 @@
4043
"react-dom": "^18.2.0",
4144
"tailwindcss": "^3.4.1",
4245
"typescript": "^5.4.2",
43-
"vite": "^5.1.5"
46+
"vite": "^5.1.6"
4447
},
4548
"peerDependencies": {
46-
"@mediapipe/tasks-text": "^0.10.9",
47-
"@sroussey/transformers": "^2.15.1"
49+
"@mediapipe/tasks-text": "^0.10.12",
50+
"@sroussey/transformers": "^2.16.0"
4851
},
4952
"engines": {
5053
"bun": "^1.0.5"

packages/cli/package.json

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -4,18 +4,18 @@
44
"version": "0.0.1",
55
"description": "Ellmers is a tool for building and running DAG pipelines of AI tasks.",
66
"scripts": {
7-
"watch": "bunx concurrently 'bun run watch-types' 'bun run watch-js'",
8-
"watch-js": "bun build --watch --target=node --sourcemap=external --external listr2 --external @sroussey/transformers --outdir ./dist ./src/lib.ts",
9-
"watch-types": "tsc --watch",
7+
"watch": "bunx concurrently -c 'auto' -n 'cli:' 'npm:watch-*'",
8+
"watch-js": "bun build --watch --target=node --sourcemap=external --external listr2 --external @sroussey/transformers --outdir ./dist ./src/lib.ts ./src/ellmers.ts",
9+
"watch-types": "tsc --watch --preserveWatchOutput",
1010
"build": "bun run build-clean && bun run build-types && bun run build-js && bun run build-types-map",
1111
"build-clean": "rm -fr dist/* tsconfig.tsbuildinfo",
12-
"build-js": "bun build --target=node --sourcemap=external --external listr2 --external @sroussey/transformers --outdir ./dist ./src/lib.ts",
12+
"build-js": "bun build --target=node --sourcemap=external --external listr2 --external @sroussey/transformers --outdir ./dist ./src/lib.ts ./src/ellmers.ts",
1313
"build-types": "tsc",
1414
"build-types-map": "tsc --declarationMap",
1515
"test": "echo \"Error: no test specified\" && exit 1"
1616
},
17-
"bin": "src/elmers.js",
18-
"module": "dist/lib.js",
17+
"bin": "./dist/elmers.js",
18+
"module": "./dist/lib.js",
1919
"exports": {
2020
".": {
2121
"import": "./dist/lib.js"
@@ -25,6 +25,6 @@
2525
"dist"
2626
],
2727
"dependencies": {
28-
"ellmers-core": "workspace:*"
28+
"ellmers-core": "workspace:packages/core"
2929
}
3030
}

packages/cli/src/TaskCLI.ts

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -19,18 +19,13 @@ import {
1919
DownloadTask,
2020
ModelUseCaseEnum,
2121
EmbeddingMultiModelTask,
22-
registerHuggingfaceLocalTasks,
23-
registerMediaPipeTfJsLocalTasks,
2422
DownloadMultiModelTask,
2523
TextRewriterMultiModelTask,
2624
TaskGraph,
2725
JsonTaskArray,
2826
JsonTask,
2927
} from "ellmers-core/server";
3028

31-
registerHuggingfaceLocalTasks();
32-
registerMediaPipeTfJsLocalTasks();
33-
3429
export function AddBaseCommands(program: Command) {
3530
program
3631
.command("download")

packages/cli/src/ellmers.ts

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,19 @@
33
import { program } from "commander";
44
import { argv } from "process";
55
import { AddBaseCommands } from "./TaskCLI";
6+
import {
7+
getProviderRegistry,
8+
registerHuggingfaceLocalTasksInMemory,
9+
registerMediaPipeTfJsLocalInMemory,
10+
} from "ellmers-core/server";
611

712
program.version("1.0.0").description("A CLI to run Ellmers.");
813

914
AddBaseCommands(program);
1015

16+
registerHuggingfaceLocalTasksInMemory();
17+
registerMediaPipeTfJsLocalInMemory();
18+
1119
await program.parseAsync(argv);
20+
21+
getProviderRegistry().stopQueues();

packages/core/package.json

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -4,17 +4,17 @@
44
"version": "0.0.1",
55
"description": "Ellmers is a tool for building and running DAG pipelines of AI tasks.",
66
"scripts": {
7-
"watch": "bunx concurrently 'bun run watch-types' 'bun run watch-browser' 'bun run watch-server-bun'",
8-
"watch-browser": "bun build --watch --target=browser --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/browser.ts",
9-
"watch-server-bun": "bun build --watch --target=bun --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/server.ts",
10-
"watch-types": "tsc --watch",
7+
"watch": "bunx concurrently -c 'auto' -n 'core:' 'npm:watch-*'",
8+
"watch-browser": "bun build --watch --target=browser --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/browser*.ts",
9+
"watch-server-bun": "bun build --watch --target=bun --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/server*.ts",
10+
"watch-types": "tsc --watch --preserveWatchOutput",
1111
"build": "bun run build-clean && bun run build-types && bun run build-browser && bun run build-server-bun && bun run build-types-map",
1212
"build-clean": "rm -fr dist/* tsconfig.tsbuildinfo",
13-
"build-browser": "bun build --target=browser --minify-whitespace --minify-syntax --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/browser.ts",
14-
"build-server-bun": "bun build --target=bun --minify-whitespace --minify-syntax --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/server.ts",
13+
"build-browser": "bun build --target=browser --minify-whitespace --minify-syntax --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/browser*.ts",
14+
"build-server-bun": "bun build --target=bun --minify-whitespace --minify-syntax --sourcemap=external --external @sroussey/transformers --outdir ./dist ./src/server*.ts",
1515
"build-types": "tsc",
1616
"build-types-map": "tsc --declarationMap",
17-
"test": "echo \"Error: no test specified\" && exit 1"
17+
"test": "bun test"
1818
},
1919
"module": "dist/server.js",
2020
"exports": {
Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
import { getProviderRegistry } from "../provider/ProviderRegistry";
2+
import { registerHuggingfaceLocalTasks } from "./local_hf";
3+
import { registerMediaPipeTfJsLocalTasks } from "./local_mp";
4+
import { InMemoryJobQueue } from "../job/InMemoryJobQueue";
5+
import { ModelProcessorEnum } from "../model/Model";
6+
import { ConcurrencyLimiter } from "../job/ConcurrencyLimiter";
7+
import { TaskInput, TaskOutput } from "../task/base/Task";
8+
9+
export async function registerHuggingfaceLocalTasksInMemory() {
10+
registerHuggingfaceLocalTasks();
11+
const ProviderRegistry = getProviderRegistry();
12+
const jobQueue = new InMemoryJobQueue<TaskInput, TaskOutput>(
13+
"local_hf",
14+
new ConcurrencyLimiter(1, 10),
15+
10
16+
);
17+
ProviderRegistry.registerQueue(ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS, jobQueue);
18+
jobQueue.start();
19+
}
20+
21+
export async function registerMediaPipeTfJsLocalInMemory() {
22+
registerMediaPipeTfJsLocalTasks();
23+
const ProviderRegistry = getProviderRegistry();
24+
const jobQueue = new InMemoryJobQueue<TaskInput, TaskOutput>(
25+
"local_media_pipe",
26+
new ConcurrencyLimiter(1, 10),
27+
10
28+
);
29+
ProviderRegistry.registerQueue(ModelProcessorEnum.MEDIA_PIPE_TFJS_MODEL, jobQueue);
30+
jobQueue.start();
31+
}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
import { registerHuggingfaceLocalTasks } from "./local_hf";
2+
import { registerMediaPipeTfJsLocalTasks } from "./local_mp";
3+
import { getProviderRegistry } from "../provider/ProviderRegistry";
4+
import { ModelProcessorEnum } from "../model/Model";
5+
import { ConcurrencyLimiter } from "../job/ConcurrencyLimiter";
6+
import { SqliteJobQueue } from "../job/SqliteJobQueue";
7+
import { getDatabase } from "../util/db_sqlite";
8+
import { TaskInput, TaskOutput } from "../task/base/Task";
9+
10+
const db = getDatabase("local.db");
11+
12+
export async function registerHuggingfaceLocalTasksSqlite() {
13+
registerHuggingfaceLocalTasks();
14+
const ProviderRegistry = getProviderRegistry();
15+
const jobQueue = new SqliteJobQueue<TaskInput, TaskOutput>(
16+
db,
17+
"local_hf",
18+
new ConcurrencyLimiter(1, 10)
19+
);
20+
ProviderRegistry.registerQueue(ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS, jobQueue);
21+
jobQueue.start();
22+
}
23+
24+
export async function registerMediaPipeTfJsLocalSqlite() {
25+
registerMediaPipeTfJsLocalTasks();
26+
const ProviderRegistry = getProviderRegistry();
27+
const jobQueue = new SqliteJobQueue<TaskInput, TaskOutput>(
28+
db,
29+
"local_media_pipe",
30+
new ConcurrencyLimiter(1, 10)
31+
);
32+
ProviderRegistry.registerQueue(ModelProcessorEnum.MEDIA_PIPE_TFJS_MODEL, jobQueue);
33+
jobQueue.start();
34+
}
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
import {
2+
HuggingFaceLocal_DownloadRun,
3+
HuggingFaceLocal_EmbeddingRun,
4+
HuggingFaceLocal_TextGenerationRun,
5+
HuggingFaceLocal_TextQuestionAnswerRun,
6+
HuggingFaceLocal_TextRewriterRun,
7+
HuggingFaceLocal_TextSummaryRun,
8+
} from "provider/local-hugging-face/HuggingFaceLocal_TaskRun";
9+
import { ModelProcessorEnum } from "../model/Model";
10+
import { getProviderRegistry } from "../provider/ProviderRegistry";
11+
import {
12+
DownloadTask,
13+
EmbeddingTask,
14+
TextGenerationTask,
15+
TextQuestionAnswerTask,
16+
TextRewriterTask,
17+
TextSummaryTask,
18+
} from "task";
19+
20+
export async function registerHuggingfaceLocalTasks() {
21+
const ProviderRegistry = getProviderRegistry();
22+
23+
ProviderRegistry.registerRunFn(
24+
DownloadTask.type,
25+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
26+
HuggingFaceLocal_DownloadRun
27+
);
28+
29+
ProviderRegistry.registerRunFn(
30+
EmbeddingTask.type,
31+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
32+
HuggingFaceLocal_EmbeddingRun
33+
);
34+
35+
ProviderRegistry.registerRunFn(
36+
TextGenerationTask.type,
37+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
38+
HuggingFaceLocal_TextGenerationRun
39+
);
40+
41+
ProviderRegistry.registerRunFn(
42+
TextRewriterTask.type,
43+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
44+
HuggingFaceLocal_TextRewriterRun
45+
);
46+
47+
ProviderRegistry.registerRunFn(
48+
TextSummaryTask.type,
49+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
50+
HuggingFaceLocal_TextSummaryRun
51+
);
52+
53+
ProviderRegistry.registerRunFn(
54+
TextQuestionAnswerTask.type,
55+
ModelProcessorEnum.LOCAL_ONNX_TRANSFORMERJS,
56+
HuggingFaceLocal_TextQuestionAnswerRun
57+
);
58+
}
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
import {
2+
MediaPipeTfJsLocal_Download,
3+
MediaPipeTfJsLocal_Embedding,
4+
} from "../provider/local-media-pipe/MediaPipeLocalTaskRun";
5+
import { ModelProcessorEnum } from "../model/Model";
6+
import { getProviderRegistry } from "../provider/ProviderRegistry";
7+
import { DownloadTask, EmbeddingTask } from "task";
8+
9+
export const registerMediaPipeTfJsLocalTasks = () => {
10+
const ProviderRegistry = getProviderRegistry();
11+
12+
ProviderRegistry.registerRunFn(
13+
DownloadTask.type,
14+
ModelProcessorEnum.MEDIA_PIPE_TFJS_MODEL,
15+
MediaPipeTfJsLocal_Download
16+
);
17+
18+
ProviderRegistry.registerRunFn(
19+
EmbeddingTask.type,
20+
ModelProcessorEnum.MEDIA_PIPE_TFJS_MODEL,
21+
MediaPipeTfJsLocal_Embedding
22+
);
23+
};

0 commit comments

Comments
 (0)