import { LogEventStatus, LogEventTypes } from '../logs/logs.types'; import { AWSSQSConfig, CloudProvisionStatus, GooglePubSubConfig, IActionRequest, IDbActionRequest, IDependencyMatrix, IStepEvent, IStepSequence, IFunctionRequest, INotificationRequest, IRequest, IStorageRequest, KafkaConfig, MessageBrokerTypes, RabbitMQConfig, RedisConfig } from './productsBuilder.types'; export { IActionRequest, INotificationRequest, IDbActionRequest, IStorageRequest, IFunctionRequest }; import { HttpMethods, Notifiers } from './enums'; import ObjectId from 'bson-objectid'; import { IParsedSample } from './inputs.types'; export interface IProcessorInputGeneric { product: string; env: string; tag: string; } export interface IProcessorInputSessionGeneric extends IProcessorInputGeneric { /** Session token in format: session_tag:jwt_token */ session?: string; cache?: string; } export interface IProcessorInput extends IProcessorInputSessionGeneric { input: Record; } export interface IRefreshTokenReq { product_tag: string; env: string; refreshToken: string; identifier: string; start_at: number; end_at: number; session_tag: string; data: string; session_id: string; /** Present when this save is rotating an existing session (refresh flow). * Backend validates this token exists before persisting the new one, atomically. */ previous_refresh_token?: string; } export interface ISessionInput extends IProcessorInputGeneric { data: Record; } export interface ISessionPayload extends IProcessorInputGeneric { token: string; } export interface ISessionRefreshPayload extends IProcessorInputGeneric { refreshToken: string; } export interface ISession extends IProcessorInputGeneric { token: string; } export interface ISessionOutput { token: string; refreshToken: string; } export interface IProcessorSequenceLevels { [level: string]: Array; } export interface IProcessingOutput { success: Array; failure: Array; waiting: Array; skipped: Array; } export interface IProcessingSuccess { event: IStepEvent; output?: Record; } export interface IProcessingFailure extends IProcessingSuccess { allow_fail: boolean; retries_left?: number; retry_at: number; error_code: string | number; reason: string; payload: IActionRequest; } export interface IProcessingWaiting extends IProcessingSuccess { dependants: Array; } export interface IProcessorResult { status: LogEventStatus; /** Execution window start time (Unix ms). Required for list/timeline and Gantt. */ start: number; /** For feature run: 'feature'. For step results: must be 'feature_step'. */ component: LogEventTypes; /** Execution window end time (Unix ms). Required for duration and timeline. */ end: number; retryable: boolean; result: IProcessingOutput; env: string; product_id: string; process_id: string; input: IProcessorInput; feature_id?: string; feature_tag?: string; product_tag?: string; workspace_id?: string; /** Step-level fields when component is feature_step (per-step success/failure log) */ step_tag?: string; /** Step kind when component is feature_step: action, notification, storage, produce, database_action, graph, vector, quota, fallback, child_feature, sleep, wait_for_signal, checkpoint */ step_type?: string; step_error?: string; step_duration_ms?: number; trace_id?: string; span_id?: string; parent_span_id?: string; feature_run_id?: string; step_run_id?: string; step_attempt?: number; function_namespace?: string; function_operation?: string; function_version?: string; function_invocation_id?: string; parent_process_id?: string; session_id?: string; session_tag?: string; session_user_id?: string; execution_kind?: string; } export interface INotificationInput { title: { key: string; }[]; body: { key: string; }[]; data: { key: string; }[]; } export interface INotificationPayload { title?: Record; body?: Record; data?: Record; } export interface IValidateNotification { input: INotificationInput; } export interface INotificationTemplate { title: string; body: string; data?: Record; } export interface IFirebaseCredential { type: string; project_id: string; private_key_id: string; private_key: string; client_email: string; client_id: string; auth_uri: string; token_uri: string; auth_provider_x509_cert_url: string; client_x509_cert_url: string; } export interface ICallbackHandler { url: string; method: HttpMethods; auth?: IActionRequest; } /** * @deprecated Use IEmailHandler from './notifications/types/notifications.types' instead * This interface is kept for backward compatibility but only supports SMTP. * For multi-provider support (SMTP, Mailgun, SendGrid, Postmark, Brevo), use the IEmailHandler from notifications.types.ts */ /** * Email handler configuration - supports multiple providers * @deprecated This interface is kept for backward compatibility. * For new code, use IEmailHandler from './notifications/types/notifications.types' */ export interface IEmailHandler { host?: string; port?: string; sender_email?: string; auth?: { user: string; pass: string; }; secure?: boolean; tls?: { rejectUnauthorized: boolean; }; provider: EmailProvider; smtp?: { host: string; port: string; sender_email: string; auth: { user: string; pass: string; }; secure: boolean; tls?: { rejectUnauthorized: boolean; }; }; mailgun?: { apiKey: string; domain: string; sender_email: string; region?: 'us' | 'eu'; baseUrl?: string; }; sendgrid?: { apiKey: string; sender_email: string; }; postmark?: { serverToken: string; sender_email: string; messageStream?: string; }; brevo?: { apiKey: string; sender_email: string; sender_name?: string; }; } export interface INotificationsHandler { type: Notifiers; credentials?: Object; databaseUrl?: string; cloud?: string; authMode?: 'cloud_connection'; } export declare enum SmsProvider { TWILIO = "twilio", NEXMO = "nexmo", PLIVO = "plivo", OTHER = "other" } export declare enum EmailProvider { SMTP = "smtp", MAILGUN = "mailgun", SENDGRID = "sendgrid", POSTMARK = "postmark", BREVO = "brevo" } export interface ISmsHandler { provider: SmsProvider; accountSid?: string; authToken?: string; apiSecret?: string; apiKey?: string; sender: string; } export interface INotificationEnv { _id?: ObjectId; slug: string; push_notifications?: INotificationsHandler; emails?: IEmailHandler; sms?: ISmsHandler; callbacks?: ICallbackHandler; } export interface IProductNotification { _id?: ObjectId; name: string; tag: string; description: string; envs: Array; messages?: IProductNotificationTemplate[]; } export interface IProductMessageBroker { _id?: ObjectId; name: string; tag: string; description: string; envs: Array; topics: Array; producers?: Array; consumers?: Array; /** Aggregate status while any env is still provisioning in the cloud */ provisionStatus?: CloudProvisionStatus; provisionError?: string; } export interface IProductMessageBrokerTopic { name: string; tag: string; description?: string; queueUrls?: [{ env_slug: string; url: string; }]; sample: Record; data?: IParsedSample[]; idempotent?: boolean; } export interface IMessageBrokerProducer { _id?: ObjectId; tag: string; topic: string; description?: string; created_at?: Date; updated_at?: Date; } export interface IMessageBrokerConsumer { _id?: ObjectId; tag: string; topic: string; description?: string; created_at?: Date; updated_at?: Date; } export interface IMessageBrokerEnvs { slug: string; type: MessageBrokerTypes; config: RedisConfig | GooglePubSubConfig | RabbitMQConfig | KafkaConfig | AWSSQSConfig; /** Set while cloud broker resources are being provisioned asynchronously */ provisionStatus?: CloudProvisionStatus; provisionError?: string; } export interface IProductNotificationTemplate { _id?: ObjectId; name: string; tag: string; description: string; push_notification?: INotificationTemplate; push_notification_data?: Array; email?: { subject: string; template: string; }; email_data?: Array; callback?: IActionRequest; callback_data?: Array; sms?: string; sms_data?: Array; } export interface INotificationData { notifications: Array; email: Array; } /** * Flat input format for simplified action execution. * Keys are automatically resolved to body, params, query, or headers * based on the action's schema definition. * * For conflicting keys (same key in multiple locations), use prefix syntax: * - 'body:id' -> goes to body.id * - 'params:id' -> goes to params.id * - 'query:id' -> goes to query.id * - 'headers:id' -> goes to headers.id * * @example * ```ts * // Flat input - auto-resolved from schema * const input: IFlatInput = { * amount: 1000, // auto-resolves to body.amount * currency: 'usd', // auto-resolves to body.currency * 'params:id': 'user_123' // explicit: goes to params.id * }; * ``` */ export type IFlatInput = Record; /** * Action input can be either structured (IActionRequest) or flat (IFlatInput) * The SDK will automatically resolve flat inputs to structured format */ export type IActionInputType = IActionRequest | IFlatInput; export interface IActionProcessorInput { env: string; product?: string; app: string; cache?: string; /** * Input can be either: * - Structured: { body: {...}, params: {...}, query: {...}, headers: {...} } * - Flat: { key: value, 'prefix:key': value } * * Flat inputs are automatically resolved based on action schema. */ input: IActionInputType; action: string; retries?: number; /** Session token in format: session_tag:jwt_token */ session?: string; /** When set, processAction skips its bootstrap API call (e.g. feature batch prefetch). */ preloadedBootstrap?: unknown; } export interface IJobProcessorInput { env: string; product: string; event: string; cache?: string; retries: number; input: IActionRequest | INotificationRequest | IDbActionRequest | IFunctionRequest | IStorageRequest | Record | IPublishRequest; start_at: number; /** Session token in format: session_tag:jwt_token */ session?: string; repeat?: IJobRepeatOptions; } interface IJobRepeatOptions { cron?: string; every?: number; limit?: number; endDate?: number | string; tz?: string; } export interface IFileReadResult extends IStorageRequest { } export interface IStorageProcessorInput { env: string; product: string; app?: string; cache?: string; input: IStorageRequest; event: string; retries?: number; /** Session token in format: session_tag:jwt_token */ session?: string; /** When set, processStorage skips its bootstrap API call (e.g. feature batch prefetch). */ preloadedBootstrap?: unknown; } export interface IFunctionProcessorInput { env: string; product: string; app?: string; cache?: string; input: IFunctionRequest; event: string; retries?: number; } export interface IPublishRequest extends IRequest { message: Record; } export interface ISubscribeRequest extends IRequest { callback: (message: object) => Promise; } export interface IMessageBrokerPublishInput { env: string; event: string; cache?: string; product: string; input: IPublishRequest; /** Session token in format: session_tag:jwt_token */ session?: string; } export interface IMessageBrokerSubscribeInput { env: string; event: string; product: string; input: ISubscribeRequest; } export interface IDBActionProcessorInput { env: string; product: string; cache?: string; input: IDbActionRequest; event: string; retries?: number; /** Session token in format: session_tag:jwt_token */ session?: string; } export interface INotificationProcessorInput { env: string; product: string; cache?: string; event: string; input: INotificationRequest; retries?: number; /** Session token in format: session_tag:jwt_token */ session?: string; /** When set, processNotification skips its bootstrap API call (e.g. feature batch prefetch). */ preloadedBootstrap?: unknown; } export interface IWebhooks { uuid: string; appEnv: string; productEnv: string; url: string; method: string; created_at?: Date; updated_at?: Date; access_tag: string; version?: string; webhook_tag: string; private_key?: string; active?: boolean; sender_workspace_id: string; receiver_workspace_id: string; app_tag: string; product_tag: string; } export interface IFileURLPayload { url: string; workspace_id: string; type: string; product: string; provider: string; process_id: string; event: string; env: string; size: number; __v?: number; _id?: string; } /** * Schedule options for dispatching operations as jobs */ export interface IDispatchSchedule { /** Start time as Unix timestamp (ms) or ISO date string */ start_at?: number | string; /** Cron expression for recurring jobs (e.g., '0 0 * * *' for daily at midnight) */ cron?: string; /** Interval in milliseconds for recurring jobs */ every?: number; /** Maximum number of times to run (for recurring jobs) */ limit?: number; /** End date as Unix timestamp (ms) or ISO date string */ endDate?: number | string; /** Timezone for cron expressions (e.g., 'America/New_York') */ tz?: string; } /** * Base dispatch input extending the original processor input with scheduling */ export interface IDispatchOptions { /** Number of retries on failure */ retries?: number; /** Schedule configuration */ schedule?: IDispatchSchedule; /** Session token in format: session_tag:jwt_token */ session?: string; /** Cache tag */ cache?: string; } /** * Dispatch input for actions - schedules an action to run as a job */ export interface IActionDispatchInput extends IDispatchOptions { env: string; product: string; app: string; action: string; input: IActionInputType; } /** * Dispatch input for database actions - schedules a DB action to run as a job */ export interface IDBActionDispatchInput extends IDispatchOptions { env: string; product: string; database: string; event: string; input: IDbActionRequest; } /** * Dispatch input for graph actions - schedules a graph action to run as a job */ export interface IGraphActionDispatchInput extends IDispatchOptions { env: string; product: string; graph: string; event: string; input: Record; } /** * Dispatch input for database operations - schedules a DB operation to run as a job */ export interface IDBOperationDispatchInput extends IDispatchOptions { env: string; product: string; database: string; /** Operation type: createOne, findMany, updateOne, deleteMany, etc. */ operation: string; input: Record; } /** * Dispatch input for graph operations - schedules a graph operation to run as a job */ export interface IGraphOperationDispatchInput extends IDispatchOptions { env: string; product: string; graph: string; /** Operation type: createNode, createRelationship, traverse, query, etc. */ operation: string; input: Record; } /** * Dispatch input for notifications - schedules a notification to be sent as a job */ export interface INotificationDispatchInput extends IDispatchOptions { env: string; product: string; notification: string; event: string; input: INotificationRequest; } /** * Dispatch input for storage operations - schedules a storage operation as a job */ export interface IStorageDispatchInput extends IDispatchOptions { env: string; product: string; storage: string; operation: 'upload' | 'download' | 'delete' | 'list' | 'getMetadata'; input: IStorageRequest; } /** * Dispatch input for message broker publish - schedules a publish as a job */ export interface IPublishDispatchInput extends IDispatchOptions { env: string; product: string; broker?: string; event: string; input: IPublishRequest; } /** * Dispatch input for features - schedules a feature to run as a job */ export interface IFeatureDispatchInput extends IDispatchOptions { env: string; product: string; tag: string; input: Record; } /** * Result returned when dispatching a job */ export interface IDispatchResult { /** Unique job ID for tracking */ job_id: string; /** Job status */ status: 'scheduled' | 'queued'; /** Scheduled start time */ scheduled_at: number; /** Whether this is a recurring job */ recurring: boolean; /** Next run time for recurring jobs */ next_run_at?: number; }