diff --git a/app/api/rpc/[command]/pull.ts b/app/api/rpc/[command]/pull.ts new file mode 100644 index 00000000..ce75a003 --- /dev/null +++ b/app/api/rpc/[command]/pull.ts @@ -0,0 +1,103 @@ +import { z } from "zod"; +import { + PullRequest, + PullResponseV1, + VersionNotSupportedResponse, +} from "replicache"; +import { Database } from "supabase/database.types"; +import { Fact } from "src/replicache"; +import postgres from "postgres"; +import { drizzle } from "drizzle-orm/postgres-js"; +import { FactWithIndexes, getClientGroup } from "src/replicache/utils"; +import { Attributes } from "src/replicache/attributes"; +import { permission_tokens } from "drizzle/schema"; +import { eq } from "drizzle-orm"; +import { makeRoute } from "../lib"; +import { Env } from "./route"; + +// First define the sub-types for V0 and V1 requests +const pullRequestV0 = z.object({ + pullVersion: z.literal(0), + schemaVersion: z.string(), + profileID: z.string(), + cookie: z.any(), // ReadonlyJSONValue + clientID: z.string(), + lastMutationID: z.number(), +}); + +// For the Cookie type used in V1 +const cookieType = z.union([ + z.null(), + z.string(), + z.number(), + z + .object({ + order: z.union([z.string(), z.number()]), + }) + .and(z.record(z.string(), z.any())), // ReadonlyJSONValue with order property +]); + +const pullRequestV1 = z.object({ + pullVersion: z.literal(1), + schemaVersion: z.string(), + profileID: z.string(), + cookie: cookieType, + clientGroupID: z.string(), +}); + +// Combined PullRequest type +const PullRequestSchema = z.union([pullRequestV0, pullRequestV1]); + +export const pull = makeRoute({ + route: "pull", + input: z.object({ pullRequest: PullRequestSchema, token_id: z.string() }), + handler: async ({ pullRequest, token_id }, { db, supabase }: Env) => { + let body = pullRequest; + if (body.pullVersion === 0) return versionNotSupported; + let [token] = await db + .select({ root_entity: permission_tokens.root_entity }) + .from(permission_tokens) + .where(eq(permission_tokens.id, token_id)); + let facts: { + attribute: string; + created_at: string; + data: any; + entity: string; + id: string; + updated_at: string | null; + version: number; + }[] = []; + let clientGroup = {}; + if (token) { + let { data } = await supabase.rpc("get_facts", { + root: token.root_entity, + }); + + clientGroup = await getClientGroup(db, body.clientGroupID); + facts = data || []; + } + + return { + cookie: Date.now(), + lastMutationIDChanges: clientGroup, + patch: [ + { op: "clear" }, + { op: "put", key: "initialized", value: true }, + ...facts.map((f) => { + return { + op: "put", + key: f.id, + value: FactWithIndexes( + f as unknown as Fact, + ), + } as const; + }), + ], + } as PullResponseV1; + }, +}); + +const versionNotSupported: VersionNotSupportedResponse = { + error: "VersionNotSupported", + versionType: "pull", +}; diff --git a/app/api/rpc/[command]/push.ts b/app/api/rpc/[command]/push.ts new file mode 100644 index 00000000..257d348c --- /dev/null +++ b/app/api/rpc/[command]/push.ts @@ -0,0 +1,111 @@ +import { PushResponse } from "replicache"; +import { serverMutationContext } from "src/replicache/serverMutationContext"; +import { mutations } from "src/replicache/mutations"; +import { eq } from "drizzle-orm"; +import { permission_token_rights, replicache_clients } from "drizzle/schema"; +import { getClientGroup } from "src/replicache/utils"; +import { makeRoute } from "../lib"; +import { z } from "zod"; +import { Env } from "./route"; + +const mutationV0Schema = z.object({ + id: z.number(), + name: z.string(), + args: z.unknown(), + timestamp: z.number(), +}); + +const mutationV1Schema = mutationV0Schema.extend({ + clientID: z.string(), +}); + +const pushRequestV0Schema = z.object({ + pushVersion: z.literal(0), + schemaVersion: z.string(), + profileID: z.string(), + clientID: z.string(), + mutations: z.array(mutationV0Schema), +}); + +const pushRequestV1Schema = z.object({ + pushVersion: z.literal(1), + schemaVersion: z.string(), + profileID: z.string(), + clientGroupID: z.string(), + mutations: z.array(mutationV1Schema), +}); + +// Combine both versions into final PushRequest schema +const pushRequestSchema = z.discriminatedUnion("pushVersion", [ + pushRequestV0Schema, + pushRequestV1Schema, +]); + +type PushRequestZ = z.infer; + +export const push = makeRoute({ + route: "push", + input: z.object({ + pushRequest: pushRequestSchema, + rootEntity: z.string(), + token: z.object({ id: z.string() }), + }), + handler: async ( + { pushRequest, rootEntity, token }, + { db, supabase }: Env, + ) => { + if (pushRequest.pushVersion !== 1) { + return { + result: { error: "VersionNotSupported", versionType: "push" } as const, + }; + } + let clientGroup = await getClientGroup(db, pushRequest.clientGroupID); + let token_rights = await db + .select() + .from(permission_token_rights) + .where(eq(permission_token_rights.token, token.id)); + for (let mutation of pushRequest.mutations) { + let lastMutationID = clientGroup[mutation.clientID] || 0; + if (mutation.id <= lastMutationID) continue; + clientGroup[mutation.clientID] = mutation.id; + let name = mutation.name as keyof typeof mutations; + if (!mutations[name]) { + continue; + } + await db.transaction(async (tx) => { + try { + await mutations[name]( + mutation.args as any, + serverMutationContext(tx, token_rights), + ); + } catch (e) { + console.log( + `Error occured while running mutation: ${name}`, + JSON.stringify(e), + JSON.stringify(mutation, null, 2), + ); + } + await tx + .insert(replicache_clients) + .values({ + client_group: pushRequest.clientGroupID, + client_id: mutation.clientID, + last_mutation: mutation.id, + }) + .onConflictDoUpdate({ + target: replicache_clients.client_id, + set: { last_mutation: mutation.id }, + }); + }); + } + + let channel = supabase.channel(`rootEntity:${rootEntity}`); + await channel.send({ + type: "broadcast", + event: "poke", + payload: { message: "poke" }, + }); + supabase.removeChannel(channel); + return { result: undefined } as const; + }, +}); diff --git a/app/api/rpc/[command]/route.ts b/app/api/rpc/[command]/route.ts new file mode 100644 index 00000000..74bb38ad --- /dev/null +++ b/app/api/rpc/[command]/route.ts @@ -0,0 +1,29 @@ +import { drizzle } from "drizzle-orm/postgres-js"; +import { makeRouter } from "../lib"; +import { push } from "./push"; +import postgres from "postgres"; +import { createClient } from "@supabase/supabase-js"; +import { Database } from "supabase/database.types"; +import { pull } from "./pull"; + +const client = postgres(process.env.DB_URL as string, { idle_timeout: 5 }); +let supabase = createClient( + process.env.NEXT_PUBLIC_SUPABASE_API_URL as string, + process.env.SUPABASE_SERVICE_ROLE_KEY as string, +); +const db = drizzle(client); + +const Env = { + supabase, + db, +}; +export type Env = typeof Env; +export type Routes = typeof Routes; +let Routes = [push, pull]; +export async function POST( + req: Request, + { params }: { params: { command: string } }, +) { + let router = makeRouter(Routes); + return router(params.command, req, Env); +} diff --git a/app/api/rpc/client.ts b/app/api/rpc/client.ts new file mode 100644 index 00000000..a403b1c8 --- /dev/null +++ b/app/api/rpc/client.ts @@ -0,0 +1,4 @@ +import { makeAPIClient } from "./lib"; +import type { Routes } from "./[command]/route"; + +export const callRPC = makeAPIClient("/api/rpc"); diff --git a/app/api/rpc/lib.ts b/app/api/rpc/lib.ts new file mode 100644 index 00000000..d9600913 --- /dev/null +++ b/app/api/rpc/lib.ts @@ -0,0 +1,104 @@ +import { ZodObject, ZodRawShape, ZodUnion, z } from "zod"; + +type Route< + Cmd extends string, + Input extends ZodObject | ZodUnion, + Result extends object, + Env extends {}, +> = { + route: Cmd; + input: Input; + handler: (msg: z.infer, env: Env, request: Request) => Promise; +}; + +type Routes = Route[]; + +export function makeAPIClient>(basePath: string) { + return async ( + route: T, + data: z.infer["input"]>, + ) => { + let result = await fetch(`${basePath}/${route}`, { + body: JSON.stringify(data), + method: "POST", + headers: { "Content-type": "application/json" }, + }); + return result.json() as Promise< + Awaited["handler"]>> + >; + }; +} + +export const makeRouter = (routes: Routes) => { + return async (route: string, request: Request, env: Env) => { + let status = 200; + let result; + switch (request.method) { + case "POST": { + let handler = routes.find((f) => f.route === route); + if (!handler) { + status = 404; + result = { error: `route ${route} not Found` }; + break; + } + + let body; + if (handler.input) + try { + body = await request.json(); + } catch (e) { + result = { error: "Request body must be valid JSON" }; + status = 400; + break; + } + + let msg = handler.input.safeParse(body); + if (!msg.success) { + status = 400; + result = msg.error; + break; + } + try { + result = (await handler.handler( + msg.data as any, + env, + request, + )) as object; + break; + } catch (e) { + console.log(e); + status = 500; + result = { + error: "An error occured while handling this request", + errorText: (e as Error).toString(), + }; + break; + } + } + default: + status = 404; + result = { error: "Only POST Supported" }; + } + + let res = new Response(JSON.stringify(result), { + status, + headers: { + "Access-Control-Allow-Credentials": "true", + "Content-type": "application/json;charset=UTF-8", + "Access-Control-Allow-Origin": "*", + "Access-Control-Allow-Methods": "GET,HEAD,POST,OPTIONS", + }, + }); + //result.headers?.forEach((h) => res.headers.append(h[0], h[1])); + return res; + }; +}; + +export function makeRoute< + Cmd extends string, + Input extends ZodObject | ZodUnion, + Result extends object, + Env extends {}, +>(d: Route) { + return d; +} diff --git a/src/replicache/index.tsx b/src/replicache/index.tsx index 6b60ee07..fcedfeeb 100644 --- a/src/replicache/index.tsx +++ b/src/replicache/index.tsx @@ -8,12 +8,11 @@ import { Replicache, WriteTransaction, } from "replicache"; -import { Pull } from "./pull"; import { mutations } from "./mutations"; import { Attributes, Data, FilterAttributes } from "./attributes"; -import { Push } from "./push"; import { clientMutationContext } from "./clientMutationContext"; import { supabaseBrowserClient } from "supabase/browserClient"; +import { callRPC } from "app/api/rpc/client"; export type Fact = { id: string; @@ -82,13 +81,23 @@ export function ReplicacheProvider(props: { mutations: pushRequest.mutations.slice(0, 250), } as PushRequest; return { - response: await Push(smolpushRequest, props.name, props.token), + response: ( + await callRPC("push", { + pushRequest: smolpushRequest, + token: props.token, + rootEntity: props.name, + }) + ).result, httpRequestInfo: { errorMessage: "", httpStatusCode: 200 }, }; }, puller: async (pullRequest) => { + let res = await callRPC("pull", { + pullRequest, + token_id: props.token.id, + }); return { - response: await Pull(pullRequest, props.token.id), + response: res, httpRequestInfo: { errorMessage: "", httpStatusCode: 200 }, }; }, diff --git a/src/replicache/pull.ts b/src/replicache/pull.ts deleted file mode 100644 index 3d967c61..00000000 --- a/src/replicache/pull.ts +++ /dev/null @@ -1,71 +0,0 @@ -"use server"; - -import { createClient } from "@supabase/supabase-js"; -import { - PullRequest, - PullResponseV1, - VersionNotSupportedResponse, -} from "replicache"; -import { Database } from "supabase/database.types"; -import { Fact } from "."; -import postgres from "postgres"; -import { drizzle } from "drizzle-orm/postgres-js"; -import { FactWithIndexes, getClientGroup } from "./utils"; -import { Attributes } from "./attributes"; -import { permission_tokens } from "drizzle/schema"; -import { eq } from "drizzle-orm"; -let supabase = createClient( - process.env.NEXT_PUBLIC_SUPABASE_API_URL as string, - process.env.SUPABASE_SERVICE_ROLE_KEY as string, -); - -const client = postgres(process.env.DB_URL as string, { idle_timeout: 5 }); -const db = drizzle(client); -export async function Pull( - body: PullRequest, - token_id: string, -): Promise { - console.log("Pull"); - if (body.pullVersion === 0) return versionNotSupported; - let [token] = await db - .select({ root_entity: permission_tokens.root_entity }) - .from(permission_tokens) - .where(eq(permission_tokens.id, token_id)); - let facts: { - attribute: string; - created_at: string; - data: any; - entity: string; - id: string; - updated_at: string | null; - version: number; - }[] = []; - let clientGroup = {}; - if (token) { - let { data } = await supabase.rpc("get_facts", { root: token.root_entity }); - - clientGroup = await getClientGroup(db, body.clientGroupID); - facts = data || []; - } - - return { - cookie: Date.now(), - lastMutationIDChanges: clientGroup, - patch: [ - { op: "clear" }, - { op: "put", key: "initialized", value: true }, - ...facts.map((f) => { - return { - op: "put", - key: f.id, - value: FactWithIndexes(f as unknown as Fact), - } as const; - }), - ], - }; -} - -const versionNotSupported: VersionNotSupportedResponse = { - error: "VersionNotSupported", - versionType: "pull", -}; diff --git a/src/replicache/push.ts b/src/replicache/push.ts deleted file mode 100644 index db82c315..00000000 --- a/src/replicache/push.ts +++ /dev/null @@ -1,76 +0,0 @@ -"use server"; -import { PushRequest, PushResponse } from "replicache"; -import { serverMutationContext } from "./serverMutationContext"; -import { mutations } from "./mutations"; -import { drizzle } from "drizzle-orm/postgres-js"; -import { eq } from "drizzle-orm"; -import postgres from "postgres"; -import { permission_token_rights, replicache_clients } from "drizzle/schema"; -import { getClientGroup } from "./utils"; -import { createClient } from "@supabase/supabase-js"; -import { Database } from "supabase/database.types"; - -const client = postgres(process.env.DB_URL as string, { idle_timeout: 5 }); -let supabase = createClient( - process.env.NEXT_PUBLIC_SUPABASE_API_URL as string, - process.env.SUPABASE_SERVICE_ROLE_KEY as string, -); -const db = drizzle(client); -export async function Push( - pushRequest: PushRequest, - rootEntity: string, - token: { id: string }, -): Promise { - console.log("Push"); - if (pushRequest.pushVersion !== 1) { - return { error: "VersionNotSupported", versionType: "push" }; - } - let clientGroup = await getClientGroup(db, pushRequest.clientGroupID); - let token_rights = await db - .select() - .from(permission_token_rights) - .where(eq(permission_token_rights.token, token.id)); - for (let mutation of pushRequest.mutations) { - let lastMutationID = clientGroup[mutation.clientID] || 0; - if (mutation.id <= lastMutationID) continue; - clientGroup[mutation.clientID] = mutation.id; - let name = mutation.name as keyof typeof mutations; - if (!mutations[name]) { - continue; - } - await db.transaction(async (tx) => { - try { - await mutations[name]( - mutation.args as any, - serverMutationContext(tx, token_rights), - ); - } catch (e) { - console.log( - `Error occured while running mutation: ${name}`, - JSON.stringify(e), - JSON.stringify(mutation, null, 2), - ); - } - await tx - .insert(replicache_clients) - .values({ - client_group: pushRequest.clientGroupID, - client_id: mutation.clientID, - last_mutation: mutation.id, - }) - .onConflictDoUpdate({ - target: replicache_clients.client_id, - set: { last_mutation: mutation.id }, - }); - }); - } - - let channel = supabase.channel(`rootEntity:${rootEntity}`); - await channel.send({ - type: "broadcast", - event: "poke", - payload: { message: "poke" }, - }); - supabase.removeChannel(channel); - return undefined; -}