mirror of
https://github.com/danny-avila/LibreChat.git
synced 2025-12-22 19:30:15 +01:00
✨feat: OAuth for Actions (#5693)
* ✨feat: OAuth for Actions * WIP: PoC flow state manager * refactor: Add identifier field to token model from action schema * chore: fix potential file type issues * ci: fix type issue with action metadata auth * fix: ensure FlowManagerOptions has a default ttl value * WIP: OAUTH actions * WIP: first pass OAuth Action * fix: standardize identifier usage in OAuth flow handling * fix: update token retrieval to include userId in query and use correct identifier * refacotr: update token retrieval to use userId for OAuth token query * feat: Tool Call Auth styling * fix: streamline token creation and add type field to token schema * refactor: cleanup OAuth flow by encrypting client credentials and ensuring oauth operations only run under condition * refactor: use encrypted credentials in OAuth callback * fix: update Token collection indexes to use expiresAt TTL index and not createdAt legacy index * refactor: enhance Token index cleanup by improving logging and removing redundant index creation logic * refactor: remove unused OAuth login route and related logic for improved clarity * refactor: replace fetch with axios for OAuth token exchange and improve error handling * refactor: better UX after authentication before oauth tool execution * refactor: implement cleanup handlers for FlowStateManager intervals to enhance resource management * refactor: encrypt OAuth tokens before storing and decrypt upon retrieval for enhanced security * refactor: enhance authentication success page with improved styling and countdown feature * refactor: add response_type parameter to OAuth redirect URI for improved compatibility * chore: update translation.json new localizations * chore: remove unused OGDialog import from OGDialogTemplate component * refactor: Actions Auth using new Dialog styling, use same component with Agents/Assistants * refactor: update removeNullishValues function to support removal of empty strings and adjust transform usage in schemas * chore: bump version of librechat-data-provider to 0.7.6991 * refactor: integrate removeNullishValues function to clean metadata before encryption in agent and assistant routes * refactor: update OAuth input fields to use 'password' type for better security * refactor: update localization placeholders for sign-in message to use double curly braces * refactor: add access_type parameter for offline access in createActionTool function * refactor: implement handleOAuthToken function for token management and encryption * feat: refresh token support * refactor: add default expiration for access token and error handling for missing token * feat: localizations for ActionAuth * refactor: set refresh token expiration to null to not expire if expiry never given * fix: prevent crash fromerror within async handleAbortError in AskController, EditController, and AgentController * feat: Action Callback URL * 🌍 i18n: Update translation.json with latest translations * refactor: handle errors in flow state checking to prevent unhandled promise rejections * fix: improve flow state concurrency to prevent multiple token creation calls * refactor: RequestExecutor to use separate axios instance * refactor: improve concurrency flows by keeping completed state until TTL expiry * refactor: increase TTL for flow state management and adjust monitoring interval * ci: mock axios instance creation in actions spec * feat: add Babel and Jest configuration files; implement FlowStateManager tests with concurrency handling * chore: add disableOAuth prop to ActionsAuth (not implemented for Assistants yet) --------- Co-authored-by: Danny Avila <danny@librechat.ai> Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com>
This commit is contained in:
parent
71c30a3640
commit
d99a9db3f6
58 changed files with 2146 additions and 1223 deletions
|
|
@ -1,4 +1,4 @@
|
|||
module.exports = {
|
||||
export default {
|
||||
collectCoverageFrom: ['src/**/*.{js,jsx,ts,tsx}', '!<rootDir>/node_modules/'],
|
||||
coveragePathIgnorePatterns: ['/node_modules/', '/dist/'],
|
||||
coverageReporters: ['text', 'cobertura'],
|
||||
|
|
@ -15,4 +15,5 @@ module.exports = {
|
|||
// },
|
||||
// },
|
||||
restoreMocks: true,
|
||||
};
|
||||
testTimeout: 15000,
|
||||
};
|
||||
|
|
@ -73,5 +73,8 @@
|
|||
"diff": "^7.0.0",
|
||||
"eventsource": "^3.0.1",
|
||||
"express": "^4.21.2"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"keyv": "^4.5.4"
|
||||
}
|
||||
}
|
||||
|
|
|
|||
152
packages/mcp/src/flow/manager.spec.ts
Normal file
152
packages/mcp/src/flow/manager.spec.ts
Normal file
|
|
@ -0,0 +1,152 @@
|
|||
import { FlowStateManager } from './manager';
|
||||
import Keyv from 'keyv';
|
||||
import type { FlowState } from './types';
|
||||
|
||||
// Create a mock class without extending Keyv
|
||||
class MockKeyv {
|
||||
private store: Map<string, FlowState<string>>;
|
||||
|
||||
constructor() {
|
||||
this.store = new Map();
|
||||
}
|
||||
|
||||
async get(key: string): Promise<FlowState<string> | undefined> {
|
||||
return this.store.get(key);
|
||||
}
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-unused-vars
|
||||
async set(key: string, value: FlowState<string>, _ttl?: number): Promise<true> {
|
||||
this.store.set(key, value);
|
||||
return true;
|
||||
}
|
||||
|
||||
async delete(key: string): Promise<boolean> {
|
||||
return this.store.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
describe('FlowStateManager', () => {
|
||||
let flowManager: FlowStateManager<string>;
|
||||
let store: MockKeyv;
|
||||
|
||||
beforeEach(() => {
|
||||
store = new MockKeyv();
|
||||
// Type assertion here since we know our mock implements the necessary methods
|
||||
flowManager = new FlowStateManager(store as unknown as Keyv, { ttl: 30000, ci: true });
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
jest.clearAllMocks();
|
||||
});
|
||||
|
||||
describe('Concurrency Tests', () => {
|
||||
it('should handle concurrent flow creation and return same result', async () => {
|
||||
const flowId = 'test-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
// Start two concurrent flow creations
|
||||
const flow1Promise = flowManager.createFlowWithHandler(flowId, type, async () => {
|
||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
||||
return 'result';
|
||||
});
|
||||
|
||||
const flow2Promise = flowManager.createFlowWithHandler(flowId, type, async () => {
|
||||
await new Promise((resolve) => setTimeout(resolve, 50));
|
||||
return 'different-result';
|
||||
});
|
||||
|
||||
// Both should resolve to the same result from the first handler
|
||||
const [result1, result2] = await Promise.all([flow1Promise, flow2Promise]);
|
||||
|
||||
expect(result1).toBe('result');
|
||||
expect(result2).toBe('result');
|
||||
});
|
||||
|
||||
it('should handle flow timeout correctly', async () => {
|
||||
const flowId = 'timeout-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
// Create flow with very short TTL
|
||||
const shortTtlManager = new FlowStateManager(store as unknown as Keyv, {
|
||||
ttl: 100,
|
||||
ci: true,
|
||||
});
|
||||
|
||||
const flowPromise = shortTtlManager.createFlow(flowId, type);
|
||||
|
||||
await expect(flowPromise).rejects.toThrow('test-type flow timed out');
|
||||
});
|
||||
|
||||
it('should maintain flow state consistency under high concurrency', async () => {
|
||||
const flowId = 'concurrent-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
// Create multiple concurrent operations
|
||||
const operations = [];
|
||||
for (let i = 0; i < 10; i++) {
|
||||
operations.push(
|
||||
flowManager.createFlowWithHandler(flowId, type, async () => {
|
||||
await new Promise((resolve) => setTimeout(resolve, Math.random() * 50));
|
||||
return `result-${i}`;
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
// All operations should resolve to the same result
|
||||
const results = await Promise.all(operations);
|
||||
const firstResult = results[0];
|
||||
results.forEach((result: string) => {
|
||||
expect(result).toBe(firstResult);
|
||||
});
|
||||
});
|
||||
|
||||
it('should handle race conditions in flow completion', async () => {
|
||||
const flowId = 'test-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
// Create initial flow
|
||||
const flowPromise = flowManager.createFlow(flowId, type);
|
||||
|
||||
// Increase delay to ensure flow is properly created
|
||||
await new Promise((resolve) => setTimeout(resolve, 500));
|
||||
|
||||
// Complete the flow
|
||||
await flowManager.completeFlow(flowId, type, 'result1');
|
||||
|
||||
const result = await flowPromise;
|
||||
expect(result).toBe('result1');
|
||||
}, 15000);
|
||||
|
||||
it('should handle concurrent flow monitoring', async () => {
|
||||
const flowId = 'test-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
// Create initial flow
|
||||
const flowPromise = flowManager.createFlow(flowId, type);
|
||||
|
||||
// Increase delay
|
||||
await new Promise((resolve) => setTimeout(resolve, 500));
|
||||
|
||||
// Complete the flow
|
||||
await flowManager.completeFlow(flowId, type, 'success');
|
||||
|
||||
const result = await flowPromise;
|
||||
expect(result).toBe('success');
|
||||
}, 15000);
|
||||
|
||||
it('should handle concurrent success and failure attempts', async () => {
|
||||
const flowId = 'race-flow';
|
||||
const type = 'test-type';
|
||||
|
||||
const flowPromise = flowManager.createFlow(flowId, type);
|
||||
|
||||
// Increase delay
|
||||
await new Promise((resolve) => setTimeout(resolve, 500));
|
||||
|
||||
// Fail the flow
|
||||
await flowManager.failFlow(flowId, type, new Error('failure'));
|
||||
|
||||
await expect(flowPromise).rejects.toThrow('failure');
|
||||
}, 15000);
|
||||
});
|
||||
});
|
||||
241
packages/mcp/src/flow/manager.ts
Normal file
241
packages/mcp/src/flow/manager.ts
Normal file
|
|
@ -0,0 +1,241 @@
|
|||
import Keyv from 'keyv';
|
||||
import type { Logger } from 'winston';
|
||||
import type { FlowState, FlowMetadata, FlowManagerOptions } from './types';
|
||||
|
||||
export class FlowStateManager<T = unknown> {
|
||||
private keyv: Keyv;
|
||||
private ttl: number;
|
||||
private logger: Logger;
|
||||
private intervals: Set<NodeJS.Timeout>;
|
||||
|
||||
private static getDefaultLogger(): Logger {
|
||||
return {
|
||||
error: console.error,
|
||||
warn: console.warn,
|
||||
info: console.info,
|
||||
debug: console.debug,
|
||||
} as Logger;
|
||||
}
|
||||
|
||||
constructor(store: Keyv, options?: FlowManagerOptions) {
|
||||
if (!options) {
|
||||
options = { ttl: 60000 * 3 };
|
||||
}
|
||||
const { ci = false, ttl, logger } = options;
|
||||
|
||||
if (!ci && !(store instanceof Keyv)) {
|
||||
throw new Error('Invalid store provided to FlowStateManager');
|
||||
}
|
||||
|
||||
this.ttl = ttl;
|
||||
this.keyv = store;
|
||||
this.logger = logger || FlowStateManager.getDefaultLogger();
|
||||
this.intervals = new Set();
|
||||
this.setupCleanupHandlers();
|
||||
}
|
||||
|
||||
private setupCleanupHandlers() {
|
||||
const cleanup = () => {
|
||||
this.logger.info('Cleaning up FlowStateManager intervals...');
|
||||
this.intervals.forEach((interval) => clearInterval(interval));
|
||||
this.intervals.clear();
|
||||
process.exit(0);
|
||||
};
|
||||
|
||||
process.on('SIGTERM', cleanup);
|
||||
process.on('SIGINT', cleanup);
|
||||
process.on('SIGQUIT', cleanup);
|
||||
process.on('SIGHUP', cleanup);
|
||||
}
|
||||
|
||||
private getFlowKey(flowId: string, type: string): string {
|
||||
return `${type}:${flowId}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new flow and waits for its completion
|
||||
*/
|
||||
async createFlow(flowId: string, type: string, metadata: FlowMetadata = {}): Promise<T> {
|
||||
const flowKey = this.getFlowKey(flowId, type);
|
||||
|
||||
let existingState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
if (existingState) {
|
||||
this.logger.debug(`[${flowKey}] Flow already exists`);
|
||||
return this.monitorFlow(flowKey, type);
|
||||
}
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 250));
|
||||
|
||||
existingState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
if (existingState) {
|
||||
this.logger.debug(`[${flowKey}] Flow exists on 2nd check`);
|
||||
return this.monitorFlow(flowKey, type);
|
||||
}
|
||||
|
||||
const initialState: FlowState = {
|
||||
type,
|
||||
status: 'PENDING',
|
||||
metadata,
|
||||
createdAt: Date.now(),
|
||||
};
|
||||
|
||||
this.logger.debug('Creating initial flow state:', flowKey);
|
||||
await this.keyv.set(flowKey, initialState, this.ttl);
|
||||
return this.monitorFlow(flowKey, type);
|
||||
}
|
||||
|
||||
private monitorFlow(flowKey: string, type: string): Promise<T> {
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
const checkInterval = 2000;
|
||||
let elapsedTime = 0;
|
||||
|
||||
const intervalId = setInterval(async () => {
|
||||
try {
|
||||
const flowState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
|
||||
if (!flowState) {
|
||||
clearInterval(intervalId);
|
||||
this.intervals.delete(intervalId);
|
||||
this.logger.error(`[${flowKey}] Flow state not found`);
|
||||
reject(new Error(`${type} Flow state not found`));
|
||||
return;
|
||||
}
|
||||
|
||||
if (flowState.status !== 'PENDING') {
|
||||
clearInterval(intervalId);
|
||||
this.intervals.delete(intervalId);
|
||||
this.logger.debug(`[${flowKey}] Flow completed`);
|
||||
|
||||
if (flowState.status === 'COMPLETED' && flowState.result !== undefined) {
|
||||
resolve(flowState.result);
|
||||
} else if (flowState.status === 'FAILED') {
|
||||
await this.keyv.delete(flowKey);
|
||||
reject(new Error(flowState.error ?? `${type} flow failed`));
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
elapsedTime += checkInterval;
|
||||
if (elapsedTime >= this.ttl) {
|
||||
clearInterval(intervalId);
|
||||
this.intervals.delete(intervalId);
|
||||
this.logger.error(
|
||||
`[${flowKey}] Flow timed out | Elapsed time: ${elapsedTime} | TTL: ${this.ttl}`,
|
||||
);
|
||||
await this.keyv.delete(flowKey);
|
||||
reject(new Error(`${type} flow timed out`));
|
||||
}
|
||||
this.logger.debug(
|
||||
`[${flowKey}] Flow state elapsed time: ${elapsedTime}, checking again...`,
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.error(`[${flowKey}] Error checking flow state:`, error);
|
||||
clearInterval(intervalId);
|
||||
this.intervals.delete(intervalId);
|
||||
reject(error);
|
||||
}
|
||||
}, checkInterval);
|
||||
|
||||
this.intervals.add(intervalId);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Completes a flow successfully
|
||||
*/
|
||||
async completeFlow(flowId: string, type: string, result: T): Promise<boolean> {
|
||||
const flowKey = this.getFlowKey(flowId, type);
|
||||
const flowState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
|
||||
if (!flowState) {
|
||||
return false;
|
||||
}
|
||||
|
||||
const updatedState: FlowState<T> = {
|
||||
...flowState,
|
||||
status: 'COMPLETED',
|
||||
result,
|
||||
completedAt: Date.now(),
|
||||
};
|
||||
|
||||
await this.keyv.set(flowKey, updatedState, this.ttl);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Marks a flow as failed
|
||||
*/
|
||||
async failFlow(flowId: string, type: string, error: Error | string): Promise<boolean> {
|
||||
const flowKey = this.getFlowKey(flowId, type);
|
||||
const flowState = (await this.keyv.get(flowKey)) as FlowState | undefined;
|
||||
|
||||
if (!flowState) {
|
||||
return false;
|
||||
}
|
||||
|
||||
const updatedState: FlowState = {
|
||||
...flowState,
|
||||
status: 'FAILED',
|
||||
error: error instanceof Error ? error.message : error,
|
||||
failedAt: Date.now(),
|
||||
};
|
||||
|
||||
await this.keyv.set(flowKey, updatedState, this.ttl);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets current flow state
|
||||
*/
|
||||
async getFlowState(flowId: string, type: string): Promise<FlowState<T> | null> {
|
||||
const flowKey = this.getFlowKey(flowId, type);
|
||||
return this.keyv.get(flowKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new flow and waits for its completion, only executing the handler if no existing flow is found
|
||||
* @param flowId - The ID of the flow
|
||||
* @param type - The type of flow
|
||||
* @param handler - Async function to execute if no existing flow is found
|
||||
* @param metadata - Optional metadata for the flow
|
||||
*/
|
||||
async createFlowWithHandler(
|
||||
flowId: string,
|
||||
type: string,
|
||||
handler: () => Promise<T>,
|
||||
metadata: FlowMetadata = {},
|
||||
): Promise<T> {
|
||||
const flowKey = this.getFlowKey(flowId, type);
|
||||
let existingState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
if (existingState) {
|
||||
this.logger.debug(`[${flowKey}] Flow already exists`);
|
||||
return this.monitorFlow(flowKey, type);
|
||||
}
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 250));
|
||||
|
||||
existingState = (await this.keyv.get(flowKey)) as FlowState<T> | undefined;
|
||||
if (existingState) {
|
||||
this.logger.debug(`[${flowKey}] Flow exists on 2nd check`);
|
||||
return this.monitorFlow(flowKey, type);
|
||||
}
|
||||
|
||||
const initialState: FlowState = {
|
||||
type,
|
||||
status: 'PENDING',
|
||||
metadata,
|
||||
createdAt: Date.now(),
|
||||
};
|
||||
this.logger.debug(`[${flowKey}] Creating initial flow state`);
|
||||
await this.keyv.set(flowKey, initialState, this.ttl);
|
||||
|
||||
try {
|
||||
const result = await handler();
|
||||
await this.completeFlow(flowId, type, result);
|
||||
return result;
|
||||
} catch (error) {
|
||||
await this.failFlow(flowId, type, error instanceof Error ? error : new Error(String(error)));
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
23
packages/mcp/src/flow/types.ts
Normal file
23
packages/mcp/src/flow/types.ts
Normal file
|
|
@ -0,0 +1,23 @@
|
|||
import type { Logger } from 'winston';
|
||||
export type FlowStatus = 'PENDING' | 'COMPLETED' | 'FAILED';
|
||||
|
||||
export interface FlowMetadata {
|
||||
[key: string]: unknown;
|
||||
}
|
||||
|
||||
export interface FlowState<T = unknown> {
|
||||
type: string;
|
||||
status: FlowStatus;
|
||||
metadata: FlowMetadata;
|
||||
createdAt: number;
|
||||
result?: T;
|
||||
error?: string;
|
||||
completedAt?: number;
|
||||
failedAt?: number;
|
||||
}
|
||||
|
||||
export interface FlowManagerOptions {
|
||||
ttl: number;
|
||||
ci?: boolean;
|
||||
logger?: Logger;
|
||||
}
|
||||
|
|
@ -1,4 +1,7 @@
|
|||
/* MCP */
|
||||
export * from './manager';
|
||||
/* Flow */
|
||||
export * from './flow/manager';
|
||||
/* types */
|
||||
export type * from './types/mcp';
|
||||
export type * from './flow/types';
|
||||
|
|
|
|||
|
|
@ -18,7 +18,7 @@
|
|||
"isolatedModules": true,
|
||||
"noEmit": true,
|
||||
"sourceMap": true,
|
||||
"baseUrl": "." // This should be the root of your package
|
||||
"baseUrl": "."
|
||||
},
|
||||
"ts-node": {
|
||||
"experimentalSpecifierResolution": "node",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue