mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-08-04 14:57:42 +00:00
refactor: Codex round 17 #3 — move scheduler service to packages/api (TS)
Per CLAUDE.md, new backend code belongs in /packages/api as TypeScript with /api a thin adapter. Extracts the scheduler service (getLimits, engineDeps factory, isOutOfBalance/isRefillEligible, recordScheduleOutcome, markScheduleRunActive, fireScheduleNow, initializeScheduleEngine) into a typed createSchedulesService(deps) factory in packages/api/src/schedules/service.ts. Behavior is preserved verbatim; api-side access (mongoose User/Balance, getAppConfig, ~/models, resolveAgentFireAccess) is injected via deps, and the engine/jobStoreShared singletons now live in the factory closure. api/server/services/Schedules/index.js is now a thin adapter that constructs those deps. Same public exports, so all existing callers are unchanged.
This commit is contained in:
parent
a6784bbd9d
commit
6f9955fe19
3 changed files with 399 additions and 298 deletions
|
|
@ -1,306 +1,29 @@
|
|||
const mongoose = require('mongoose');
|
||||
const {
|
||||
Permissions,
|
||||
PermissionTypes,
|
||||
getRefillEligibilityDate,
|
||||
} = require('librechat-data-provider');
|
||||
const { logger, runAsSystem, tenantStorage } = require('@librechat/data-schemas');
|
||||
const { resolveAgentFireAccess } = require('./access');
|
||||
const {
|
||||
fireSchedule,
|
||||
getBalanceConfig,
|
||||
startScheduleEngine,
|
||||
generateShortLivedToken,
|
||||
buildBalanceUpdateFields,
|
||||
getAppConfigOptionsFromUser,
|
||||
DEFAULT_SCHEDULE_LIMITS,
|
||||
SCHEDULE_FIRE_TOKEN_TTL,
|
||||
SCHEDULE_FIRE_SCOPE,
|
||||
} = require('@librechat/api');
|
||||
const { createSchedulesService } = require('@librechat/api');
|
||||
const { getAppConfig } = require('~/server/services/Config/app');
|
||||
const { resolveAgentFireAccess } = require('./access');
|
||||
const methods = require('~/models');
|
||||
|
||||
/**
|
||||
* Resolves schedule limits, honoring per-principal (role/user) config overrides
|
||||
* when a user is supplied (routes pass req.user, the fire path passes the owner).
|
||||
* @returns {Promise<import('@librechat/api').ScheduleLimits>}
|
||||
*/
|
||||
async function getLimits(user) {
|
||||
const appConfig = user
|
||||
? await getAppConfig(getAppConfigOptionsFromUser(user))
|
||||
: await getAppConfig();
|
||||
const config = appConfig?.interfaceConfig?.schedules;
|
||||
// Disabled config is a hard stop: the engine must not keep firing existing
|
||||
// schedules after an admin turns the feature off.
|
||||
if (config === false) {
|
||||
return { ...DEFAULT_SCHEDULE_LIMITS, enabled: false };
|
||||
}
|
||||
if (config == null || typeof config === 'boolean') {
|
||||
return DEFAULT_SCHEDULE_LIMITS;
|
||||
}
|
||||
return {
|
||||
enabled: config.use !== false,
|
||||
maxPerUser: config.maxPerUser ?? DEFAULT_SCHEDULE_LIMITS.maxPerUser,
|
||||
minIntervalMinutes: config.minIntervalMinutes ?? DEFAULT_SCHEDULE_LIMITS.minIntervalMinutes,
|
||||
autoDisableAfterFailures:
|
||||
config.autoDisableAfterFailures ?? DEFAULT_SCHEDULE_LIMITS.autoDisableAfterFailures,
|
||||
fireConcurrency: config.fireConcurrency ?? DEFAULT_SCHEDULE_LIMITS.fireConcurrency,
|
||||
};
|
||||
}
|
||||
|
||||
const MANUAL_RUN_LEASE_MS = 5 * 60 * 1000;
|
||||
|
||||
// Whether every engine replica observes the same jobs. The standard backend is a
|
||||
// single process (its one engine sees all its jobs), so it defaults true and keeps
|
||||
// full reconciliation; a clustered backend with private in-memory stores passes
|
||||
// false unless Redis-backed (see initializeScheduleEngine / experimental.js).
|
||||
let jobStoreShared = true;
|
||||
|
||||
/**
|
||||
* Whether a refill would top up this zero-credit balance record right now,
|
||||
* mirroring the chat balance check's auto-refill eligibility (record-based).
|
||||
* @param {{ autoRefillEnabled?: boolean, refillAmount?: number, refillIntervalValue?: number, refillIntervalUnit?: import('librechat-data-provider').RefillIntervalUnit, lastRefill?: Date } | null | undefined} record
|
||||
* @returns {boolean}
|
||||
*/
|
||||
function isRefillEligible(record) {
|
||||
if (record?.autoRefillEnabled !== true) {
|
||||
return false;
|
||||
}
|
||||
if (!(typeof record.refillAmount === 'number' && record.refillAmount > 0)) {
|
||||
return false;
|
||||
}
|
||||
const lastRefillDate = new Date(record.lastRefill ?? 0);
|
||||
if (Number.isNaN(lastRefillDate.getTime())) {
|
||||
return true;
|
||||
}
|
||||
// Mirror checkBalanceRecord's fallbacks exactly (interval 0 / 'days' when a
|
||||
// partially-synced record is missing them) so we never pre-skip a record the
|
||||
// interactive chat balance check would have refilled.
|
||||
return (
|
||||
new Date() >=
|
||||
getRefillEligibilityDate(
|
||||
lastRefillDate,
|
||||
record.refillIntervalValue ?? 0,
|
||||
record.refillIntervalUnit ?? 'days',
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
/** @type {import('@librechat/api').ScheduleEngineDeps} */
|
||||
const engineDeps = {
|
||||
const service = createSchedulesService({
|
||||
methods,
|
||||
getLimits,
|
||||
getUserContext: async (userId) => {
|
||||
const user = await mongoose.models.User.findById(userId).select('_id tenantId role').lean();
|
||||
if (user == null) {
|
||||
return null;
|
||||
}
|
||||
return { id: user._id.toString(), tenantId: user.tenantId, role: user.role };
|
||||
},
|
||||
hasScheduleAccess: async (user) => {
|
||||
const role = await methods.getRoleByName(user.role);
|
||||
return role?.permissions?.[PermissionTypes.SCHEDULES]?.[Permissions.USE] === true;
|
||||
},
|
||||
isOutOfBalance: async (user) => {
|
||||
const appConfig = await getAppConfig(getAppConfigOptionsFromUser(user));
|
||||
const balanceConfig = getBalanceConfig(appConfig);
|
||||
if (balanceConfig?.enabled !== true) {
|
||||
return false;
|
||||
}
|
||||
const Balance = mongoose.models.Balance;
|
||||
let record = await Balance.findOne({ user: user.id }).lean();
|
||||
// Initialize/sync the record exactly as the chat's balance middleware would,
|
||||
// so a new user's startBalance is applied before we read it (avoids skipping
|
||||
// a schedule that an interactive chat would have allowed).
|
||||
if (balanceConfig.startBalance != null) {
|
||||
const updateFields = buildBalanceUpdateFields(balanceConfig, record, user.id);
|
||||
if (Object.keys(updateFields).length > 0) {
|
||||
record = await Balance.findOneAndUpdate(
|
||||
{ user: user.id },
|
||||
{ $set: updateFields },
|
||||
{ upsert: true, new: true },
|
||||
).lean();
|
||||
}
|
||||
}
|
||||
const credits = record?.tokenCredits ?? 0;
|
||||
if (credits > 0) {
|
||||
return false;
|
||||
}
|
||||
// At/below zero: an auto-refill user is only spared a pre-skip when a refill
|
||||
// would actually fire now (mirrors the chat balance check's eligibility). If
|
||||
// they aren't eligible yet, or the refill settings are incomplete, pre-skip as
|
||||
// a balance skip — otherwise the zero-credit fire reaches the chat, is rejected
|
||||
// there, and records a generic error that walks the schedule toward
|
||||
// too_many_failures instead of skipped_balance/insufficient_balance.
|
||||
if (balanceConfig.autoRefillEnabled === true && isRefillEligible(record)) {
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
},
|
||||
// Mirrors the loopback chat route's authorization (role AGENTS:USE + resource
|
||||
// VIEW with the manage:agents bypass); shared with the create/update precheck
|
||||
// so the two never diverge.
|
||||
agentAccess: (agentId, user) => resolveAgentFireAccess(agentId, user),
|
||||
resolveFiles: async (fileIds, user) => {
|
||||
const files = await methods.getFiles(
|
||||
{ file_id: { $in: fileIds }, user: user.id },
|
||||
null,
|
||||
'-text',
|
||||
);
|
||||
return (files ?? []).map((file) => ({
|
||||
file_id: file.file_id,
|
||||
filepath: file.filepath,
|
||||
filename: file.filename,
|
||||
type: file.type,
|
||||
height: file.height,
|
||||
width: file.width,
|
||||
source: file.source,
|
||||
}));
|
||||
},
|
||||
mintFireToken: (userId) =>
|
||||
generateShortLivedToken(userId, SCHEDULE_FIRE_TOKEN_TTL, { scope: SCHEDULE_FIRE_SCOPE }),
|
||||
getSelfUrl: () =>
|
||||
process.env.SCHEDULES_SELF_URL ?? `http://127.0.0.1:${process.env.PORT ?? 3080}`,
|
||||
runInTenantContext: (user, fn) =>
|
||||
tenantStorage.run({ tenantId: user.tenantId, userId: user.id }, fn),
|
||||
getJobStatus: async (conversationId) => {
|
||||
const { GenerationJobManager } = require('@librechat/api');
|
||||
const job = await GenerationJobManager.getJobStore()?.getJob(conversationId);
|
||||
return job?.status ?? null;
|
||||
},
|
||||
clearReconciledJob: async (conversationId) => {
|
||||
const { GenerationJobManager } = require('@librechat/api');
|
||||
await GenerationJobManager.getJobStore()?.deleteJob(conversationId);
|
||||
},
|
||||
isJobStoreShared: () => jobStoreShared,
|
||||
// Counted in system scope so the cap is GLOBAL — a per-owner (tenant-scoped)
|
||||
// count would let multiple tenants collectively exceed fireConcurrency.
|
||||
countActiveRunsGlobal: () => runAsSystem(() => methods.countActiveRuns()),
|
||||
};
|
||||
|
||||
/** @type {ReturnType<typeof startScheduleEngine> | undefined} */
|
||||
let engine;
|
||||
|
||||
async function initializeScheduleEngine(options) {
|
||||
if (engine != null) {
|
||||
return engine;
|
||||
}
|
||||
// A clustered backend passes isJobStoreShared=false (unless Redis-backed) so the
|
||||
// reconciler skips job-status checks it can't trust across workers.
|
||||
if (options?.isJobStoreShared != null) {
|
||||
jobStoreShared = options.isJobStoreShared;
|
||||
}
|
||||
// Explicitly build the Schedule/ScheduleRun indexes first — the unique
|
||||
// idempotency index and TTL retention index would otherwise never exist when
|
||||
// MONGO_AUTO_INDEX is disabled (the production default). If this fails the
|
||||
// unique {scheduleId, scheduledFor} guard may be absent, so leave the engine
|
||||
// DISABLED rather than firing without duplicate protection — the app still
|
||||
// runs; schedules simply don't fire until an operator resolves the index.
|
||||
try {
|
||||
await runAsSystem(() => methods.ensureScheduleIndexes());
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
'[schedules] index creation failed — scheduler NOT started (fires need the unique idempotency index):',
|
||||
err,
|
||||
);
|
||||
return undefined;
|
||||
}
|
||||
engine = startScheduleEngine(engineDeps);
|
||||
return engine;
|
||||
}
|
||||
|
||||
/**
|
||||
* Manual run-now fire. Acquires the schedule lease to serialize concurrent
|
||||
* run-now clicks (and to block against a background engine claim), then fires
|
||||
* in manual mode so the next automatic occurrence is left untouched. Returns
|
||||
* null when the lease is already held (a run is in progress).
|
||||
*/
|
||||
async function fireScheduleNow(schedule, limits) {
|
||||
const acquired = await methods.acquireManualRunLease(
|
||||
schedule.id,
|
||||
schedule.user,
|
||||
MANUAL_RUN_LEASE_MS,
|
||||
);
|
||||
if (!acquired) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
return await fireSchedule(engineDeps, schedule, limits, new Date(), { manual: true });
|
||||
} catch (err) {
|
||||
await methods.releaseLease(schedule.id).catch(() => undefined);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
const OUTCOME_RETRY_ATTEMPTS = 3;
|
||||
|
||||
/**
|
||||
* Completion hook: called from the agents controller finalize paths when the
|
||||
* request carried a scheduleId. The caller deletes the job (`completeJob`) right
|
||||
* after, destroying the only evidence the reconciler could use — so a transient
|
||||
* Mongo failure here is RETRIED (bounded) before giving up, and the failure is
|
||||
* surfaced to the caller (returns false) so it can keep the job when it matters.
|
||||
* @returns {Promise<boolean>} true when the outcome was recorded.
|
||||
*/
|
||||
async function recordScheduleOutcome({ scheduleId, scheduledFor, status, conversationId, error }) {
|
||||
if (!scheduleId || !scheduledFor) {
|
||||
return true;
|
||||
}
|
||||
for (let attempt = 1; attempt <= OUTCOME_RETRY_ATTEMPTS; attempt++) {
|
||||
try {
|
||||
// Resolve the owner's limits so auto-disable uses the same per-principal
|
||||
// threshold as the fire path (not the global default).
|
||||
const schedule = await methods.getScheduleById(scheduleId);
|
||||
const owner = schedule ? await engineDeps.getUserContext(schedule.user) : null;
|
||||
const limits = await getLimits(owner ?? undefined);
|
||||
await methods.recordRunOutcome({
|
||||
scheduleId,
|
||||
scheduledFor: new Date(scheduledFor),
|
||||
status,
|
||||
conversationId,
|
||||
error,
|
||||
autoDisableAfterFailures: limits.autoDisableAfterFailures,
|
||||
});
|
||||
return true;
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
`[schedules] failed to record run outcome (attempt ${attempt}/${OUTCOME_RETRY_ATTEMPTS}):`,
|
||||
err,
|
||||
);
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Moves a paused scheduled run back to `started` when its HITL resume claim
|
||||
* succeeds, so overlap/capacity (which key on `started`) count the resuming
|
||||
* generation as active and a second run for the same schedule can't start
|
||||
* concurrently. Best-effort: a failure just leaves it `requires_action` (the
|
||||
* terminal hook still records the outcome).
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
async function markScheduleRunActive(scheduleId, scheduledFor) {
|
||||
if (!scheduleId || !scheduledFor) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await methods.transitionRunStatus(
|
||||
scheduleId,
|
||||
new Date(scheduledFor),
|
||||
'requires_action',
|
||||
'started',
|
||||
);
|
||||
} catch (err) {
|
||||
logger.error('[schedules] failed to mark resumed run active:', err);
|
||||
}
|
||||
}
|
||||
getAppConfig,
|
||||
findUserById: (userId) =>
|
||||
mongoose.models.User.findById(userId).select('_id tenantId role').lean(),
|
||||
findBalance: (userId) => mongoose.models.Balance.findOne({ user: userId }).lean(),
|
||||
upsertBalance: (userId, fields) =>
|
||||
mongoose.models.Balance.findOneAndUpdate(
|
||||
{ user: userId },
|
||||
{ $set: fields },
|
||||
{ upsert: true, new: true },
|
||||
).lean(),
|
||||
resolveAgentFireAccess,
|
||||
});
|
||||
|
||||
module.exports = {
|
||||
getLimits,
|
||||
engineDeps,
|
||||
fireScheduleNow,
|
||||
recordScheduleOutcome,
|
||||
markScheduleRunActive,
|
||||
initializeScheduleEngine,
|
||||
getLimits: service.getLimits,
|
||||
engineDeps: service.engineDeps,
|
||||
fireScheduleNow: service.fireScheduleNow,
|
||||
recordScheduleOutcome: service.recordScheduleOutcome,
|
||||
markScheduleRunActive: service.markScheduleRunActive,
|
||||
initializeScheduleEngine: service.initializeScheduleEngine,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -51,6 +51,7 @@ export * from './prompts';
|
|||
export * from './projects';
|
||||
/* Skills */
|
||||
export * from './schedules';
|
||||
export * from './schedules/service';
|
||||
export * from './skills';
|
||||
export * from './favorites';
|
||||
/* Endpoints */
|
||||
|
|
|
|||
377
packages/api/src/schedules/service.ts
Normal file
377
packages/api/src/schedules/service.ts
Normal file
|
|
@ -0,0 +1,377 @@
|
|||
import { logger, runAsSystem, tenantStorage } from '@librechat/data-schemas';
|
||||
import { getRefillEligibilityDate, Permissions, PermissionTypes } from 'librechat-data-provider';
|
||||
import type { ScheduleMethods, AppConfig, IBalance } from '@librechat/data-schemas';
|
||||
import type { Types } from 'mongoose';
|
||||
import type {
|
||||
ScheduleEngineDeps,
|
||||
ScheduleLimits,
|
||||
ScheduleUserContext,
|
||||
FireableSchedule,
|
||||
FireResult,
|
||||
} from './types';
|
||||
import type { BalanceUpdateFields } from '../types/balance';
|
||||
import type { GetAppConfigOptions } from '../app/service';
|
||||
import { generateShortLivedToken, SCHEDULE_FIRE_SCOPE } from '../crypto/jwt';
|
||||
import { GenerationJobManager } from '../stream/GenerationJobManager';
|
||||
import { buildBalanceUpdateFields } from '../middleware/balance';
|
||||
import { fireSchedule, SCHEDULE_FIRE_TOKEN_TTL } from './fire';
|
||||
import { getAppConfigOptionsFromUser } from '../app/service';
|
||||
import { DEFAULT_SCHEDULE_LIMITS } from './types';
|
||||
import { getBalanceConfig } from '../app/config';
|
||||
import { startScheduleEngine } from './engine';
|
||||
|
||||
/** Recordable terminal/paused run outcome, as accepted by `recordRunOutcome`. */
|
||||
type ScheduleRunOutcomeStatus = Parameters<ScheduleMethods['recordRunOutcome']>[0]['status'];
|
||||
|
||||
export interface RecordScheduleOutcomeInput {
|
||||
scheduleId?: string;
|
||||
scheduledFor?: string | Date;
|
||||
status: ScheduleRunOutcomeStatus;
|
||||
conversationId?: string;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Api-side dependencies the schedules service needs injected: model methods,
|
||||
* config/balance access, and the owner-scoped agent access check. Everything
|
||||
* else (job store, tenant context, token minting) lives in `@librechat/api` and
|
||||
* is imported directly.
|
||||
*/
|
||||
export interface SchedulesServiceDeps {
|
||||
methods: ScheduleMethods & {
|
||||
getRoleByName: (
|
||||
role?: string,
|
||||
) => Promise<{ permissions?: Record<string, Record<string, boolean | undefined>> } | null>;
|
||||
getFiles: (
|
||||
filter: unknown,
|
||||
sort: unknown,
|
||||
select: unknown,
|
||||
) => Promise<Array<{
|
||||
file_id: string;
|
||||
filepath: string;
|
||||
filename: string;
|
||||
type: string;
|
||||
height?: number;
|
||||
width?: number;
|
||||
source: string;
|
||||
}> | null>;
|
||||
};
|
||||
getAppConfig: (options?: GetAppConfigOptions) => Promise<AppConfig | undefined>;
|
||||
findUserById: (
|
||||
userId: string | Types.ObjectId,
|
||||
) => Promise<{ _id: Types.ObjectId; tenantId?: string; role?: string } | null>;
|
||||
findBalance: (userId: string) => Promise<IBalance | null>;
|
||||
upsertBalance: (userId: string, fields: BalanceUpdateFields) => Promise<IBalance | null>;
|
||||
resolveAgentFireAccess: (
|
||||
agentId: string,
|
||||
user: ScheduleUserContext,
|
||||
) => Promise<'ok' | 'missing' | 'forbidden'>;
|
||||
}
|
||||
|
||||
export interface SchedulesService {
|
||||
getLimits: (user?: ScheduleUserContext) => Promise<ScheduleLimits>;
|
||||
engineDeps: ScheduleEngineDeps;
|
||||
fireScheduleNow: (
|
||||
schedule: FireableSchedule,
|
||||
limits: ScheduleLimits,
|
||||
) => Promise<FireResult | null>;
|
||||
recordScheduleOutcome: (input: RecordScheduleOutcomeInput) => Promise<boolean>;
|
||||
markScheduleRunActive: (scheduleId: string, scheduledFor: string | Date) => Promise<void>;
|
||||
initializeScheduleEngine: (options?: {
|
||||
isJobStoreShared?: boolean;
|
||||
}) => Promise<ReturnType<typeof startScheduleEngine> | undefined>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds the scheduler service around api-side dependencies. Each call owns its
|
||||
* own engine singleton and job-store-shared flag, so state never leaks between
|
||||
* instances.
|
||||
*/
|
||||
export function createSchedulesService(deps: SchedulesServiceDeps): SchedulesService {
|
||||
const { methods } = deps;
|
||||
|
||||
/**
|
||||
* Resolves schedule limits, honoring per-principal (role/user) config overrides
|
||||
* when a user is supplied (routes pass req.user, the fire path passes the owner).
|
||||
*/
|
||||
async function getLimits(user?: ScheduleUserContext): Promise<ScheduleLimits> {
|
||||
const appConfig = user
|
||||
? await deps.getAppConfig(getAppConfigOptionsFromUser(user))
|
||||
: await deps.getAppConfig();
|
||||
const config = appConfig?.interfaceConfig?.schedules;
|
||||
// Disabled config is a hard stop: the engine must not keep firing existing
|
||||
// schedules after an admin turns the feature off.
|
||||
if (config === false) {
|
||||
return { ...DEFAULT_SCHEDULE_LIMITS, enabled: false };
|
||||
}
|
||||
if (config == null || typeof config === 'boolean') {
|
||||
return DEFAULT_SCHEDULE_LIMITS;
|
||||
}
|
||||
return {
|
||||
enabled: config.use !== false,
|
||||
maxPerUser: config.maxPerUser ?? DEFAULT_SCHEDULE_LIMITS.maxPerUser,
|
||||
minIntervalMinutes: config.minIntervalMinutes ?? DEFAULT_SCHEDULE_LIMITS.minIntervalMinutes,
|
||||
autoDisableAfterFailures:
|
||||
config.autoDisableAfterFailures ?? DEFAULT_SCHEDULE_LIMITS.autoDisableAfterFailures,
|
||||
fireConcurrency: config.fireConcurrency ?? DEFAULT_SCHEDULE_LIMITS.fireConcurrency,
|
||||
};
|
||||
}
|
||||
|
||||
const MANUAL_RUN_LEASE_MS = 5 * 60 * 1000;
|
||||
|
||||
// Whether every engine replica observes the same jobs. The standard backend is a
|
||||
// single process (its one engine sees all its jobs), so it defaults true and keeps
|
||||
// full reconciliation; a clustered backend with private in-memory stores passes
|
||||
// false unless Redis-backed (see initializeScheduleEngine / experimental.js).
|
||||
let jobStoreShared = true;
|
||||
|
||||
/**
|
||||
* Whether a refill would top up this zero-credit balance record right now,
|
||||
* mirroring the chat balance check's auto-refill eligibility (record-based).
|
||||
*/
|
||||
function isRefillEligible(record: IBalance | null | undefined): boolean {
|
||||
if (record?.autoRefillEnabled !== true) {
|
||||
return false;
|
||||
}
|
||||
if (!(typeof record.refillAmount === 'number' && record.refillAmount > 0)) {
|
||||
return false;
|
||||
}
|
||||
const lastRefillDate = new Date(record.lastRefill ?? 0);
|
||||
if (Number.isNaN(lastRefillDate.getTime())) {
|
||||
return true;
|
||||
}
|
||||
// Mirror checkBalanceRecord's fallbacks exactly (interval 0 / 'days' when a
|
||||
// partially-synced record is missing them) so we never pre-skip a record the
|
||||
// interactive chat balance check would have refilled.
|
||||
return (
|
||||
new Date() >=
|
||||
getRefillEligibilityDate(
|
||||
lastRefillDate,
|
||||
record.refillIntervalValue ?? 0,
|
||||
record.refillIntervalUnit ?? 'days',
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
const engineDeps: ScheduleEngineDeps = {
|
||||
methods,
|
||||
getLimits,
|
||||
getUserContext: async (userId) => {
|
||||
const user = await deps.findUserById(userId);
|
||||
if (user == null) {
|
||||
return null;
|
||||
}
|
||||
return { id: user._id.toString(), tenantId: user.tenantId, role: user.role };
|
||||
},
|
||||
hasScheduleAccess: async (user) => {
|
||||
const role = await methods.getRoleByName(user.role);
|
||||
return role?.permissions?.[PermissionTypes.SCHEDULES]?.[Permissions.USE] === true;
|
||||
},
|
||||
isOutOfBalance: async (user) => {
|
||||
const appConfig = await deps.getAppConfig(getAppConfigOptionsFromUser(user));
|
||||
const balanceConfig = getBalanceConfig(appConfig);
|
||||
if (balanceConfig?.enabled !== true) {
|
||||
return false;
|
||||
}
|
||||
let record = await deps.findBalance(user.id);
|
||||
// Initialize/sync the record exactly as the chat's balance middleware would,
|
||||
// so a new user's startBalance is applied before we read it (avoids skipping
|
||||
// a schedule that an interactive chat would have allowed).
|
||||
if (balanceConfig.startBalance != null) {
|
||||
const updateFields = buildBalanceUpdateFields(balanceConfig, record, user.id);
|
||||
if (Object.keys(updateFields).length > 0) {
|
||||
record = await deps.upsertBalance(user.id, updateFields);
|
||||
}
|
||||
}
|
||||
const credits = record?.tokenCredits ?? 0;
|
||||
if (credits > 0) {
|
||||
return false;
|
||||
}
|
||||
// At/below zero: an auto-refill user is only spared a pre-skip when a refill
|
||||
// would actually fire now (mirrors the chat balance check's eligibility). If
|
||||
// they aren't eligible yet, or the refill settings are incomplete, pre-skip as
|
||||
// a balance skip — otherwise the zero-credit fire reaches the chat, is rejected
|
||||
// there, and records a generic error that walks the schedule toward
|
||||
// too_many_failures instead of skipped_balance/insufficient_balance.
|
||||
if (balanceConfig.autoRefillEnabled === true && isRefillEligible(record)) {
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
},
|
||||
// Mirrors the loopback chat route's authorization (role AGENTS:USE + resource
|
||||
// VIEW with the manage:agents bypass); shared with the create/update precheck
|
||||
// so the two never diverge.
|
||||
agentAccess: (agentId, user) => deps.resolveAgentFireAccess(agentId, user),
|
||||
resolveFiles: async (fileIds, user) => {
|
||||
const files = await methods.getFiles(
|
||||
{ file_id: { $in: fileIds }, user: user.id },
|
||||
null,
|
||||
'-text',
|
||||
);
|
||||
return (files ?? []).map((file) => ({
|
||||
file_id: file.file_id,
|
||||
filepath: file.filepath,
|
||||
filename: file.filename,
|
||||
type: file.type,
|
||||
height: file.height,
|
||||
width: file.width,
|
||||
source: file.source,
|
||||
}));
|
||||
},
|
||||
mintFireToken: (userId) =>
|
||||
generateShortLivedToken(userId, SCHEDULE_FIRE_TOKEN_TTL, { scope: SCHEDULE_FIRE_SCOPE }),
|
||||
getSelfUrl: () =>
|
||||
process.env.SCHEDULES_SELF_URL ?? `http://127.0.0.1:${process.env.PORT ?? 3080}`,
|
||||
runInTenantContext: (user, fn) =>
|
||||
tenantStorage.run({ tenantId: user.tenantId, userId: user.id }, fn),
|
||||
getJobStatus: async (conversationId) => {
|
||||
const job = await GenerationJobManager.getJobStore()?.getJob(conversationId);
|
||||
return job?.status ?? null;
|
||||
},
|
||||
clearReconciledJob: async (conversationId) => {
|
||||
await GenerationJobManager.getJobStore()?.deleteJob(conversationId);
|
||||
},
|
||||
isJobStoreShared: () => jobStoreShared,
|
||||
// Counted in system scope so the cap is GLOBAL — a per-owner (tenant-scoped)
|
||||
// count would let multiple tenants collectively exceed fireConcurrency.
|
||||
countActiveRunsGlobal: () => runAsSystem(() => methods.countActiveRuns()),
|
||||
};
|
||||
|
||||
let engine: ReturnType<typeof startScheduleEngine> | undefined;
|
||||
|
||||
async function initializeScheduleEngine(options?: {
|
||||
isJobStoreShared?: boolean;
|
||||
}): Promise<ReturnType<typeof startScheduleEngine> | undefined> {
|
||||
if (engine != null) {
|
||||
return engine;
|
||||
}
|
||||
// A clustered backend passes isJobStoreShared=false (unless Redis-backed) so the
|
||||
// reconciler skips job-status checks it can't trust across workers.
|
||||
if (options?.isJobStoreShared != null) {
|
||||
jobStoreShared = options.isJobStoreShared;
|
||||
}
|
||||
// Explicitly build the Schedule/ScheduleRun indexes first — the unique
|
||||
// idempotency index and TTL retention index would otherwise never exist when
|
||||
// MONGO_AUTO_INDEX is disabled (the production default). If this fails the
|
||||
// unique {scheduleId, scheduledFor} guard may be absent, so leave the engine
|
||||
// DISABLED rather than firing without duplicate protection — the app still
|
||||
// runs; schedules simply don't fire until an operator resolves the index.
|
||||
try {
|
||||
await runAsSystem(() => methods.ensureScheduleIndexes());
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
'[schedules] index creation failed — scheduler NOT started (fires need the unique idempotency index):',
|
||||
err,
|
||||
);
|
||||
return undefined;
|
||||
}
|
||||
engine = startScheduleEngine(engineDeps);
|
||||
return engine;
|
||||
}
|
||||
|
||||
/**
|
||||
* Manual run-now fire. Acquires the schedule lease to serialize concurrent
|
||||
* run-now clicks (and to block against a background engine claim), then fires
|
||||
* in manual mode so the next automatic occurrence is left untouched. Returns
|
||||
* null when the lease is already held (a run is in progress).
|
||||
*/
|
||||
async function fireScheduleNow(
|
||||
schedule: FireableSchedule,
|
||||
limits: ScheduleLimits,
|
||||
): Promise<FireResult | null> {
|
||||
const acquired = await methods.acquireManualRunLease(
|
||||
schedule.id,
|
||||
schedule.user,
|
||||
MANUAL_RUN_LEASE_MS,
|
||||
);
|
||||
if (!acquired) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
return await fireSchedule(engineDeps, schedule, limits, new Date(), { manual: true });
|
||||
} catch (err) {
|
||||
await methods.releaseLease(schedule.id).catch(() => undefined);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
const OUTCOME_RETRY_ATTEMPTS = 3;
|
||||
|
||||
/**
|
||||
* Completion hook: called from the agents controller finalize paths when the
|
||||
* request carried a scheduleId. The caller deletes the job (`completeJob`) right
|
||||
* after, destroying the only evidence the reconciler could use — so a transient
|
||||
* Mongo failure here is RETRIED (bounded) before giving up, and the failure is
|
||||
* surfaced to the caller (returns false) so it can keep the job when it matters.
|
||||
*/
|
||||
async function recordScheduleOutcome({
|
||||
scheduleId,
|
||||
scheduledFor,
|
||||
status,
|
||||
conversationId,
|
||||
error,
|
||||
}: RecordScheduleOutcomeInput): Promise<boolean> {
|
||||
if (!scheduleId || !scheduledFor) {
|
||||
return true;
|
||||
}
|
||||
for (let attempt = 1; attempt <= OUTCOME_RETRY_ATTEMPTS; attempt++) {
|
||||
try {
|
||||
// Resolve the owner's limits so auto-disable uses the same per-principal
|
||||
// threshold as the fire path (not the global default).
|
||||
const schedule = await methods.getScheduleById(scheduleId);
|
||||
const owner = schedule ? await engineDeps.getUserContext(schedule.user) : null;
|
||||
const limits = await getLimits(owner ?? undefined);
|
||||
await methods.recordRunOutcome({
|
||||
scheduleId,
|
||||
scheduledFor: new Date(scheduledFor),
|
||||
status,
|
||||
conversationId,
|
||||
error,
|
||||
autoDisableAfterFailures: limits.autoDisableAfterFailures,
|
||||
});
|
||||
return true;
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
`[schedules] failed to record run outcome (attempt ${attempt}/${OUTCOME_RETRY_ATTEMPTS}):`,
|
||||
err,
|
||||
);
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Moves a paused scheduled run back to `started` when its HITL resume claim
|
||||
* succeeds, so overlap/capacity (which key on `started`) count the resuming
|
||||
* generation as active and a second run for the same schedule can't start
|
||||
* concurrently. Best-effort: a failure just leaves it `requires_action` (the
|
||||
* terminal hook still records the outcome).
|
||||
*/
|
||||
async function markScheduleRunActive(
|
||||
scheduleId: string,
|
||||
scheduledFor: string | Date,
|
||||
): Promise<void> {
|
||||
if (!scheduleId || !scheduledFor) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await methods.transitionRunStatus(
|
||||
scheduleId,
|
||||
new Date(scheduledFor),
|
||||
'requires_action',
|
||||
'started',
|
||||
);
|
||||
} catch (err) {
|
||||
logger.error('[schedules] failed to mark resumed run active:', err);
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
getLimits,
|
||||
engineDeps,
|
||||
fireScheduleNow,
|
||||
recordScheduleOutcome,
|
||||
markScheduleRunActive,
|
||||
initializeScheduleEngine,
|
||||
};
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue