import { Queue, Job, QueueOptions, JobOptions, AdvancedSettings } from 'bull'; import BaseService = require('./core'); import Task = require('../support/task'); import ProcessRunner from '../utils/processRunner'; import RedisMessager, { AddMessageOptions, ConsumeMessageOptions } from '../utils/redisMessager'; import { ChainConfig, TokenConfig } from "../models"; import * as http from 'http'; import * as https from 'https'; import * as Router from '@koa/router'; export as namespace services declare interface JobDef { name: string // data数据 fileName: string prefix?: string chainKey?: string // task实例 instance: Task } declare interface JobQueueOptions { tasks?: JobDef[] settings?: AdvancedSettings } declare interface JobToRun { /** 任务队列名称 */ name: string /** 任务子名称名称 */ subName?: string /** 任务数据 */ data?: object /** 任务参数 */ options?: JobOptions } declare class JobQueueService extends BaseService { constructor(services: any) initialize(opts: JobQueueOptions): Promise /** * 获取一个队列 * @param taskName */ fetchQueue(taskName: String, opts?: QueueOptions): Promise /** * 注册任务 * @param tasks */ registerJobQueues(tasks: JobDef[]): Promise /** * 注册任务 * @param task */ registerJobQueue(task: JobDef): Promise /** * 注册方法 * @param name * @param method * @param opts */ registerMethod(name: string, method: Function, opts?: AdvancedSettings): Promise /** * 注册方法 * @param method * @param opts */ registerMethod(method: string, opts?: AdvancedSettings): Promise /** * 启动/重启任务 */ startOrReloadJobs(): Promise /** * 查询正在running的jobs * @param taskName */ runningJobs(taskName: string): Promise /** * 创建循环任务 * @param interval * @param taskName * @param data * @param options */ every(interval: number | string, taskName: string, data?: object, options?: JobOptions): Promise /** * 创建计划任务 * @param when * @param taskName * @param data * @param options */ schedule(when: number | string | Date, taskName: string, data?: object, options?: JobOptions): Promise /** * 创建单次任务 * @param job 任务说明 */ add(job: JobToRun): Promise } interface AddMsgOptions extends AddMessageOptions { /** 预创建的 group */ group?: string /** uid是否要和 group 组合 */ attachGroup?: boolean } interface ConsumeMsgOptions { /** 当 method 为 string 时有用 */ namespace?: string /** uid是否要和 group 组合 */ attachGroup?: boolean /** 同名 msg 处理的间隔 */ sameInterval?: number /** 一批获取多少个 */ msgCount?: number /** 每次获取的等待时间 */ msgIdleTime?: number /** 成功后执行的代码 */ onSucceed?: Function /** 失败后执行的代码 */ onFailure?: Function } declare class MsgQueueService extends BaseService { constructor(services: any) initialize(opts: any): Promise /** * 添加消息 * @param msgKey 消息key * @param msgs 消息内容 * @param opts 参数 */ addMessages(msgKey: string, msgs: object[], opts?: AddMsgOptions): Promise /** * 消费消息 * @param msgKey 消息key * @param group 消费组 * @param method 处理 data 的 method 方法或执行函数 * @param opts 参数 */ consumeMessages(msgKey: string, group: string, method: string | Function, opts?: ConsumeMsgOptions): Promise } declare type ServerCfg = { host: string, http: { port: number, disabled: boolean }, https: { port: number, disabled: boolean, key: string, cert: string, ca: string } } declare interface HttpOptions { server?: ServerCfg listenManually: boolean /** @default 500 */ defaultErrorStatus?: number } declare class HttpService extends BaseService { constructor(services: any) initialize(opts: HttpOptions): Promise listen(): Promise port?: number server?: http.Server portSSL?: number serverSSL?: https.Server } declare interface ExpressOptions extends HttpOptions { routes: (app: any) => void } declare class ExpressService extends HttpService { initialize(opts: ExpressOptions): Promise } declare interface KoaOptions extends HttpOptions { router: Router } declare class KoaService extends HttpService { initialize(opts: KoaOptions): Promise } declare type LocaleData = { code: number | string, category: string, message: string } declare interface LocaleCodeOptions { isHost: Boolean localePath?: string } declare class ErrorCodeService extends BaseService { constructor(services: any) initialize(opts: LocaleCodeOptions): Promise getErrorInfo(code: number, locale: string): Promise } declare class ActivityService extends BaseService { constructor(services: any) initialize(opts: any): Promise /** * 加载并安装本地化文档 * @param localePath */ setupActivityLocales(localePath: string): Promise /** * 获取活动日志 * @param activity 日志实例 * @param locale 本地化码 */ getActivityLog(activity: string | models.ActivityDocument, locale?: string): Promise /** * 添加日志参数 * @param activityId 活动实例 ID * @param params 参数列表 */ pushActivityLogParams(activityId: string, params: string | string[]): Promise /** * 添加 api 日志 * @param moduleName 模块名,global/wallet * @param name 活动名称 * @param operator 操作者 * @param operatorRole 操作角色 * @param input * @param logParams */ startApiLog(moduleName: string, name: string, operatorType: string, operator: string, operatorRole: string, resourceName?: string, resourceId?: string, input?: models.ActivityDataInput, logParams?: string[]): Promise /** * 结束 api 日志 * @param activityId 活动实例 ID * @param output * @param params 参数列表 */ finishApiLog(activityId: string, output?: models.ActivityDataOutput, logParams?: string[]): Promise /** * 添加事件类活动日志 * @param moduleName 模块名,global/wallet * @param name 活动名称 * @param operator 操作者 * @param logParams */ createEventActivity(moduleName: string, name: string, operatorType: string, operator: string, operatorRole: string, resourceName?: string, resourceId?: string, eventContent: models.ActivityEventContent, logParams?: string[]): Promise /** * 解析Operator相关信息的工具方法 * @param isSystemInvoke 是否是系统内部调用 * @param extra 相关信息 */ parseOperatorDetail(isSystemInvoke: boolean, extra: models.ActivityOperatorInfoInput): Promise } declare interface RPCMethodDefine { name: string, func?: Function, encryptResult?: boolean } declare interface JSONRPCOptions { acceptMethods: string[] noAuth?: boolean signerId?: string signer?: Buffer | string verifier?: Buffer | string sort?: string hash?: string encode?: string authWithTimestamp?: boolean withoutTimestamp?: boolean } declare interface JSONRPCServerOptions extends JSONRPCOptions { host?: string port?: number } declare class JSONRpcService extends BaseService { public host: string public port: number constructor(services: any) initialize(opts: JSONRPCServerOptions): Promise /** * 设置可接受rpc方法 * @param acceptMethods 方法名 */ setAcceptMethods(acceptMethods: string[] | RPCMethodDefine[]): void; /** * 添加可接受的RPC方法 * @param methodName 方法名 * @param methodFunc 可选,方法执行的代码 */ addAcceptableMethod(methodName: string, methodFunc?: Function, encryptResult?: boolean): void; /** * 移除可接受的RPC方法 * @param methodName 方法名 */ removeAcceptableMethod(methodName: string): void; /** * 请求RPC地址 * @param ws 请求的客户端 * @param methodName 方法名 * @param args 参数 */ requestJSONRPC(ws: any, methodName: string, args: any): Promise; } declare interface RequestOptions extends JSONRPCOptions { acceptNamespace?: string } declare class JSONRpcClientService extends BaseService { constructor(services: any) initialize(opts: JSONRPCOptions): Promise /** * 设置可接受rpc方法 * @param acceptMethods 方法名 */ setAcceptMethods(acceptMethods: string[]): void; /** * 加入ws RPC服务器 * @param url * @param opts */ joinRPCServer(url: string, opts?: RequestOptions): Promise /** * 关闭ws RPC服务器 * @param url */ closeRPCServer(url: string): Promise /** * 检测连接状态 * @param url */ getClientReadyState(url: string): number /** * 请求JSONRPC */ requestJSONRPC(url: string, methodName: string, args: any, opts?: RequestOptions): Promise } declare interface InternalRPCOptions { namespace: string port?: number } declare class InternalRpcService extends BaseService { constructor(services: any) initialize(opts: InternalRPCOptions): Promise /** * 注册本服务的方法 * @param {String} methodName * @param {Function} [func=undefined] */ registerRPCMethod(method: string, methodFunc?: Function): Promise; /** * 注册一堆本服务的方法 * @param methods 方法定义 */ registerRPCMethods(methods: string[] | RPCMethodDefine[]): Promise; /** * 调用rpc方法 * @param {String} namespace * @param {String} method * @param {any} params */ invokeRPCMethod(namespace: string, method: string, params: any): Promise } declare interface ScriptOptions { onStartup: Function; onExit: Function; } declare class ScriptService extends BaseService { constructor(services: any) initialize(opts: ScriptOptions): Promise /** * 运行脚本 */ runScript(name: string): Promise } declare interface AsyncPlanOptions { processEvery: number; } /** * 该services依赖JobQueueService */ declare class AsyncPlanService extends BaseService { constructor(services: any) initialize(opts: AsyncPlanOptions): Promise; } declare interface StartOptions { /** 进程模式 */ mode?: 'app' | 'task' /** 进程参数 */ param: string /** 进程任务名称 */ task?: string /** 进程负责的多个任务名称 */ jobs?: string | string[] /** 退出进程的超时 */ timeout?: number /** 启动运行路径 */ cwd?: string /** 启动脚本 */ script?: string /** 是否使用最新的 node */ useLatestNode?: boolean /** 启动node附加的指令 */ nodeArgs?: string | string[] /** 是否启动cluster模式 */ cluster?: boolean /** cluster模式下启动的进程数 */ instances?: number /** 合并日志 */ mergeLogs?: boolean } declare interface PM2Desc { name: string worker_id: number pid?: number status: string restarts: number unstable_restarts: number } declare interface ProcessDesc extends PM2Desc { uuid: string, monit: any, uptime: number, env: { JP_MODE?: string, JP_PARAM?: string, JP_TASK?: string, JP_JOBS?: string, NODE_ENV?: string } } /** * pm2服务的参数 */ declare interface Pm2Options { masterKey?: string } /** * 该服务将使用pm2进行进程管理 */ declare class Pm2Service extends BaseService { constructor(services: any) public processPrefix: string; initialize(opts: Pm2Options): Promise; /** 启动进程 */ start(opts: StartOptions): Promise; /** 重启进程 */ restart(nameOrId: string): Promise; /** 关闭进程 */ stop(nameOrId: string, isDelete: boolean): Promise; /** 列出指定信息 */ info(nameOrId: string): Promise; /** 列出信息 */ /** * 列出进程信息 * @param excludeSelfControlled 是否反选 * @param forceIncludeAll 是否列出所有 */ list(excludeSelfControlled?: boolean, forceIncludeAll?: boolean): Promise; } declare class ProcessService extends BaseService { constructor(services: any) initialize(opts: undefined): Promise; /** * 创建或获取一个常驻的子进程 * @param name 子进程唯一别名 * @param execPath 子进程路径 * @param env 子进程环境变量 * @param cwd 可选,运行目录,默认为脚本目录 */ forkNamedProcess(name: string, execPath: string, env: { [key: string]: string }, cwd?: string): ProcessRunner; /** * 重启进程 * @param name 子进程唯一别名 */ restartNamedProcess(name: string): Promise /** * 调用子进程方法(jsonrpc形式) * @param name 子进程唯一别名 * @param method 函数名 * @param params 参数 */ requestProcess(name: string, method: string, params: any): Promise; } declare interface SocketIOOptions { timeout?: number; disableInternal?: boolean; adapter?: { type: string }; } declare class SocketIOService extends BaseService { constructor(services: any) initialize(opts: SocketIOOptions): Promise /** * 调用内部跨进程方法 * @param namespace inteneralKey,通常为chain的Key * @param methodName 方法名 * @param args 参数 */ invokeInternalMethod(namespace: string, methodName: string, args: any): Promise /** * 调用内部跨进程方法 * @param namespace inteneralKey,通常为chain的Key * @param socketId socket的sid,需再次拼接修改为socket.id * @param methodName 方法名 * @param args 参数 */ invokeInternalMethodById(namespace: string, socketId: string, methodName: string, args: any): Promise /** * 广播内部跨进程方法 * @param namespace inteneralKey,通常为chain的Key * @param methodName 方法名 * @paramargs 参数 */ broadcastInternalMethod(namespace: string, methodName: string, args: any): Promise } declare interface SocketIOWorkerOptions { namespaces: string[] } declare class SocketIOWorkerService extends BaseService { constructor(services: any) initialize(opts: SocketIOWorkerOptions): Promise } declare interface ConfigOptions { isHost: boolean } declare type GeneralConfig = { id: string, [key: string]: any } declare class ConfigService extends BaseService { constructor(services: any) initialize(opts: ConfigOptions): Promise // 便捷查询方法 /** * 获取实时的区块链配置 * @param keyOrNameOrCoreType */ loadChainCfg(keyOrNameOrCoreType: string): Promise /** * 读取默认配置中的token配置信息 * @param chainKey * @param tokenNameOrAssetIdOrContract */ loadCoinCfg(chain: string, tokenNameOrAssetIdOrContract: string): Promise /** * 获取实时的全部链名称 */ loadAllChainNames(includeDisabled?: boolean): Promise /** * 获取实时的可用coinNames * @param chainKey */ loadAllCoinNames(chain: string, includeDisabled?: boolean): Promise // 通用方法 /** * 从数据库中读取配置,若该配置不存在,则从文件中读取并保存到数据库 * @param cfgPath 目录名 * @param key 子目录名 * @param parent */ loadConfig(cfgPath: string, key: string, parent?: string): Promise; /** * 从数据库中读取path相同的全部配置,同时也从文件夹中读取全部路径 * @param cfgPath * @param parent */ loadConfigKeys(cfgPath: string, parent?: string, includeDisabled?: boolean): Promise; // 写入型方法无默认实现 /** * 设置是否自动保存 * @param value */ setAutoSaveWhenLoad(value: boolean): Promise; /** * 设置path + key的别名目录 * 对于loadConfig来说,只取最后一个被设置的别名目录 * 对于loadConfigKeys来说,别名目录 + config目录下的结果都将累加到最终结果 * @param path * @param key * @param aliasPath */ setAliasConfigPath(path: string, key: string, aliasPath: string): Promise; /** * 保存配置修改 * @param cfgPath 目录名 * @param key 子目录名 * @param modJson 配置修改Json,需Merge * @param disabled 是否禁用 * @param parent */ saveConfig(cfgPath: string, key: string, modJson?: object, disabled?: boolean, parent?: string): Promise; /** * 从数据库中删除配置,该配置必须是customized的配置 * @param cfgPath 目录名 * @param key 子目录名 * @param parent */ deleteConfig(cfgPath: string, key: string, parent?: string): Promise; } declare interface ConsulOptions { url?: string } declare type KeyValueMeta = { // consul 属性 Key: string, Value: string, // base64 of JSON string // 锁定的 sessionid Session?: string, // locking session id // 一些我也搞不清楚的 consul 参数 Flags: number, LockIndex: number, CreateIndex: number, ModifyIndex: number, // 其他我不知道的 [key: string]: string } declare type ServiceData = { host: string port: number meta?: KeyValueMeta status: 'passing' | 'warning' | 'critical' } type TTLMethodType = string | (() => boolean) | Promise<(() => boolean)> declare class ConsulService extends BaseService { constructor(services: any); initialize(opts: ConsulOptions): Promise /** * 设置 KV 值,监控 * @param key * @param json */ watchKeyValue(key: string, json: Object):Promise /** * 设置 KV 值 * @param key * @param json */ updateKeyValue(key: string, json: Object): Promise /** * 获取 kv 值 * @param key */ getKeyValue(key: string): Promise<{ meta: KeyValueMeta, value: Object }> /** * 注册服务到consul * @param serviceName * @param port * @param meta */ registerService(serviceName: string, port: number, meta?: KeyValueMeta, ttlCheckMethod?: TTLMethodType): Promise /** * 移除服务 * @param serviceName */ deregisterService(serviceName: string): Promise /** * 直接返回consul的services数据 * @param serviceName */ listServices(serviceName: string): Promise /** * 等待到服务发现未知 * @param serviceName * @param timeout 等待超时时间 */ waitForService(serviceName: string, timeout?: number): Promise /** * 获取服务信息 * @param serviceName * @param waitForService 是否等待 */ getServiceData(serviceName: string, waitForService?: boolean): Promise } declare interface ConsulKVMonitorOptions { /** 监控频率 */ interval?: number /** 是否为主服务(处理监控触发内容的服务) */ isHost?: boolean /** 作为主服务时,执行的代码 */ monitorTriggerMethod?: string /** 作为主服务时,可并发执行的数量(默认 5) */ monitorTriggerThreads?: number } declare class ConsulKVMonitorService extends BaseService { constructor(services: any); initialize(opts: ConsulKVMonitorOptions): Promise /** * 注册监控 key * @param key */ registerMonitorKey(key: string): Promise /** * 注销监控 key * @param key */ unregisterMonitorKey(key: string): Promise } declare interface StatServiceInitializeOptions { interval: number isMain: boolean redisOpts: Object } declare type StatsEventData = { host: string processType: string scope: string event: string action: string timestamp: number } declare class StatService extends BaseService { constructor(services: any); initialize(opts: StatServiceInitializeOptions): Promise /** * 记录具体打点事件 * @param scope [enum] consts.STAT_POINT_RECORD_SCOPE * @param event [enum] consts.STAT_POINT_RECORD_EVENT * @param action * @param data * @param timestamp 13位时间戳 * @param host */ recordPoint(scope: string, event: string, action: string, data: string | string[], timestamp: number, host: string): Promise /** * 记录汇总后事件 * @param events [object] */ recordStatsEvent(events: any[]): Promise ackStatsStatistic(start:number,end:number):Promise getRuler(): Promise setRuler(ruler:number):Promise getProcessUids(start:number,end:number):Promise ackProcessUids(start:number,end:number):Promise getEvents(start:number,end:number):Promise ackEvents(start:number,end:number):Promise getFristEvents(index:number):Promise readAllKeys():Promise saveNewKeys(keys:string[]):Promise }