import cookie from "@fastify/cookie"; import helmet from "@fastify/helmet"; import rateLimit from "@fastify/rate-limit"; import { ledgerThemeAccents, type LedgerEntry, type LedgerThemeId, type SyncOperation } from "@cents/domain"; import Fastify, { type FastifyReply, type FastifyRequest } from "fastify"; import { createUserWithPersonalLedger, PERSONAL_LEDGER_NAME } from "./accounts.js"; import { initializeDatabase, pool } from "./db.js"; import { clearSessionCookie, createId, createSecret, hashPassword, hashSecret, INVITATION_TTL_MS, normalizeName, SESSION_TTL_MS, sessionCookieName, setSessionCookie, validPassword, verifyPassword, } from "./security.js"; type User = { id: string; name: string }; type MemberRole = "owner" | "member"; const server = Fastify({ logger: true, trustProxy: true }); await server.register(cookie); await server.register(helmet, { contentSecurityPolicy: false }); await server.register(rateLimit, { max: 120, timeWindow: "1 minute" }); await initializeDatabase(); server.addHook("onRequest", async (request, reply) => { if (!["POST", "PATCH", "PUT", "DELETE"].includes(request.method)) return; const origin = request.headers.origin; if (!origin) return; try { const configuredOrigin = process.env.PUBLIC_ORIGIN ? new URL(process.env.PUBLIC_ORIGIN).origin : null; if (new URL(origin).host !== request.headers.host && new URL(origin).origin !== configuredOrigin) { return reply.code(403).send({ error: "不允许跨站提交" }); } } catch { return reply.code(403).send({ error: "请求来源无效" }); } }); server.addHook("onSend", async (request, reply, payload) => { if ( request.url.startsWith("/api/auth/") || request.url.startsWith("/api/invitations/") || request.url.startsWith("/api/ledger-invitations/") || request.url.startsWith("/api/sync/") || request.url === "/api/me" ) { reply.header("Cache-Control", "no-store"); } return payload; }); function uniqueViolation(error: unknown) { return typeof error === "object" && error !== null && "code" in error && error.code === "23505"; } async function currentUser(request: FastifyRequest, reply?: FastifyReply): Promise { const token = request.cookies[sessionCookieName]; if (!token) return null; const tokenHash = hashSecret(token); const result = await pool.query( `SELECT u.id, u.name FROM sessions s JOIN users u ON u.id = s.user_id WHERE s.token_hash = $1 AND s.revoked_at IS NULL AND s.expires_at > now()`, [tokenHash], ); const user = result.rows[0]; if (!user) return null; await pool.query( "UPDATE sessions SET last_seen_at = now(), expires_at = $1 WHERE token_hash = $2", [new Date(Date.now() + SESSION_TTL_MS), tokenHash], ); if (reply) setSessionCookie(reply, token); return user; } async function requireUser(request: FastifyRequest, reply: FastifyReply) { const user = await currentUser(request, reply); if (!user) { clearSessionCookie(reply); await reply.code(401).send({ error: "请先登录" }); return null; } return user; } async function ledgerRole(userId: string, ledgerId: string): Promise { const result = await pool.query<{ role: MemberRole }>( `SELECT role FROM ledger_members WHERE ledger_id = $1 AND user_id = $2 AND removed_at IS NULL`, [ledgerId, userId], ); return result.rows[0]?.role ?? null; } async function createSession(userId: string) { const token = createSecret(); await pool.query( "INSERT INTO sessions (id, user_id, token_hash, expires_at) VALUES ($1, $2, $3, $4)", [createId(), userId, hashSecret(token), new Date(Date.now() + SESSION_TTL_MS)], ); return token; } function invitationStatus(invitation: { expiresAt: Date | string; acceptedAt: Date | string | null; revokedAt: Date | string | null }) { if (invitation.revokedAt) return "revoked"; if (invitation.acceptedAt) return "accepted"; if (new Date(invitation.expiresAt).getTime() <= Date.now()) return "expired"; return "valid"; } const currencies = new Set(["CNY", "USD", "EUR", "JPY", "HKD"]); function validLedgerTheme(value: unknown): value is LedgerThemeId { return typeof value === "string" && Object.hasOwn(ledgerThemeAccents, value); } function ledgerThemeFromColor(color: unknown): LedgerThemeId { if (typeof color !== "string") return "jade"; return (Object.entries(ledgerThemeAccents) as Array<[LedgerThemeId, string]>) .find(([, accent]) => accent.toLowerCase() === color.toLowerCase())?.[0] ?? "jade"; } function validDate(value: unknown) { return typeof value === "string" && !Number.isNaN(Date.parse(value)); } function validLedgerIds(value: unknown): value is string[] { return Array.isArray(value) && value.length > 0 && value.length <= 20 && value.every((ledgerId) => typeof ledgerId === "string" && ledgerId.length > 0 && ledgerId.length <= 100) && new Set(value).size === value.length; } function validEntry(value: unknown): value is LedgerEntry { if (!value || typeof value !== "object") return false; const entry = value as Partial; return typeof entry.id === "string" && entry.id.length > 0 && entry.id.length <= 100 && typeof entry.ownerId === "string" && entry.ownerId.length > 0 && entry.ownerId.length <= 100 && validLedgerIds(entry.ledgerIds) && (entry.type === "expense" || entry.type === "income") && Number.isSafeInteger(entry.amount) && entry.amount! > 0 && entry.amount! <= 2_147_483_647 && typeof entry.currency === "string" && currencies.has(entry.currency) && typeof entry.baseCurrency === "string" && currencies.has(entry.baseCurrency) && Number.isSafeInteger(entry.baseAmount) && entry.baseAmount! > 0 && entry.baseAmount! <= 2_147_483_647 && typeof entry.exchangeRate === "string" && entry.exchangeRate.length > 0 && entry.exchangeRate.length <= 32 && (entry.exchangeRateSource === "manual" || entry.exchangeRateSource === "system") && typeof entry.categoryId === "string" && entry.categoryId.length > 0 && entry.categoryId.length <= 100 && typeof entry.note === "string" && entry.note.length <= 500 && validDate(entry.occurredAt) && validDate(entry.createdAt) && validDate(entry.updatedAt) && (entry.deletedAt === null || validDate(entry.deletedAt)) && Number.isSafeInteger(entry.version) && entry.version! > 0; } function validSyncOperation(value: unknown): value is SyncOperation { if (!value || typeof value !== "object") return false; const operation = value as Partial; const payload = operation.payload; return typeof operation.id === "string" && operation.id.length > 0 && operation.id.length <= 100 && operation.entity === "entry" && (operation.action === "create" || operation.action === "update" || operation.action === "delete") && validEntry(payload) && operation.entityId === payload.id && validLedgerIds(operation.ledgerIds) && operation.ledgerIds.length === payload.ledgerIds.length && operation.ledgerIds.every((ledgerId) => payload.ledgerIds.includes(ledgerId)) && (operation.action !== "delete" || payload.deletedAt !== null); } type SyncCursor = { serverUpdatedAt: string; id: string }; function encodeSyncCursor(cursor: SyncCursor) { return Buffer.from(JSON.stringify(cursor)).toString("base64url"); } function decodeSyncCursor(value: string): SyncCursor | null { try { const cursor = JSON.parse(Buffer.from(value, "base64url").toString("utf8")) as Partial; if (!validDate(cursor.serverUpdatedAt) || typeof cursor.id !== "string" || cursor.id.length > 100) return null; return { serverUpdatedAt: cursor.serverUpdatedAt!, id: cursor.id }; } catch { return null; } } server.get("/health", async () => ({ ok: true, service: "cents-api" })); server.get("/api/auth/session", async (request, reply) => ({ user: await currentUser(request, reply) })); server.post<{ Body: { name?: string; password?: string } }>( "/api/auth/login", { config: { rateLimit: { max: 10, timeWindow: "1 minute" } } }, async (request, reply) => { const nameKey = normalizeName(request.body.name ?? ""); const password = request.body.password ?? ""; const result = await pool.query<{ id: string; name: string; passwordHash: string | null }>( `SELECT id, name, password_hash AS "passwordHash" FROM users WHERE name_key = $1`, [nameKey], ); const account = result.rows[0]; const valid = account?.passwordHash ? await verifyPassword(account.passwordHash, password).catch(() => false) : (await hashPassword(password || "invalid-password"), false); if (!account || !valid) return reply.code(401).send({ error: "姓名或密码错误" }); const token = await createSession(account.id); setSessionCookie(reply, token); return { user: { id: account.id, name: account.name } }; }, ); server.post("/api/auth/logout", async (request, reply) => { const token = request.cookies[sessionCookieName]; if (token) await pool.query("UPDATE sessions SET revoked_at = now() WHERE token_hash = $1", [hashSecret(token)]); clearSessionCookie(reply); return { ok: true }; }); server.patch<{ Body: { name?: string } }>("/api/me", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const name = request.body.name?.trim() ?? ""; if (!name || name.length > 40) return reply.code(400).send({ error: "姓名应为 1 至 40 个字符" }); try { const result = await pool.query( "UPDATE users SET name = $1, name_key = $2, updated_at = now() WHERE id = $3 RETURNING id, name", [name, normalizeName(name), user.id], ); return { user: result.rows[0] }; } catch (error) { if (uniqueViolation(error)) return reply.code(409).send({ error: "该姓名已被使用" }); throw error; } }); server.get("/api/ledgers", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const result = await pool.query( `SELECT l.id, l.name, l.color, l.theme, l.default_currency AS "defaultCurrency", l.created_at AS "createdAt", l.updated_at AS "updatedAt", l.archived_at AS "archivedAt", m.role, (l.personal_owner_id IS NOT NULL) AS "isPersonal" FROM ledgers l JOIN ledger_members m ON m.ledger_id = l.id WHERE m.user_id = $1 AND m.removed_at IS NULL AND l.archived_at IS NULL ORDER BY (l.personal_owner_id = $1) DESC, l.updated_at DESC`, [user.id], ); return { ledgers: result.rows }; }); server.post<{ Body: { name?: string; color?: string; theme?: string; defaultCurrency?: string } }>( "/api/ledgers", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const name = request.body.name?.trim() ?? ""; if (request.body.theme !== undefined && !validLedgerTheme(request.body.theme)) { return reply.code(400).send({ error: "账本主题无效" }); } const theme = request.body.theme ?? ledgerThemeFromColor(request.body.color); const color = ledgerThemeAccents[theme]; const defaultCurrency = request.body.defaultCurrency ?? "CNY"; if (!name || name.length > 40) return reply.code(400).send({ error: "账本名称应为 1 至 40 个字符" }); if (!/^#[0-9a-f]{6}$/i.test(color)) return reply.code(400).send({ error: "账本颜色无效" }); if (!["CNY", "USD", "EUR", "JPY", "HKD"].includes(defaultCurrency)) { return reply.code(400).send({ error: "默认币种无效" }); } const client = await pool.connect(); try { await client.query("BEGIN"); const ledgerId = createId(); const result = await client.query( `INSERT INTO ledgers (id, name, color, theme, default_currency) VALUES ($1, $2, $3, $4, $5) RETURNING id, name, color, theme, default_currency AS "defaultCurrency", created_at AS "createdAt", updated_at AS "updatedAt", archived_at AS "archivedAt"`, [ledgerId, name, color, theme, defaultCurrency], ); await client.query( "INSERT INTO ledger_members (ledger_id, user_id, role) VALUES ($1, $2, 'owner')", [ledgerId, user.id], ); await client.query("COMMIT"); return reply.code(201).send({ ledger: { ...result.rows[0], role: "owner" } }); } catch (error) { await client.query("ROLLBACK"); throw error; } finally { client.release(); } }, ); server.get<{ Params: { ledgerId: string } }>("/api/ledgers/:ledgerId/members", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; if (!await ledgerRole(user.id, request.params.ledgerId)) return reply.code(403).send({ error: "无权访问该账本" }); const result = await pool.query( `SELECT u.id, u.name, m.role, m.joined_at AS "joinedAt" FROM ledger_members m JOIN users u ON u.id = m.user_id WHERE m.ledger_id = $1 AND m.removed_at IS NULL ORDER BY CASE m.role WHEN 'owner' THEN 0 ELSE 1 END, m.joined_at`, [request.params.ledgerId], ); return { members: result.rows }; }); server.patch<{ Params: { ledgerId: string }; Body: { name?: string; color?: string; theme?: string; defaultCurrency?: string }; }>("/api/ledgers/:ledgerId", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; if (await ledgerRole(user.id, request.params.ledgerId) !== "owner") { return reply.code(403).send({ error: "只有账本拥有者可以修改设置" }); } const personalLedger = await pool.query<{ isPersonal: boolean }>( `SELECT (personal_owner_id IS NOT NULL) AS "isPersonal" FROM ledgers WHERE id = $1`, [request.params.ledgerId], ); const name = request.body.name?.trim() ?? ""; if (!name || name.length > 40) return reply.code(400).send({ error: "账本名称无效" }); if (request.body.theme !== undefined && !validLedgerTheme(request.body.theme)) { return reply.code(400).send({ error: "账本主题无效" }); } const theme = request.body.theme ?? ledgerThemeFromColor(request.body.color); const color = ledgerThemeAccents[theme]; const result = await pool.query( `UPDATE ledgers SET name = $1, color = $2, theme = $3, default_currency = $4, updated_at = now() WHERE id = $5 RETURNING id, name, color, theme, default_currency AS "defaultCurrency", created_at AS "createdAt", updated_at AS "updatedAt", archived_at AS "archivedAt"`, [personalLedger.rows[0]?.isPersonal ? PERSONAL_LEDGER_NAME : name, color, theme, request.body.defaultCurrency ?? "CNY", request.params.ledgerId], ); return { ledger: result.rows[0] }; }); // Application invitations create accounts but never grant access to an existing ledger. server.get<{ Params: { key: string } }>("/api/invitations/:key", async (request, reply) => { const user = await currentUser(request, reply); if (user) return { user, alreadyLoggedIn: true }; if (request.params.key.length < 32) return reply.code(404).send({ error: "邀请链接无效" }); const result = await pool.query( `SELECT expires_at AS "expiresAt", accepted_at AS "acceptedAt", revoked_at AS "revokedAt" FROM app_invitations WHERE key_hash = $1`, [hashSecret(request.params.key)], ); const invitation = result.rows[0]; if (!invitation) return reply.code(404).send({ error: "邀请链接无效" }); return { invitation: { ...invitation, status: invitationStatus(invitation) }, user: null }; }); server.post<{ Params: { key: string }; Body: { name?: string; password?: string } }>( "/api/invitations/:key/accept", { config: { rateLimit: { max: 10, timeWindow: "1 minute" } } }, async (request, reply) => { const existingUser = await currentUser(request, reply); if (existingUser) return { user: existingUser, created: false }; const name = request.body.name?.trim() ?? ""; const password = request.body.password ?? ""; if (!name || name.length > 40) return reply.code(400).send({ error: "请填写 1 至 40 个字符的姓名" }); if (!validPassword(password)) return reply.code(400).send({ error: "密码应为 8 至 128 个字符" }); const client = await pool.connect(); try { await client.query("BEGIN"); const result = await client.query<{ id: string; expiresAt: Date; acceptedAt: Date | null; revokedAt: Date | null; }>( `SELECT id, expires_at AS "expiresAt", accepted_at AS "acceptedAt", revoked_at AS "revokedAt" FROM app_invitations WHERE key_hash = $1 FOR UPDATE`, [hashSecret(request.params.key)], ); const invitation = result.rows[0]; const status = invitation ? invitationStatus(invitation) : "missing"; if (status !== "valid") { await client.query("ROLLBACK"); const code = status === "accepted" ? 409 : status === "missing" ? 404 : 410; return reply.code(code).send({ error: status === "accepted" ? "邀请已被使用" : status === "expired" ? "邀请已过期" : status === "revoked" ? "邀请已撤销" : "邀请链接无效" }); } const user = await createUserWithPersonalLedger(client, name, password); await client.query( "UPDATE app_invitations SET accepted_by = $1, accepted_at = now() WHERE id = $2", [user.id, invitation!.id], ); const sessionToken = createSecret(); await client.query( "INSERT INTO sessions (id, user_id, token_hash, expires_at) VALUES ($1, $2, $3, $4)", [createId(), user.id, hashSecret(sessionToken), new Date(Date.now() + SESSION_TTL_MS)], ); await client.query("COMMIT"); setSessionCookie(reply, sessionToken); return { user, created: true }; } catch (error) { await client.query("ROLLBACK"); if (uniqueViolation(error)) return reply.code(409).send({ error: "该姓名已被使用" }); throw error; } finally { client.release(); } }, ); server.post<{ Params: { ledgerId: string } }>("/api/ledgers/:ledgerId/invitations", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const personalLedger = await pool.query("SELECT 1 FROM ledgers WHERE id = $1 AND personal_owner_id IS NOT NULL", [request.params.ledgerId]); if (personalLedger.rowCount) return reply.code(403).send({ error: "个人账本不能共享" }); if (await ledgerRole(user.id, request.params.ledgerId) !== "owner") { return reply.code(403).send({ error: "只有账本拥有者可以邀请成员" }); } const key = createSecret(); const expiresAt = new Date(Date.now() + INVITATION_TTL_MS); await pool.query( `INSERT INTO ledger_invitations (id, ledger_id, key_hash, role, created_by, expires_at) VALUES ($1, $2, $3, 'member', $4, $5)`, [createId(), request.params.ledgerId, hashSecret(key), user.id, expiresAt], ); const origin = process.env.PUBLIC_ORIGIN ?? `${request.protocol}://${request.headers.host}`; return { invitation: { key, url: `${origin}/join-ledger?key=${encodeURIComponent(key)}`, expiresAt } }; }); server.get<{ Params: { key: string } }>("/api/ledger-invitations/:key", async (request, reply) => { const result = await pool.query( `SELECT i.expires_at AS "expiresAt", i.accepted_at AS "acceptedAt", i.revoked_at AS "revokedAt", l.id AS "ledgerId", l.name AS "ledgerName", u.name AS "inviterName" FROM ledger_invitations i JOIN ledgers l ON l.id = i.ledger_id LEFT JOIN users u ON u.id = i.created_by WHERE i.key_hash = $1`, [hashSecret(request.params.key)], ); const invitation = result.rows[0]; if (!invitation) return reply.code(404).send({ error: "账本邀请无效" }); return { invitation: { ...invitation, status: invitationStatus(invitation) } }; }); server.post<{ Params: { key: string } }>("/api/ledger-invitations/:key/accept", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const client = await pool.connect(); try { await client.query("BEGIN"); const result = await client.query<{ id: string; ledgerId: string; role: MemberRole; expiresAt: Date; acceptedAt: Date | null; revokedAt: Date | null; personalOwnerId: string | null; }>( `SELECT i.id, i.ledger_id AS "ledgerId", i.role, i.expires_at AS "expiresAt", i.accepted_at AS "acceptedAt", i.revoked_at AS "revokedAt", l.personal_owner_id AS "personalOwnerId" FROM ledger_invitations i JOIN ledgers l ON l.id = i.ledger_id WHERE i.key_hash = $1 FOR UPDATE OF i`, [hashSecret(request.params.key)], ); const invitation = result.rows[0]; if (invitation?.personalOwnerId) { await client.query("ROLLBACK"); return reply.code(403).send({ error: "个人账本不能共享" }); } const status = invitation ? invitationStatus(invitation) : "missing"; if (status !== "valid") { await client.query("ROLLBACK"); return reply.code(status === "accepted" ? 409 : status === "missing" ? 404 : 410).send({ error: "账本邀请不可用" }); } await client.query( `INSERT INTO ledger_members (ledger_id, user_id, role) VALUES ($1, $2, $3) ON CONFLICT (ledger_id, user_id) DO UPDATE SET removed_at = NULL, joined_at = now()`, [invitation!.ledgerId, user.id, invitation!.role], ); await client.query( `UPDATE entries e SET server_updated_at = now() WHERE EXISTS ( SELECT 1 FROM entry_ledgers el WHERE el.entry_id = e.id AND el.ledger_id = $1 AND el.unlinked_at IS NULL )`, [invitation!.ledgerId], ); await client.query( "UPDATE ledger_invitations SET accepted_by = $1, accepted_at = now() WHERE id = $2", [user.id, invitation!.id], ); await client.query("COMMIT"); return { ledgerId: invitation!.ledgerId }; } catch (error) { await client.query("ROLLBACK"); throw error; } finally { client.release(); } }); server.post<{ Body: { operations?: unknown[] } }>("/api/sync/push", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const operations = request.body.operations ?? []; if (!Array.isArray(operations) || operations.length > 200 || !operations.every(validSyncOperation)) { return reply.code(400).send({ error: "同步数据无效" }); } const client = await pool.connect(); const acceptedOperationIds: string[] = []; try { await client.query("BEGIN"); for (const operation of operations) { const alreadyProcessed = await client.query( "SELECT 1 FROM entry_sync_operations WHERE id = $1", [operation.id], ); if (alreadyProcessed.rowCount) { acceptedOperationIds.push(operation.id); continue; } const existing = await client.query<{ ownerId: string; canAccess: boolean }>( `SELECT e.owner_id AS "ownerId", (e.owner_id = $2 OR EXISTS ( SELECT 1 FROM entry_ledgers el JOIN ledger_members m ON m.ledger_id = el.ledger_id WHERE el.entry_id = e.id AND el.unlinked_at IS NULL AND m.user_id = $2 AND m.removed_at IS NULL )) AS "canAccess" FROM entries e WHERE e.id = $1`, [operation.entityId, user.id], ); if (existing.rows[0] && !existing.rows[0].canAccess) { await client.query("ROLLBACK"); return reply.code(403).send({ error: "无权修改该流水" }); } const inaccessibleTargets = await client.query<{ count: number }>( `SELECT count(*)::int AS count FROM unnest($1::text[]) AS target(ledger_id) WHERE NOT EXISTS ( SELECT 1 FROM ledger_members m WHERE m.ledger_id = target.ledger_id AND m.user_id = $2 AND m.removed_at IS NULL ) AND NOT ($3::boolean AND EXISTS ( SELECT 1 FROM entry_ledgers el WHERE el.entry_id = $4 AND el.ledger_id = target.ledger_id AND el.unlinked_at IS NULL ))`, [operation.ledgerIds, user.id, existing.rows[0]?.ownerId === user.id, operation.entityId], ); if (inaccessibleTargets.rows[0]?.count) { await client.query("ROLLBACK"); return reply.code(403).send({ error: "无权关联其中一个账本" }); } const entryOwnerId = existing.rows[0]?.ownerId ?? user.id; const personalLedger = await client.query<{ id: string }>( "SELECT id FROM ledgers WHERE personal_owner_id = $1", [entryOwnerId], ); if (!personalLedger.rows[0]) throw new Error("用户缺少个人账本"); const targetLedgerIds = [...new Set([...operation.ledgerIds, personalLedger.rows[0].id])]; const entry = operation.payload; const writeResult = await client.query( `INSERT INTO entries ( id, owner_id, type, amount, currency, base_currency, base_amount, exchange_rate, exchange_rate_source, category_id, note, occurred_at, created_by, updated_by, created_at, updated_at, deleted_at, version ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18 ) ON CONFLICT (id) DO UPDATE SET type = EXCLUDED.type, amount = EXCLUDED.amount, currency = EXCLUDED.currency, base_currency = EXCLUDED.base_currency, base_amount = EXCLUDED.base_amount, exchange_rate = EXCLUDED.exchange_rate, exchange_rate_source = EXCLUDED.exchange_rate_source, category_id = EXCLUDED.category_id, note = EXCLUDED.note, occurred_at = EXCLUDED.occurred_at, updated_by = EXCLUDED.updated_by, updated_at = EXCLUDED.updated_at, deleted_at = EXCLUDED.deleted_at, version = EXCLUDED.version, server_updated_at = now() WHERE entries.version < EXCLUDED.version OR (entries.version = EXCLUDED.version AND entries.updated_at <= EXCLUDED.updated_at)`, [ entry.id, user.id, entry.type, entry.amount, entry.currency, entry.baseCurrency, entry.baseAmount, entry.exchangeRate, entry.exchangeRateSource, entry.categoryId, entry.note, entry.occurredAt, user.id, user.id, entry.createdAt, entry.updatedAt, entry.deletedAt, entry.version, ], ); if (operation.action !== "delete" && writeResult.rowCount) { if (existing.rows[0]?.ownerId === user.id) { await client.query( `UPDATE entry_ledgers SET unlinked_at = now(), updated_at = now() WHERE entry_id = $1 AND unlinked_at IS NULL AND NOT (ledger_id = ANY($2::text[]))`, [entry.id, targetLedgerIds], ); } else if (existing.rows[0]) { await client.query( `UPDATE entry_ledgers el SET unlinked_at = now(), updated_at = now() FROM ledger_members m WHERE el.entry_id = $1 AND el.unlinked_at IS NULL AND el.ledger_id = m.ledger_id AND m.user_id = $2 AND m.removed_at IS NULL AND NOT (el.ledger_id = ANY($3::text[]))`, [entry.id, user.id, targetLedgerIds], ); } await client.query( `INSERT INTO entry_ledgers (entry_id, ledger_id) SELECT $1, unnest($2::text[]) ON CONFLICT (entry_id, ledger_id) DO UPDATE SET unlinked_at = NULL, updated_at = now()`, [entry.id, targetLedgerIds], ); } await client.query("UPDATE entries SET server_updated_at = now() WHERE id = $1", [entry.id]); await client.query( "INSERT INTO entry_sync_operations (id, user_id) VALUES ($1, $2)", [operation.id, user.id], ); acceptedOperationIds.push(operation.id); } await client.query("COMMIT"); return { acceptedOperationIds }; } catch (error) { await client.query("ROLLBACK"); throw error; } finally { client.release(); } }); server.get<{ Querystring: { cursor?: string } }>("/api/sync/pull", async (request, reply) => { const user = await requireUser(request, reply); if (!user) return; const cursor = request.query.cursor ? decodeSyncCursor(request.query.cursor) : null; if (request.query.cursor && !cursor) return reply.code(400).send({ error: "同步游标无效" }); const result = await pool.query( `SELECT e.id, e.owner_id AS "ownerId", ARRAY( SELECT el.ledger_id FROM entry_ledgers el LEFT JOIN ledger_members visible_member ON visible_member.ledger_id = el.ledger_id AND visible_member.user_id = $1 AND visible_member.removed_at IS NULL WHERE el.entry_id = e.id AND el.unlinked_at IS NULL AND (e.owner_id = $1 OR visible_member.user_id IS NOT NULL) ORDER BY el.ledger_id ) AS "ledgerIds", e.type, e.amount, e.currency, e.base_currency AS "baseCurrency", e.base_amount AS "baseAmount", e.exchange_rate AS "exchangeRate", e.exchange_rate_source AS "exchangeRateSource", e.category_id AS "categoryId", e.note, e.occurred_at AS "occurredAt", e.created_by AS "createdBy", e.updated_by AS "updatedBy", e.created_at AS "createdAt", e.updated_at AS "updatedAt", e.deleted_at AS "deletedAt", e.version, entry_change.changed_at AS "serverUpdatedAt" FROM entries e CROSS JOIN LATERAL ( SELECT GREATEST( e.server_updated_at, COALESCE(MAX(el.updated_at) FILTER ( WHERE e.owner_id = $1 OR change_member.user_id IS NOT NULL ), e.server_updated_at) ) AS changed_at FROM entry_ledgers el LEFT JOIN ledger_members change_member ON change_member.ledger_id = el.ledger_id AND change_member.user_id = $1 AND change_member.removed_at IS NULL WHERE el.entry_id = e.id ) entry_change WHERE (e.owner_id = $1 OR EXISTS ( SELECT 1 FROM entry_ledgers accessible_link JOIN ledger_members accessible_member ON accessible_member.ledger_id = accessible_link.ledger_id WHERE accessible_link.entry_id = e.id AND accessible_member.user_id = $1 AND accessible_member.removed_at IS NULL )) AND ($2::timestamptz IS NULL OR entry_change.changed_at > $2::timestamptz OR (entry_change.changed_at = $2::timestamptz AND e.id > $3)) ORDER BY entry_change.changed_at, e.id LIMIT 501`, [user.id, cursor?.serverUpdatedAt ?? null, cursor?.id ?? ""], ); const pageRows = result.rows.slice(0, 500); const lastRow = pageRows.at(-1); const entries = pageRows.map(({ serverUpdatedAt: _serverUpdatedAt, ...entry }) => entry); return { entries, cursor: lastRow ? encodeSyncCursor({ serverUpdatedAt: lastRow.serverUpdatedAt.toISOString(), id: lastRow.id }) : request.query.cursor ?? "", hasMore: result.rows.length > 500, }; }); const port = Number(process.env.PORT ?? 3000); const host = process.env.HOST ?? "0.0.0.0"; await server.listen({ port, host });