1- import { v } from "convex/values" ;
1+ import { type Infer , v } from "convex/values" ;
22import {
33 internalAction ,
44 internalMutation ,
55 internalQuery ,
66 mutation ,
77 query ,
88} from "./_generated/server.js" ;
9+ import type { MutationCtx } from "./_generated/server.js" ;
910import { components , internal } from "./_generated/api.js" ;
1011import { omit , withoutSystemFields } from "convex-helpers" ;
1112import { WorkOS , type Event as WorkOSEvent } from "@workos-inc/node" ;
@@ -18,32 +19,117 @@ const eventWorkpool = new Workpool(components.eventWorkpool, {
1819 maxParallelism : 1 ,
1920} ) ;
2021
21- const vEvent = v . object ( {
22+ export const vEvent = v . object ( {
2223 id : v . string ( ) ,
2324 createdAt : v . string ( ) ,
2425 event : v . string ( ) ,
2526 data : v . record ( v . string ( ) , v . any ( ) ) ,
2627 context : v . optional ( v . record ( v . string ( ) , v . any ( ) ) ) ,
2728} ) ;
2829
29- export const enqueueWebhookEvent = mutation ( {
30+ async function processEventHandler (
31+ ctx : MutationCtx ,
32+ args : {
33+ event : Infer < typeof vEvent > ;
34+ logLevel ?: "DEBUG" ;
35+ onEventHandle ?: string ;
36+ }
37+ ) {
38+ if ( args . logLevel === "DEBUG" ) {
39+ console . log ( "processing event" , args . event ) ;
40+ }
41+ const dbEvent = await ctx . db
42+ . query ( "events" )
43+ . withIndex ( "eventId" , ( q ) => q . eq ( "eventId" , args . event . id ) )
44+ . unique ( ) ;
45+ if ( dbEvent ) {
46+ console . log ( "event already processed" , args . event . id ) ;
47+ return ;
48+ }
49+ await ctx . db . insert ( "events" , {
50+ eventId : args . event . id ,
51+ event : args . event . event ,
52+ updatedAt : args . event . data . updatedAt as string | undefined ,
53+ } ) ;
54+ const event = args . event as WorkOSEvent ;
55+ switch ( event . event ) {
56+ case "user.created" : {
57+ const data = omit ( event . data , [ "object" ] ) ;
58+ const existingUser = await ctx . db
59+ . query ( "users" )
60+ . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
61+ . unique ( ) ;
62+ if ( existingUser ) {
63+ console . warn ( "user already exists" , data . id ) ;
64+ break ;
65+ }
66+ await ctx . db . insert ( "users" , data ) ;
67+ break ;
68+ }
69+ case "user.updated" : {
70+ const data = omit ( event . data , [ "object" ] ) ;
71+ const user = await ctx . db
72+ . query ( "users" )
73+ . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
74+ . unique ( ) ;
75+ if ( ! user ) {
76+ console . error ( "user not found" , data . id ) ;
77+ break ;
78+ }
79+ if ( user . updatedAt >= data . updatedAt ) {
80+ console . warn ( `user already updated for event ${ event . id } , skipping` ) ;
81+ break ;
82+ }
83+ await ctx . db . patch ( user . _id , data ) ;
84+ break ;
85+ }
86+ case "user.deleted" : {
87+ const data = omit ( event . data , [ "object" ] ) ;
88+ const user = await ctx . db
89+ . query ( "users" )
90+ . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
91+ . unique ( ) ;
92+ if ( ! user ) {
93+ console . warn ( "user not found" , data . id ) ;
94+ break ;
95+ }
96+ await ctx . db . delete ( user . _id ) ;
97+ break ;
98+ }
99+ }
100+ if ( args . onEventHandle ) {
101+ await ctx . runMutation ( args . onEventHandle as FunctionHandle < "mutation" > , {
102+ event : args . event . event ,
103+ data : args . event . data ,
104+ } ) ;
105+ }
106+ }
107+
108+ export const onWebhookEvent = mutation ( {
30109 args : {
31110 apiKey : v . string ( ) ,
32- eventId : v . string ( ) ,
33- event : v . string ( ) ,
34- updatedAt : v . optional ( v . string ( ) ) ,
111+ event : vEvent ,
35112 onEventHandle : v . optional ( v . string ( ) ) ,
36113 eventTypes : v . optional ( v . array ( v . string ( ) ) ) ,
37114 logLevel : v . optional ( v . literal ( "DEBUG" ) ) ,
38115 } ,
116+ returns : v . null ( ) ,
39117 handler : async ( ctx , args ) => {
40- await eventWorkpool . cancelAll ( ctx ) ;
41- await eventWorkpool . enqueueAction ( ctx , internal . lib . updateEvents , {
42- apiKey : args . apiKey ,
43- onEventHandle : args . onEventHandle ,
44- eventTypes : args . eventTypes ,
45- logLevel : args . logLevel ,
46- } ) ;
118+ const isCreateEvent = args . event . event . endsWith ( ".created" ) ;
119+
120+ if ( isCreateEvent ) {
121+ // Process create events immediately
122+ await processEventHandler ( ctx , args ) ;
123+ } else {
124+ // Enqueue update/delete events to workpool
125+ await eventWorkpool . enqueueAction ( ctx , internal . lib . updateEvents , {
126+ apiKey : args . apiKey ,
127+ onEventHandle : args . onEventHandle ,
128+ eventTypes : args . eventTypes ,
129+ logLevel : args . logLevel ,
130+ } ) ;
131+ }
132+ return null ;
47133 } ,
48134} ) ;
49135
@@ -67,6 +153,9 @@ export const updateEvents = internalAction({
67153 logLevel : v . optional ( v . literal ( "DEBUG" ) ) ,
68154 } ,
69155 handler : async ( ctx , args ) => {
156+ // Cancel other pending workpool jobs since this run will
157+ // process all available events from the WorkOS API.
158+ await eventWorkpool . cancelAll ( ctx ) ;
70159 const workos = new WorkOS ( args . apiKey ) ;
71160 const cursor = await ctx . runQuery ( internal . lib . getCursor ) ;
72161 let nextCursor = cursor ?? undefined ;
@@ -107,74 +196,7 @@ export const processEvent = internalMutation({
107196 onEventHandle : v . optional ( v . string ( ) ) ,
108197 } ,
109198 handler : async ( ctx , args ) => {
110- if ( args . logLevel === "DEBUG" ) {
111- console . log ( "processing event" , args . event ) ;
112- }
113- const dbEvent = await ctx . db
114- . query ( "events" )
115- . withIndex ( "eventId" , ( q ) => q . eq ( "eventId" , args . event . id ) )
116- . unique ( ) ;
117- if ( dbEvent ) {
118- console . log ( "event already processed" , args . event . id ) ;
119- return ;
120- }
121- await ctx . db . insert ( "events" , {
122- eventId : args . event . id ,
123- event : args . event . event ,
124- updatedAt : args . event . data . updatedAt ,
125- } ) ;
126- const event = args . event as WorkOSEvent ;
127- switch ( event . event ) {
128- case "user.created" : {
129- const data = omit ( event . data , [ "object" ] ) ;
130- const existingUser = await ctx . db
131- . query ( "users" )
132- . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
133- . unique ( ) ;
134- if ( existingUser ) {
135- console . warn ( "user already exists" , data . id ) ;
136- break ;
137- }
138- await ctx . db . insert ( "users" , data ) ;
139- break ;
140- }
141- case "user.updated" : {
142- const data = omit ( event . data , [ "object" ] ) ;
143- const user = await ctx . db
144- . query ( "users" )
145- . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
146- . unique ( ) ;
147- if ( ! user ) {
148- console . error ( "user not found" , data . id ) ;
149- break ;
150- }
151- if ( user . updatedAt >= data . updatedAt ) {
152- console . warn ( `user already updated for event ${ event . id } , skipping` ) ;
153- break ;
154- }
155- await ctx . db . patch ( user . _id , data ) ;
156- break ;
157- }
158- case "user.deleted" : {
159- const data = omit ( event . data , [ "object" ] ) ;
160- const user = await ctx . db
161- . query ( "users" )
162- . withIndex ( "id" , ( q ) => q . eq ( "id" , data . id ) )
163- . unique ( ) ;
164- if ( ! user ) {
165- console . warn ( "user not found" , data . id ) ;
166- break ;
167- }
168- await ctx . db . delete ( user . _id ) ;
169- break ;
170- }
171- }
172- if ( args . onEventHandle ) {
173- await ctx . runMutation ( args . onEventHandle as FunctionHandle < "mutation" > , {
174- event : args . event . event ,
175- data : args . event . data ,
176- } ) ;
177- }
199+ await processEventHandler ( ctx , args ) ;
178200 } ,
179201} ) ;
180202
0 commit comments