diff --git a/prisma/d1/00007_add_user_device_alias.sql b/prisma/d1/00007_add_user_device_alias.sql new file mode 100644 index 0000000..bd1f918 --- /dev/null +++ b/prisma/d1/00007_add_user_device_alias.sql @@ -0,0 +1,12 @@ +CREATE TABLE "user_device_alias" ( + "device_id" TEXT NOT NULL, + "user_id" INTEGER NOT NULL, + "bound_at" DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + "bind_source" TEXT NOT NULL CHECK ("bind_source" IN ('signup', 'login', 'backfill')) +); + +CREATE UNIQUE INDEX "user_device_alias_device_id_user_id_key" + ON "user_device_alias"("device_id", "user_id"); + +CREATE INDEX "user_device_alias_user_id_bound_at_idx" + ON "user_device_alias"("user_id", "bound_at"); diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 52d5e7c..4caeb62 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -33,6 +33,16 @@ model slax_user { @@index([invite_code]) } +model user_device_alias { + device_id String + user_id Int + bound_at DateTime @default(now()) + bind_source String + + @@unique([device_id, user_id]) + @@index([user_id, bound_at]) +} + model slax_bookmark { id Int @id @default(autoincrement()) title String @default("") diff --git a/src/domain/orchestrator/sync.ts b/src/domain/orchestrator/sync.ts index c784f46..a894340 100644 --- a/src/domain/orchestrator/sync.ts +++ b/src/domain/orchestrator/sync.ts @@ -4,6 +4,7 @@ import { ContextManager } from '../../utils/context' import { UserService } from '../user' import { SignJWT } from 'jose' import { DBSyncBatchOperation } from '../../infra/repository/dbSyncBatch' +import type { SyncCommitObserver } from '../../infra/repository/dbSyncBatch' import { QueueClient, queueRetryParseMessage, callbackType } from '../../infra/queue/queueClient' import { parserType, URLPolicie } from '../../utils/urlPolicie' import { markType } from '../../infra/repository/dbMark' @@ -177,7 +178,8 @@ export class SyncOrchestrator { } } - const result = await this.dbSyncBatch.executeOrderedOperations(orderedOperations) + const observer = this.syncCommitObserver(ctx) + const result = observer ? await this.dbSyncBatch.executeOrderedOperations(orderedOperations, observer) : await this.dbSyncBatch.executeOrderedOperations(orderedOperations) for (const newBookmark of result) { await this.sendRetryParseEvent(ctx, newBookmark) } @@ -188,6 +190,10 @@ export class SyncOrchestrator { } } + protected syncCommitObserver(_ctx: ContextManager): SyncCommitObserver | undefined { + return undefined + } + public processUserBookmarkCommentChange(change: SyncChangeItem, userId: number, operations: OrderedSyncOperation[]) { if (!change.data) return diff --git a/src/domain/user.ts b/src/domain/user.ts index 0c49c79..bcb72ff 100644 --- a/src/domain/user.ts +++ b/src/domain/user.ts @@ -1,4 +1,5 @@ import { ContextManager } from '../utils/context' +import { resolveDeviceId } from '../utils/eventContext' import { Auth } from '../utils/jwt' import { NeedCreateUsernameError, RegisterUserError, SaveReportError, UnauthorizedError, UserNotFoundError } from '../const/err' import { platformBindType, userInfoPO } from '../infra/repository/dbUser' @@ -210,6 +211,8 @@ export class UserService { // await new KVClient(env.KV).delete.USER_INFO(regInfo.id) + this.bindRequestDevice(ctx, request, regInfo.id, isFirstRegister ? 'signup' : 'login') + return { token, user_id: signUserId.toString(), @@ -220,6 +223,12 @@ export class UserService { } } + public bindRequestDevice(ctx: ContextManager, request: Request, userId: number, source: 'signup' | 'login') { + const deviceId = resolveDeviceId(request) + if (!deviceId || !Number.isInteger(userId) || userId < 1) return + ctx.execution.waitUntil(this.userRepo.bindUserDeviceAlias(deviceId, userId, source).catch(error => console.error('[events] failed to bind user device:', error))) + } + /** * 用户信息接口 * @param ctx diff --git a/src/infra/queue/queueClient.ts b/src/infra/queue/queueClient.ts index 89f510d..fa52573 100644 --- a/src/infra/queue/queueClient.ts +++ b/src/infra/queue/queueClient.ts @@ -1,6 +1,7 @@ import { container, injectable } from '../../decorators/di' import { ContextManager } from '../../utils/context' import { parserType } from '../../utils/urlPolicie' +import type { EventRequestContext } from '../../utils/eventContext' export enum callbackType { NOT_CALLBACK = 0, @@ -19,6 +20,7 @@ export interface parseMessage { } export interface importBookmarkMessage { + eventContext?: EventRequestContext type: string id: number data: any[] diff --git a/src/infra/repository/dbSyncBatch.ts b/src/infra/repository/dbSyncBatch.ts index ed8fda4..949b965 100644 --- a/src/infra/repository/dbSyncBatch.ts +++ b/src/infra/repository/dbSyncBatch.ts @@ -18,15 +18,36 @@ import { markType } from './dbMark' export type prismaTx = Omit, '$connect' | '$disconnect' | '$on' | '$transaction' | '$extends'> export type executeFunction = (tx: prismaTx, operation: OrderedSyncOperation) => Promise<{ bookmarkId: number; targetUrl: string; userId: number } | null | void> +export interface SyncEntitySnapshot { + uuid: string + user_id: number + type: number + user_bookmark_uuid?: string + archive_status?: number + is_starred?: boolean + deleted_at?: Date | null + is_deleted?: boolean + bookmark?: { target_url: string } +} + +export interface CommittedSyncMutation { + operation: OrderedSyncOperation + before: SyncEntitySnapshot | null + after: SyncEntitySnapshot | null +} + +export type SyncCommitObserver = (mutations: CommittedSyncMutation[], committedAt: string) => Promise | void + @injectable() export class DBSyncBatchOperation { constructor(@inject(PRISIMA_HYPERDRIVE_CLIENT) public prismaHyperdrive: LazyInstance) {} /** execute */ - public async executeOrderedOperations(operations: OrderedSyncOperation[]): Promise<{ bookmarkId: number; targetUrl: string; userId: number }[]> { + public async executeOrderedOperations(operations: OrderedSyncOperation[], observer?: SyncCommitObserver): Promise<{ bookmarkId: number; targetUrl: string; userId: number }[]> { if (operations.length === 0) return [] const newBookmarks: { bookmarkId: number; targetUrl: string; userId: number }[] = [] + const mutations: CommittedSyncMutation[] = [] const executeMap: Record = { create_tag: this.executeCreateTag, create_bookmark: this.executeCreateBookmark, @@ -41,14 +62,44 @@ export class DBSyncBatchOperation { await this.prismaHyperdrive().$transaction(async tx => { for (const operation of operations) { + const before = observer ? await this.eventSnapshot(tx, operation) : null const result = await executeMap[operation.type].bind(this)(tx, operation) if (result) newBookmarks.push(result) + if (observer) { + const after = await this.eventSnapshot(tx, operation) + if (before || after) mutations.push({ operation, before, after }) + } } }) + // A rolled-back batch never reaches this observer. Telemetry cannot fail an already-committed sync. + if (observer) { + try { + await observer(mutations, new Date().toISOString()) + } catch (error) { + console.error('[events] sync observer failed:', error) + } + } + return newBookmarks } + private async eventSnapshot(tx: prismaTx, operation: OrderedSyncOperation): Promise { + if (['create_bookmark', 'update_bookmark', 'delete_bookmark'].includes(operation.type) && 'bookmarkUuid' in operation) { + return tx.sr_user_bookmark.findFirst({ + where: { uuid: operation.bookmarkUuid, user_id: operation.userId }, + select: { uuid: true, user_id: true, type: true, archive_status: true, is_starred: true, deleted_at: true, bookmark: { select: { target_url: true } } } + }) + } + if (operation.type === 'create_comment' || operation.type === 'delete_comment') { + return tx.sr_bookmark_comment.findFirst({ + where: { uuid: operation.commentUuid, user_id: operation.userId }, + select: { uuid: true, user_id: true, user_bookmark_uuid: true, type: true, is_deleted: true } + }) + } + return null + } + /** create tag */ public async executeCreateTag(tx: prismaTx, operation: OrderedSyncOperation): Promise { if (operation.type !== 'create_tag') return diff --git a/src/infra/repository/dbUser.ts b/src/infra/repository/dbUser.ts index a7f8463..dab4a3d 100644 --- a/src/infra/repository/dbUser.ts +++ b/src/infra/repository/dbUser.ts @@ -76,6 +76,20 @@ export class UserRepo { @inject(PRISIMA_HYPERDRIVE_CLIENT) private prismaPg: LazyInstance ) {} + public async bindUserDeviceAlias(deviceId: string, userId: number, bindSource: 'signup' | 'login') { + if (!deviceId || userId < 1) return + await this.prisma().user_device_alias.upsert({ + where: { device_id_user_id: { device_id: deviceId, user_id: userId } }, + create: { + device_id: deviceId, + user_id: userId, + bind_source: bindSource, + bound_at: new Date() + }, + update: {} + }) + } + public async getInfoByEmail(email: string): Promise { if (!email) return null let res = await this.prismaPg().sr_user.findFirst({ where: { email: email } }) diff --git a/src/utils/eventContext.ts b/src/utils/eventContext.ts new file mode 100644 index 0000000..6a0d374 --- /dev/null +++ b/src/utils/eventContext.ts @@ -0,0 +1,28 @@ +export type BookmarkEventSource = 'extension_save' | 'share_sheet' | 'web_button' | 'import' | 'api' | 'cli' | 'onboarding_gs' + +/** Serializable request context, carried unchanged into asynchronous work. */ +export interface EventRequestContext { + device_id: string + platform: string + locale: string + client_version: string + source: BookmarkEventSource + ua?: string +} + +export const resolveDeviceId = (request: Request): string => { + const explicit = request.headers.get('x-device-id')?.trim() + if (explicit) return explicit + for (const item of (request.headers.get('cookie') || '').split(';')) { + const separator = item.indexOf('=') + if (separator < 0 || item.slice(0, separator).trim() !== '_su') continue + try { + return decodeURIComponent(item.slice(separator + 1).trim()).trim() + } catch { + return item.slice(separator + 1).trim() + } + } + return '' +} + +export const normalizeEventLocale = (locale: string): string => locale.trim().toLowerCase().replaceAll('_', '-').split('-')[0] || 'en'