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
12 changes: 12 additions & 0 deletions prisma/d1/00007_add_user_device_alias.sql
Original file line number Diff line number Diff line change
@@ -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");
10 changes: 10 additions & 0 deletions prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -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("")
Expand Down
8 changes: 7 additions & 1 deletion src/domain/orchestrator/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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)
}
Expand All @@ -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

Expand Down
9 changes: 9 additions & 0 deletions src/domain/user.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down Expand Up @@ -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(),
Expand All @@ -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
Expand Down
2 changes: 2 additions & 0 deletions src/infra/queue/queueClient.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -19,6 +20,7 @@ export interface parseMessage {
}

export interface importBookmarkMessage {
eventContext?: EventRequestContext
type: string
id: number
data: any[]
Expand Down
53 changes: 52 additions & 1 deletion src/infra/repository/dbSyncBatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,36 @@ import { markType } from './dbMark'
export type prismaTx = Omit<HyperdrivePrismaClient<Prisma.PrismaClientOptions>, '$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> | void

@injectable()
export class DBSyncBatchOperation {
constructor(@inject(PRISIMA_HYPERDRIVE_CLIENT) public prismaHyperdrive: LazyInstance<HyperdrivePrismaClient>) {}

/** 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<string, executeFunction> = {
create_tag: this.executeCreateTag,
create_bookmark: this.executeCreateBookmark,
Expand All @@ -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<SyncEntitySnapshot | null> {
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<void> {
if (operation.type !== 'create_tag') return
Expand Down
14 changes: 14 additions & 0 deletions src/infra/repository/dbUser.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,20 @@ export class UserRepo {
@inject(PRISIMA_HYPERDRIVE_CLIENT) private prismaPg: LazyInstance<HyperdrivePrismaClient>
) {}

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<userInfoPO | null> {
if (!email) return null
let res = await this.prismaPg().sr_user.findFirst({ where: { email: email } })
Expand Down
28 changes: 28 additions & 0 deletions src/utils/eventContext.ts
Original file line number Diff line number Diff line change
@@ -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'
Loading