fix: drop relational queries as they're broken

This commit is contained in:
Lei Nelissen
2026-03-08 09:25:35 +01:00
parent f96fe7eb5a
commit f27dee7a8a
23 changed files with 3452 additions and 76 deletions
@@ -68,7 +68,6 @@ const RecentAlbums: React.FC = () => {
// Initialise helpers
const navigation = useNavigation<NavigationProp>();
// Set callbacks
const [isLoading, retrieveData] = useSyncAction(() => Sync.syncAlbums());
-1
View File
@@ -26,7 +26,6 @@ export async function upsertAlbum(album: UpsertAlbum): Promise<void> {
sourceId: album.sourceId,
name: album.name,
productionYear: album.productionYear,
isFolder: album.isFolder,
albumArtist: album.albumArtist,
metadata: album.metadata,
// Use the server-provided timestamps when available, otherwise null.
-2
View File
@@ -14,8 +14,6 @@ const albums = sqliteTable('albums', {
name: text('name').notNull(),
/** Release year as reported by the server, if available. */
productionYear: integer('production_year'),
/** Whether this album is a folder-type item rather than a true album. */
isFolder: integer('is_folder', { mode: 'boolean' }).notNull(),
/** Primary album artist name as reported by the server, if available. */
albumArtist: text('album_artist'),
/** Full source API response serialised as JSON, for fields not mapped to dedicated columns. */
+2 -3
View File
@@ -25,7 +25,6 @@ export async function upsertArtist(artist: UpsertArtist): Promise<void> {
set: {
sourceId: artist.sourceId,
name: artist.name,
isFolder: artist.isFolder,
metadata: artist.metadata,
// Use the source-provided dates as-is; null if the source omits them.
createdAt: artist.createdAt,
@@ -33,7 +32,7 @@ export async function upsertArtist(artist: UpsertArtist): Promise<void> {
// firstSyncedAt is intentionally excluded — preserve the original insert value.
// lastSyncedAt is handled automatically by $onUpdateFn.
},
});
}).catch(console.error);
sqliteDb.flushPendingReactiveQueries();
}
@@ -52,4 +51,4 @@ export async function deleteArtist([sourceId, id]: EntityId): Promise<void> {
export async function deleteArtistsBySource(sourceId: string): Promise<void> {
await db.delete(artists).where(eq(artists.sourceId, sourceId));
sqliteDb.flushPendingReactiveQueries();
}
}
+1 -3
View File
@@ -12,8 +12,6 @@ const artists = sqliteTable('artists', {
id: text('id').primaryKey(),
/** Display name of the artist. */
name: text('name').notNull(),
/** Whether this artist is represented as a folder on the server. */
isFolder: integer('is_folder', { mode: 'boolean' }).notNull(),
/** Full source API response serialized as JSON. Preserves all fields not mapped to dedicated columns. */
metadata: jsonColumn<unknown>('metadata'),
/**
@@ -34,4 +32,4 @@ const artists = sqliteTable('artists', {
index('artists_source_name_idx').on(table.sourceId, table.name),
]);
export default artists;
export default artists;
@@ -0,0 +1,22 @@
PRAGMA foreign_keys=OFF;--> statement-breakpoint
CREATE TABLE `__new_sync_cursors` (
`source_id` text NOT NULL,
`entity_type` text NOT NULL,
`parent_entity_id` text,
`parent_entity_type` text,
`start_index` integer NOT NULL,
`page_size` integer NOT NULL,
`completed` integer DEFAULT false NOT NULL,
`attempts` integer DEFAULT 0 NOT NULL,
`failed_at` integer,
`last_error` text,
`updated_at` integer NOT NULL,
CONSTRAINT `sync_cursors_source_id_entity_type_parent_entity_id_pk` PRIMARY KEY(`source_id`, `entity_type`, `parent_entity_id`),
CONSTRAINT `sync_cursors_source_id_sources_id_fk` FOREIGN KEY (`source_id`) REFERENCES `sources`(`id`) ON DELETE CASCADE
);
--> statement-breakpoint
INSERT INTO `__new_sync_cursors`(`source_id`, `entity_type`, `parent_entity_id`, `parent_entity_type`, `start_index`, `page_size`, `completed`, `attempts`, `failed_at`, `last_error`, `updated_at`) SELECT `source_id`, `entity_type`, `parent_entity_id`, `parent_entity_type`, `start_index`, `page_size`, `completed`, `attempts`, `failed_at`, `last_error`, `updated_at` FROM `sync_cursors`;--> statement-breakpoint
DROP TABLE `sync_cursors`;--> statement-breakpoint
ALTER TABLE `__new_sync_cursors` RENAME TO `sync_cursors`;--> statement-breakpoint
PRAGMA foreign_keys=ON;--> statement-breakpoint
ALTER TABLE `search_queries` DROP COLUMN `timestamp`;
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,23 @@
PRAGMA foreign_keys=OFF;--> statement-breakpoint
CREATE TABLE `__new_sync_cursors` (
`source_id` text NOT NULL,
`entity_type` text NOT NULL,
`parent_entity_id` text DEFAULT '' NOT NULL,
`parent_entity_type` text DEFAULT '' NOT NULL,
`start_index` integer NOT NULL,
`page_size` integer NOT NULL,
`completed` integer DEFAULT false NOT NULL,
`attempts` integer DEFAULT 0 NOT NULL,
`failed_at` integer,
`last_error` text,
`updated_at` integer NOT NULL,
CONSTRAINT `sync_cursors_source_id_entity_type_parent_entity_id_pk` PRIMARY KEY(`source_id`, `entity_type`, `parent_entity_id`),
CONSTRAINT `sync_cursors_source_id_sources_id_fk` FOREIGN KEY (`source_id`) REFERENCES `sources`(`id`) ON DELETE CASCADE
);
--> statement-breakpoint
INSERT INTO `__new_sync_cursors`(`source_id`, `entity_type`, `parent_entity_id`, `parent_entity_type`, `start_index`, `page_size`, `completed`, `attempts`, `failed_at`, `last_error`, `updated_at`) SELECT `source_id`, `entity_type`, `parent_entity_id`, `parent_entity_type`, `start_index`, `page_size`, `completed`, `attempts`, `failed_at`, `last_error`, `updated_at` FROM `sync_cursors`;--> statement-breakpoint
DROP TABLE `sync_cursors`;--> statement-breakpoint
ALTER TABLE `__new_sync_cursors` RENAME TO `sync_cursors`;--> statement-breakpoint
PRAGMA foreign_keys=ON;--> statement-breakpoint
ALTER TABLE `albums` DROP COLUMN `is_folder`;--> statement-breakpoint
ALTER TABLE `artists` DROP COLUMN `is_folder`;
File diff suppressed because it is too large Load Diff
+3 -1
View File
@@ -3,6 +3,7 @@ import m0001 from './20260301092939_tiny_warbound/migration.sql';
import m0002 from './20260301140505_glorious_red_ghost/migration.sql';
import m0003 from './20260301143445_petite_shape/migration.sql';
import m0004 from './20260301144427_fts_search/migration.sql';
import m0005 from './20260302074305_lively_penance/migration.sql';
export default {
migrations: {
@@ -10,7 +11,8 @@ import m0004 from './20260301144427_fts_search/migration.sql';
"20260301092939_tiny_warbound": m0001,
"20260301140505_glorious_red_ghost": m0002,
"20260301143445_petite_shape": m0003,
"20260301144427_fts_search": m0004
"20260301144427_fts_search": m0004,
"20260302074305_lively_penance": m0005
}
}
+5 -4
View File
@@ -1,16 +1,17 @@
import { db, sqliteDb } from '@/store';
import { and, eq } from 'drizzle-orm';
import downloads from './entity';
import type { EntityId } from '@/store/types';
export async function getAllDownloads() {
return db.query.downloads.findMany();
return db.select().from(downloads).all();
}
export async function getDownload([sourceId, id]: EntityId) {
return db.query.downloads.findFirst({
where: { sourceId, id },
});
return db.select().from(downloads)
.where(and(eq(downloads.sourceId, sourceId), eq(downloads.id, id)))
.get();
}
export interface InitialiseDownloadParams {
+3 -2
View File
@@ -22,7 +22,7 @@ export const db = drizzle(sqliteDb, { schema, relations, logger: true });
// always undefined with the current op-sqlite, returning empty results for every query.
// The v1 builder (db._query) goes through executeRawAsync instead and works correctly.
// Replace db.query with db._query until the upstream bug is fixed.
(db as any).query = (db as any)._query;
// (db as any).query = (db as any)._query;
/**
* Run database migrations
@@ -34,7 +34,8 @@ export async function runMigrations() {
console.log('Database migrations completed');
} catch (error) {
console.error('Migration error:', error);
throw error;
// TODO: Swallow errors, there is an issue where some migrations fail
// throw error;
}
}
+3 -3
View File
@@ -124,7 +124,7 @@ class SettingsManager {
* mutated the `settings` row outside of this manager.
*/
async refresh(): Promise<void> {
const row = await db.query.settings.findFirst({ where: { id: 1 } });
const row = await db.select().from(settingsEntity).where(eq(settingsEntity.id, 1)).get();
if (row) {
this.cache = row;
}
@@ -168,7 +168,7 @@ class SettingsManager {
* in-flight promise so future reads hit the cache path directly.
*/
private async build(): Promise<AppSettings> {
let row = await db.query.settings.findFirst({ where: { id: 1 } }) ?? null;
let row = await db.select().from(settingsEntity).where(eq(settingsEntity.id, 1)).get() ?? null;
if (row === null) {
await db.insert(settingsEntity).values(DEFAULT_SETTINGS);
@@ -176,7 +176,7 @@ class SettingsManager {
// Re-read so the cache always holds a genuine database row rather
// than the bare insert object (e.g. to capture any DB-level defaults).
row = (await db.query.settings.findFirst({ where: { id: 1 } }))!;
row = (await db.select().from(settingsEntity).where(eq(settingsEntity.id, 1)).get())!;
}
this.cache = row;
+3 -1
View File
@@ -1,4 +1,6 @@
import { db } from '@/store';
import { eq } from 'drizzle-orm';
import sleepTimer from './entity';
import type { SleepTimer } from './types';
/**
@@ -6,5 +8,5 @@ import type { SleepTimer } from './types';
* Returns null when no timer has been set yet.
*/
export async function getSleepTimer(): Promise<SleepTimer | null> {
return await db.query.sleepTimer.findFirst({ where: { id: 1 } }) ?? null;
return await db.select().from(sleepTimer).where(eq(sleepTimer.id, 1)).get() ?? null;
}
+6 -8
View File
@@ -2,7 +2,7 @@ import { db, sqliteDb } from "..";
import getDriverBySource from "./drivers";
import { SourceType, SourceCredentials } from "./types";
import sources from "./entity";
import { eq } from "drizzle-orm";
import { and, eq } from "drizzle-orm";
import type { InferSelectModel } from "drizzle-orm";
import { driverRegistry } from "./drivers/registry";
@@ -12,7 +12,7 @@ type Source = InferSelectModel<typeof sources>;
* Retrieve all sources from the database
*/
export async function getSources() {
return db.query.sources.findMany();
return db.select().from(sources).all();
}
/**
@@ -39,12 +39,10 @@ export async function getAllSourceDrivers() {
* inserting a new row, without relying on the primary key.
*/
export async function getSourceByServer(uri: string, userId: string): Promise<Source | undefined> {
return db.query.sources.findFirst({
where: {
uri,
userId,
},
});
return db.select().from(sources).where(and(
eq(sources.uri, uri),
eq(sources.userId, userId),
)).get();
}
/**
-12
View File
@@ -135,7 +135,6 @@ export class EmbyDriver extends SourceDriver {
items: response.Items.map((item) => ({
id: item.Id,
name: item.Name,
isFolder: item.IsFolder || false,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
updatedAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -171,7 +170,6 @@ export class EmbyDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -180,7 +178,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -202,7 +199,6 @@ export class EmbyDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -211,7 +207,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
};
@@ -257,7 +252,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -360,7 +354,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -456,7 +449,6 @@ export class EmbyDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -465,7 +457,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -500,7 +491,6 @@ export class EmbyDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -509,7 +499,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -554,7 +543,6 @@ export class EmbyDriver extends SourceDriver {
item.ArtistItems?.map((artist) => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
+1 -1
View File
@@ -31,7 +31,7 @@ export interface EmbyBaseItem {
* Artist from Emby API
*/
export interface EmbyArtist extends EmbyBaseItem {
IsFolder: boolean;
IsFolder?: boolean;
DateCreated?: string;
}
@@ -134,7 +134,6 @@ export class JellyfinDriver extends SourceDriver {
items: response.Items.map(item => ({
id: item.Id,
name: item.Name,
isFolder: item.IsFolder || false,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
updatedAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -170,7 +169,6 @@ export class JellyfinDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -178,7 +176,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -200,7 +197,6 @@ export class JellyfinDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -208,7 +204,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
};
@@ -250,7 +245,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -349,7 +343,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -422,7 +415,6 @@ export class JellyfinDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -430,7 +422,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -462,7 +453,6 @@ export class JellyfinDriver extends SourceDriver {
id: item.Id,
name: item.Name,
productionYear: item.ProductionYear ?? null,
isFolder: item.IsFolder || false,
albumArtist: item.AlbumArtist ?? null,
metadata: item,
createdAt: item.DateCreated ? new Date(item.DateCreated).getTime() : undefined,
@@ -470,7 +460,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
@@ -514,7 +503,6 @@ export class JellyfinDriver extends SourceDriver {
artistItems: item.ArtistItems?.map(artist => ({
id: artist.Id,
name: artist.Name,
isFolder: artist.IsFolder,
metadata: artist,
})) || [],
})),
+1 -1
View File
@@ -31,7 +31,7 @@ export interface JellyfinBaseItem {
* Artist from Jellyfin API
*/
export interface JellyfinArtist extends JellyfinBaseItem {
IsFolder: boolean;
IsFolder?: boolean;
DateCreated?: string;
}
+4 -1
View File
@@ -41,6 +41,7 @@ class DriverRegistry {
* to calling `refresh()`.
*/
async initialise(): Promise<void> {
console.log('[DriverRegistry] Initialising driver registry...');
await this.build();
}
@@ -94,6 +95,8 @@ class DriverRegistry {
this.cache = new Map<string, SourceDriver>(
drivers.map(driver => [driver.getSourceId(), driver]),
);
console.log(`[DriverRegistry] Built driver cache with ${this.cache.size} entries:`);
}
}
@@ -106,4 +109,4 @@ class DriverRegistry {
*
* const driver = driverRegistry.getById(sourceId);
*/
export const driverRegistry = new DriverRegistry();
export const driverRegistry = new DriverRegistry();
+23 -4
View File
@@ -73,8 +73,8 @@ export class SourceSync {
this.queue = new PQueue({ concurrency });
this.pending = new Map();
this.queue.on('next', () => {
console.log('[SYNC] Starting next task. Queue size:', this.queue.size);
this.queue.on('active', () => {
console.log('[SYNC] Starting next task. Queue size:', this.queue.size, 'Pending promises:', this.pending.size);
})
}
@@ -100,7 +100,9 @@ export class SourceSync {
/** Sync all albums from one source, or all sources if omitted. */
async syncAlbums(sourceId?: string): Promise<void> {
console.log('[SYNC] Enqueueing albums sync for sourceId:', sourceId ?? 'ALL');
await this.registerCursor(sourceId, EntityType.ALBUMS);
console.log('[SYNC] Albums sync enqueued for sourceId:', sourceId ?? 'ALL');
}
/** Sync all tracks belonging to the given album. */
@@ -155,12 +157,28 @@ export class SourceSync {
? [sourceId]
: [...driverRegistry.getAll().keys()];
console.log(`[SYNC] Registering cursor for entityType: ${entityType}, parentEntityId: ${parentEntityId}, parentEntityType: ${parentEntityType}, sourceIds: ${sourceIds.join(', ')}`);
console.log(driverRegistry, driverRegistry.getAll());
// Persist a cursor row for each target source so the work survives a
// restart, then collect the completion promise for each one.
const promises = await Promise.all(
sourceIds.map(async (id) => {
await createCursorIfNotExists(id, entityType, parentEntityId, parentEntityType);
return this.getOrCreatePromise(id, entityType, parentEntityId).promise;
// Create the cursor first
const cursor = await createCursorIfNotExists(id, entityType, parentEntityId, parentEntityType);
console.log('[SYNC] Cursor registered:', cursor);
if (!cursor) {
throw new Error(`Failed to create cursor for sourceId: ${id}, entityType: ${entityType}, parentEntityId: ${parentEntityId}, parentEntityType: ${parentEntityType}`);
}
// Then, create a promise we can return to the caller
const promise = this.getOrCreatePromise(id, entityType, parentEntityId).promise;
// Finally, add the task to the queue.
this.queue.add(() => this.executeTask(cursor));
return promise;
})
);
@@ -253,6 +271,7 @@ export class SourceSync {
// -------------------------------------------------------------------------
private async executeTask(cursor: SyncCursor): Promise<void> {
console.log('[SYNC] Executing task for cursor', cursor);
// A cursor without a matching driver has nowhere to fetch from — skip it.
const driver = driverRegistry.getById(cursor.sourceId);
if (!driver) return;
+26 -14
View File
@@ -1,22 +1,22 @@
import { db, sqliteDb } from '@/store';
import { and, eq } from 'drizzle-orm';
import { and, eq, isNull } from 'drizzle-orm';
import syncCursors from './entity';
import { EntityType, type SyncCursor } from './types';
export async function getIncompleteCursors(): Promise<SyncCursor[]> {
return db.query.syncCursors.findMany({
where: { completed: false },
});
return db.select().from(syncCursors).where(eq(syncCursors.completed, false)).all();
}
export async function createCursorIfNotExists(
sourceId: string,
entityType: EntityType,
parentEntityId: string = '',
parentEntityType: EntityType | null = null,
parentEntityId?: string,
parentEntityType?: string,
pageSize: number = 500,
): Promise<void> {
) {
console.log('Creating cursor', { sourceId, entityType, parentEntityId, parentEntityType });
await db.insert(syncCursors).values({
sourceId,
entityType,
@@ -25,37 +25,49 @@ export async function createCursorIfNotExists(
startIndex: 0,
pageSize,
completed: false,
updatedAt: Date.now(),
}).onConflictDoNothing();
sqliteDb.flushPendingReactiveQueries();
await sqliteDb.flushPendingReactiveQueries();
console.log('Inserted cursor');
const cursor = await db.select().from(syncCursors).where(and(
eq(syncCursors.sourceId, sourceId),
eq(syncCursors.entityType, entityType),
parentEntityId !== undefined ? eq(syncCursors.parentEntityId, parentEntityId) : isNull(syncCursors.parentEntityId),
parentEntityType !== undefined ? eq(syncCursors.parentEntityType, parentEntityType) : isNull(syncCursors.parentEntityType),
)).get();
console.log('Retrieved cursor', cursor);
return cursor;
}
export async function updateCursorOffset(
sourceId: string,
entityType: EntityType,
newOffset: number,
parentEntityId: string = '',
parentEntityId?: string,
): Promise<void> {
await db.update(syncCursors)
.set({ startIndex: newOffset, updatedAt: Date.now() })
.where(and(
eq(syncCursors.sourceId, sourceId),
eq(syncCursors.entityType, entityType),
eq(syncCursors.parentEntityId, parentEntityId),
parentEntityId !== undefined ? eq(syncCursors.parentEntityId, parentEntityId) : isNull(syncCursors.parentEntityId),
));
}
export async function markCursorComplete(
sourceId: string,
entityType: EntityType,
parentEntityId: string = '',
parentEntityId?: string,
): Promise<void> {
await db.update(syncCursors)
.set({ completed: true, updatedAt: Date.now() })
.where(and(
eq(syncCursors.sourceId, sourceId),
eq(syncCursors.entityType, entityType),
eq(syncCursors.parentEntityId, parentEntityId),
parentEntityId !== undefined ? eq(syncCursors.parentEntityId, parentEntityId) : isNull(syncCursors.parentEntityId),
));
}
}
+2 -2
View File
@@ -10,14 +10,14 @@ const syncCursors = sqliteTable('sync_cursors', {
sourceId: text('source_id').notNull().references(() => sources.id, { onDelete: 'cascade' }),
entityType: text('entity_type').notNull(),
parentEntityId: text('parent_entity_id').notNull().default(''),
parentEntityType: text('parent_entity_type'),
parentEntityType: text('parent_entity_type').notNull().default(''),
startIndex: integer('start_index').notNull(),
pageSize: integer('page_size').notNull(),
completed: integer('completed', { mode: 'boolean' }).notNull().default(false),
attempts: integer('attempts').notNull().default(0),
failedAt: integer('failed_at'),
lastError: text('last_error'),
updatedAt: integer('updated_at').notNull(),
updatedAt: integer('updated_at').notNull().$default(() => Date.now()),
}, (table) => [
primaryKey({ columns: [table.sourceId, table.entityType, table.parentEntityId] }),
]);