From 3b9b2dcd11babcb8793feacf816a1bc5ced01af5 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 12 Mar 2026 16:46:34 +0000 Subject: [PATCH] feat(cron): add scheduled task execution system with AI-callable tools Implement a cron/scheduling system for NeuroLink that allows AI models to create, manage, and execute scheduled tasks via generate() calls. New files (src/lib/cron/): - types.ts: Type definitions for tasks, schedules, stores, backends - taskStore.ts: InMemoryTaskStore + RedisTaskStore with atomic operations - schedulerBackend.ts: NodeTimeoutScheduler using setTimeout/setInterval/croner - cronManager.ts: Core orchestrator for task lifecycle and execution - cronTools.ts: 4 AI-callable tools (create, list, cancel, getStatus) - index.ts: Public exports Features: - Three schedule types: "at" (one-shot), "every" (interval), "cron" (expression) - Two session modes: "isolated" (fresh context) and "same-session" (shared) - Instance-scoped CronManager (no module-global state, multi-instance safe) - Atomic incrementRunCountIfUnderLimit() to prevent maxRuns race conditions - Cryptographically secure task IDs via crypto.randomBytes() - Pluggable scheduler backend interface (Node.js default, extensible) - Pluggable persistence (in-memory default, optional Redis) - Concurrency control via p-limit - Env-based control: NEUROLINK_DISABLE_CRON_TOOLS, NEUROLINK_CRON_STORE, NEUROLINK_CRON_MAX_CONCURRENT (validated before use) - Registered as built-in tools via BaseProvider with NeuroLink instance binding - Graceful shutdown in both shutdown() and dispose() paths https://claude.ai/code/session_01LGQEq8JtRqf6QnrrrkvXfo --- package.json | 17 +- pnpm-lock.yaml | 49 ++++- src/lib/agent/directTools.ts | 21 ++ src/lib/core/baseProvider.ts | 18 +- src/lib/cron/cronManager.ts | 297 +++++++++++++++++++++++++++ src/lib/cron/cronTools.ts | 338 +++++++++++++++++++++++++++++++ src/lib/cron/index.ts | 39 ++++ src/lib/cron/schedulerBackend.ts | 217 ++++++++++++++++++++ src/lib/cron/taskStore.ts | 248 +++++++++++++++++++++++ src/lib/cron/types.ts | 233 +++++++++++++++++++++ src/lib/neurolink.ts | 111 +++++++++- src/lib/types/configTypes.ts | 5 + src/lib/types/index.ts | 16 ++ src/lib/utils/toolUtils.ts | 13 ++ 14 files changed, 1606 insertions(+), 16 deletions(-) create mode 100644 src/lib/cron/cronManager.ts create mode 100644 src/lib/cron/cronTools.ts create mode 100644 src/lib/cron/index.ts create mode 100644 src/lib/cron/schedulerBackend.ts create mode 100644 src/lib/cron/taskStore.ts create mode 100644 src/lib/cron/types.ts diff --git a/package.json b/package.json index ade26c423..684aacfc3 100644 --- a/package.json +++ b/package.json @@ -198,6 +198,7 @@ "adm-zip": "^0.5.16", "ai": "4.3.19", "chalk": "^5.6.2", + "croner": "^10.0.1", "csv-parser": "^3.2.0", "dotenv": "^16.6.1", "exceljs": "^4.4.0", @@ -237,21 +238,21 @@ "@opentelemetry/sdk-trace-node": "^2.0.0" }, "optionalDependencies": { + "@fastify/cors": "^11.2.0", + "@fastify/rate-limit": "^10.3.0", "@hono/node-server": "^1.13.0", + "@koa/cors": "^5.0.0", + "@koa/router": "^15.3.0", "canvas": "^3.2.0", + "cors": "^2.8.5", + "express": "^5.1.0", "express-rate-limit": "^7.4.0", + "fastify": "^5.7.2", "ffmpeg-static": "^5.3.0", "ffprobe-static": "^3.1.0", - "sharp": "^0.34.5", - "@fastify/rate-limit": "^10.3.0", - "fastify": "^5.7.2", - "@fastify/cors": "^11.2.0", - "@koa/cors": "^5.0.0", - "@koa/router": "^15.3.0", "koa": "^3.1.1", "koa-bodyparser": "^4.4.1", - "express": "^5.1.0", - "cors": "^2.8.5" + "sharp": "^0.34.5" }, "devDependencies": { "@actions/core": "^2.0.2", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 6a20132b3..ba536097c 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -113,6 +113,9 @@ importers: chalk: specifier: ^5.6.2 version: 5.6.2 + croner: + specifier: ^10.0.1 + version: 10.0.1 csv-parser: specifier: ^3.2.0 version: 3.2.0 @@ -873,24 +876,28 @@ packages: engines: {node: '>=14.21.3'} cpu: [arm64] os: [linux] + libc: [musl] '@biomejs/cli-linux-arm64@2.2.5': resolution: {integrity: sha512-5DjiiDfHqGgR2MS9D+AZ8kOfrzTGqLKywn8hoXpXXlJXIECGQ32t+gt/uiS2XyGBM2XQhR6ztUvbjZWeccFMoQ==} engines: {node: '>=14.21.3'} cpu: [arm64] os: [linux] + libc: [glibc] '@biomejs/cli-linux-x64-musl@2.2.5': resolution: {integrity: sha512-AVqLCDb/6K7aPNIcxHaTQj01sl1m989CJIQFQEaiQkGr2EQwyOpaATJ473h+nXDUuAcREhccfRpe/tu+0wu0eQ==} engines: {node: '>=14.21.3'} cpu: [x64] os: [linux] + libc: [musl] '@biomejs/cli-linux-x64@2.2.5': resolution: {integrity: sha512-fq9meKm1AEXeAWan3uCg6XSP5ObA6F/Ovm89TwaMiy1DNIwdgxPkNwxlXJX8iM6oRbFysYeGnT0OG8diCWb9ew==} engines: {node: '>=14.21.3'} cpu: [x64] os: [linux] + libc: [glibc] '@biomejs/cli-win32-arm64@2.2.5': resolution: {integrity: sha512-xaOIad4wBambwJa6mdp1FigYSIF9i7PCqRbvBqtIi9y29QtPVQ13sDGtUnsRoe6SjL10auMzQ6YAe+B3RpZXVg==} @@ -1472,89 +1479,105 @@ packages: resolution: {integrity: sha512-excjX8DfsIcJ10x1Kzr4RcWe1edC9PquDRRPx3YVCvQv+U5p7Yin2s32ftzikXojb1PIFc/9Mt28/y+iRklkrw==} cpu: [arm64] os: [linux] + libc: [glibc] '@img/sharp-libvips-linux-arm@1.2.4': resolution: {integrity: sha512-bFI7xcKFELdiNCVov8e44Ia4u2byA+l3XtsAj+Q8tfCwO6BQ8iDojYdvoPMqsKDkuoOo+X6HZA0s0q11ANMQ8A==} cpu: [arm] os: [linux] + libc: [glibc] '@img/sharp-libvips-linux-ppc64@1.2.4': resolution: {integrity: sha512-FMuvGijLDYG6lW+b/UvyilUWu5Ayu+3r2d1S8notiGCIyYU/76eig1UfMmkZ7vwgOrzKzlQbFSuQfgm7GYUPpA==} cpu: [ppc64] os: [linux] + libc: [glibc] '@img/sharp-libvips-linux-riscv64@1.2.4': resolution: {integrity: sha512-oVDbcR4zUC0ce82teubSm+x6ETixtKZBh/qbREIOcI3cULzDyb18Sr/Wcyx7NRQeQzOiHTNbZFF1UwPS2scyGA==} cpu: [riscv64] os: [linux] + libc: [glibc] '@img/sharp-libvips-linux-s390x@1.2.4': resolution: {integrity: sha512-qmp9VrzgPgMoGZyPvrQHqk02uyjA0/QrTO26Tqk6l4ZV0MPWIW6LTkqOIov+J1yEu7MbFQaDpwdwJKhbJvuRxQ==} cpu: [s390x] os: [linux] + libc: [glibc] '@img/sharp-libvips-linux-x64@1.2.4': resolution: {integrity: sha512-tJxiiLsmHc9Ax1bz3oaOYBURTXGIRDODBqhveVHonrHJ9/+k89qbLl0bcJns+e4t4rvaNBxaEZsFtSfAdquPrw==} cpu: [x64] os: [linux] + libc: [glibc] '@img/sharp-libvips-linuxmusl-arm64@1.2.4': resolution: {integrity: sha512-FVQHuwx1IIuNow9QAbYUzJ+En8KcVm9Lk5+uGUQJHaZmMECZmOlix9HnH7n1TRkXMS0pGxIJokIVB9SuqZGGXw==} cpu: [arm64] os: [linux] + libc: [musl] '@img/sharp-libvips-linuxmusl-x64@1.2.4': resolution: {integrity: sha512-+LpyBk7L44ZIXwz/VYfglaX/okxezESc6UxDSoyo2Ks6Jxc4Y7sGjpgU9s4PMgqgjj1gZCylTieNamqA1MF7Dg==} cpu: [x64] os: [linux] + libc: [musl] '@img/sharp-linux-arm64@0.34.5': resolution: {integrity: sha512-bKQzaJRY/bkPOXyKx5EVup7qkaojECG6NLYswgktOZjaXecSAeCWiZwwiFf3/Y+O1HrauiE3FVsGxFg8c24rZg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] + libc: [glibc] '@img/sharp-linux-arm@0.34.5': resolution: {integrity: sha512-9dLqsvwtg1uuXBGZKsxem9595+ujv0sJ6Vi8wcTANSFpwV/GONat5eCkzQo/1O6zRIkh0m/8+5BjrRr7jDUSZw==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm] os: [linux] + libc: [glibc] '@img/sharp-linux-ppc64@0.34.5': resolution: {integrity: sha512-7zznwNaqW6YtsfrGGDA6BRkISKAAE1Jo0QdpNYXNMHu2+0dTrPflTLNkpc8l7MUP5M16ZJcUvysVWWrMefZquA==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [ppc64] os: [linux] + libc: [glibc] '@img/sharp-linux-riscv64@0.34.5': resolution: {integrity: sha512-51gJuLPTKa7piYPaVs8GmByo7/U7/7TZOq+cnXJIHZKavIRHAP77e3N2HEl3dgiqdD/w0yUfiJnII77PuDDFdw==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [riscv64] os: [linux] + libc: [glibc] '@img/sharp-linux-s390x@0.34.5': resolution: {integrity: sha512-nQtCk0PdKfho3eC5MrbQoigJ2gd1CgddUMkabUj+rBevs8tZ2cULOx46E7oyX+04WGfABgIwmMC0VqieTiR4jg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [s390x] os: [linux] + libc: [glibc] '@img/sharp-linux-x64@0.34.5': resolution: {integrity: sha512-MEzd8HPKxVxVenwAa+JRPwEC7QFjoPWuS5NZnBt6B3pu7EG2Ge0id1oLHZpPJdn3OQK+BQDiw9zStiHBTJQQQQ==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] + libc: [glibc] '@img/sharp-linuxmusl-arm64@0.34.5': resolution: {integrity: sha512-fprJR6GtRsMt6Kyfq44IsChVZeGN97gTD331weR1ex1c1rypDEABN6Tm2xa1wE6lYb5DdEnk03NZPqA7Id21yg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] + libc: [musl] '@img/sharp-linuxmusl-x64@0.34.5': resolution: {integrity: sha512-Jg8wNT1MUzIvhBFxViqrEhWDGzqymo3sV7z7ZsaWbZNDLXRJZoRGrjulp60YYtV4wfY8VIKcWidjojlLcWrd8Q==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] + libc: [musl] '@img/sharp-wasm32@0.34.5': resolution: {integrity: sha512-OdWTEiVkY2PHwqkbBI8frFxQQFekHaSSkUIJkwzclWZe64O1X4UlUjqqqLaPbUpMOQk6FBu/HtlGXNblIs0huw==} @@ -1737,30 +1760,35 @@ packages: engines: {node: '>= 10'} cpu: [arm64] os: [linux] + libc: [glibc] '@napi-rs/canvas-linux-arm64-musl@0.1.80': resolution: {integrity: sha512-1XbCOz/ymhj24lFaIXtWnwv/6eFHXDrjP0jYkc6iHQ9q8oXKzUX1Lc6bu+wuGiLhGh2GS/2JlfORC5ZcXimRcg==} engines: {node: '>= 10'} cpu: [arm64] os: [linux] + libc: [musl] '@napi-rs/canvas-linux-riscv64-gnu@0.1.80': resolution: {integrity: sha512-XTzR125w5ZMs0lJcxRlS1K3P5RaZ9RmUsPtd1uGt+EfDyYMu4c6SEROYsxyatbbu/2+lPe7MPHOO/0a0x7L/gw==} engines: {node: '>= 10'} cpu: [riscv64] os: [linux] + libc: [glibc] '@napi-rs/canvas-linux-x64-gnu@0.1.80': resolution: {integrity: sha512-BeXAmhKg1kX3UCrJsYbdQd3hIMDH/K6HnP/pG2LuITaXhXBiNdh//TVVVVCBbJzVQaV5gK/4ZOCMrQW9mvuTqA==} engines: {node: '>= 10'} cpu: [x64] os: [linux] + libc: [glibc] '@napi-rs/canvas-linux-x64-musl@0.1.80': resolution: {integrity: sha512-x0XvZWdHbkgdgucJsRxprX/4o4sEed7qo9rCQA9ugiS9qE2QvP0RIiEugtZhfLH3cyI+jIRFJHV4Fuz+1BHHMg==} engines: {node: '>= 10'} cpu: [x64] os: [linux] + libc: [musl] '@napi-rs/canvas-win32-x64-msvc@0.1.80': resolution: {integrity: sha512-Z8jPsM6df5V8B1HrCHB05+bDiCxjE9QA//3YrkKIdVDEwn5RKaqOxCJDRJkl48cJbylcrJbW4HxZbTte8juuPg==} @@ -2717,56 +2745,67 @@ packages: resolution: {integrity: sha512-PsNAbcyv9CcecAUagQefwX8fQn9LQ4nZkpDboBOttmyffnInRy8R8dSg6hxxl2Re5QhHBf6FYIDhIj5v982ATQ==} cpu: [arm] os: [linux] + libc: [glibc] '@rollup/rollup-linux-arm-musleabihf@4.52.5': resolution: {integrity: sha512-Fw4tysRutyQc/wwkmcyoqFtJhh0u31K+Q6jYjeicsGJJ7bbEq8LwPWV/w0cnzOqR2m694/Af6hpFayLJZkG2VQ==} cpu: [arm] os: [linux] + libc: [musl] '@rollup/rollup-linux-arm64-gnu@4.52.5': resolution: {integrity: sha512-a+3wVnAYdQClOTlyapKmyI6BLPAFYs0JM8HRpgYZQO02rMR09ZcV9LbQB+NL6sljzG38869YqThrRnfPMCDtZg==} cpu: [arm64] os: [linux] + libc: [glibc] '@rollup/rollup-linux-arm64-musl@4.52.5': resolution: {integrity: sha512-AvttBOMwO9Pcuuf7m9PkC1PUIKsfaAJ4AYhy944qeTJgQOqJYJ9oVl2nYgY7Rk0mkbsuOpCAYSs6wLYB2Xiw0Q==} cpu: [arm64] os: [linux] + libc: [musl] '@rollup/rollup-linux-loong64-gnu@4.52.5': resolution: {integrity: sha512-DkDk8pmXQV2wVrF6oq5tONK6UHLz/XcEVow4JTTerdeV1uqPeHxwcg7aFsfnSm9L+OO8WJsWotKM2JJPMWrQtA==} cpu: [loong64] os: [linux] + libc: [glibc] '@rollup/rollup-linux-ppc64-gnu@4.52.5': resolution: {integrity: sha512-W/b9ZN/U9+hPQVvlGwjzi+Wy4xdoH2I8EjaCkMvzpI7wJUs8sWJ03Rq96jRnHkSrcHTpQe8h5Tg3ZzUPGauvAw==} cpu: [ppc64] os: [linux] + libc: [glibc] '@rollup/rollup-linux-riscv64-gnu@4.52.5': resolution: {integrity: sha512-sjQLr9BW7R/ZiXnQiWPkErNfLMkkWIoCz7YMn27HldKsADEKa5WYdobaa1hmN6slu9oWQbB6/jFpJ+P2IkVrmw==} cpu: [riscv64] os: [linux] + libc: [glibc] '@rollup/rollup-linux-riscv64-musl@4.52.5': resolution: {integrity: sha512-hq3jU/kGyjXWTvAh2awn8oHroCbrPm8JqM7RUpKjalIRWWXE01CQOf/tUNWNHjmbMHg/hmNCwc/Pz3k1T/j/Lg==} cpu: [riscv64] os: [linux] + libc: [musl] '@rollup/rollup-linux-s390x-gnu@4.52.5': resolution: {integrity: sha512-gn8kHOrku8D4NGHMK1Y7NA7INQTRdVOntt1OCYypZPRt6skGbddska44K8iocdpxHTMMNui5oH4elPH4QOLrFQ==} cpu: [s390x] os: [linux] + libc: [glibc] '@rollup/rollup-linux-x64-gnu@4.52.5': resolution: {integrity: sha512-hXGLYpdhiNElzN770+H2nlx+jRog8TyynpTVzdlc6bndktjKWyZyiCsuDAlpd+j+W+WNqfcyAWz9HxxIGfZm1Q==} cpu: [x64] os: [linux] + libc: [glibc] '@rollup/rollup-linux-x64-musl@4.52.5': resolution: {integrity: sha512-arCGIcuNKjBoKAXD+y7XomR9gY6Mw7HnFBv5Rw7wQRvwYLR7gBAgV7Mb2QTyjXfTveBNFAtPt46/36vV9STLNg==} cpu: [x64] os: [linux] + libc: [musl] '@rollup/rollup-openharmony-arm64@4.52.5': resolution: {integrity: sha512-QoFqB6+/9Rly/RiPjaomPLmR/13cgkIGfA40LHly9zcH1S0bN2HVFYk3a1eAyHQyjs3ZJYlXvIGtcCs5tko9Cw==} @@ -4304,6 +4343,10 @@ packages: resolution: {integrity: sha512-NT7w2JVU7DFroFdYkeq8cywxrgjPHWkdX1wjpRQXPX5Asews3tA+Ght6lddQO5Mkumffp3X7GEqku3epj2toIw==} engines: {node: '>= 10'} + croner@10.0.1: + resolution: {integrity: sha512-ixNtAJndqh173VQ4KodSdJEI6nuioBWI0V1ITNKhZZsO0pEMoDxz539T4FTTbSZ/xIOSuDnzxLVRqBVSvPNE2g==} + engines: {node: '>=18.0'} + cross-spawn@7.0.6: resolution: {integrity: sha512-uV2QOWP2nWzsy2aMp8aRibhi9dlzF5Hgh5SHaB9OiTGEyDTiJJyx0uy51QXdyWbtAHNua4XJzUKca3OzKUd3vA==} engines: {node: '>= 8'} @@ -9829,8 +9872,8 @@ snapshots: '@aws-sdk/client-bedrock-runtime': 3.901.0 '@aws-sdk/client-sagemaker': 3.901.0 '@aws-sdk/client-sagemaker-runtime': 3.901.0 - '@aws-sdk/credential-provider-node': 3.901.0 - '@aws-sdk/types': 3.901.0 + '@aws-sdk/credential-provider-node': 3.972.14 + '@aws-sdk/types': 3.973.4 '@google-cloud/text-to-speech': 5.8.1(encoding@0.1.13) '@google-cloud/vertexai': 1.10.0(encoding@0.1.13) '@google/genai': 1.34.0(@modelcontextprotocol/sdk@1.26.0(@cfworker/json-schema@4.1.1)(zod@3.25.76)) @@ -13310,6 +13353,8 @@ snapshots: crc-32: 1.2.2 readable-stream: 3.6.2 + croner@10.0.1: {} + cross-spawn@7.0.6: dependencies: path-key: 3.1.1 diff --git a/src/lib/agent/directTools.ts b/src/lib/agent/directTools.ts index b93a712d1..b76850992 100644 --- a/src/lib/agent/directTools.ts +++ b/src/lib/agent/directTools.ts @@ -10,6 +10,9 @@ import * as path from "path"; import { logger } from "../utils/logger.js"; import { VertexAI } from "@google-cloud/vertexai"; import { CSVProcessor } from "../utils/csvProcessor.js"; +import { shouldDisableCronTools } from "../utils/toolUtils.js"; +import { createCronTools } from "../cron/cronTools.js"; +import type { CronManager } from "../cron/cronManager.js"; // Runtime Google Search tool creation - bypasses TypeScript strict typing function createGoogleSearchTools() { @@ -896,3 +899,21 @@ export function validateToolStructure(): boolean { return false; } } + +// --- Cron Tools Integration --- + +/** + * Get cron tools bound to a specific CronManager getter. + * Returns empty object if disabled via env var. + * + * @param getCronManager - Function that returns the CronManager instance + * from the owning NeuroLink instance. This avoids module-global state + * and supports multiple NeuroLink instances correctly. + */ +// eslint-disable-next-line @typescript-eslint/no-explicit-any +export function getCronTools(getCronManager?: () => CronManager | undefined): Record { + if (shouldDisableCronTools()) { + return {}; + } + return createCronTools(getCronManager ?? (() => undefined)); +} diff --git a/src/lib/core/baseProvider.ts b/src/lib/core/baseProvider.ts index f2a28fae4..0dc5891f4 100644 --- a/src/lib/core/baseProvider.ts +++ b/src/lib/core/baseProvider.ts @@ -2,7 +2,7 @@ import { generateText } from "ai"; import type { CoreMessage, LanguageModelV1, Tool } from "ai"; import { SpanKind, SpanStatusCode } from "@opentelemetry/api"; import { tracers } from "../telemetry/tracers.js"; -import { directAgentTools } from "../agent/directTools.js"; +import { directAgentTools, getCronTools } from "../agent/directTools.js"; import type { AIProviderName } from "../constants/enums.js"; import { IMAGE_GENERATION_MODELS } from "../core/constants.js"; import type { EvaluationData } from "../index.js"; @@ -59,10 +59,10 @@ export abstract class BaseProvider implements AIProvider { protected readonly defaultTimeout: number = 30000; // 30 seconds protected middlewareOptions?: MiddlewareFactoryOptions; // TODO: Implement global level middlewares that can be used - // Tools are conditionally included based on centralized configuration - protected readonly directTools = shouldDisableBuiltinTools() - ? {} - : directAgentTools; + // Tools are conditionally included based on centralized configuration. + // Cron tools are merged in the constructor after this.neurolink is assigned, + // binding to the correct NeuroLink instance (avoids module-global state). + protected directTools: Record = {}; protected mcpTools?: Record; // MCP tools loaded dynamically when available protected customTools?: Map; // Custom tools from registerTool() protected toolExecutor?: ( @@ -92,6 +92,14 @@ export abstract class BaseProvider implements AIProvider { this.neurolink = neurolink; this.middlewareOptions = middleware; + // Initialize direct tools (built-in + cron) bound to this instance's NeuroLink + if (!shouldDisableBuiltinTools()) { + const cronTools = getCronTools( + () => this.neurolink?.getCronManager?.(), + ); + this.directTools = { ...directAgentTools, ...cronTools }; + } + // Initialize composition modules this.messageBuilder = new MessageBuilder(this.providerName, this.modelName); this.streamHandler = new StreamHandler(this.providerName, this.modelName); diff --git a/src/lib/cron/cronManager.ts b/src/lib/cron/cronManager.ts new file mode 100644 index 000000000..fbed121c5 --- /dev/null +++ b/src/lib/cron/cronManager.ts @@ -0,0 +1,297 @@ +/** + * CronManager - Core orchestrator for NeuroLink scheduled tasks + * + * Manages task lifecycle: creation, scheduling, execution, and cleanup. + * Executes tasks by calling a TaskExecutor (typically neurolink.generate()). + * Supports isolated and same-session execution modes. + */ + +import { randomBytes } from "crypto"; +import pLimit from "p-limit"; +import { logger } from "../utils/logger.js"; +import { NodeTimeoutScheduler } from "./schedulerBackend.js"; +import { InMemoryTaskStore, RedisTaskStore } from "./taskStore.js"; +import type { + CreateTaskOptions, + CronManagerConfig, + ScheduledTask, + SchedulerBackend, + TaskExecutor, + TaskFilter, + TaskRunResult, + TaskStore, +} from "./types.js"; + +/** + * CronManager orchestrates all scheduled task operations. + * + * Lifecycle: + * 1. User/AI creates a task via createTask() + * 2. Task is persisted in the TaskStore + * 3. Task is registered with the SchedulerBackend + * 4. When triggered, the TaskExecutor runs the prompt via generate() + * 5. Results are stored in the task's run history + */ +export class CronManager { + private store: TaskStore; + private scheduler: SchedulerBackend; + private executor: TaskExecutor | null = null; + private concurrencyLimit: ReturnType; + private config: Required< + Pick + > & + CronManagerConfig; + private isShutdown = false; + + constructor(config?: CronManagerConfig) { + this.config = { + enabled: config?.enabled ?? true, + maxConcurrentRuns: config?.maxConcurrentRuns ?? 3, + maxRunHistory: config?.maxRunHistory ?? 50, + ...config, + }; + + this.concurrencyLimit = pLimit(this.config.maxConcurrentRuns); + + // Initialize store + if (this.config.store === "redis" && this.config.redisConfig) { + this.store = new RedisTaskStore(this.config.redisConfig); + } else { + this.store = new InMemoryTaskStore(); + } + + // Initialize scheduler (Node.js timeout by default) + this.scheduler = new NodeTimeoutScheduler(); + + logger.debug("[CronManager] Initialized", { + enabled: this.config.enabled, + store: this.config.store || "memory", + maxConcurrent: this.config.maxConcurrentRuns, + }); + } + + /** + * Set the task executor function. + * This is called by NeuroLink after construction to wire up generate(). + */ + setExecutor(executor: TaskExecutor): void { + this.executor = executor; + } + + /** + * Create and schedule a new task. + */ + async createTask(options: CreateTaskOptions): Promise { + if (!this.config.enabled) { + throw new Error("Cron system is disabled"); + } + if (this.isShutdown) { + throw new Error("CronManager is shut down"); + } + + const task: ScheduledTask = { + id: generateTaskId(), + name: options.name || `task-${Date.now()}`, + schedule: options.schedule, + prompt: options.prompt, + sessionMode: options.sessionMode || "isolated", + provider: options.provider, + model: options.model, + maxRuns: options.maxRuns, + status: "active", + creatorSessionId: options.creatorSessionId, + runCount: 0, + runs: [], + maxRunHistory: this.config.maxRunHistory, + createdAt: Date.now(), + updatedAt: Date.now(), + }; + + // Compute next run time + task.nextRunAt = this.scheduler.getNextRunTime(task); + + // Persist + await this.store.save(task); + + // Schedule + this.scheduler.schedule(task, () => this.executeTask(task.id)); + + logger.info("[CronManager] Task created", { + taskId: task.id, + name: task.name, + scheduleType: task.schedule.type, + sessionMode: task.sessionMode, + nextRunAt: task.nextRunAt + ? new Date(task.nextRunAt).toISOString() + : "immediate", + }); + + return task; + } + + /** + * Cancel a scheduled task. + */ + async cancelTask(taskId: string): Promise { + const task = await this.store.get(taskId); + if (!task) { + return false; + } + + this.scheduler.cancel(taskId); + await this.store.updateStatus(taskId, "cancelled"); + + logger.info("[CronManager] Task cancelled", { taskId, name: task.name }); + return true; + } + + /** + * Get a task by ID. + */ + async getTask(taskId: string): Promise { + const task = await this.store.get(taskId); + if (task) { + // Refresh next run time + task.nextRunAt = this.scheduler.getNextRunTime(task); + } + return task; + } + + /** + * List tasks with optional filtering. + */ + async listTasks(filter?: TaskFilter): Promise { + const tasks = await this.store.list(filter); + // Refresh next run times + for (const task of tasks) { + if (task.status === "active") { + task.nextRunAt = this.scheduler.getNextRunTime(task); + } + } + return tasks; + } + + /** + * Gracefully shutdown the cron manager. + */ + async shutdown(): Promise { + this.isShutdown = true; + this.scheduler.shutdown(); + await this.store.shutdown(); + logger.debug("[CronManager] Shutdown complete"); + } + + // --- Private execution logic --- + + private async executeTask(taskId: string): Promise { + await this.concurrencyLimit(async () => { + const task = await this.store.get(taskId); + if (!task || task.status !== "active") { + return; + } + + if (!this.executor) { + logger.warn("[CronManager] No executor set, skipping task", { taskId }); + return; + } + + // Atomic check-and-increment for maxRuns to prevent race conditions + // where concurrent executions both pass the limit check. + if (task.maxRuns !== undefined) { + const newRunCount = await this.store.incrementRunCountIfUnderLimit(taskId, task.maxRuns); + if (newRunCount === undefined) { + this.scheduler.cancel(taskId); + await this.store.updateStatus(taskId, "completed"); + logger.info("[CronManager] Task completed (maxRuns reached)", { + taskId, + runCount: task.runCount, + maxRuns: task.maxRuns, + }); + return; + } + } + + const runNumber = task.runCount + 1; + const sessionId = + task.sessionMode === "isolated" + ? `cron:${taskId}:${runNumber}` + : task.creatorSessionId || `cron:${taskId}`; + + const run: TaskRunResult = { + runId: `${taskId}-run-${runNumber}`, + runNumber, + startedAt: Date.now(), + status: "running", + sessionId, + }; + + logger.info("[CronManager] Executing task", { + taskId, + name: task.name, + runNumber, + sessionMode: task.sessionMode, + sessionId, + }); + + try { + const result = await this.executor(task, sessionId); + + run.status = "completed"; + run.completedAt = Date.now(); + run.durationMs = run.completedAt - run.startedAt; + run.responseText = result.responseText; + run.tokenUsage = result.tokenUsage; + + logger.info("[CronManager] Task run completed", { + taskId, + runNumber, + durationMs: run.durationMs, + }); + } catch (error) { + run.status = "failed"; + run.completedAt = Date.now(); + run.durationMs = run.completedAt - run.startedAt; + run.error = error instanceof Error ? error.message : String(error); + + logger.error("[CronManager] Task run failed", { + taskId, + runNumber, + error: run.error, + }); + } + + // Store run result + await this.store.addRunResult(taskId, run); + + // Update next run time + const updatedTask = await this.store.get(taskId); + if (updatedTask) { + updatedTask.nextRunAt = this.scheduler.getNextRunTime(updatedTask); + await this.store.save(updatedTask); + } + + // For "at" (one-shot) tasks, mark as completed after execution + if (task.schedule.type === "at") { + await this.store.updateStatus(taskId, "completed"); + } + + // Check if maxRuns reached after this run (using fresh data from store) + if (task.maxRuns !== undefined && updatedTask) { + if (updatedTask.runCount >= task.maxRuns) { + this.scheduler.cancel(taskId); + await this.store.updateStatus(taskId, "completed"); + } + } + }); + } +} + +/** + * Generate a cryptographically secure unique task ID. + * Uses crypto.randomBytes() instead of Math.random() for security. + */ +function generateTaskId(): string { + const timestamp = Date.now().toString(36); + const random = randomBytes(6).toString("hex"); + return `cron_${timestamp}_${random}`; +} diff --git a/src/lib/cron/cronTools.ts b/src/lib/cron/cronTools.ts new file mode 100644 index 000000000..2c669c452 --- /dev/null +++ b/src/lib/cron/cronTools.ts @@ -0,0 +1,338 @@ +/** + * Cron Tool Definitions for NeuroLink + * + * Four AI-callable tools for scheduled task management: + * - createScheduledTask: Create a new scheduled task + * - listScheduledTasks: List tasks with optional filtering + * - cancelScheduledTask: Cancel a scheduled task + * - getScheduledTaskStatus: Get detailed status of a task + * + * These tools follow the same pattern as directTools.ts (Vercel AI SDK tool()). + * They require a CronManager reference injected at registration time. + */ + +import { tool } from "ai"; +import { z } from "zod"; +import type { CronManager } from "./cronManager.js"; +import type { SessionMode, TaskFilter } from "./types.js"; + +/** + * Create cron tools bound to a specific CronManager instance. + * This factory pattern allows tools to reference the manager without global state. + */ +export function createCronTools(getCronManager: () => CronManager | undefined) { + return { + createScheduledTask: tool({ + description: + "Create a new scheduled task that will execute an AI prompt at specified times. " + + 'Supports three schedule types: "at" (one-shot at specific time), "every" (recurring interval), ' + + 'and "cron" (cron expression). Tasks can run in "isolated" mode (fresh context each run) ' + + 'or "same-session" mode (shared conversation context).', + parameters: z.object({ + scheduleType: z + .enum(["at", "every", "cron"]) + .describe( + 'Schedule type: "at" for one-shot, "every" for interval, "cron" for cron expression', + ), + scheduleValue: z + .string() + .describe( + 'Schedule value: ISO 8601 timestamp for "at" (e.g., "2024-12-25T09:00:00Z"), ' + + 'interval for "every" (e.g., "30s", "5m", "1h", or milliseconds), ' + + 'cron expression for "cron" (e.g., "*/5 * * * *" for every 5 minutes, "0 9 * * *" for daily at 9am)', + ), + prompt: z + .string() + .describe("The AI prompt to execute at each scheduled run"), + sessionMode: z + .enum(["isolated", "same-session"]) + .optional() + .default("isolated") + .describe( + '"isolated" = fresh session each run (no history), "same-session" = shared conversation context', + ), + taskName: z + .string() + .optional() + .describe("Human-readable name for the task"), + provider: z + .string() + .optional() + .describe( + 'AI provider override (e.g., "openai", "anthropic", "google")', + ), + model: z + .string() + .optional() + .describe('Model override (e.g., "gpt-4o", "claude-sonnet-4-20250514")'), + maxRuns: z + .number() + .optional() + .describe( + "Maximum number of executions (omit for unlimited recurring tasks)", + ), + timezone: z + .string() + .optional() + .describe( + 'IANA timezone for cron expressions (e.g., "America/New_York", "Asia/Kolkata")', + ), + }), + execute: async ({ + scheduleType, + scheduleValue, + prompt, + sessionMode, + taskName, + provider, + model, + maxRuns, + timezone, + }) => { + const manager = getCronManager(); + if (!manager) { + return { + success: false, + error: + "Cron system is not initialized. Ensure NEUROLINK_DISABLE_CRON_TOOLS is not set to true.", + }; + } + + try { + // Parse scheduleValue for "every" type - could be number string + let value: string | number = scheduleValue; + if (scheduleType === "every") { + const num = Number(scheduleValue); + if (!isNaN(num)) { + value = num; + } + } + + const task = await manager.createTask({ + schedule: { + type: scheduleType, + value, + timezone, + }, + prompt, + sessionMode: sessionMode as SessionMode, + name: taskName, + provider, + model, + maxRuns, + }); + + return { + success: true, + taskId: task.id, + name: task.name, + scheduleType: task.schedule.type, + scheduleValue: String(task.schedule.value), + sessionMode: task.sessionMode, + status: task.status, + nextRunAt: task.nextRunAt + ? new Date(task.nextRunAt).toISOString() + : "immediate", + message: `Scheduled task "${task.name}" created successfully. Task ID: ${task.id}`, + }; + } catch (error) { + return { + success: false, + error: error instanceof Error ? error.message : String(error), + }; + } + }, + }), + + listScheduledTasks: tool({ + description: + "List all scheduled tasks with optional filtering by status. " + + "Shows task details including schedule, status, run count, and next execution time.", + parameters: z.object({ + status: z + .enum(["active", "paused", "completed", "failed", "cancelled"]) + .optional() + .describe("Filter tasks by status"), + limit: z + .number() + .optional() + .default(20) + .describe("Maximum number of tasks to return (default: 20)"), + }), + execute: async ({ status, limit }) => { + const manager = getCronManager(); + if (!manager) { + return { + success: false, + error: "Cron system is not initialized.", + }; + } + + try { + const filter: TaskFilter = { limit }; + if (status) { + filter.status = status; + } + + const tasks = await manager.listTasks(filter); + + return { + success: true, + count: tasks.length, + tasks: tasks.map((t) => ({ + taskId: t.id, + name: t.name, + scheduleType: t.schedule.type, + scheduleValue: String(t.schedule.value), + sessionMode: t.sessionMode, + status: t.status, + runCount: t.runCount, + maxRuns: t.maxRuns, + nextRunAt: t.nextRunAt + ? new Date(t.nextRunAt).toISOString() + : undefined, + createdAt: new Date(t.createdAt).toISOString(), + lastRunStatus: + t.runs.length > 0 ? t.runs[0].status : undefined, + lastRunAt: + t.runs.length > 0 + ? new Date(t.runs[0].startedAt).toISOString() + : undefined, + })), + }; + } catch (error) { + return { + success: false, + error: error instanceof Error ? error.message : String(error), + }; + } + }, + }), + + cancelScheduledTask: tool({ + description: + "Cancel a scheduled task by its task ID. The task will stop executing and be marked as cancelled.", + parameters: z.object({ + taskId: z + .string() + .describe("The task ID to cancel (returned from createScheduledTask)"), + }), + execute: async ({ taskId }) => { + const manager = getCronManager(); + if (!manager) { + return { + success: false, + error: "Cron system is not initialized.", + }; + } + + try { + const cancelled = await manager.cancelTask(taskId); + if (cancelled) { + return { + success: true, + taskId, + message: `Task ${taskId} has been cancelled successfully.`, + }; + } else { + return { + success: false, + error: `Task ${taskId} not found.`, + }; + } + } catch (error) { + return { + success: false, + error: error instanceof Error ? error.message : String(error), + }; + } + }, + }), + + getScheduledTaskStatus: tool({ + description: + "Get detailed status of a scheduled task including run history, next execution time, " + + "and execution results. Use this to monitor task health and review past runs.", + parameters: z.object({ + taskId: z + .string() + .describe("The task ID to get status for"), + includeRuns: z + .boolean() + .optional() + .default(true) + .describe("Include run history in response (default: true)"), + maxRuns: z + .number() + .optional() + .default(5) + .describe("Maximum number of recent runs to include (default: 5)"), + }), + execute: async ({ taskId, includeRuns, maxRuns }) => { + const manager = getCronManager(); + if (!manager) { + return { + success: false, + error: "Cron system is not initialized.", + }; + } + + try { + const task = await manager.getTask(taskId); + if (!task) { + return { + success: false, + error: `Task ${taskId} not found.`, + }; + } + + const result: Record = { + success: true, + taskId: task.id, + name: task.name, + status: task.status, + scheduleType: task.schedule.type, + scheduleValue: String(task.schedule.value), + timezone: task.schedule.timezone, + sessionMode: task.sessionMode, + provider: task.provider, + model: task.model, + runCount: task.runCount, + maxRuns: task.maxRuns, + nextRunAt: task.nextRunAt + ? new Date(task.nextRunAt).toISOString() + : undefined, + createdAt: new Date(task.createdAt).toISOString(), + updatedAt: new Date(task.updatedAt).toISOString(), + }; + + if (includeRuns) { + result.recentRuns = task.runs.slice(0, maxRuns).map((r) => ({ + runId: r.runId, + runNumber: r.runNumber, + status: r.status, + startedAt: new Date(r.startedAt).toISOString(), + completedAt: r.completedAt + ? new Date(r.completedAt).toISOString() + : undefined, + durationMs: r.durationMs, + responsePreview: r.responseText + ? r.responseText.substring(0, 200) + + (r.responseText.length > 200 ? "..." : "") + : undefined, + error: r.error, + tokenUsage: r.tokenUsage, + })); + } + + return result; + } catch (error) { + return { + success: false, + error: error instanceof Error ? error.message : String(error), + }; + } + }, + }), + }; +} diff --git a/src/lib/cron/index.ts b/src/lib/cron/index.ts new file mode 100644 index 000000000..e1a0323c6 --- /dev/null +++ b/src/lib/cron/index.ts @@ -0,0 +1,39 @@ +/** + * NeuroLink Cron/Scheduling System + * + * Provides scheduled task execution for AI prompts with support for: + * - Three schedule types: one-shot ("at"), interval ("every"), cron expressions ("cron") + * - Two session modes: isolated (fresh context) and same-session (shared context) + * - Pluggable scheduler backends (Node.js timeout default, extensible for RabbitMQ, webhooks) + * - Pluggable persistence (in-memory default, optional Redis) + * - AI-callable tools for task management (create, list, cancel, status) + */ + +// Core manager +export { CronManager } from "./cronManager.js"; + +// Scheduler backends +export { NodeTimeoutScheduler, parseInterval } from "./schedulerBackend.js"; + +// Task stores +export { InMemoryTaskStore, RedisTaskStore } from "./taskStore.js"; + +// Tool factory +export { createCronTools } from "./cronTools.js"; + +// Types +export type { + CronManagerConfig, + CreateTaskOptions, + ScheduledTask, + Schedule, + ScheduleType, + SchedulerBackend, + SessionMode, + TaskExecutor, + TaskFilter, + TaskRunResult, + TaskStatus, + TaskStore, + RunStatus, +} from "./types.js"; diff --git a/src/lib/cron/schedulerBackend.ts b/src/lib/cron/schedulerBackend.ts new file mode 100644 index 000000000..0da0a479d --- /dev/null +++ b/src/lib/cron/schedulerBackend.ts @@ -0,0 +1,217 @@ +/** + * Scheduler Backends for NeuroLink Cron System + * + * Provides the timer management layer. The default NodeTimeoutScheduler + * uses native Node.js setTimeout/setInterval. The interface is designed + * to be extended with RabbitMQ, webhook, or other backends in the future. + */ + +import { Cron } from "croner"; +import type { ScheduledTask, SchedulerBackend } from "./types.js"; + +/** + * Node.js native timeout/interval-based scheduler. + * + * - "at": setTimeout to fire at a specific timestamp + * - "every": setInterval with fixed interval + * - "cron": croner library for cron expression parsing + setTimeout chain + */ +export class NodeTimeoutScheduler implements SchedulerBackend { + /** Active timer handles keyed by task ID */ + private timers: Map = new Map(); + /** Active croner instances keyed by task ID */ + private cronJobs: Map = new Map(); + + schedule(task: ScheduledTask, callback: () => Promise): void { + // Cancel any existing timer for this task + this.cancel(task.id); + + switch (task.schedule.type) { + case "at": + this.scheduleAt(task, callback); + break; + case "every": + this.scheduleEvery(task, callback); + break; + case "cron": + this.scheduleCron(task, callback); + break; + } + } + + cancel(taskId: string): void { + const timer = this.timers.get(taskId); + if (timer) { + clearTimeout(timer); + clearInterval(timer); + this.timers.delete(taskId); + } + + const cronJob = this.cronJobs.get(taskId); + if (cronJob) { + cronJob.stop(); + this.cronJobs.delete(taskId); + } + } + + cancelAll(): void { + this.timers.forEach((timer) => { + clearTimeout(timer); + clearInterval(timer); + }); + this.timers.clear(); + + this.cronJobs.forEach((cronJob) => { + cronJob.stop(); + }); + this.cronJobs.clear(); + } + + getNextRunTime(task: ScheduledTask): number | undefined { + switch (task.schedule.type) { + case "at": { + const targetTime = + typeof task.schedule.value === "string" + ? new Date(task.schedule.value).getTime() + : task.schedule.value; + return targetTime > Date.now() ? targetTime : undefined; + } + case "every": { + const intervalMs = + typeof task.schedule.value === "number" + ? task.schedule.value + : parseInterval(task.schedule.value); + return Date.now() + intervalMs; + } + case "cron": { + const cronJob = this.cronJobs.get(task.id); + if (cronJob) { + const next = cronJob.nextRun(); + return next ? next.getTime() : undefined; + } + // Compute without an active job + const tempCron = new Cron(String(task.schedule.value), { + timezone: task.schedule.timezone, + paused: true, + }); + const next = tempCron.nextRun(); + tempCron.stop(); + return next ? next.getTime() : undefined; + } + default: + return undefined; + } + } + + shutdown(): void { + this.cancelAll(); + } + + // --- Private scheduling methods --- + + private scheduleAt( + task: ScheduledTask, + callback: () => Promise, + ): void { + const targetTime = + typeof task.schedule.value === "string" + ? new Date(task.schedule.value).getTime() + : task.schedule.value; + + const delayMs = targetTime - Date.now(); + if (delayMs <= 0) { + // Already past - execute immediately + void callback(); + return; + } + + const timer = setTimeout(() => { + this.timers.delete(task.id); + void callback(); + }, delayMs); + + // Prevent the timer from keeping the process alive + if (timer.unref) { + timer.unref(); + } + this.timers.set(task.id, timer); + } + + private scheduleEvery( + task: ScheduledTask, + callback: () => Promise, + ): void { + const intervalMs = + typeof task.schedule.value === "number" + ? task.schedule.value + : parseInterval(String(task.schedule.value)); + + if (intervalMs <= 0) { + throw new Error(`Invalid interval: ${task.schedule.value}`); + } + + const timer = setInterval(() => { + void callback(); + }, intervalMs); + + if (timer.unref) { + timer.unref(); + } + this.timers.set(task.id, timer); + } + + private scheduleCron( + task: ScheduledTask, + callback: () => Promise, + ): void { + const cronExpression = String(task.schedule.value); + + const job = new Cron(cronExpression, { + timezone: task.schedule.timezone, + protect: true, // Prevent overlapping runs + }, () => { + void callback(); + }); + + this.cronJobs.set(task.id, job); + } +} + +/** + * Parse a human-readable interval string to milliseconds. + * Supports: "30s", "5m", "1h", "1d", or raw milliseconds as string. + */ +export function parseInterval(value: string): number { + const trimmed = value.trim().toLowerCase(); + + // Raw number = milliseconds + const rawNum = Number(trimmed); + if (!isNaN(rawNum) && rawNum > 0) { + return rawNum; + } + + const match = trimmed.match(/^(\d+(?:\.\d+)?)\s*(ms|s|m|h|d)$/); + if (!match) { + throw new Error( + `Invalid interval format: "${value}". Use "30s", "5m", "1h", "1d", or milliseconds.`, + ); + } + + const num = parseFloat(match[1]); + const unit = match[2]; + + switch (unit) { + case "ms": + return num; + case "s": + return num * 1000; + case "m": + return num * 60 * 1000; + case "h": + return num * 60 * 60 * 1000; + case "d": + return num * 24 * 60 * 60 * 1000; + default: + throw new Error(`Unknown time unit: ${unit}`); + } +} diff --git a/src/lib/cron/taskStore.ts b/src/lib/cron/taskStore.ts new file mode 100644 index 000000000..dd12d1c12 --- /dev/null +++ b/src/lib/cron/taskStore.ts @@ -0,0 +1,248 @@ +/** + * Task Persistence Stores for NeuroLink Cron System + * + * Provides two implementations: + * - InMemoryTaskStore: Default, Map-based in-memory storage + * - RedisTaskStore: Optional Redis-backed persistence for production use + */ + +import type { + ScheduledTask, + TaskFilter, + TaskRunResult, + TaskStatus, + TaskStore, +} from "./types.js"; + +/** + * In-memory task store using Map. + * Tasks are lost on process restart. Suitable for development + * and session-scoped scheduling. + */ +export class InMemoryTaskStore implements TaskStore { + private tasks: Map = new Map(); + + async save(task: ScheduledTask): Promise { + this.tasks.set(task.id, { ...task }); + } + + async get(taskId: string): Promise { + const task = this.tasks.get(taskId); + return task ? { ...task } : undefined; + } + + async list(filter?: TaskFilter): Promise { + let tasks = Array.from(this.tasks.values()); + + if (filter?.status) { + tasks = tasks.filter((t) => t.status === filter.status); + } + if (filter?.sessionMode) { + tasks = tasks.filter((t) => t.sessionMode === filter.sessionMode); + } + if (filter?.limit) { + tasks = tasks.slice(0, filter.limit); + } + + return tasks.map((t) => ({ ...t })); + } + + async delete(taskId: string): Promise { + return this.tasks.delete(taskId); + } + + async updateStatus(taskId: string, status: TaskStatus): Promise { + const task = this.tasks.get(taskId); + if (task) { + task.status = status; + task.updatedAt = Date.now(); + } + } + + async addRunResult(taskId: string, run: TaskRunResult): Promise { + const task = this.tasks.get(taskId); + if (task) { + task.runs.unshift(run); + // Trim run history + if (task.runs.length > task.maxRunHistory) { + task.runs = task.runs.slice(0, task.maxRunHistory); + } + task.updatedAt = Date.now(); + } + } + + async incrementRunCountIfUnderLimit(taskId: string, maxRuns: number): Promise { + const task = this.tasks.get(taskId); + if (!task) { + return undefined; + } + if (task.runCount >= maxRuns) { + return undefined; + } + task.runCount = task.runCount + 1; + task.updatedAt = Date.now(); + return task.runCount; + } + + async shutdown(): Promise { + this.tasks.clear(); + } +} + +/** + * Redis-backed task store for production persistence. + * Tasks survive process restarts. Uses the existing Redis config pattern. + */ +export class RedisTaskStore implements TaskStore { + private keyPrefix: string; + private redisClient: RedisLikeClient | null = null; + private config: RedisTaskStoreConfig; + + constructor(config: RedisTaskStoreConfig) { + this.config = config; + this.keyPrefix = config.keyPrefix || "neurolink:cron:"; + } + + private async getClient(): Promise { + if (this.redisClient) { + return this.redisClient; + } + + // Dynamic import to avoid requiring redis as a dependency + const { createClient } = await import("redis"); + + const url = + this.config.url || + `redis://${this.config.host || "localhost"}:${this.config.port || 6379}`; + + this.redisClient = createClient({ + url, + password: this.config.password, + }) as unknown as RedisLikeClient; + + await this.redisClient.connect(); + return this.redisClient; + } + + private taskKey(taskId: string): string { + return `${this.keyPrefix}task:${taskId}`; + } + + private indexKey(): string { + return `${this.keyPrefix}task-index`; + } + + async save(task: ScheduledTask): Promise { + const client = await this.getClient(); + const key = this.taskKey(task.id); + await client.set(key, JSON.stringify(task)); + await client.sAdd(this.indexKey(), task.id); + } + + async get(taskId: string): Promise { + const client = await this.getClient(); + const data = await client.get(this.taskKey(taskId)); + return data ? (JSON.parse(data) as ScheduledTask) : undefined; + } + + async list(filter?: TaskFilter): Promise { + const client = await this.getClient(); + const ids = await client.sMembers(this.indexKey()); + const tasks: ScheduledTask[] = []; + + for (const id of ids) { + const data = await client.get(this.taskKey(id)); + if (data) { + const task = JSON.parse(data) as ScheduledTask; + if (filter?.status && task.status !== filter.status) { + continue; + } + if (filter?.sessionMode && task.sessionMode !== filter.sessionMode) { + continue; + } + tasks.push(task); + } + } + + if (filter?.limit) { + return tasks.slice(0, filter.limit); + } + return tasks; + } + + async delete(taskId: string): Promise { + const client = await this.getClient(); + const deleted = await client.del(this.taskKey(taskId)); + await client.sRem(this.indexKey(), taskId); + return deleted > 0; + } + + async updateStatus(taskId: string, status: TaskStatus): Promise { + const task = await this.get(taskId); + if (task) { + task.status = status; + task.updatedAt = Date.now(); + await this.save(task); + } + } + + async addRunResult(taskId: string, run: TaskRunResult): Promise { + const task = await this.get(taskId); + if (task) { + task.runs.unshift(run); + if (task.runs.length > task.maxRunHistory) { + task.runs = task.runs.slice(0, task.maxRunHistory); + } + task.updatedAt = Date.now(); + await this.save(task); + } + } + + async incrementRunCountIfUnderLimit(taskId: string, maxRuns: number): Promise { + // For Redis, this is a read-modify-write. In a multi-process environment, + // consider using Redis WATCH/MULTI for true atomicity. + const task = await this.get(taskId); + if (!task) { + return undefined; + } + if (task.runCount >= maxRuns) { + return undefined; + } + task.runCount = task.runCount + 1; + task.updatedAt = Date.now(); + await this.save(task); + return task.runCount; + } + + async shutdown(): Promise { + if (this.redisClient) { + await this.redisClient.quit(); + this.redisClient = null; + } + } +} + +/** + * Minimal Redis client interface to avoid tight coupling to any specific Redis library + */ +interface RedisLikeClient { + connect(): Promise; + quit(): Promise; + get(key: string): Promise; + set(key: string, value: string): Promise; + del(key: string): Promise; + sAdd(key: string, member: string): Promise; + sRem(key: string, member: string): Promise; + sMembers(key: string): Promise; +} + +/** + * Redis task store configuration + */ +type RedisTaskStoreConfig = { + url?: string; + host?: string; + port?: number; + password?: string; + keyPrefix?: string; +}; diff --git a/src/lib/cron/types.ts b/src/lib/cron/types.ts new file mode 100644 index 000000000..97347dead --- /dev/null +++ b/src/lib/cron/types.ts @@ -0,0 +1,233 @@ +/** + * NeuroLink Cron/Scheduling System - Type Definitions + * + * Defines types for scheduled task management including task definitions, + * scheduler backends, persistence stores, and configuration. + */ + +/** + * Schedule type - how a task is triggered + * - "at": One-shot execution at a specific time (ISO 8601) + * - "every": Fixed-interval repetition (milliseconds) + * - "cron": Standard cron expression (5 or 6 field) + */ +export type ScheduleType = "at" | "every" | "cron"; + +/** + * Session mode - how task execution relates to conversation context + * - "isolated": Fresh session per execution, no conversation history + * - "same-session": Uses the creating session's conversation context + */ +export type SessionMode = "isolated" | "same-session"; + +/** + * Task lifecycle status + */ +export type TaskStatus = + | "active" + | "paused" + | "completed" + | "failed" + | "cancelled"; + +/** + * Individual run status + */ +export type RunStatus = "running" | "completed" | "failed"; + +/** + * Schedule definition + */ +export type Schedule = { + /** Schedule type */ + type: ScheduleType; + /** + * Schedule value: + * - For "at": ISO 8601 timestamp string (e.g., "2024-12-25T09:00:00Z") + * - For "every": interval in milliseconds (e.g., 60000 for 1 minute) + * - For "cron": cron expression (e.g., "0 9 * * *" for daily at 9am) + */ + value: string | number; + /** IANA timezone for cron expressions (e.g., "America/New_York") */ + timezone?: string; +}; + +/** + * Result of a single task execution + */ +export type TaskRunResult = { + /** Unique run identifier */ + runId: string; + /** Run number (1-based) */ + runNumber: number; + /** Start time (Unix ms) */ + startedAt: number; + /** End time (Unix ms) */ + completedAt?: number; + /** Duration in ms */ + durationMs?: number; + /** Run status */ + status: RunStatus; + /** AI-generated response text */ + responseText?: string; + /** Token usage for the run */ + tokenUsage?: { promptTokens?: number; completionTokens?: number; totalTokens?: number }; + /** Error message if failed */ + error?: string; + /** Session ID used for this run */ + sessionId: string; +}; + +/** + * Full scheduled task definition + */ +export type ScheduledTask = { + /** Unique task identifier */ + id: string; + /** Human-readable task name */ + name: string; + /** Schedule configuration */ + schedule: Schedule; + /** Prompt to send to the AI provider */ + prompt: string; + /** Session execution mode */ + sessionMode: SessionMode; + /** AI provider override (e.g., "openai", "anthropic") */ + provider?: string; + /** Model override (e.g., "gpt-4o", "claude-sonnet-4-20250514") */ + model?: string; + /** Maximum number of runs (undefined = unlimited for recurring) */ + maxRuns?: number; + /** Current task status */ + status: TaskStatus; + /** Session ID of the creator (used for same-session mode) */ + creatorSessionId?: string; + /** Total number of completed runs */ + runCount: number; + /** Run history (most recent first) */ + runs: TaskRunResult[]; + /** Maximum run history to keep */ + maxRunHistory: number; + /** Next scheduled execution time (Unix ms) */ + nextRunAt?: number; + /** Task creation time (Unix ms) */ + createdAt: number; + /** Last update time (Unix ms) */ + updatedAt: number; +}; + +/** + * Options for creating a new scheduled task + */ +export type CreateTaskOptions = { + /** Schedule configuration */ + schedule: Schedule; + /** Prompt to send to the AI provider */ + prompt: string; + /** Session execution mode (default: "isolated") */ + sessionMode?: SessionMode; + /** Human-readable task name */ + name?: string; + /** AI provider override */ + provider?: string; + /** Model override */ + model?: string; + /** Maximum number of runs */ + maxRuns?: number; + /** Session ID of the creator */ + creatorSessionId?: string; +}; + +/** + * Filter options for listing tasks + */ +export type TaskFilter = { + /** Filter by status */ + status?: TaskStatus; + /** Filter by session mode */ + sessionMode?: SessionMode; + /** Maximum number of results */ + limit?: number; +}; + +/** + * Interface for scheduler backends (timer management) + * Implementations: NodeTimeoutScheduler (default), future: RabbitMQ, webhook + */ +export interface SchedulerBackend { + /** Schedule a task for execution */ + schedule(task: ScheduledTask, callback: () => Promise): void; + /** Cancel a scheduled task */ + cancel(taskId: string): void; + /** Cancel all scheduled tasks */ + cancelAll(): void; + /** Get next run time for a task (Unix ms) */ + getNextRunTime(task: ScheduledTask): number | undefined; + /** Shutdown the scheduler backend */ + shutdown(): void; +} + +/** + * Interface for task persistence stores + * Implementations: InMemoryTaskStore (default), RedisTaskStore (optional) + */ +export interface TaskStore { + /** Save or update a task */ + save(task: ScheduledTask): Promise; + /** Get a task by ID */ + get(taskId: string): Promise; + /** List tasks with optional filter */ + list(filter?: TaskFilter): Promise; + /** Delete a task */ + delete(taskId: string): Promise; + /** Update task status */ + updateStatus(taskId: string, status: TaskStatus): Promise; + /** Add a run result to a task */ + addRunResult(taskId: string, run: TaskRunResult): Promise; + /** + * Atomically check if runCount < maxRuns and increment runCount. + * Returns the new runCount if successful, or undefined if the limit was already reached. + * This prevents race conditions where concurrent executions both pass the maxRuns check. + */ + incrementRunCountIfUnderLimit(taskId: string, maxRuns: number): Promise; + /** Shutdown the store */ + shutdown(): Promise; +} + +/** + * Cron manager configuration + */ +export type CronManagerConfig = { + /** Enable/disable the cron system (default: true) */ + enabled?: boolean; + /** Maximum concurrent task executions (default: 3) */ + maxConcurrentRuns?: number; + /** Store backend: "memory" or "redis" (default: "memory") */ + store?: "memory" | "redis"; + /** Redis configuration (required if store is "redis") */ + redisConfig?: { + url?: string; + host?: string; + port?: number; + password?: string; + keyPrefix?: string; + }; + /** Maximum run history per task (default: 50) */ + maxRunHistory?: number; + /** Default provider for task execution */ + defaultProvider?: string; + /** Default model for task execution */ + defaultModel?: string; +}; + +/** + * Task execution callback - called by CronManager to execute a task + * This is the function that calls neurolink.generate() with the task's prompt + */ +export type TaskExecutor = ( + task: ScheduledTask, + sessionId: string, +) => Promise<{ + responseText?: string; + tokenUsage?: { promptTokens?: number; completionTokens?: number; totalTokens?: number }; +}>; diff --git a/src/lib/neurolink.ts b/src/lib/neurolink.ts index e227a9d07..5f61eca46 100644 --- a/src/lib/neurolink.ts +++ b/src/lib/neurolink.ts @@ -188,6 +188,9 @@ import { ATTR } from "./telemetry/attributes.js"; import { getWorkflow } from "./workflow/core/workflowRegistry.js"; import { runWorkflow } from "./workflow/core/workflowRunner.js"; import type { WorkflowConfig } from "./workflow/types.js"; +import { CronManager } from "./cron/cronManager.js"; +import type { CronManagerConfig } from "./cron/types.js"; +import { shouldDisableCronTools } from "./utils/toolUtils.js"; /** * Check if an error is a non-retryable provider error that should immediately @@ -628,6 +631,7 @@ export class NeuroLink { * @throws {Error} When HITL configuration is invalid (if enabled) */ private observabilityConfig?: ObservabilityConfig; + private cronManager?: CronManager; constructor(config?: NeurolinkConstructorConfig) { this.toolRegistry = config?.toolRegistry || new MCPToolRegistry(); @@ -673,6 +677,7 @@ export class NeuroLink { ); this.registerFileTools(); this.registerMemoryRetrievalTools(); + this.initializeCronManager(config?.cron); this.initializeLangfuse( constructorId, constructorStartTime, @@ -1335,6 +1340,84 @@ Current user's request: ${currentInput}`; }); } + /** + * Initialize the CronManager for scheduled task execution. + * Wires up the task executor to call this.generate() with the task's prompt. + */ + private initializeCronManager(cronConfig?: CronManagerConfig): void { + if (shouldDisableCronTools()) { + logger.debug("[NeuroLink] Cron tools disabled via environment/config"); + return; + } + + const envStore = process.env.NEUROLINK_CRON_STORE; + const envMaxConcurrent = process.env.NEUROLINK_CRON_MAX_CONCURRENT; + + // Validate environment variables before use + const validStores = ["memory", "redis"]; + const resolvedStore = (envStore && validStores.includes(envStore)) + ? (envStore as "memory" | "redis") + : cronConfig?.store || "memory"; + + const parsedMaxConcurrent = envMaxConcurrent + ? Number.parseInt(envMaxConcurrent, 10) + : undefined; + const resolvedMaxConcurrent = (parsedMaxConcurrent != null && Number.isInteger(parsedMaxConcurrent) && parsedMaxConcurrent > 0) + ? parsedMaxConcurrent + : cronConfig?.maxConcurrentRuns; + + const config: CronManagerConfig = { + ...cronConfig, + enabled: true, + store: resolvedStore, + maxConcurrentRuns: resolvedMaxConcurrent, + maxRunHistory: cronConfig?.maxRunHistory, + redisConfig: cronConfig?.redisConfig, + defaultProvider: cronConfig?.defaultProvider, + defaultModel: cronConfig?.defaultModel, + }; + + this.cronManager = new CronManager(config); + + // Wire up the executor: when a task fires, call this.generate() + this.cronManager.setExecutor(async (task, sessionId) => { + const generateOptions = { + input: { text: task.prompt }, + provider: task.provider || config.defaultProvider, + model: task.model || config.defaultModel, + sessionId, + } as Record; + + const result = await this.generate( + generateOptions as import("./types/generateTypes.js").GenerateOptions, + ); + + return { + responseText: result.content, + tokenUsage: result.usage + ? { + promptTokens: result.usage.input, + completionTokens: result.usage.output, + totalTokens: result.usage.total, + } + : undefined, + }; + }); + + logger.debug("[NeuroLink] CronManager initialized", { + store: config.store, + maxConcurrent: config.maxConcurrentRuns, + }); + } + + /** + * Get the CronManager instance for this NeuroLink instance. + * Used by BaseProvider to bind cron tools to the correct instance scope. + */ + getCronManager(): CronManager | undefined { + return this.cronManager; + } + /** * Initialize Langfuse observability for AI operations tracking */ @@ -2177,6 +2260,15 @@ Current user's request: ${currentInput}`; logger.warn("[NeuroLink] OpenTelemetry shutdown failed:", error); } + if (this.cronManager) { + try { + await this.cronManager.shutdown(); + logger.debug("[NeuroLink] CronManager shutdown completed"); + } catch (error) { + logger.warn("[NeuroLink] CronManager shutdown failed:", error); + } + } + if (this.externalServerManager) { try { await this.externalServerManager.shutdown(); @@ -8871,7 +8963,24 @@ Current user's request: ${currentInput}`; logger.warn("[NeuroLink] Error shutting down OpenTelemetry:", error); } - // 2. Shutdown external MCP server connections + // 2. Shutdown CronManager (timers, store connections) + if (this.cronManager) { + try { + logger.debug("[NeuroLink] Shutting down CronManager..."); + await this.cronManager.shutdown(); + this.cronManager = undefined; + logger.debug("[NeuroLink] CronManager shutdown successfully"); + } catch (error) { + const err = + error instanceof Error + ? error + : new Error(`CronManager shutdown error: ${String(error)}`); + cleanupErrors.push(err); + logger.warn("[NeuroLink] Error shutting down CronManager:", error); + } + } + + // 3. Shutdown external MCP server connections if (this.externalServerManager) { try { logger.debug("[NeuroLink] Shutting down external MCP servers..."); diff --git a/src/lib/types/configTypes.ts b/src/lib/types/configTypes.ts index 44dd9a3bd..66f60352b 100644 --- a/src/lib/types/configTypes.ts +++ b/src/lib/types/configTypes.ts @@ -6,6 +6,7 @@ import { MCPToolRegistry } from "../mcp/toolRegistry.js"; import type { HITLConfig } from "../types/hitlTypes.js"; import type { ConversationMemoryConfig } from "./conversation.js"; +import type { CronManagerConfig } from "../cron/types.js"; import type { ObservabilityConfig } from "./observability.js"; /** @@ -28,6 +29,7 @@ export type NeurolinkConstructorConfig = { conversationMemory?: Partial; enableOrchestration?: boolean; hitl?: HITLConfig; + cron?: CronManagerConfig; toolRegistry?: MCPToolRegistry; observability?: ObservabilityConfig; }; @@ -127,6 +129,8 @@ export type ToolConfig = { maxToolsPerProvider?: number; /** Whether MCP tools should be enabled */ enableMCPTools?: boolean; + /** Whether cron/scheduling tools should be disabled */ + disableCronTools?: boolean; }; /** @@ -215,6 +219,7 @@ export const DEFAULT_CONFIG: NeuroLinkConfig = { allowCustomTools: true, maxToolsPerProvider: 100, enableMCPTools: true, + disableCronTools: false, }, configVersion: "3.0.1", }; diff --git a/src/lib/types/index.ts b/src/lib/types/index.ts index ea247c7bf..d7d330e2c 100644 --- a/src/lib/types/index.ts +++ b/src/lib/types/index.ts @@ -22,6 +22,22 @@ export type { RetryConfig, ToolConfig, } from "./configTypes.js"; +// Cron/Scheduling types +export type { + CronManagerConfig, + CreateTaskOptions, + ScheduledTask, + Schedule, + ScheduleType, + SchedulerBackend, + SessionMode, + TaskExecutor, + TaskFilter, + TaskRunResult, + TaskStatus, + TaskStore, + RunStatus, +} from "../cron/types.js"; // External MCP types export type { ExternalMCPConfigValidation, diff --git a/src/lib/utils/toolUtils.ts b/src/lib/utils/toolUtils.ts index 49fd1c4ea..12652b208 100644 --- a/src/lib/utils/toolUtils.ts +++ b/src/lib/utils/toolUtils.ts @@ -69,3 +69,16 @@ export function getMaxToolsPerProvider(toolConfig?: ToolConfig): number { return 100; // Default } + +/** + * Check if cron/scheduling tools should be disabled + * @param toolConfig - Optional tool configuration + * @returns true if cron tools should be disabled + */ +export function shouldDisableCronTools(toolConfig?: ToolConfig): boolean { + if (toolConfig?.disableCronTools !== undefined) { + return toolConfig.disableCronTools; + } + + return process.env.NEUROLINK_DISABLE_CRON_TOOLS === "true"; +}