Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
9a46710
feat(workflow): distributed execution with multi-tenant isolation and…
kojiwakayama Feb 7, 2026
a84b3fd
chore: remove unused files
kojiwakayama Feb 8, 2026
bdf7058
refactor: enforce internal module boundaries with #veryfront/ import …
kojiwakayama Feb 8, 2026
ccc876b
chore: refactor
kojiwakayama Feb 9, 2026
0098a32
refactor: align data resolution and harden module fetch/runtime paths
kojiwakayama Feb 9, 2026
670370b
refactor: split large modules into focused, cohesive files
kojiwakayama Feb 9, 2026
a007e9f
refactor: extract SSR vf-module resolver and tmp path helpers
kojiwakayama Feb 9, 2026
bae6fd9
refactor: extract read in-flight dedupe policy
kojiwakayama Feb 9, 2026
fad6495
refactor: extract SSR preflight import filtering
kojiwakayama Feb 9, 2026
f6f72eb
refactor: split read operation fetch pipeline and fallback tests
kojiwakayama Feb 9, 2026
62dd2a0
refactor: extract extensionless path resolution in read ops
kojiwakayama Feb 9, 2026
1be355f
refactor: modularize cache layer, tighten barrels, remove dead code
kojiwakayama Feb 9, 2026
3582784
refactor: extract stat index and resolution helpers
kojiwakayama Feb 9, 2026
83f215b
refactor: isolate stat API search circuit breaker
kojiwakayama Feb 9, 2026
0b1f487
refactor: isolate SSR cache import recovery path
kojiwakayama Feb 9, 2026
78ca51d
refactor: extract SSR cross-project import flow
kojiwakayama Feb 9, 2026
04d34e3
refactor: extract veryfront adapter content-context helpers
kojiwakayama Feb 9, 2026
8d31764
refactor: fix minor issues from parallel agent work
kojiwakayama Feb 9, 2026
77ba4b3
Merge remote-tracking branch 'origin/main' into chore/refactor
kojiwakayama Feb 9, 2026
dfa435f
chore: fix test
kojiwakayama Feb 9, 2026
e7beacc
Merge branch 'main' into chore/refactor
kojiwakayama Feb 9, 2026
3d42dd1
fix: address PR review comments
kojiwakayama Feb 9, 2026
4eeac27
fix: remove unused TIMEOUT_ERROR import and dead isLocalProject variable
kojiwakayama Feb 9, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
1 change: 1 addition & 0 deletions cli/commands/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,3 +19,4 @@ export { startCommand } from "./start/index.ts";
export { serveCommand } from "./serve/index.ts";
export { doctorCommand } from "./doctor/index.ts";
export { installCommand, uninstallCommand } from "./install/index.ts";
export { workerCommand } from "./worker/index.ts";
48 changes: 48 additions & 0 deletions cli/commands/worker/command-help.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
import type { CommandHelp } from "../../help/types.ts";

export const workerHelp: CommandHelp = {
name: "worker",
description: "Start workflow job worker",
usage: "veryfront worker [options]",
options: [
{
flag: "--redis-url <url>",
description: "Redis connection URL",
default: "redis://localhost:6379",
},
{
flag: "-c, --concurrency <number>",
description: "Maximum concurrent jobs",
default: "3",
},
{
flag: "--poll-interval <ms>",
description: "Poll interval in milliseconds",
default: "5000",
},
{
flag: "--stalled-threshold <ms>",
description: "Time before a run is considered stalled",
default: "60000",
},
{
flag: "-e, --executor <type>",
description: "Job executor type (process | k8s)",
default: "process",
},
{
flag: "--entrypoint <path>",
description: "Path to job entrypoint script",
default: "./workflow-job.ts",
},
{
flag: "--debug",
description: "Enable debug logging",
},
],
examples: [
"veryfront worker",
"veryfront worker --redis-url redis://prod:6379 --concurrency 5",
"veryfront worker --entrypoint ./src/jobs/workflow-runner.ts --debug",
],
};
114 changes: 114 additions & 0 deletions cli/commands/worker/command.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
/**
* Worker command - Start workflow job manager
*
* Polls Redis for pending/stalled workflow runs and executes them
* as isolated processes. Supports multi-tenant execution: each job
* runs with its own tenant context captured at workflow creation time.
*/

import { cliLogger } from "#cli/utils";
import { exitProcess, registerTerminationSignals, showLogo } from "#cli/utils";
import type { WorkerArgs } from "./handler.ts";

export interface WorkerOptions extends WorkerArgs {}

