mirror of
https://github.com/bknd-io/bknd/
synced 2026-08-03 08:36:01 +00:00
initial refactor
This commit is contained in:
@@ -1,15 +1,16 @@
|
||||
import { type Static, transformObject } from "core/utils";
|
||||
import { transformObject } from "core/utils";
|
||||
import { Flow, HttpTrigger } from "flows";
|
||||
import { Hono } from "hono";
|
||||
import { Module } from "modules/Module";
|
||||
import { TASKS, flowsConfigSchema } from "./flows-schema";
|
||||
import type { s } from "core/object/schema";
|
||||
|
||||
export type AppFlowsSchema = Static<typeof flowsConfigSchema>;
|
||||
export type AppFlowsSchema = s.Static<typeof flowsConfigSchema>;
|
||||
export type TAppFlowSchema = AppFlowsSchema["flows"][number];
|
||||
export type TAppFlowTriggerSchema = TAppFlowSchema["trigger"];
|
||||
export type { TAppFlowTaskSchema } from "./flows-schema";
|
||||
|
||||
export class AppFlows extends Module<typeof flowsConfigSchema> {
|
||||
export class AppFlows extends Module<AppFlowsSchema> {
|
||||
private flows: Record<string, Flow> = {};
|
||||
|
||||
getSchema() {
|
||||
@@ -80,6 +81,8 @@ export class AppFlows extends Module<typeof flowsConfigSchema> {
|
||||
this.setBuilt();
|
||||
}
|
||||
|
||||
// @todo: fix this
|
||||
// @ts-expect-error
|
||||
override toJSON() {
|
||||
return {
|
||||
...this.config,
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import { Const, type Static, StringRecord, transformObject } from "core/utils";
|
||||
import { transformObject } from "core/utils";
|
||||
import { TaskMap, TriggerMap } from "flows";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
export const TASKS = {
|
||||
...TaskMap,
|
||||
@@ -10,77 +9,59 @@ export const TASKS = {
|
||||
export const TRIGGERS = TriggerMap;
|
||||
|
||||
const taskSchemaObject = transformObject(TASKS, (task, name) => {
|
||||
return Type.Object(
|
||||
return s.strictObject(
|
||||
{
|
||||
type: Const(name),
|
||||
type: s.literal(name),
|
||||
params: task.cls.schema,
|
||||
},
|
||||
{ title: String(name), additionalProperties: false },
|
||||
{ title: String(name) },
|
||||
);
|
||||
});
|
||||
const taskSchema = Type.Union(Object.values(taskSchemaObject));
|
||||
export type TAppFlowTaskSchema = Static<typeof taskSchema>;
|
||||
const taskSchema = s.anyOf(Object.values(taskSchemaObject));
|
||||
export type TAppFlowTaskSchema = s.Static<typeof taskSchema>;
|
||||
|
||||
const triggerSchemaObject = transformObject(TRIGGERS, (trigger, name) => {
|
||||
return Type.Object(
|
||||
return s.strictObject(
|
||||
{
|
||||
type: Const(name),
|
||||
config: trigger.cls.schema,
|
||||
type: s.literal(name),
|
||||
config: trigger.cls.schema.optional(),
|
||||
},
|
||||
{ title: String(name), additionalProperties: false },
|
||||
{ title: String(name) },
|
||||
);
|
||||
});
|
||||
const triggerSchema = s.anyOf(Object.values(triggerSchemaObject));
|
||||
export type TAppFlowTriggerSchema = s.Static<typeof triggerSchema>;
|
||||
|
||||
const connectionSchema = Type.Object({
|
||||
source: Type.String(),
|
||||
target: Type.String(),
|
||||
config: Type.Object(
|
||||
{
|
||||
condition: Type.Optional(
|
||||
Type.Union([
|
||||
Type.Object(
|
||||
{ type: Const("success") },
|
||||
{ additionalProperties: false, title: "success" },
|
||||
),
|
||||
Type.Object(
|
||||
{ type: Const("error") },
|
||||
{ additionalProperties: false, title: "error" },
|
||||
),
|
||||
Type.Object(
|
||||
{ type: Const("matches"), path: Type.String(), value: Type.String() },
|
||||
{ additionalProperties: false, title: "matches" },
|
||||
),
|
||||
]),
|
||||
),
|
||||
max_retries: Type.Optional(Type.Number()),
|
||||
},
|
||||
{ default: {}, additionalProperties: false },
|
||||
),
|
||||
const connectionSchema = s.strictObject({
|
||||
source: s.string(),
|
||||
target: s.string(),
|
||||
config: s
|
||||
.strictObject({
|
||||
condition: s.anyOf([
|
||||
s.strictObject({ type: s.literal("success") }, { title: "success" }),
|
||||
s.strictObject({ type: s.literal("error") }, { title: "error" }),
|
||||
s.strictObject(
|
||||
{ type: s.literal("matches"), path: s.string(), value: s.string() },
|
||||
{ title: "matches" },
|
||||
),
|
||||
]),
|
||||
max_retries: s.number(),
|
||||
})
|
||||
.partial(),
|
||||
});
|
||||
|
||||
// @todo: rework to have fixed ids per task and connections (and preferrably arrays)
|
||||
// causes issues with canvas
|
||||
export const flowSchema = Type.Object(
|
||||
{
|
||||
trigger: Type.Union(Object.values(triggerSchemaObject)),
|
||||
tasks: Type.Optional(StringRecord(Type.Union(Object.values(taskSchemaObject)))),
|
||||
connections: Type.Optional(StringRecord(connectionSchema)),
|
||||
start_task: Type.Optional(Type.String()),
|
||||
responding_task: Type.Optional(Type.String()),
|
||||
},
|
||||
{
|
||||
additionalProperties: false,
|
||||
},
|
||||
);
|
||||
export type TAppFlowSchema = Static<typeof flowSchema>;
|
||||
export const flowSchema = s.strictObject({
|
||||
trigger: s.anyOf(Object.values(triggerSchemaObject)),
|
||||
tasks: s.record(s.anyOf(Object.values(taskSchemaObject))).optional(),
|
||||
connections: s.record(connectionSchema).optional(),
|
||||
start_task: s.string().optional(),
|
||||
responding_task: s.string().optional(),
|
||||
});
|
||||
export type TAppFlowSchema = s.Static<typeof flowSchema>;
|
||||
|
||||
export const flowsConfigSchema = Type.Object(
|
||||
{
|
||||
basepath: Type.String({ default: "/api/flows" }),
|
||||
flows: StringRecord(flowSchema, { default: {} }),
|
||||
},
|
||||
{
|
||||
default: {},
|
||||
additionalProperties: false,
|
||||
},
|
||||
);
|
||||
export const flowsConfigSchema = s.strictObject({
|
||||
basepath: s.string({ default: "/api/flows" }),
|
||||
flows: s.record(flowSchema, { default: {} }),
|
||||
});
|
||||
|
||||
@@ -2,19 +2,15 @@ import type { EventManager } from "core/events";
|
||||
import type { Flow } from "../Flow";
|
||||
import { Trigger } from "./Trigger";
|
||||
import { $console } from "core";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
export class EventTrigger extends Trigger<typeof EventTrigger.schema> {
|
||||
override type = "event";
|
||||
|
||||
static override schema = Type.Composite([
|
||||
Trigger.schema,
|
||||
Type.Object({
|
||||
event: Type.String(),
|
||||
// add match
|
||||
}),
|
||||
]);
|
||||
static override schema = s.strictObject({
|
||||
event: s.string(),
|
||||
...Trigger.schema.properties,
|
||||
});
|
||||
|
||||
override async register(flow: Flow, emgr: EventManager<any>) {
|
||||
if (!emgr.eventExists(this.config.event)) {
|
||||
|
||||
@@ -1,23 +1,19 @@
|
||||
import { StringEnum } from "core/utils";
|
||||
import type { Context, Hono } from "hono";
|
||||
import type { Flow } from "../Flow";
|
||||
import { Trigger } from "./Trigger";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
const httpMethods = ["GET", "POST", "PUT", "PATCH", "DELETE"] as const;
|
||||
|
||||
export class HttpTrigger extends Trigger<typeof HttpTrigger.schema> {
|
||||
override type = "http";
|
||||
|
||||
static override schema = Type.Composite([
|
||||
Trigger.schema,
|
||||
Type.Object({
|
||||
path: Type.String({ pattern: "^/.*$" }),
|
||||
method: StringEnum(httpMethods, { default: "GET" }),
|
||||
response_type: StringEnum(["json", "text", "html"], { default: "json" }),
|
||||
}),
|
||||
]);
|
||||
static override schema = s.strictObject({
|
||||
path: s.string({ pattern: "^/.*$" }),
|
||||
method: s.string({ enum: httpMethods, default: "GET" }),
|
||||
response_type: s.string({ enum: ["json", "text", "html"], default: "json" }),
|
||||
...Trigger.schema.properties,
|
||||
});
|
||||
|
||||
override async register(flow: Flow, hono: Hono<any>) {
|
||||
const method = this.config.method.toLowerCase() as any;
|
||||
|
||||
@@ -1,20 +1,18 @@
|
||||
import { type Static, StringEnum, parse } from "core/utils";
|
||||
import type { Execution } from "../Execution";
|
||||
import type { Flow } from "../Flow";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s, parse } from "core/object/schema";
|
||||
|
||||
export class Trigger<Schema extends typeof Trigger.schema = typeof Trigger.schema> {
|
||||
// @todo: remove this
|
||||
executions: Execution[] = [];
|
||||
type = "manual";
|
||||
config: Static<Schema>;
|
||||
config: s.Static<Schema>;
|
||||
|
||||
static schema = Type.Object({
|
||||
mode: StringEnum(["sync", "async"], { default: "async" }),
|
||||
static schema = s.strictObject({
|
||||
mode: s.string({ enum: ["sync", "async"], default: "async" }),
|
||||
});
|
||||
|
||||
constructor(config?: Partial<Static<Schema>>) {
|
||||
constructor(config?: Partial<s.Static<Schema>>) {
|
||||
const schema = (this.constructor as typeof Trigger).schema;
|
||||
// @ts-ignore for now
|
||||
this.config = parse(schema, config ?? {});
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
import type { StaticDecode, TSchema } from "@sinclair/typebox";
|
||||
import { BkndError, SimpleRenderer } from "core";
|
||||
import { type Static, type TObject, Value, parse, ucFirst } from "core/utils";
|
||||
//import { BkndError, SimpleRenderer } from "core";
|
||||
import { BkndError } from "core/errors";
|
||||
|
||||
import { s, parse } from "core/object/schema";
|
||||
import type { InputsMap } from "../flows/Execution";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { SimpleRenderer } from "core/template/SimpleRenderer";
|
||||
|
||||
//type InstanceOf<T> = T extends new (...args: any) => infer R ? R : never;
|
||||
|
||||
export type TaskResult<Output = any> = {
|
||||
@@ -16,7 +17,10 @@ export type TaskResult<Output = any> = {
|
||||
|
||||
export type TaskRenderProps<T extends Task = Task> = any;
|
||||
|
||||
export function dynamic<Type extends TSchema>(
|
||||
// @todo: CURRENT WORKAROUND
|
||||
export const dynamic = <S extends s.Schema>(a: S, b?: any) => null as unknown as S;
|
||||
|
||||
/* export function dynamic<Type extends TSchema>(
|
||||
type: Type,
|
||||
parse?: (val: any | string) => Static<Type>,
|
||||
) {
|
||||
@@ -51,23 +55,23 @@ export function dynamic<Type extends TSchema>(
|
||||
// @ts-ignore
|
||||
.Encode((val) => val)
|
||||
);
|
||||
}
|
||||
} */
|
||||
|
||||
export abstract class Task<Params extends TObject = TObject, Output = unknown> {
|
||||
export abstract class Task<Params extends s.Schema = s.Schema, Output = unknown> {
|
||||
abstract type: string;
|
||||
name: string;
|
||||
|
||||
/**
|
||||
* The schema of the task's parameters.
|
||||
*/
|
||||
static schema = Type.Object({});
|
||||
static schema = s.any();
|
||||
|
||||
/**
|
||||
* The task's parameters.
|
||||
*/
|
||||
_params: Static<Params>;
|
||||
_params: s.Static<Params>;
|
||||
|
||||
constructor(name: string, params?: Static<Params>) {
|
||||
constructor(name: string, params?: s.Static<Params>) {
|
||||
if (typeof name !== "string") {
|
||||
throw new Error(`Task name must be a string, got ${typeof name}`);
|
||||
}
|
||||
@@ -81,7 +85,7 @@ export abstract class Task<Params extends TObject = TObject, Output = unknown> {
|
||||
if (
|
||||
schema === Task.schema &&
|
||||
typeof params !== "undefined" &&
|
||||
Object.keys(params).length > 0
|
||||
Object.keys(params || {}).length > 0
|
||||
) {
|
||||
throw new Error(
|
||||
`Task "${name}" has no schema defined but params passed: ${JSON.stringify(params)}`,
|
||||
@@ -93,18 +97,18 @@ export abstract class Task<Params extends TObject = TObject, Output = unknown> {
|
||||
}
|
||||
|
||||
get params() {
|
||||
return this._params as StaticDecode<Params>;
|
||||
return this._params as s.StaticCoerced<Params>;
|
||||
}
|
||||
|
||||
protected clone(name: string, params: Static<Params>): Task {
|
||||
protected clone(name: string, params: s.Static<Params>): Task {
|
||||
return new (this.constructor as any)(name, params);
|
||||
}
|
||||
|
||||
static async resolveParams<S extends TSchema>(
|
||||
static async resolveParams<S extends s.Schema>(
|
||||
schema: S,
|
||||
params: any,
|
||||
inputs: object = {},
|
||||
): Promise<StaticDecode<S>> {
|
||||
): Promise<s.StaticCoerced<S>> {
|
||||
const newParams: any = {};
|
||||
const renderer = new SimpleRenderer(inputs, { renderKeys: true });
|
||||
|
||||
@@ -134,7 +138,8 @@ export abstract class Task<Params extends TObject = TObject, Output = unknown> {
|
||||
newParams[key] = value;
|
||||
}
|
||||
|
||||
return Value.Decode(schema, newParams);
|
||||
return schema.coerce(newParams);
|
||||
//return Value.Decode(schema, newParams);
|
||||
}
|
||||
|
||||
private async cloneWithResolvedParams(_inputs: Map<string, any>) {
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
import { StringEnum } from "core/utils";
|
||||
import { Task, dynamic } from "../Task";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
const FetchMethods = ["GET", "POST", "PUT", "PATCH", "DELETE"];
|
||||
|
||||
@@ -11,24 +9,22 @@ export class FetchTask<Output extends Record<string, any>> extends Task<
|
||||
> {
|
||||
type = "fetch";
|
||||
|
||||
static override schema = Type.Object({
|
||||
url: Type.String({
|
||||
static override schema = s.strictObject({
|
||||
url: s.string({
|
||||
pattern: "^(http|https)://",
|
||||
}),
|
||||
method: Type.Optional(dynamic(StringEnum(FetchMethods, { default: "GET" }))),
|
||||
headers: Type.Optional(
|
||||
dynamic(
|
||||
Type.Array(
|
||||
Type.Object({
|
||||
key: Type.String(),
|
||||
value: Type.String(),
|
||||
}),
|
||||
),
|
||||
JSON.parse,
|
||||
method: dynamic(s.string({ enum: FetchMethods, default: "GET" })).optional(),
|
||||
headers: dynamic(
|
||||
s.array(
|
||||
s.strictObject({
|
||||
key: s.string(),
|
||||
value: s.string(),
|
||||
}),
|
||||
),
|
||||
),
|
||||
body: Type.Optional(dynamic(Type.String())),
|
||||
normal: Type.Optional(dynamic(Type.Number(), Number.parseInt)),
|
||||
JSON.parse,
|
||||
).optional(),
|
||||
body: dynamic(s.string()).optional(),
|
||||
normal: dynamic(s.number(), Number.parseInt).optional(),
|
||||
});
|
||||
|
||||
protected getBody(): string | undefined {
|
||||
|
||||
@@ -1,13 +1,12 @@
|
||||
import { Task } from "../Task";
|
||||
import { $console } from "core";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
export class LogTask extends Task<typeof LogTask.schema> {
|
||||
type = "log";
|
||||
|
||||
static override schema = Type.Object({
|
||||
delay: Type.Number({ default: 10 }),
|
||||
static override schema = s.strictObject({
|
||||
delay: s.number({ default: 10 }),
|
||||
});
|
||||
|
||||
async execute() {
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
import { Task } from "../Task";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
export class RenderTask<Output extends Record<string, any>> extends Task<
|
||||
typeof RenderTask.schema,
|
||||
@@ -8,8 +7,8 @@ export class RenderTask<Output extends Record<string, any>> extends Task<
|
||||
> {
|
||||
type = "render";
|
||||
|
||||
static override schema = Type.Object({
|
||||
render: Type.String(),
|
||||
static override schema = s.strictObject({
|
||||
render: s.string(),
|
||||
});
|
||||
|
||||
async execute() {
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import { Flow } from "../../flows/Flow";
|
||||
import { Task, dynamic } from "../Task";
|
||||
import * as tbbox from "@sinclair/typebox";
|
||||
const { Type } = tbbox;
|
||||
import { s } from "core/object/schema";
|
||||
|
||||
export class SubFlowTask<Output extends Record<string, any>> extends Task<
|
||||
typeof SubFlowTask.schema,
|
||||
@@ -9,10 +8,10 @@ export class SubFlowTask<Output extends Record<string, any>> extends Task<
|
||||
> {
|
||||
type = "subflow";
|
||||
|
||||
static override schema = Type.Object({
|
||||
flow: Type.Any(),
|
||||
input: Type.Optional(dynamic(Type.Any(), JSON.parse)),
|
||||
loop: Type.Optional(Type.Boolean()),
|
||||
static override schema = s.strictObject({
|
||||
flow: s.any(),
|
||||
input: dynamic(s.any(), JSON.parse).optional(),
|
||||
loop: s.boolean().optional(),
|
||||
});
|
||||
|
||||
async execute() {
|
||||
|
||||
Reference in New Issue
Block a user