Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
"scripts": {
"build": "bun run --cwd packages/tauri-release build && bun run --cwd packages/auth build && bun run --cwd packages/brand build && bun run --cwd packages/flags build && bun run --cwd packages/billing build && bun run --cwd packages/edge-shared build && bun run --cwd packages/sync build",
"test": "vitest run",
"test:supabase": "supabase test db supabase/tests/database --local && bun run supabase/tests/welcome-credits-concurrency.ts && bun test packages/billing/test/checkout-lifecycle.contract.test.ts",
"test:supabase": "supabase test db supabase/tests/database --local && bun run supabase/tests/welcome-credits-concurrency.ts && bun test packages/billing/test/checkout-lifecycle.contract.test.ts && bun test packages/sync/test/lww.contract.test.ts",
"test:coverage": "vitest run --coverage",
"typecheck": "tsc -p tsconfig.json --noEmit",
"check": "bun run typecheck && bun run test && bun run test:coverage && bun run build"
Expand Down
129 changes: 57 additions & 72 deletions packages/sync/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,17 +33,8 @@ export interface SupabaseQueryResult<T> {

export interface SupabaseQueryBuilderLike<T = unknown> extends PromiseLike<SupabaseQueryResult<T>> {
select(columns?: string): SupabaseQueryBuilderLike<T>;
update(values: unknown): SupabaseQueryBuilderLike<T>;
upsert(
values: unknown,
options?: {
onConflict?: string;
ignoreDuplicates?: boolean;
},
): SupabaseQueryBuilderLike<T>;
eq(column: string, value: unknown): SupabaseQueryBuilderLike<T>;
is(column: string, value: unknown): SupabaseQueryBuilderLike<T>;
lt(column: string, value: unknown): SupabaseQueryBuilderLike<T>;
order(
column: string,
options?: {
Expand All @@ -55,6 +46,7 @@ export interface SupabaseQueryBuilderLike<T = unknown> extends PromiseLike<Supab

export interface SupabaseClientLike {
from<T = unknown>(table: string): SupabaseQueryBuilderLike<T>;
rpc(fn: string, args?: Record<string, unknown>): unknown;
}

export interface PushGroupOptions {
Expand Down Expand Up @@ -171,25 +163,20 @@ export async function pushGroup({
const validGroup = validateFieldGroup(group);
const validUpdatedAt = validateUpdatedAt(updatedAt);
const sanitizedPayload = sanitizeSyncPayload(validGroup, payload);
const row = {

const { written, row } = await applyLwwWrite(supabase, {
user_id: validUserId,
field_group: validGroup,
payload: sanitizedPayload,
updated_at: validUpdatedAt,
deleted_at: null,
};

const written = await writeRowWithLww(supabase, row);
if (written !== null) {
return { status: "written", record: toSyncRecord(written) };
}
});

const server = await selectStoredGroup(supabase, validUserId, validGroup);
if (server === null) {
throw new SyncError("database_error", "Failed to read stale sync group");
if (written) {
return { status: "written", record: toSyncRecord(row) };
}

return { status: "stale", record: server.deleted_at === null ? toSyncRecord(server) : null };
return { status: "stale", record: row.deleted_at === null ? toSyncRecord(row) : null };
}

export async function pullAll({ supabase, userId }: PullAllOptions): Promise<SyncRecord[]> {
Expand Down Expand Up @@ -225,23 +212,19 @@ export async function deleteGroup({
const validGroup = validateFieldGroup(group);
const validUpdatedAt = validateUpdatedAt(updatedAt);

const written = await writeRowWithLww(supabase, {
const { written, row } = await applyLwwWrite(supabase, {
user_id: validUserId,
field_group: validGroup,
payload: {},
updated_at: validUpdatedAt,
deleted_at: validUpdatedAt,
});
if (written !== null) {
return { status: "deleted" };
}

const server = await selectStoredGroup(supabase, validUserId, validGroup);
if (server === null) {
throw new SyncError("database_error", "Failed to read stale sync group");
if (written) {
return { status: "deleted" };
}

return { status: "stale", record: server.deleted_at === null ? toSyncRecord(server) : null };
return { status: "stale", record: row.deleted_at === null ? toSyncRecord(row) : null };
}

async function selectGroup(
Expand All @@ -263,58 +246,60 @@ async function selectGroup(
return data === null ? null : toSyncRecord(data);
}

async function selectStoredGroup(
interface LwwWriteResult {
written: boolean;
row: SettingsSyncRow;
}

/**
* Calls the durable `settings_sync_lww_write` Postgres function (see
* supabase/migrations/20260720000000_settings_sync_lww.sql), which locks the
* (user_id, field_group) row, compares `row.updated_at` against it, and
* performs whichever of insert/update/no-op wins, all inside one durable
* operation. It always returns the row that is now authoritative, so callers
* never need a follow-up read to learn the outcome.
*/
async function applyLwwWrite(
supabase: SupabaseClientLike,
userId: string,
group: SyncFieldGroup,
): Promise<SettingsSyncRow | null> {
const { data, error } = await supabase
.from<SettingsSyncRow>("settings_sync")
.select(SYNC_COLUMNS)
.eq("user_id", userId)
.eq("field_group", group)
.maybeSingle();
row: SettingsSyncRow,
): Promise<LwwWriteResult> {
const { data, error } = (await supabase.rpc("settings_sync_lww_write", {
target_user: row.user_id,
target_group: row.field_group,
new_payload: row.payload,
new_updated_at: row.updated_at,
new_deleted_at: row.deleted_at,
})) as SupabaseQueryResult<unknown>;
if (error !== null) {
throw databaseError("Failed to pull stored sync group", error);
throw databaseError("Failed to durably write sync group", error);
}

return data;
return toLwwWriteResult(data);
}

async function writeRowWithLww(
supabase: SupabaseClientLike,
row: SettingsSyncRow,
): Promise<SettingsSyncRow | null> {
for (let attempt = 0; attempt < 2; attempt += 1) {
const updated = await supabase
.from<SettingsSyncRow>("settings_sync")
.update(row)
.eq("user_id", row.user_id)
.eq("field_group", row.field_group)
.lt("updated_at", row.updated_at)
.select(SYNC_COLUMNS)
.maybeSingle();
if (updated.error !== null) {
throw databaseError("Failed to update sync group", updated.error);
}
if (updated.data !== null) {
return updated.data;
}

const inserted = await supabase
.from<SettingsSyncRow>("settings_sync")
.upsert(row, { onConflict: "user_id,field_group", ignoreDuplicates: true })
.select(SYNC_COLUMNS)
.maybeSingle();
if (inserted.error !== null) {
throw databaseError("Failed to insert sync group", inserted.error);
}
if (inserted.data !== null) {
return inserted.data;
}
function toLwwWriteResult(data: unknown): LwwWriteResult {
if (
!isRecord(data) ||
typeof data.written !== "boolean" ||
typeof data.user_id !== "string" ||
typeof data.field_group !== "string" ||
typeof data.updated_at !== "string" ||
!isJson(data.payload) ||
(data.deleted_at !== null && typeof data.deleted_at !== "string")
) {
throw new SyncError("database_error", "settings_sync_lww_write returned an invalid result");
}

return null;
return {
written: data.written,
row: {
user_id: data.user_id,
field_group: data.field_group,
payload: data.payload,
updated_at: data.updated_at,
deleted_at: data.deleted_at,
},
};
}

function walkPayload(group: SyncFieldGroup, value: Json, path: string[]): void {
Expand Down
99 changes: 99 additions & 0 deletions packages/sync/test/lww-fixtures.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
// Local-Postgres contract-lane fixtures. Mirrors the `SUPABASE_DB_URL` /
// localhost-only guard used by `packages/billing/test/lifecycle-fixtures.ts`
// and `supabase/tests/welcome-credits-concurrency.ts` so this lane can never
// accidentally target a non-disposable database, and gracefully reports why
// it is skipped when no local Supabase is running.
import { SQL } from "bun";

export const DEFAULT_DATABASE_URL = "postgresql://postgres:postgres@127.0.0.1:54322/postgres";

export function resolveLocalDatabaseUrl(): string | null {
const raw = process.env.SUPABASE_DB_URL ?? DEFAULT_DATABASE_URL;
let parsed: URL;
try {
parsed = new URL(raw);
} catch {
return null;
}
const isDisposableLocalDatabase =
["127.0.0.1", "localhost", "[::1]"].includes(parsed.hostname) &&
parsed.port === "54322" &&
parsed.pathname === "/postgres" &&
parsed.username === "postgres";
return isDisposableLocalDatabase ? raw : null;
}

/**
* Connects to the local Supabase Postgres started by `supabase start`. Returns
* `null` (never throws) when the database is unreachable so the contract lane
* can skip cleanly instead of failing CI runs that have no local Postgres.
*/
export async function connectToLocalDatabase(): Promise<SQL | null> {
const databaseUrl = resolveLocalDatabaseUrl();
if (databaseUrl === null) {
return null;
}

const sql = new SQL(databaseUrl, { max: 10 });
try {
await sql.unsafe("select 1");
} catch {
await sql.close().catch(() => {});
return null;
}
return sql;
}

export function uniqueId(prefix: string): string {
return `${prefix}_${crypto.randomUUID().replaceAll("-", "")}`;
}

export async function insertAuthUser(
sql: Pick<SQL, "unsafe">,
user: { id: string; email: string },
): Promise<void> {
await sql.unsafe(
`insert into auth.users (
id, aud, role, email, raw_app_meta_data, raw_user_meta_data, is_anonymous, created_at, updated_at
) values ($1, 'authenticated', 'authenticated', $2, '{}'::jsonb, '{}'::jsonb, false, now(), now())`,
[user.id, user.email],
);
}

/**
* Cleans up rows the contract lane cannot roll back transactionally: used
* only by the true-concurrency tests, which need two independent
* connections/transactions racing for real, so they commit instead of
* rolling back. Deleting the auth user cascades its settings_sync rows.
*/
export async function cleanupSyncFixtures(
sql: Pick<SQL, "unsafe" | "array">,
fixtures: { userIds: string[] },
): Promise<void> {
if (fixtures.userIds.length === 0) {
return;
}
await sql.unsafe("delete from auth.users where id = any($1::uuid[])", [
sql.array(fixtures.userIds, "uuid"),
]);
}

/**
* Runs `run` inside a transaction that always rolls back, so each contract
* test is fully isolated without needing bespoke per-table cleanup. Only the
* true-concurrency tests opt out of this helper, since they need two
* independently committing connections.
*/
export async function withRollback(sql: SQL, run: (tx: SQL) => Promise<void>): Promise<void> {
const rollbackSentinel = Symbol("contract-lane-rollback");
try {
await sql.begin(async (tx) => {
await run(tx);
throw rollbackSentinel;
});
} catch (error) {
if (error !== rollbackSentinel) {
throw error;
}
}
}
Loading