export async function workerCommand(options: WorkerOptions): Promise<void> {
showLogo();

const { WorkflowJobManager } = await import(
"../../../src/workflow/worker/job-manager.ts"
);
const { ProcessJobExecutor } = await import(
"../../../src/workflow/worker/executors/process.ts"
);
const { RedisBackend } = await import(
"../../../src/workflow/backends/redis.ts"
);

cliLogger.info("Starting workflow worker...");
cliLogger.info(` Redis: ${options.redisUrl}`);
cliLogger.info(` Executor: ${options.executor}`);
cliLogger.info(` Concurrency: ${options.concurrency}`);
cliLogger.info(` Poll: ${options.pollInterval}ms`);

// Initialize Redis backend
const backend = new RedisBackend({
url: options.redisUrl,
debug: options.debug,
});

if (backend.initialize) {
await backend.initialize();
}

// Create job executor
// The entrypoint script runs inside each spawned process.
// It reads WORKFLOW_RUN_ID + TENANT_* env vars, discovers workflows
// from the user's project, and executes the matching one.
const entrypointPath = options.entrypoint ?? "./workflow-job.ts";

if (options.executor === "k8s") {
cliLogger.error(
"K8s executor requires custom configuration. Use --executor process for local dev, " +
"or configure K8sJobExecutor programmatically for production.",
);
exitProcess(1);
return;
}

const executor = new ProcessJobExecutor({
entrypointPath,
env: {
REDIS_URL: options.redisUrl,
},
debug: options.debug,
});

// Create and start job manager
const manager = new WorkflowJobManager({
backend,
executor,
pollInterval: options.pollInterval,
maxConcurrentJobs: options.concurrency,
stalledThreshold: options.stalledThreshold,
debug: options.debug,
});

await manager.start();

cliLogger.info(
`Workflow worker started (manager: ${manager.getManagerId()})`,
);
cliLogger.info("Polling for workflow jobs...\n");

// Graceful shutdown
let shuttingDown = false;
const shutdown = async (signal: "SIGINT" | "SIGTERM"): Promise<void> => {
if (shuttingDown) return;
shuttingDown = true;

cliLogger.info(`\nReceived ${signal}, shutting down worker...`);

try {
await manager.stop();
await backend.destroy();

const stats = manager.getStats();
cliLogger.info("Worker stopped.");
cliLogger.info(
` Jobs: ${stats.jobsCreated} created, ${stats.jobsCompleted} completed, ${stats.jobsFailed} failed`,
);
} catch (error) {
cliLogger.warn("Error during shutdown:", error);
} finally {
exitProcess(0);
}
};

registerTerminationSignals((signal) => {
void shutdown(signal);
});

// Keep alive
await new Promise(() => {});
}
31 changes: 31 additions & 0 deletions cli/commands/worker/handler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import { z } from "zod";
import { createArgParser, parseArgsOrThrow } from "#cli/shared/args";
import type { ParsedArgs } from "#cli/shared/types";

const WorkerArgsSchema = z.object({
redisUrl: z.string().default("redis://localhost:6379"),
concurrency: z.number().default(3),
pollInterval: z.number().default(5000),
stalledThreshold: z.number().default(60000),
executor: z.enum(["process", "k8s"]).default("process"),
entrypoint: z.string().optional(),
debug: z.boolean().default(false),
});

export type WorkerArgs = z.infer<typeof WorkerArgsSchema>;

export const parseWorkerArgs = createArgParser(WorkerArgsSchema, {
redisUrl: { keys: ["redis-url", "redis"], type: "string" },
concurrency: { keys: ["concurrency", "c"], type: "number" },
pollInterval: { keys: ["poll-interval"], type: "number" },
stalledThreshold: { keys: ["stalled-threshold"], type: "number" },
executor: { keys: ["executor", "e"], type: "string" },
entrypoint: { keys: ["entrypoint"], type: "string" },
debug: { keys: ["debug"], type: "boolean" },
});

export async function handleWorkerCommand(args: ParsedArgs): Promise<void> {
const opts = parseArgsOrThrow(parseWorkerArgs, "worker", args);
const { workerCommand } = await import("./command.ts");
await workerCommand(opts);
}
8 changes: 8 additions & 0 deletions cli/commands/worker/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
/**
* Worker command - Start workflow job manager
*/

export { workerCommand } from "./command.ts";
export type { WorkerOptions } from "./command.ts";
export { handleWorkerCommand } from "./handler.ts";
export { workerHelp } from "./command-help.ts";
2 changes: 2 additions & 0 deletions cli/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import { handleServeCommand } from "./commands/serve/handler.ts";
import { handleStartCommand } from "./commands/start/handler.ts";
import { handleStudioCommand } from "./commands/studio/handler.ts";
import { handleUpCommand } from "./commands/up/index.ts";
import { handleWorkerCommand } from "./commands/worker/handler.ts";
import { login, logout, whoami } from "./auth/index.ts";
import { parseLoginMethod } from "./auth/utils.ts";
import { showCommandHelp, showMainHelp } from "./help/index.ts";
Expand Down Expand Up @@ -74,6 +75,7 @@ const commands: Record<string, (args: ParsedArgs) => Promise<void>> = {
"mcp": handleMCPCommand,
"issues": handleIssuesCommand,
"start": handleStartCommand,
"worker": handleWorkerCommand,
};

