diff --git a/bun.lockb b/bun.lockb index cb23ae1e..624291ce 100755 Binary files a/bun.lockb and b/bun.lockb differ diff --git a/packages/core/types.ts b/packages/core/types.ts index c466d33f..86dca83e 100644 --- a/packages/core/types.ts +++ b/packages/core/types.ts @@ -35,6 +35,45 @@ export type SendSmsParams = { readonly message: string, } +export type JobStatus = 'completed' | 'failed' | 'queued' | 'active' | 'cancelled' + +export type Job = { + readonly jobId: string + readonly status: JobStatus + readonly type: TType +} + +export type Cron = { + readonly recurrencePattern: string + readonly worker: () => Promise + readonly errorHandler?: (error: unknown) => void +} + +export type Queue = { + readonly worker: () => Promise + readonly errorHandler?: (error: unknown) => void +} + +export interface StartJobParams { + readonly jobId?: string + readonly type: TType + readonly data: Zemble.QueueRegistry[TType] + readonly delayMs?: number + readonly priority?: number +} + +type StartCronParams = { + readonly type: string + readonly data: Record + readonly recurrencePattern: string +} + +export type IStandardCronService = (options: StartCronParams) => Promise + +export type IStandardRemoveFromQueueService = (jobId: string) => Promise> | Job + +export type IStandardAddToQueueService = (options: StartJobParams) => Promise> | Job + export type IStandardSendSmsService = (options: SendSmsParams) => Promise export abstract class IStandardKeyValueService { @@ -257,6 +296,10 @@ declare global { // readonly UnknownToken: Record } + interface QueueRegistry { + + } + interface PushTokenRegistry { } diff --git a/packages/cron-in-memory/README.md b/packages/cron-in-memory/README.md new file mode 100644 index 00000000..899754d7 --- /dev/null +++ b/packages/cron-in-memory/README.md @@ -0,0 +1,13 @@ +# cron-in-memory + +## Develop + +```bash +bun dev +``` + +## Test + +```bash +bun test +``` \ No newline at end of file diff --git a/packages/cron-in-memory/codegen.ts b/packages/cron-in-memory/codegen.ts new file mode 100644 index 00000000..1e6cc7ed --- /dev/null +++ b/packages/cron-in-memory/codegen.ts @@ -0,0 +1,9 @@ +import defaultConfig from '@zemble/graphql/codegen' + +import type { CodegenConfig } from '@graphql-codegen/cli' + +const config: CodegenConfig = { + ...defaultConfig, +} + +export default config diff --git a/packages/cron-in-memory/graphql/Mutation/randomNumber.test.ts b/packages/cron-in-memory/graphql/Mutation/randomNumber.test.ts new file mode 100644 index 00000000..a252872d --- /dev/null +++ b/packages/cron-in-memory/graphql/Mutation/randomNumber.test.ts @@ -0,0 +1,23 @@ +import { createTestApp } from '@zemble/core/test-utils' +import { + it, expect, +} from 'bun:test' + +import plugin from '../../plugin' +import { graphql } from '../client.generated' + +const randomNumberMutation = graphql(` + mutation RandomNumber { + randomNumber + } +`) + +it('Should return a number', async () => { + const app = await createTestApp(plugin) + + const response = await app.gqlRequest(randomNumberMutation, {}) + expect(response.data).toEqual({ + + randomNumber: expect.any(Number), + }) +}) diff --git a/packages/cron-in-memory/graphql/Mutation/randomNumber.ts b/packages/cron-in-memory/graphql/Mutation/randomNumber.ts new file mode 100644 index 00000000..80b510ff --- /dev/null +++ b/packages/cron-in-memory/graphql/Mutation/randomNumber.ts @@ -0,0 +1,11 @@ +import type { MutationResolvers } from '../schema.generated' + +const randomNumber: MutationResolvers['randomNumber'] = (_, __, { pubsub }) => { + const randomNumber = Math.floor(Math.random() * 1000) + + pubsub.publish('randomNumber', randomNumber) + + return randomNumber +} + +export default randomNumber diff --git a/packages/cron-in-memory/graphql/Query/hello.test.ts b/packages/cron-in-memory/graphql/Query/hello.test.ts new file mode 100644 index 00000000..32e109cc --- /dev/null +++ b/packages/cron-in-memory/graphql/Query/hello.test.ts @@ -0,0 +1,21 @@ +import { createTestApp } from '@zemble/core/test-utils' +import { it, expect } from 'bun:test' + +import plugin from '../../plugin' +import { graphql } from '../client.generated' + +const HelloWorldQuery = graphql(` + query Hello { + hello + } +`) + +it('Should return world!', async () => { + const app = await createTestApp(plugin) + + const response = await app.gqlRequest(HelloWorldQuery, {}) + + expect(response.data).toEqual({ + hello: 'world!', + }) +}) diff --git a/packages/cron-in-memory/graphql/Query/hello.ts b/packages/cron-in-memory/graphql/Query/hello.ts new file mode 100644 index 00000000..49793ad9 --- /dev/null +++ b/packages/cron-in-memory/graphql/Query/hello.ts @@ -0,0 +1,5 @@ +import type { QueryResolvers } from '../schema.generated' + +const hello: QueryResolvers['hello'] = () => 'world!' + +export default hello diff --git a/packages/cron-in-memory/graphql/Subscription/countdown.ts b/packages/cron-in-memory/graphql/Subscription/countdown.ts new file mode 100644 index 00000000..e0cf4d68 --- /dev/null +++ b/packages/cron-in-memory/graphql/Subscription/countdown.ts @@ -0,0 +1,20 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +const countdown: SubscriptionResolvers['countdown'] = { + // This will return the value on every 1 sec until it reaches 0 + // eslint-disable-next-line object-shorthand + subscribe: async function* (_, { from }, { logger }) { + // eslint-disable-next-line no-plusplus + for (let i = from; i >= 0; i--) { + // eslint-disable-next-line no-await-in-loop + await new Promise((resolve) => { + logger.info('countdown', { countdown: i }) + setTimeout(resolve, 1000) + }) + yield { countdown: i } + } + }, + resolve: (payload: unknown) => (payload as { readonly countdown: number}).countdown, +} + +export default countdown diff --git a/packages/cron-in-memory/graphql/Subscription/randomNumber.ts b/packages/cron-in-memory/graphql/Subscription/randomNumber.ts new file mode 100644 index 00000000..2959ba0b --- /dev/null +++ b/packages/cron-in-memory/graphql/Subscription/randomNumber.ts @@ -0,0 +1,15 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +const randomNumber: SubscriptionResolvers['randomNumber'] = { + // subscribe to the randomNumber event + subscribe: (_, __, { pubsub }) => { + console.log('subscribing to randomNumber') + return pubsub.subscribe('randomNumber') + }, + resolve: (payload: number) => { + console.log('resolving randomNumber', payload) + return payload + }, +} + +export default randomNumber diff --git a/packages/cron-in-memory/graphql/Subscription/tick.ts b/packages/cron-in-memory/graphql/Subscription/tick.ts new file mode 100644 index 00000000..ce7e549d --- /dev/null +++ b/packages/cron-in-memory/graphql/Subscription/tick.ts @@ -0,0 +1,22 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +let initialized = false +const initializeOnce = (pubsub: Zemble.PubSubType) => { + if (initialized) return + initialized = true + setInterval(() => { + pubsub.publish('tick', Date.now()) + }, 1000) +} + +const tick: SubscriptionResolvers['tick'] = { + // subscribe to the tick event + subscribe: (_, __, { pubsub }) => { + initializeOnce(pubsub) + console.log('subscribing to tick') + return pubsub.subscribe('tick') + }, + resolve: (payload: number) => payload, +} + +export default tick diff --git a/packages/cron-in-memory/graphql/schema.generated.ts b/packages/cron-in-memory/graphql/schema.generated.ts new file mode 100644 index 00000000..c9298623 --- /dev/null +++ b/packages/cron-in-memory/graphql/schema.generated.ts @@ -0,0 +1,154 @@ +// @ts-nocheck +import '@zemble/core' +import type { GraphQLResolveInfo } from 'graphql'; +export type Maybe = T | null | undefined; +export type InputMaybe = T | null | undefined; +export type Exact = { [K in keyof T]: T[K] }; +export type MakeOptional = Omit & { [SubKey in K]?: Maybe }; +export type MakeMaybe = Omit & { [SubKey in K]: Maybe }; +export type MakeEmpty = { [_ in K]?: never }; +export type Incremental = T | { [P in keyof T]?: P extends ' $fragmentName' | '__typename' ? T[P] : never }; +export type RequireFields = Omit & { [P in K]-?: NonNullable }; +/** All built-in and custom scalars, mapped to their actual values */ +export type Scalars = { + ID: { input: string; output: string; } + String: { input: string; output: string; } + Boolean: { input: boolean; output: boolean; } + Int: { input: number; output: number; } + Float: { input: number; output: number; } +}; + +export type Mutation = { + readonly __typename?: 'Mutation'; + readonly randomNumber: Scalars['Int']['output']; +}; + +export type Query = { + readonly __typename?: 'Query'; + readonly hello: Scalars['String']['output']; +}; + +export type Subscription = { + readonly __typename?: 'Subscription'; + readonly countdown: Scalars['Int']['output']; + readonly randomNumber: Scalars['Int']['output']; + readonly tick: Scalars['Float']['output']; +}; + + +export type SubscriptionCountdownArgs = { + from: Scalars['Int']['input']; +}; + +export type WithIndex = TObject & Record; +export type ResolversObject = WithIndex; + +export type ResolverTypeWrapper = Promise | T; + + +export type ResolverWithResolve = { + resolve: ResolverFn; +}; +export type Resolver = ResolverFn | ResolverWithResolve; + +export type ResolverFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => Promise | TResult; + +export type SubscriptionSubscribeFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => AsyncIterable | Promise>; + +export type SubscriptionResolveFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => TResult | Promise; + +export interface SubscriptionSubscriberObject { + subscribe: SubscriptionSubscribeFn<{ [key in TKey]: TResult }, TParent, TContext, TArgs>; + resolve?: SubscriptionResolveFn; +} + +export interface SubscriptionResolverObject { + subscribe: SubscriptionSubscribeFn; + resolve: SubscriptionResolveFn; +} + +export type SubscriptionObject = + | SubscriptionSubscriberObject + | SubscriptionResolverObject; + +export type SubscriptionResolver = + | ((...args: any[]) => SubscriptionObject) + | SubscriptionObject; + +export type TypeResolveFn = ( + parent: TParent, + context: TContext, + info: GraphQLResolveInfo +) => Maybe | Promise>; + +export type IsTypeOfResolverFn = (obj: T, context: TContext, info: GraphQLResolveInfo) => boolean | Promise; + +export type NextResolverFn = () => Promise; + +export type DirectiveResolverFn = ( + next: NextResolverFn, + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => TResult | Promise; + + + +/** Mapping between all available schema types and the resolvers types */ +export type ResolversTypes = ResolversObject<{ + Boolean: ResolverTypeWrapper; + Float: ResolverTypeWrapper; + Int: ResolverTypeWrapper; + Mutation: ResolverTypeWrapper<{}>; + Query: ResolverTypeWrapper<{}>; + String: ResolverTypeWrapper; + Subscription: ResolverTypeWrapper<{}>; +}>; + +/** Mapping between all available schema types and the resolvers parents */ +export type ResolversParentTypes = ResolversObject<{ + Boolean: Scalars['Boolean']['output']; + Float: Scalars['Float']['output']; + Int: Scalars['Int']['output']; + Mutation: {}; + Query: {}; + String: Scalars['String']['output']; + Subscription: {}; +}>; + +export type MutationResolvers = ResolversObject<{ + randomNumber?: Resolver; +}>; + +export type QueryResolvers = ResolversObject<{ + hello?: Resolver; +}>; + +export type SubscriptionResolvers = ResolversObject<{ + countdown?: SubscriptionResolver>; + randomNumber?: SubscriptionResolver; + tick?: SubscriptionResolver; +}>; + +export type Resolvers = ResolversObject<{ + Mutation?: MutationResolvers; + Query?: QueryResolvers; + Subscription?: SubscriptionResolvers; +}>; + diff --git a/packages/cron-in-memory/graphql/schema.graphql b/packages/cron-in-memory/graphql/schema.graphql new file mode 100644 index 00000000..aea15a76 --- /dev/null +++ b/packages/cron-in-memory/graphql/schema.graphql @@ -0,0 +1,13 @@ +type Query { + hello: String! +} + +type Mutation { + randomNumber: Int! +} + +type Subscription { + countdown(from: Int!): Int! + tick: Float! + randomNumber: Int! +} diff --git a/packages/cron-in-memory/package.json b/packages/cron-in-memory/package.json new file mode 100644 index 00000000..3a3c2060 --- /dev/null +++ b/packages/cron-in-memory/package.json @@ -0,0 +1,38 @@ +{ + "name": "@zemble/cron-in-memory", + "version": "0.0.1", + "description": "", + "type": "module", + "keywords": [ + "zemble", + "zemble-plugin", + "@zemble" + ], + "dependencies": { + "@zemble/bun": "workspace:*", + "@zemble/core": "workspace:*", + "@zemble/graphql": "workspace:*", + "@zemble/routes": "workspace:*", + "cron-schedule": "^5.0.1" + }, + "scripts": { + "test": "bun test", + "dev": "zemble-dev plugin.ts", + "typecheck": "tsc --noEmit", + "codegen": "graphql-codegen" + }, + "devDependencies": { + "@types/bun": "*", + "@graphql-codegen/add": "^5.0.0", + "@graphql-codegen/cli": "^5.0.0", + "@graphql-codegen/client-preset": "^4.1.0", + "@graphql-codegen/typescript": "^4.0.1", + "@graphql-codegen/typescript-resolvers": "^4.0.1", + "@tsconfig/bun": "^1.0.1" + }, + "peerDependencies": { + "typescript": "^5.3.3" + }, + "module": "plugin.ts", + "main": "plugin.ts" +} diff --git a/packages/cron-in-memory/plugin.ts b/packages/cron-in-memory/plugin.ts new file mode 100644 index 00000000..016ea124 --- /dev/null +++ b/packages/cron-in-memory/plugin.ts @@ -0,0 +1,51 @@ +import { Plugin, type Cron } from '@zemble/core' +import GraphQL from '@zemble/graphql' +import { parseCronExpression } from 'cron-schedule' +import { IntervalBasedCronScheduler } from 'cron-schedule/schedulers/interval-based.js' + +interface Config extends Zemble.GlobalConfig { + readonly crons: ReadonlyArray +} + +const plugin = new Plugin( + import.meta.dir, + { + dependencies: [{ plugin: GraphQL }], + middleware: ({ config, self }) => { + const scheduler = new IntervalBasedCronScheduler(1000) + config.crons.forEach((cronConfig) => { + const { recurrencePattern, errorHandler, worker } = cronConfig + const cron = parseCronExpression(recurrencePattern) + scheduler.registerTask(cron, worker, { + isOneTimeTask: false, + errorHandler: (error) => { + if (errorHandler) { + errorHandler(error) + } else { + self.providers.logger.error(error) + } + }, + }) + }) + }, + }, +) + +plugin.configure({ + crons: [ + { + recurrencePattern: '*/5 * * * * *', + worker: async () => { + console.log('Hello from in-memory cron') + }, + }, + { + recurrencePattern: '*/5 * * * * *', + worker: async () => { + throw new Error('Error from in-memory cron') + }, + }, + ], +}) + +export default plugin diff --git a/packages/cron-in-memory/routes/hello-world.ts b/packages/cron-in-memory/routes/hello-world.ts new file mode 100644 index 00000000..5208d963 --- /dev/null +++ b/packages/cron-in-memory/routes/hello-world.ts @@ -0,0 +1,3 @@ +export default (ctx: Zemble.RouteContext) => ctx.json({ + message: 'Hello, world!', +}) diff --git a/packages/cron-in-memory/routes/hello/world.html b/packages/cron-in-memory/routes/hello/world.html new file mode 100644 index 00000000..644783de --- /dev/null +++ b/packages/cron-in-memory/routes/hello/world.html @@ -0,0 +1,8 @@ + + + Hello World + + +

Hello World

+ + \ No newline at end of file diff --git a/packages/cron-in-memory/tsconfig.json b/packages/cron-in-memory/tsconfig.json new file mode 100644 index 00000000..0037f214 --- /dev/null +++ b/packages/cron-in-memory/tsconfig.json @@ -0,0 +1,4 @@ +{ + "extends": "../../tsconfig.json", + "exclude": [ "**/*.generated/*.ts" ] +} diff --git a/packages/queue-in-memory/README.md b/packages/queue-in-memory/README.md new file mode 100644 index 00000000..4fb467de --- /dev/null +++ b/packages/queue-in-memory/README.md @@ -0,0 +1,13 @@ +# queue-in-memory + +## Develop + +```bash +bun dev +``` + +## Test + +```bash +bun test +``` \ No newline at end of file diff --git a/packages/queue-in-memory/codegen.ts b/packages/queue-in-memory/codegen.ts new file mode 100644 index 00000000..1e6cc7ed --- /dev/null +++ b/packages/queue-in-memory/codegen.ts @@ -0,0 +1,9 @@ +import defaultConfig from '@zemble/graphql/codegen' + +import type { CodegenConfig } from '@graphql-codegen/cli' + +const config: CodegenConfig = { + ...defaultConfig, +} + +export default config diff --git a/packages/queue-in-memory/graphql/Mutation/randomNumber.test.ts b/packages/queue-in-memory/graphql/Mutation/randomNumber.test.ts new file mode 100644 index 00000000..a252872d --- /dev/null +++ b/packages/queue-in-memory/graphql/Mutation/randomNumber.test.ts @@ -0,0 +1,23 @@ +import { createTestApp } from '@zemble/core/test-utils' +import { + it, expect, +} from 'bun:test' + +import plugin from '../../plugin' +import { graphql } from '../client.generated' + +const randomNumberMutation = graphql(` + mutation RandomNumber { + randomNumber + } +`) + +it('Should return a number', async () => { + const app = await createTestApp(plugin) + + const response = await app.gqlRequest(randomNumberMutation, {}) + expect(response.data).toEqual({ + + randomNumber: expect.any(Number), + }) +}) diff --git a/packages/queue-in-memory/graphql/Mutation/randomNumber.ts b/packages/queue-in-memory/graphql/Mutation/randomNumber.ts new file mode 100644 index 00000000..80b510ff --- /dev/null +++ b/packages/queue-in-memory/graphql/Mutation/randomNumber.ts @@ -0,0 +1,11 @@ +import type { MutationResolvers } from '../schema.generated' + +const randomNumber: MutationResolvers['randomNumber'] = (_, __, { pubsub }) => { + const randomNumber = Math.floor(Math.random() * 1000) + + pubsub.publish('randomNumber', randomNumber) + + return randomNumber +} + +export default randomNumber diff --git a/packages/queue-in-memory/graphql/Query/hello.test.ts b/packages/queue-in-memory/graphql/Query/hello.test.ts new file mode 100644 index 00000000..32e109cc --- /dev/null +++ b/packages/queue-in-memory/graphql/Query/hello.test.ts @@ -0,0 +1,21 @@ +import { createTestApp } from '@zemble/core/test-utils' +import { it, expect } from 'bun:test' + +import plugin from '../../plugin' +import { graphql } from '../client.generated' + +const HelloWorldQuery = graphql(` + query Hello { + hello + } +`) + +it('Should return world!', async () => { + const app = await createTestApp(plugin) + + const response = await app.gqlRequest(HelloWorldQuery, {}) + + expect(response.data).toEqual({ + hello: 'world!', + }) +}) diff --git a/packages/queue-in-memory/graphql/Query/hello.ts b/packages/queue-in-memory/graphql/Query/hello.ts new file mode 100644 index 00000000..49793ad9 --- /dev/null +++ b/packages/queue-in-memory/graphql/Query/hello.ts @@ -0,0 +1,5 @@ +import type { QueryResolvers } from '../schema.generated' + +const hello: QueryResolvers['hello'] = () => 'world!' + +export default hello diff --git a/packages/queue-in-memory/graphql/Subscription/countdown.ts b/packages/queue-in-memory/graphql/Subscription/countdown.ts new file mode 100644 index 00000000..e0cf4d68 --- /dev/null +++ b/packages/queue-in-memory/graphql/Subscription/countdown.ts @@ -0,0 +1,20 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +const countdown: SubscriptionResolvers['countdown'] = { + // This will return the value on every 1 sec until it reaches 0 + // eslint-disable-next-line object-shorthand + subscribe: async function* (_, { from }, { logger }) { + // eslint-disable-next-line no-plusplus + for (let i = from; i >= 0; i--) { + // eslint-disable-next-line no-await-in-loop + await new Promise((resolve) => { + logger.info('countdown', { countdown: i }) + setTimeout(resolve, 1000) + }) + yield { countdown: i } + } + }, + resolve: (payload: unknown) => (payload as { readonly countdown: number}).countdown, +} + +export default countdown diff --git a/packages/queue-in-memory/graphql/Subscription/randomNumber.ts b/packages/queue-in-memory/graphql/Subscription/randomNumber.ts new file mode 100644 index 00000000..2959ba0b --- /dev/null +++ b/packages/queue-in-memory/graphql/Subscription/randomNumber.ts @@ -0,0 +1,15 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +const randomNumber: SubscriptionResolvers['randomNumber'] = { + // subscribe to the randomNumber event + subscribe: (_, __, { pubsub }) => { + console.log('subscribing to randomNumber') + return pubsub.subscribe('randomNumber') + }, + resolve: (payload: number) => { + console.log('resolving randomNumber', payload) + return payload + }, +} + +export default randomNumber diff --git a/packages/queue-in-memory/graphql/Subscription/tick.ts b/packages/queue-in-memory/graphql/Subscription/tick.ts new file mode 100644 index 00000000..ce7e549d --- /dev/null +++ b/packages/queue-in-memory/graphql/Subscription/tick.ts @@ -0,0 +1,22 @@ +import type { SubscriptionResolvers } from '../schema.generated' + +let initialized = false +const initializeOnce = (pubsub: Zemble.PubSubType) => { + if (initialized) return + initialized = true + setInterval(() => { + pubsub.publish('tick', Date.now()) + }, 1000) +} + +const tick: SubscriptionResolvers['tick'] = { + // subscribe to the tick event + subscribe: (_, __, { pubsub }) => { + initializeOnce(pubsub) + console.log('subscribing to tick') + return pubsub.subscribe('tick') + }, + resolve: (payload: number) => payload, +} + +export default tick diff --git a/packages/queue-in-memory/graphql/schema.generated.ts b/packages/queue-in-memory/graphql/schema.generated.ts new file mode 100644 index 00000000..c9298623 --- /dev/null +++ b/packages/queue-in-memory/graphql/schema.generated.ts @@ -0,0 +1,154 @@ +// @ts-nocheck +import '@zemble/core' +import type { GraphQLResolveInfo } from 'graphql'; +export type Maybe = T | null | undefined; +export type InputMaybe = T | null | undefined; +export type Exact = { [K in keyof T]: T[K] }; +export type MakeOptional = Omit & { [SubKey in K]?: Maybe }; +export type MakeMaybe = Omit & { [SubKey in K]: Maybe }; +export type MakeEmpty = { [_ in K]?: never }; +export type Incremental = T | { [P in keyof T]?: P extends ' $fragmentName' | '__typename' ? T[P] : never }; +export type RequireFields = Omit & { [P in K]-?: NonNullable }; +/** All built-in and custom scalars, mapped to their actual values */ +export type Scalars = { + ID: { input: string; output: string; } + String: { input: string; output: string; } + Boolean: { input: boolean; output: boolean; } + Int: { input: number; output: number; } + Float: { input: number; output: number; } +}; + +export type Mutation = { + readonly __typename?: 'Mutation'; + readonly randomNumber: Scalars['Int']['output']; +}; + +export type Query = { + readonly __typename?: 'Query'; + readonly hello: Scalars['String']['output']; +}; + +export type Subscription = { + readonly __typename?: 'Subscription'; + readonly countdown: Scalars['Int']['output']; + readonly randomNumber: Scalars['Int']['output']; + readonly tick: Scalars['Float']['output']; +}; + + +export type SubscriptionCountdownArgs = { + from: Scalars['Int']['input']; +}; + +export type WithIndex = TObject & Record; +export type ResolversObject = WithIndex; + +export type ResolverTypeWrapper = Promise | T; + + +export type ResolverWithResolve = { + resolve: ResolverFn; +}; +export type Resolver = ResolverFn | ResolverWithResolve; + +export type ResolverFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => Promise | TResult; + +export type SubscriptionSubscribeFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => AsyncIterable | Promise>; + +export type SubscriptionResolveFn = ( + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => TResult | Promise; + +export interface SubscriptionSubscriberObject { + subscribe: SubscriptionSubscribeFn<{ [key in TKey]: TResult }, TParent, TContext, TArgs>; + resolve?: SubscriptionResolveFn; +} + +export interface SubscriptionResolverObject { + subscribe: SubscriptionSubscribeFn; + resolve: SubscriptionResolveFn; +} + +export type SubscriptionObject = + | SubscriptionSubscriberObject + | SubscriptionResolverObject; + +export type SubscriptionResolver = + | ((...args: any[]) => SubscriptionObject) + | SubscriptionObject; + +export type TypeResolveFn = ( + parent: TParent, + context: TContext, + info: GraphQLResolveInfo +) => Maybe | Promise>; + +export type IsTypeOfResolverFn = (obj: T, context: TContext, info: GraphQLResolveInfo) => boolean | Promise; + +export type NextResolverFn = () => Promise; + +export type DirectiveResolverFn = ( + next: NextResolverFn, + parent: TParent, + args: TArgs, + context: TContext, + info: GraphQLResolveInfo +) => TResult | Promise; + + + +/** Mapping between all available schema types and the resolvers types */ +export type ResolversTypes = ResolversObject<{ + Boolean: ResolverTypeWrapper; + Float: ResolverTypeWrapper; + Int: ResolverTypeWrapper; + Mutation: ResolverTypeWrapper<{}>; + Query: ResolverTypeWrapper<{}>; + String: ResolverTypeWrapper; + Subscription: ResolverTypeWrapper<{}>; +}>; + +/** Mapping between all available schema types and the resolvers parents */ +export type ResolversParentTypes = ResolversObject<{ + Boolean: Scalars['Boolean']['output']; + Float: Scalars['Float']['output']; + Int: Scalars['Int']['output']; + Mutation: {}; + Query: {}; + String: Scalars['String']['output']; + Subscription: {}; +}>; + +export type MutationResolvers = ResolversObject<{ + randomNumber?: Resolver; +}>; + +export type QueryResolvers = ResolversObject<{ + hello?: Resolver; +}>; + +export type SubscriptionResolvers = ResolversObject<{ + countdown?: SubscriptionResolver>; + randomNumber?: SubscriptionResolver; + tick?: SubscriptionResolver; +}>; + +export type Resolvers = ResolversObject<{ + Mutation?: MutationResolvers; + Query?: QueryResolvers; + Subscription?: SubscriptionResolvers; +}>; + diff --git a/packages/queue-in-memory/graphql/schema.graphql b/packages/queue-in-memory/graphql/schema.graphql new file mode 100644 index 00000000..aea15a76 --- /dev/null +++ b/packages/queue-in-memory/graphql/schema.graphql @@ -0,0 +1,13 @@ +type Query { + hello: String! +} + +type Mutation { + randomNumber: Int! +} + +type Subscription { + countdown(from: Int!): Int! + tick: Float! + randomNumber: Int! +} diff --git a/packages/queue-in-memory/package.json b/packages/queue-in-memory/package.json new file mode 100644 index 00000000..e8c60c81 --- /dev/null +++ b/packages/queue-in-memory/package.json @@ -0,0 +1,38 @@ +{ + "name": "@zemble/queue-in-memory", + "version": "0.0.1", + "description": "", + "type": "module", + "keywords": [ + "zemble", + "zemble-plugin", + "@zemble" + ], + "dependencies": { + "@zemble/core": "workspace:*", + "@zemble/bun": "workspace:*", + "@zemble/utils": "workspace:*", + "@zemble/graphql": "workspace:*" + }, + "scripts": { + "test": "bun test", + "dev": "zemble-dev plugin.ts", + "testing-it-out": "zemble-dev testing-it-out.ts", + "typecheck": "tsc --noEmit", + "codegen": "graphql-codegen" + }, + "devDependencies": { + "@types/bun": "*", + "@graphql-codegen/add": "^5.0.0", + "@graphql-codegen/cli": "^5.0.0", + "@graphql-codegen/client-preset": "^4.1.0", + "@graphql-codegen/typescript": "^4.0.1", + "@graphql-codegen/typescript-resolvers": "^4.0.1", + "@tsconfig/bun": "^1.0.1" + }, + "peerDependencies": { + "typescript": "^5.3.3" + }, + "module": "plugin.ts", + "main": "plugin.ts" +} diff --git a/packages/queue-in-memory/plugin.ts b/packages/queue-in-memory/plugin.ts new file mode 100644 index 00000000..8714108b --- /dev/null +++ b/packages/queue-in-memory/plugin.ts @@ -0,0 +1,145 @@ +import { + Plugin, setupProvider, type IStandardAddToQueueService, type IStandardRemoveFromQueueService, type JobStatus, +} from '@zemble/core' +import GraphQL from '@zemble/graphql' +import { concurrencyExecutionLock } from '@zemble/utils/executionLock' + +type Worker = (queueConfig: Zemble.QueueRegistry[TType]) => Promise | void + +interface Config extends Zemble.GlobalConfig { + readonly queues: Record + }> +} + +declare global { + // eslint-disable-next-line @typescript-eslint/no-namespace + namespace Zemble { + interface Providers { + readonly addToQueue: IStandardAddToQueueService + readonly removeFromQueue: IStandardRemoveFromQueueService + } + + interface MiddlewareConfig { + readonly ['@zemble/queue-in-memory']?: undefined + } + } +} + +export interface JobData { + readonly jobId: string + readonly type: TType + readonly data: Zemble.QueueRegistry[TType] + readonly priority: number + readonly earliestExecution?: Date +} + +const plugin = new Plugin( + import.meta.dir, + { + dependencies: [{ plugin: GraphQL }], + middleware: async ({ app }) => { + const jobTimers = new Map() + const jobStatus = new Map() + // eslint-disable-next-line no-spaced-func + const workers = new Map void)>() + // eslint-disable-next-line functional/prefer-readonly-type + const executionQueues = new Map>() + let jobIdCounter = 0 + + function createWorker(type: TType, concurrency: number, worker: Worker) { + return concurrencyExecutionLock(async () => { + const executionQueue = executionQueues.get(type)! + // eslint-disable-next-line functional/immutable-data + const pickedJob = executionQueue.find((j) => !j.earliestExecution || j.earliestExecution.valueOf() < Date.now()) + + if (!pickedJob) { + throw new Error('No job found') + } + + const { jobId, data } = pickedJob + executionQueues.set(type, executionQueue.filter((j) => j.jobId !== jobId)) + jobStatus.set(jobId, 'active') + try { + await worker(data) + jobStatus.set(jobId, 'completed') + } catch (error) { + jobStatus.set(jobId, 'failed') + } + }, concurrency ?? 1) + } + + await setupProvider({ + app, + // eslint-disable-next-line unicorn/consistent-function-scoping + initializeProvider: () => (cfg) => { + const queue = plugin.config.queues[cfg.type] + // eslint-disable-next-line no-plusplus + const jobId = `${cfg.type}:${cfg.jobId ?? (jobIdCounter++).toString()}` + + if (!queue) { + throw new Error(`Queue ${cfg.type} not found`) + } + + const executionQueue = executionQueues.get(cfg.type) ?? [] + + executionQueues.set(cfg.type, [ + ...executionQueue, { + jobId, + data: cfg.data, + type: cfg.type, + priority: cfg.priority ?? -Date.now(), + earliestExecution: cfg.delayMs ? new Date(Date.now() + cfg.delayMs) : undefined, + }, + ].sort((a, b) => a.priority - b.priority)) + + if (!workers.has(cfg.type)) { + workers.set(cfg.type, createWorker(cfg.type, queue.concurrency ?? 1, queue.worker)) + } + + jobTimers.set(jobId, setTimeout(workers.get(cfg.type)!, cfg.delayMs ?? 0)) + jobStatus.set(jobId, 'queued') + + return { + jobId, + status: jobStatus.get(jobId)!, + type: cfg.type, + } + }, + middlewareKey: '@zemble/queue-in-memory', + providerKey: 'addToQueue', + }) + + await setupProvider({ + app, + // eslint-disable-next-line unicorn/consistent-function-scoping + initializeProvider: () => (jobId) => { + const status = jobStatus.get(jobId) + if (!status) { + throw new Error('job with id not found') + } + const [type] = jobId.split(':') + if (status === 'queued') { + jobStatus.set(jobId, 'cancelled') + const executionQueue = executionQueues.get(type!)! + executionQueues.set(type!, executionQueue.filter((j) => j.jobId !== jobId)) + + const timer = jobTimers.get(jobId) + clearTimeout(timer) + } + + return { + jobId, + status, + type: type as keyof Zemble.QueueRegistry, + } + }, + middlewareKey: '@zemble/queue-in-memory', + providerKey: 'removeFromQueue', + }) + }, + }, +) + +export default plugin diff --git a/packages/queue-in-memory/routes/hello-world.ts b/packages/queue-in-memory/routes/hello-world.ts new file mode 100644 index 00000000..5208d963 --- /dev/null +++ b/packages/queue-in-memory/routes/hello-world.ts @@ -0,0 +1,3 @@ +export default (ctx: Zemble.RouteContext) => ctx.json({ + message: 'Hello, world!', +}) diff --git a/packages/queue-in-memory/routes/hello/world.html b/packages/queue-in-memory/routes/hello/world.html new file mode 100644 index 00000000..644783de --- /dev/null +++ b/packages/queue-in-memory/routes/hello/world.html @@ -0,0 +1,8 @@ + + + Hello World + + +

Hello World

+ + \ No newline at end of file diff --git a/packages/queue-in-memory/testing-it-out.ts b/packages/queue-in-memory/testing-it-out.ts new file mode 100644 index 00000000..40b07234 --- /dev/null +++ b/packages/queue-in-memory/testing-it-out.ts @@ -0,0 +1,35 @@ +import { wait } from '@zemble/utils' + +import plugin from './plugin' + +interface MyLibreQueueType { + readonly patientId: string +} + +declare global { + // eslint-disable-next-line @typescript-eslint/no-namespace + namespace Zemble { + interface QueueRegistry { + readonly ['process-libre']: MyLibreQueueType + } + } +} + +plugin.configure({ + queues: { + 'process-libre': { + worker: async (data) => { + console.log('queueConfig.data', data) + await wait(1000) + console.log('done', data) + }, + }, + }, +}) + +setTimeout(() => { + void plugin.providers.addToQueue({ data: { patientId: 'sdf' }, type: 'process-libre' }) + void plugin.providers.addToQueue({ data: { patientId: 'string2' }, type: 'process-libre' }) +}, 1000) + +export default plugin diff --git a/packages/queue-in-memory/tsconfig.json b/packages/queue-in-memory/tsconfig.json new file mode 100644 index 00000000..0037f214 --- /dev/null +++ b/packages/queue-in-memory/tsconfig.json @@ -0,0 +1,4 @@ +{ + "extends": "../../tsconfig.json", + "exclude": [ "**/*.generated/*.ts" ] +} diff --git a/packages/utils/executionLock.test.ts b/packages/utils/executionLock.test.ts new file mode 100644 index 00000000..0ec3b193 --- /dev/null +++ b/packages/utils/executionLock.test.ts @@ -0,0 +1,86 @@ +import { + describe, expect, it, jest, +} from 'bun:test' + +import { concurrencyExecutionLock, singleExecutionLock } from './executionLock' +import wait from './wait' + +describe('singleExecutionLock', () => { + it('should run the function only once', async () => { + const fn = jest.fn(async () => wait(10)) + const fnLocked = singleExecutionLock(fn) + + await Promise.all([ + fnLocked(), + fnLocked(), + fnLocked(), + ]) + + expect(fn).toHaveBeenCalledTimes(1) + }) + + it('should run the function more if in sequence', async () => { + const fn = jest.fn() + const fnLocked = singleExecutionLock(fn) + + await fnLocked() + await fnLocked() + await fnLocked() + + expect(fn).toHaveBeenCalledTimes(3) + }) +}) + +describe('concurrencyExecutionLock', () => { + it('should run all functions', async () => { + const fn = jest.fn() + const fnLocked = concurrencyExecutionLock(fn) + + await fnLocked() + await fnLocked() + await fnLocked() + + expect(fn).toHaveBeenCalledTimes(3) + }) + + it('should run the function more if in parallel with timeout', async () => { + const start = Date.now() + const fn = jest.fn() + const fnLocked = concurrencyExecutionLock(fn, 3, 100) + + await Promise.all([ + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + ]) + + const timeConsumed = Date.now() - start + + expect(fn).toHaveBeenCalledTimes(6) + expect(timeConsumed).toBeGreaterThanOrEqual(100) + expect(timeConsumed).toBeLessThan(150) + }) + + it('should run the function more if in parallel with timeout', async () => { + const start = Date.now() + const fn = jest.fn() + const fnLocked = concurrencyExecutionLock(fn, 6, 100) + + await Promise.all([ + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + fnLocked(), + ]) + + const timeConsumed = Date.now() - start + + expect(fn).toHaveBeenCalledTimes(6) + expect(timeConsumed).toBeLessThan(50) + }) +}) diff --git a/packages/utils/executionLock.ts b/packages/utils/executionLock.ts new file mode 100644 index 00000000..a0727e79 --- /dev/null +++ b/packages/utils/executionLock.ts @@ -0,0 +1,38 @@ +/* eslint-disable no-plusplus */ + +import wait from './wait' + +// eslint-disable-next-line @typescript-eslint/no-explicit-any +type AnyFunction = (...args: readonly any[]) => any; + +export function singleExecutionLock(fn: T) { + let isRunning = false + return async (...args: Parameters): Promise | undefined> => { + if (isRunning) { + return undefined + } + isRunning = true + try { + return await fn(...args) + } finally { + isRunning = false + } + } +} + +export function concurrencyExecutionLock(fn: T, concurrency = 1, timeout = 1000) { + let running = 0 + const self = async (...args: Parameters): Promise> => { + if (running >= concurrency) { + await wait(timeout) + return self(...args) + } + running++ + try { + return await fn(...args) + } finally { + running-- + } + } + return self +}