@trigger.dev/supabase: You can now trigger on multiple database events in the same job
This commit is contained in:
@@ -0,0 +1,5 @@
|
|||||||
|
---
|
||||||
|
"@trigger.dev/supabase": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
You can now trigger on multiple database events in the same job
|
||||||
@@ -38,8 +38,8 @@ client.defineJob({
|
|||||||
supabase,
|
supabase,
|
||||||
},
|
},
|
||||||
run: async (payload, io, ctx) => {
|
run: async (payload, io, ctx) => {
|
||||||
const { data: users, error } = await io.supabase.runTask("find-users", async (db) => {
|
const { data: todos, error } = await io.supabase.runTask("find-todos", async (db) => {
|
||||||
return db.from("users").select("*");
|
return db.from("todos").select("*");
|
||||||
});
|
});
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
@@ -60,8 +60,8 @@ client.defineJob({
|
|||||||
supabase,
|
supabase,
|
||||||
},
|
},
|
||||||
run: async (payload, io, ctx) => {
|
run: async (payload, io, ctx) => {
|
||||||
const users = await io.supabase.runTask("find-users", async (db) => {
|
const todos = await io.supabase.runTask("find-todos", async (db) => {
|
||||||
const { data, error } = await db.from("users").select("*");
|
const { data, error } = await db.from("todos").select("*");
|
||||||
|
|
||||||
if (error) throw error;
|
if (error) throw error;
|
||||||
|
|
||||||
@@ -100,15 +100,15 @@ client.defineJob({
|
|||||||
supabase,
|
supabase,
|
||||||
},
|
},
|
||||||
run: async (payload, io, ctx) => {
|
run: async (payload, io, ctx) => {
|
||||||
const users = await io.supabase.runTask("find-users", async (db) => {
|
const todos = await io.supabase.runTask("find-todos", async (db) => {
|
||||||
const { data, error } = await db.from("users").select("*");
|
const { data, error } = await db.from("todos").select("*");
|
||||||
|
|
||||||
if (error) throw error;
|
if (error) throw error;
|
||||||
|
|
||||||
return data;
|
return data;
|
||||||
});
|
});
|
||||||
|
|
||||||
// users is now typed as User[] instead of any[]
|
// todos is now typed as Todo[] instead of any[]
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
```
|
```
|
||||||
|
|||||||
@@ -126,7 +126,7 @@ client.defineJob({
|
|||||||
id: "supabase-trigger",
|
id: "supabase-trigger",
|
||||||
name: "Supabase Trigger",
|
name: "Supabase Trigger",
|
||||||
trigger: db.onInserted({
|
trigger: db.onInserted({
|
||||||
table: "users",
|
table: "todos",
|
||||||
}),
|
}),
|
||||||
run: async (payload, io, ctx) => {
|
run: async (payload, io, ctx) => {
|
||||||
// payload is the database webhook body (see https://supabase.com/docs/guides/database/webhooks#payload)
|
// payload is the database webhook body (see https://supabase.com/docs/guides/database/webhooks#payload)
|
||||||
@@ -137,20 +137,6 @@ client.defineJob({
|
|||||||
You can add additional filters to the trigger by passing a `filter` object:
|
You can add additional filters to the trigger by passing a `filter` object:
|
||||||
|
|
||||||
```ts
|
```ts
|
||||||
client.defineJob({
|
|
||||||
id: "supabase-trigger",
|
|
||||||
name: "Supabase Trigger",
|
|
||||||
trigger: db.onUpdated({
|
|
||||||
table: "users",
|
|
||||||
filter: {
|
|
||||||
country: ["USA", "Canada"], // This will only trigger the job if the user.country is USA or Canada
|
|
||||||
},
|
|
||||||
}),
|
|
||||||
run: async (payload, io, ctx) => {
|
|
||||||
// payload is the database webhook body (see https://supabase.com/docs/guides/database/webhooks#payload)
|
|
||||||
},
|
|
||||||
});
|
|
||||||
|
|
||||||
client.defineJob({
|
client.defineJob({
|
||||||
id: "supabase-trigger",
|
id: "supabase-trigger",
|
||||||
name: "Supabase Trigger",
|
name: "Supabase Trigger",
|
||||||
@@ -172,6 +158,31 @@ client.defineJob({
|
|||||||
});
|
});
|
||||||
```
|
```
|
||||||
|
|
||||||
|
You can also listen for multiple different events using the `on` trigger:
|
||||||
|
|
||||||
|
```ts
|
||||||
|
client.defineJob({
|
||||||
|
id: "supabase-trigger",
|
||||||
|
name: "Supabase Trigger",
|
||||||
|
trigger: db.on({
|
||||||
|
table: "todos",
|
||||||
|
events: ["INSERT", "UPDATE"] // Trigger on both insert and update events
|
||||||
|
filter: {
|
||||||
|
record: {
|
||||||
|
is_completed: [false],
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}),
|
||||||
|
run: async (payload, io, ctx) => {
|
||||||
|
if (payload.type === "INSERT") {
|
||||||
|
// payload will be typed as the INSERT payload
|
||||||
|
} else {
|
||||||
|
// payload will be typed as the UPDATE payload
|
||||||
|
}
|
||||||
|
},
|
||||||
|
});
|
||||||
|
```
|
||||||
|
|
||||||
<Note>
|
<Note>
|
||||||
We will only create at most 1 database webhook per table, to limit resource usage when writing to
|
We will only create at most 1 database webhook per table, to limit resource usage when writing to
|
||||||
your database. This means we cannot support scoping updated triggers to specific columns.
|
your database. This means we cannot support scoping updated triggers to specific columns.
|
||||||
@@ -196,10 +207,10 @@ client.defineJob({
|
|||||||
id: "supabase-trigger",
|
id: "supabase-trigger",
|
||||||
name: "Supabase Trigger",
|
name: "Supabase Trigger",
|
||||||
trigger: db.onUpdated({
|
trigger: db.onUpdated({
|
||||||
table: "users",
|
table: "todos",
|
||||||
}),
|
}),
|
||||||
run: async (payload, io, ctx) => {
|
run: async (payload, io, ctx) => {
|
||||||
// payload.record and payload.old_record are now correctly typed to match the users table
|
// payload.record and payload.old_record are now correctly typed to match the todos table
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
```
|
```
|
||||||
|
|||||||
@@ -1,14 +1,19 @@
|
|||||||
import { TriggerClient, eventTrigger } from "@trigger.dev/sdk";
|
import { TriggerClient } from "@trigger.dev/sdk";
|
||||||
import { createExpressServer } from "@trigger.dev/express";
|
import { createExpressServer } from "@trigger.dev/express";
|
||||||
import { z } from "zod";
|
import { Supabase, SupabaseManagement } from "@trigger.dev/supabase";
|
||||||
import { SupabaseManagement } from "@trigger.dev/supabase";
|
|
||||||
|
|
||||||
const supabaseManagement = new SupabaseManagement({
|
const supabaseManagement = new SupabaseManagement({
|
||||||
id: "supabase-management",
|
id: "supabase-management",
|
||||||
apiKey: process.env["SUPABASE_API_KEY"]!,
|
apiKey: process.env["SUPABASE_API_KEY"]!,
|
||||||
});
|
});
|
||||||
|
|
||||||
const db = supabaseManagement.db(process.env["SUPABASE_ID"]!);
|
const triggers = supabaseManagement.db<Database>(process.env["SUPABASE_ID"]!);
|
||||||
|
|
||||||
|
const supabase = new Supabase({
|
||||||
|
id: "supabase",
|
||||||
|
supabaseKey: process.env["SUPABASE_SERVICE_ROLE_KEY"]!,
|
||||||
|
supabaseUrl: process.env["SUPABASE_URL"]!,
|
||||||
|
});
|
||||||
|
|
||||||
export const client = new TriggerClient({
|
export const client = new TriggerClient({
|
||||||
id: "job-catalog",
|
id: "job-catalog",
|
||||||
@@ -24,8 +29,8 @@ client.defineJob({
|
|||||||
id: "supabase-management-example-1",
|
id: "supabase-management-example-1",
|
||||||
name: "Supabase Management Example 1",
|
name: "Supabase Management Example 1",
|
||||||
version: "0.1.0",
|
version: "0.1.0",
|
||||||
trigger: db.onInserted({
|
trigger: triggers.onInserted({
|
||||||
table: "users",
|
table: "todos",
|
||||||
}),
|
}),
|
||||||
run: async (payload, io, ctx) => {},
|
run: async (payload, io, ctx) => {},
|
||||||
});
|
});
|
||||||
@@ -34,8 +39,86 @@ client.defineJob({
|
|||||||
id: "supabase-management-example-2",
|
id: "supabase-management-example-2",
|
||||||
name: "Supabase Management Example 2",
|
name: "Supabase Management Example 2",
|
||||||
version: "0.1.0",
|
version: "0.1.0",
|
||||||
trigger: db.onUpdated({
|
trigger: triggers.onUpdated({
|
||||||
table: "users",
|
table: "todos",
|
||||||
}),
|
}),
|
||||||
run: async (payload, io, ctx) => {},
|
run: async (payload, io, ctx) => {},
|
||||||
});
|
});
|
||||||
|
|
||||||
|
client.defineJob({
|
||||||
|
id: "supabase-management-example-on",
|
||||||
|
name: "Supabase Management Example On",
|
||||||
|
version: "0.1.0",
|
||||||
|
trigger: triggers.on({
|
||||||
|
table: "todos",
|
||||||
|
events: ["INSERT", "UPDATE"],
|
||||||
|
}),
|
||||||
|
integrations: {
|
||||||
|
supabase,
|
||||||
|
},
|
||||||
|
run: async (payload, io, ctx) => {
|
||||||
|
const user = await io.supabase.runTask("fetch-user", async (db) => {
|
||||||
|
const { data, error } = await db.auth.admin.getUserById(payload.record.user_id);
|
||||||
|
|
||||||
|
if (error) {
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
|
||||||
|
return data.user;
|
||||||
|
});
|
||||||
|
|
||||||
|
return user;
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
export type Json = string | number | boolean | null | { [key: string]: Json | undefined } | Json[];
|
||||||
|
|
||||||
|
export interface Database {
|
||||||
|
public: {
|
||||||
|
Tables: {
|
||||||
|
todos: {
|
||||||
|
Row: {
|
||||||
|
id: number;
|
||||||
|
inserted_at: string;
|
||||||
|
is_complete: boolean | null;
|
||||||
|
task: string | null;
|
||||||
|
user_id: string;
|
||||||
|
};
|
||||||
|
Insert: {
|
||||||
|
id?: number;
|
||||||
|
inserted_at?: string;
|
||||||
|
is_complete?: boolean | null;
|
||||||
|
task?: string | null;
|
||||||
|
user_id: string;
|
||||||
|
};
|
||||||
|
Update: {
|
||||||
|
id?: number;
|
||||||
|
inserted_at?: string;
|
||||||
|
is_complete?: boolean | null;
|
||||||
|
task?: string | null;
|
||||||
|
user_id?: string;
|
||||||
|
};
|
||||||
|
Relationships: [
|
||||||
|
{
|
||||||
|
foreignKeyName: "todos_user_id_fkey";
|
||||||
|
columns: ["user_id"];
|
||||||
|
referencedRelation: "users";
|
||||||
|
referencedColumns: ["id"];
|
||||||
|
},
|
||||||
|
];
|
||||||
|
};
|
||||||
|
};
|
||||||
|
Views: {
|
||||||
|
[_ in never]: never;
|
||||||
|
};
|
||||||
|
Functions: {
|
||||||
|
[_ in never]: never;
|
||||||
|
};
|
||||||
|
Enums: {
|
||||||
|
[_ in never]: never;
|
||||||
|
};
|
||||||
|
CompositeTypes: {
|
||||||
|
[_ in never]: never;
|
||||||
|
};
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|||||||
@@ -34,6 +34,91 @@ class SupabaseDatabase<Database = any> {
|
|||||||
private projectRef: string
|
private projectRef: string
|
||||||
) {}
|
) {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The function `on` creates a trigger for when a record is inserted, updated, or deleted on a
|
||||||
|
* specific table in a database schema.
|
||||||
|
* @param params - The `params` parameter is an object that contains the following properties:
|
||||||
|
* @param params.table - The `table` property is a string that specifies the name of the table
|
||||||
|
* that the trigger will be created for.
|
||||||
|
* @param params.events - The `events` property is an array of events that specifies the events
|
||||||
|
* that the trigger will be called for. The events that can be specified are `INSERT`, `UPDATE`, or `DELETE`.
|
||||||
|
* By default, the trigger will be called for all events.
|
||||||
|
* @param params.schema - The `schema` property is a string that specifies the name of the schema
|
||||||
|
* that the trigger will be created for. If the schema is not specified, the default schema will
|
||||||
|
* be used. (public)
|
||||||
|
* @param params.filter - The `filter` property is an object that specifies the filter that will
|
||||||
|
* be used to determine if the trigger should be called. If the filter is not specified, the
|
||||||
|
* trigger will be called for all records.
|
||||||
|
*
|
||||||
|
* @example
|
||||||
|
*
|
||||||
|
* ```ts
|
||||||
|
* const supabase = new SupabaseManagement({ id: "supabase" });
|
||||||
|
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
|
||||||
|
*
|
||||||
|
* client.defineJob({
|
||||||
|
* trigger: database.on({
|
||||||
|
* table: "todos",
|
||||||
|
* events: ["INSERTED", "UPDATED"],
|
||||||
|
* schema: "public",
|
||||||
|
* filter: {
|
||||||
|
* record: { is_completed: [false] },
|
||||||
|
* },
|
||||||
|
* }),
|
||||||
|
* })
|
||||||
|
* ```
|
||||||
|
*/
|
||||||
|
on<
|
||||||
|
SchemaName extends string & keyof Database = "public" extends keyof Database
|
||||||
|
? "public"
|
||||||
|
: string & keyof Database,
|
||||||
|
Schema extends GenericSchema = Database[SchemaName] extends GenericSchema
|
||||||
|
? Database[SchemaName]
|
||||||
|
: any,
|
||||||
|
TTableName extends string & keyof Schema["Tables"] = string & keyof Schema["Tables"],
|
||||||
|
TTable extends Schema["Tables"][TTableName] = Schema["Tables"][TTableName],
|
||||||
|
TEvents extends WebhookEvents[] = ["INSERT", "UPDATE", "DELETE"],
|
||||||
|
>(params: { table: TTableName; events?: TEvents; schema?: SchemaName; filter?: EventFilter }) {
|
||||||
|
return createTrigger<Prettify<UnionPayloads<TEvents, TTableName, SchemaName, TTable["Row"]>>>(
|
||||||
|
this.integration.source,
|
||||||
|
{
|
||||||
|
event: params.events ?? ["INSERT", "UPDATE", "DELETE"],
|
||||||
|
projectRef: this.projectRef,
|
||||||
|
...params,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The function `onInserted` creates a trigger for when a new record is inserted into a specific
|
||||||
|
* table in a database schema.
|
||||||
|
* @param params - The `params` parameter is an object that contains the following properties:
|
||||||
|
* @param params.table - The `table` property is a string that specifies the name of the table
|
||||||
|
* that the trigger will be created for.
|
||||||
|
* @param params.schema - The `schema` property is a string that specifies the name of the schema
|
||||||
|
* that the trigger will be created for. If the schema is not specified, the default schema will
|
||||||
|
* be used. (public)
|
||||||
|
* @param params.filter - The `filter` property is an object that specifies the filter that will
|
||||||
|
* be used to determine if the trigger should be called. If the filter is not specified, the
|
||||||
|
* trigger will be called for all records.
|
||||||
|
*
|
||||||
|
* @example
|
||||||
|
*
|
||||||
|
* ```ts
|
||||||
|
* const supabase = new SupabaseManagement({ id: "supabase" });
|
||||||
|
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
|
||||||
|
*
|
||||||
|
* client.defineJob({
|
||||||
|
* trigger: database.onInserted({
|
||||||
|
* table: "todos",
|
||||||
|
* schema: "public",
|
||||||
|
* filter: {
|
||||||
|
* record: { is_completed: [false] },
|
||||||
|
* },
|
||||||
|
* }),
|
||||||
|
* })
|
||||||
|
* ```
|
||||||
|
*/
|
||||||
onInserted<
|
onInserted<
|
||||||
SchemaName extends string & keyof Database = "public" extends keyof Database
|
SchemaName extends string & keyof Database = "public" extends keyof Database
|
||||||
? "public"
|
? "public"
|
||||||
@@ -57,6 +142,37 @@ class SupabaseDatabase<Database = any> {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The function `onUpdated` creates a trigger for when a new record is updated on a specific
|
||||||
|
* table in a database schema.
|
||||||
|
* @param params - The `params` parameter is an object that contains the following properties:
|
||||||
|
* @param params.table - The `table` property is a string that specifies the name of the table
|
||||||
|
* that the trigger will be created for.
|
||||||
|
* @param params.schema - The `schema` property is a string that specifies the name of the schema
|
||||||
|
* that the trigger will be created for. If the schema is not specified, the default schema will
|
||||||
|
* be used. (public)
|
||||||
|
* @param params.filter - The `filter` property is an object that specifies the filter that will
|
||||||
|
* be used to determine if the trigger should be called. If the filter is not specified, the
|
||||||
|
* trigger will be called for all records.
|
||||||
|
*
|
||||||
|
* @example
|
||||||
|
*
|
||||||
|
* ```ts
|
||||||
|
* const supabase = new SupabaseManagement({ id: "supabase" });
|
||||||
|
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
|
||||||
|
*
|
||||||
|
* client.defineJob({
|
||||||
|
* trigger: database.onUpdated({
|
||||||
|
* table: "todos",
|
||||||
|
* schema: "public",
|
||||||
|
* filter: {
|
||||||
|
* record: { completed: [true] },
|
||||||
|
* old_record: { completed: [false] },
|
||||||
|
* },
|
||||||
|
* }),
|
||||||
|
* })
|
||||||
|
* ```
|
||||||
|
*/
|
||||||
onUpdated<
|
onUpdated<
|
||||||
SchemaName extends string & keyof Database = "public" extends keyof Database
|
SchemaName extends string & keyof Database = "public" extends keyof Database
|
||||||
? "public"
|
? "public"
|
||||||
@@ -80,6 +196,36 @@ class SupabaseDatabase<Database = any> {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The function `onDeleted` creates a trigger for when a new record is deleted from a specific
|
||||||
|
* table in a database schema.
|
||||||
|
* @param params - The `params` parameter is an object that contains the following properties:
|
||||||
|
* @param params.table - The `table` property is a string that specifies the name of the table
|
||||||
|
* that the trigger will be created for.
|
||||||
|
* @param params.schema - The `schema` property is a string that specifies the name of the schema
|
||||||
|
* that the trigger will be created for. If the schema is not specified, the default schema will
|
||||||
|
* be used. (public)
|
||||||
|
* @param params.filter - The `filter` property is an object that specifies the filter that will
|
||||||
|
* be used to determine if the trigger should be called. If the filter is not specified, the
|
||||||
|
* trigger will be called for all records.
|
||||||
|
*
|
||||||
|
* @example
|
||||||
|
*
|
||||||
|
* ```ts
|
||||||
|
* const supabase = new SupabaseManagement({ id: "supabase" });
|
||||||
|
* const database = supabase.database<Database>("https://<project-id>.supabase.co");
|
||||||
|
*
|
||||||
|
* client.defineJob({
|
||||||
|
* trigger: database.onDeleted({
|
||||||
|
* table: "todos",
|
||||||
|
* schema: "public",
|
||||||
|
* filter: {
|
||||||
|
* old_record: { is_completed: [true] },
|
||||||
|
* },
|
||||||
|
* }),
|
||||||
|
* })
|
||||||
|
* ```
|
||||||
|
*/
|
||||||
onDeleted<
|
onDeleted<
|
||||||
SchemaName extends string & keyof Database = "public" extends keyof Database
|
SchemaName extends string & keyof Database = "public" extends keyof Database
|
||||||
? "public"
|
? "public"
|
||||||
@@ -175,9 +321,44 @@ type WebhookEventSource = ReturnType<typeof createWebhookEventSource>;
|
|||||||
|
|
||||||
type WebhookEvents = "INSERT" | "UPDATE" | "DELETE";
|
type WebhookEvents = "INSERT" | "UPDATE" | "DELETE";
|
||||||
|
|
||||||
|
type WebhookEventPayloads<
|
||||||
|
TTableName extends string,
|
||||||
|
TSchemaName extends string = "public",
|
||||||
|
TRecord = any,
|
||||||
|
> = {
|
||||||
|
INSERT: {
|
||||||
|
table: TTableName;
|
||||||
|
record: Prettify<TRecord>;
|
||||||
|
type: "INSERT";
|
||||||
|
schema: TSchemaName;
|
||||||
|
old_record: null;
|
||||||
|
};
|
||||||
|
UPDATE: {
|
||||||
|
table: TTableName;
|
||||||
|
record: Prettify<TRecord>;
|
||||||
|
type: "UPDATE";
|
||||||
|
schema: TSchemaName;
|
||||||
|
old_record: Prettify<TRecord>;
|
||||||
|
};
|
||||||
|
DELETE: {
|
||||||
|
table: TTableName;
|
||||||
|
record: null;
|
||||||
|
type: "DELETE";
|
||||||
|
schema: TSchemaName;
|
||||||
|
old_record: Prettify<TRecord>;
|
||||||
|
};
|
||||||
|
};
|
||||||
|
|
||||||
|
type UnionPayloads<
|
||||||
|
T extends WebhookEvents[],
|
||||||
|
TTableName extends string,
|
||||||
|
TSchemaName extends string = "public",
|
||||||
|
TRecord = any,
|
||||||
|
> = WebhookEventPayloads<TTableName, TSchemaName, TRecord>[T[number]];
|
||||||
|
|
||||||
function createTrigger<TEvent extends any>(
|
function createTrigger<TEvent extends any>(
|
||||||
source: WebhookEventSource,
|
source: WebhookEventSource,
|
||||||
params: { event: WebhookEvents; filter?: EventFilter } & {
|
params: { event: WebhookEvents | WebhookEvents[]; filter?: EventFilter } & {
|
||||||
projectRef: string;
|
projectRef: string;
|
||||||
table: string;
|
table: string;
|
||||||
schema?: string;
|
schema?: string;
|
||||||
@@ -190,7 +371,7 @@ function createTrigger<TEvent extends any>(
|
|||||||
icon: "supabase",
|
icon: "supabase",
|
||||||
filter: {
|
filter: {
|
||||||
...params.filter,
|
...params.filter,
|
||||||
type: [params.event],
|
type: typeof params.event === "string" ? [params.event] : params.event,
|
||||||
schema: [params.schema ?? "public"],
|
schema: [params.schema ?? "public"],
|
||||||
},
|
},
|
||||||
properties: [],
|
properties: [],
|
||||||
|
|||||||
Reference in New Issue
Block a user