/**
Expand Down
16 changes: 15 additions & 1 deletion deno.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@
"./agent": "./src/agent/index.ts",
"./tool": "./src/tool/index.ts",
"./workflow": "./src/workflow/index.ts",
"./workflow/worker": "./src/workflow/worker/index.ts",
"./workflow/claude-code": "./src/workflow/claude-code/index.ts",
"./workflow/claude-code/react": "./src/workflow/claude-code/react/index.ts",
"./workflow/discovery": "./src/workflow/discovery/index.ts",
"./prompt": "./src/prompt/index.ts",
"./resource": "./src/resource/index.ts",
"./mcp": "./src/mcp/index.ts",
Expand Down Expand Up @@ -56,6 +60,10 @@
"veryfront/workflow/executor": "./src/workflow/executor/index.ts",
"veryfront/workflow/blob": "./src/workflow/blob/index.ts",
"veryfront/workflow/runtime": "./src/workflow/runtime/agent-registry.ts",
"veryfront/workflow/worker": "./src/workflow/worker/index.ts",
"veryfront/workflow/claude-code": "./src/workflow/claude-code/index.ts",
"veryfront/workflow/claude-code/react": "./src/workflow/claude-code/react/index.ts",
"veryfront/workflow/discovery": "./src/workflow/discovery/index.ts",
"veryfront/utils/box": "./src/utils/box.ts",
"veryfront/utils/case-utils": "./src/utils/case-utils.ts",
"veryfront/utils/constants/server": "./src/utils/constants/server.ts",
Expand Down Expand Up @@ -130,6 +138,10 @@
"#veryfront/utils": "./src/utils/index.ts",
"#veryfront/workflow": "./src/workflow/index.ts",
"#veryfront/workflow/react": "./src/workflow/react/index.ts",
"#veryfront/workflow/worker": "./src/workflow/worker/index.ts",
"#veryfront/workflow/claude-code": "./src/workflow/claude-code/index.ts",
"#veryfront/workflow/claude-code/react": "./src/workflow/claude-code/react/index.ts",
"#veryfront/workflow/discovery": "./src/workflow/discovery/index.ts",
"#veryfront/schemas": "./src/schemas/index.ts",
"#veryfront/http/responses": "./src/platform/compat/http/responses.ts",
"#veryfront/compat/console": "./src/platform/compat/console/index.ts",
Expand All @@ -141,6 +153,7 @@
"#veryfront/platform/": "./src/platform/",
"#veryfront/components": "./src/react/components/index.ts",
"#veryfront/": "./src/",
"#deno-config": "./deno.json",
"std/": "jsr:@std/",
"#std/path": "jsr:@std/path",
"#std/path.ts": "jsr:@std/path",
Expand Down Expand Up @@ -215,6 +228,7 @@
"@ai-sdk/react": "npm:@ai-sdk/react@3.0.35",
"@ai-sdk/openai": "https://esm.sh/@ai-sdk/openai@2.0.1",
"@ai-sdk/anthropic": "https://esm.sh/@ai-sdk/anthropic@2.0.1",
"@anthropic-ai/claude-agent-sdk": "npm:@anthropic-ai/claude-agent-sdk@0.2.37",
"tailwindcss": "https://esm.sh/tailwindcss@4.1.8",
"tailwindcss/plugin": "https://esm.sh/tailwindcss@4.1.8/plugin",
"tailwindcss/defaultTheme": "https://esm.sh/tailwindcss@4.1.8/defaultTheme",
Expand Down Expand Up @@ -280,13 +294,13 @@
"docs:check-links": "deno run -A scripts/lint/check-doc-links.ts",
"lint:ban-console": "deno run --allow-read scripts/lint/ban-console.ts",
"lint:ban-deep-imports": "deno run --allow-read scripts/lint/ban-deep-imports.ts",
"lint:imports": "deno run --allow-read scripts/lint/no-cross-boundary-relative-imports.ts",
"lint:ban-internal-root-imports": "deno run --allow-read scripts/lint/ban-internal-root-imports.ts",
"lint:cli-boundary": "deno run --allow-read scripts/lint/enforce-cli-boundary.ts",
"validate:architecture": "deno run --allow-read scripts/lint/validate-architecture.ts",
"lint:check-awaits": "deno run --allow-read scripts/lint/check-unawaited-promises.ts",
"dupes": "deno run --allow-read scripts/lint/find-duplicate-functions.ts",
"lint:platform": "deno run --allow-read scripts/lint/lint-platform-agnostic.ts",
"test:smoke": "deno run --allow-all src/platform/compat/smoke-test.ts",
"test:cross-runtime": "deno run --allow-all src/platform/compat/cross-runtime.test.ts",
"test:node": "node ./tests/node/run-tests.mjs 'src/**/*.test.ts'",
"test:bun": "node ./tests/bun/run-tests.mjs src/",
Expand Down
Loading
Loading