cents/services/api/src/server.ts

708 lines
31 KiB
TypeScript

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<User | null> {
const token = request.cookies[sessionCookieName];
if (!token) return null;
const tokenHash = hashSecret(token);
const result = await pool.query<User>(
`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<MemberRole | null> {
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<LedgerEntry>;
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<SyncOperation>;
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<SyncCursor>;
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<User>(
"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<LedgerEntry & { serverUpdatedAt: Date }>(
`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 });