import { useEffect, useState } from "react"; import { Alert, App as AntApp, Button, Checkbox, Form, Input, InputNumber, Modal, Select, Space, Table, Tag, Typography, } from "antd"; import { PlusOutlined, ReloadOutlined, PlayCircleOutlined, PauseCircleOutlined, ApiOutlined, TableOutlined, EyeOutlined, DeleteOutlined, } from "@ant-design/icons"; import { AgentAccount, BindCode, Session, SyncBinding, SyncChannel, SyncInspectResult, SyncPreviewResult, createSyncChannel, deleteSyncChannel, dropSyncTable, getSyncCheckpointMeta, inspectSyncChannel, listAgents, listApps, listSyncBindings, listSyncChannels, listBindCodes, createBindCode, revokeBindCode, previewSyncTable, reconcileSyncChannel, restoreSyncChannel, startSyncChannel, stopSyncChannel, testSyncEndpoints, updateSyncChannel, } from "./api"; type EndpointForm = { driver: string; dsn: string; tables: string; }; type FormValues = { name: string; direction: string; conflict_policy: string; poll_interval_ms: number; agent_id?: number | null; app_slug?: string; local: EndpointForm; remote: EndpointForm; pk_columns: string; // table:pk,table:pk }; function parsePKs(s: string): Record { const out: Record = {}; for (const part of (s || "").split(",")) { const [t, p] = part.split(":").map((x) => x.trim()); if (t && p) out[t] = p; } return out; } /** DSN 脱敏;file 库尽量只露文件名 */ function formatDsnHint(dsn: string): string { const masked = (dsn || "").replace(/:[^:@/]+@/, ":***@"); const fileMatch = masked.match(/(?:^|[\\/])([^\\/?#]+\.db)(?:\?|$)/i); if (fileMatch) return fileMatch[1]; if (masked.length > 56) return `${masked.slice(0, 56)}…`; return masked; } function bindingLocalLabel(b: SyncBinding): string { return (b.database_name || "").trim() || b.local_database_id; } function bindingOnlineLabel(b: SyncBinding): string { return (b.display_name || "").trim() || b.online_db_id; } /** 联调烟雾 Binding:占位 id,易被误认为「库名对不上」 */ function isSmokeBinding(b: SyncBinding): boolean { const note = (b.note || "").toLowerCase(); if (note.includes("smoke") || note.includes("user-jwt")) return true; const local = (b.local_database_id || "").toLowerCase(); const online = (b.online_db_id || "").toLowerCase(); return local.startsWith("local_") || online.startsWith("online_"); } /** 产品:与本地一致的库名优先,勿用落库文件名当主标题 */ function channelOnlineTitle(ch: SyncChannel, bindings: SyncBinding[]): string { const linked = bindings.filter((b) => b.channel_id === ch.id); const localNames = linked.map((b) => (b.database_name || "").trim()).filter(Boolean); if (localNames.length) return localNames[0]; const onlineNames = linked.map((b) => (b.display_name || "").trim()).filter(Boolean); if (onlineNames.length) return onlineNames[0]; if ((ch.name || "").trim()) return ch.name.trim(); return ch.remote?.driver || "线上库"; } function channelOnlineIds(ch: SyncChannel, bindings: SyncBinding[]): string { const linked = bindings.filter((b) => b.channel_id === ch.id); const ids = linked.map((b) => b.online_db_id).filter(Boolean); const uniq = [...new Set(ids)]; if (!uniq.length) return ""; // 有可读本地名时副文案带上 local↔online id,便于对账 const localName = linked.map((b) => (b.database_name || "").trim()).find(Boolean); if (localName && linked[0]?.local_database_id) { return `${linked[0].local_database_id} ↔ ${uniq.join(" · ")}`; } return uniq.join(" · "); } function toChannelBody(v: FormValues, id?: string): Partial { return { id, name: v.name, direction: v.direction, conflict_policy: v.conflict_policy, poll_interval_ms: v.poll_interval_ms || 500, agent_id: v.agent_id || undefined, app_slug: (v.app_slug || "").trim() || undefined, local: { driver: v.local.driver, dsn: v.local.dsn, tables: (v.local.tables || "") .split(",") .map((x) => x.trim()) .filter(Boolean), }, remote: { driver: v.remote.driver, dsn: v.remote.dsn, tables: (v.remote.tables || "") .split(",") .map((x) => x.trim()) .filter(Boolean), }, pk_columns: parsePKs(v.pk_columns), }; } export function SyncPage(props: { session: Session; busy: boolean; setBusy: (v: boolean) => void; setError: (v: string) => void; setInfo: (v: string) => void; }) { const { session, busy, setBusy, setError, setInfo } = props; const { message } = AntApp.useApp(); const [items, setItems] = useState([]); const [bindings, setBindings] = useState([]); const [bindCodes, setBindCodes] = useState([]); const [agents, setAgents] = useState([]); const [appOptions, setAppOptions] = useState<{ value: string; label: string }[]>([]); const [loading, setLoading] = useState(false); const [open, setOpen] = useState(false); const [editing, setEditing] = useState(null); const [showSmokeBindings, setShowSmokeBindings] = useState(false); const [filterAgentId, setFilterAgentId] = useState(undefined); const [inspectOpen, setInspectOpen] = useState(false); const [inspectCh, setInspectCh] = useState(null); const [inspectSide, setInspectSide] = useState<"remote" | "local">("remote"); const [inspectData, setInspectData] = useState(null); const [inspectLoading, setInspectLoading] = useState(false); const [previewOpen, setPreviewOpen] = useState(false); const [previewData, setPreviewData] = useState(null); const [previewLoading, setPreviewLoading] = useState(false); const [form] = Form.useForm(); const visibleBindings = showSmokeBindings ? bindings : bindings.filter((b) => !isSmokeBinding(b)); const smokeHidden = bindings.length - visibleBindings.length; const visibleChannels = filterAgentId ? items.filter((c) => c.agent_id === filterAgentId || agents.find((a) => a.agent_id === filterAgentId)?.channel_id === c.id) : items; async function refresh() { setLoading(true); try { const [ch, bind, ag, apps, codes] = await Promise.all([ listSyncChannels(session), listSyncBindings(session).catch(() => ({ items: [] as SyncBinding[] })), listAgents(session).catch(() => ({ items: [] as AgentAccount[] })), listApps(session).catch(() => ({ items: [] as { slug: string; name: string; status: string }[] })), listBindCodes(session).catch(() => ({ items: [] as BindCode[] })), ]); setItems(ch.items || []); setBindings(bind.items || []); setAgents(ag.items || []); setBindCodes(codes.items || []); setAppOptions( (apps.items || []) .filter((a) => a.status === "published" || !a.status) .map((a) => ({ value: a.slug, label: `${a.name || a.slug} (${a.slug})` })) ); } catch (e: any) { const msg = e.message || String(e); setError(msg); message.error(msg); } finally { setLoading(false); } } useEffect(() => { void refresh(); // eslint-disable-next-line react-hooks/exhaustive-deps }, [session.accessToken]); function openCreate() { setEditing(null); form.setFieldsValue({ name: "本地 B ↔ 线上 A", direction: "local_to_remote", conflict_policy: "lww_source", poll_interval_ms: 500, agent_id: undefined, app_slug: undefined, local: { driver: "sqlite", dsn: "file:./data/local.db", tables: "article" }, remote: { driver: "postgres", dsn: "postgres://user:pass@127.0.0.1:5432/app?sslmode=disable", tables: "article", }, pk_columns: "article:id", }); setOpen(true); } function openEdit(ch: SyncChannel) { setEditing(ch); form.setFieldsValue({ name: ch.name, direction: ch.direction, conflict_policy: ch.conflict_policy, poll_interval_ms: ch.poll_interval_ms, agent_id: ch.agent_id || undefined, app_slug: ch.app_slug || undefined, local: { driver: ch.local.driver, dsn: ch.local.dsn, tables: (ch.local.tables || []).join(","), }, remote: { driver: ch.remote.driver, dsn: ch.remote.dsn, tables: (ch.remote.tables || []).join(","), }, pk_columns: Object.entries(ch.pk_columns || {}) .map(([t, p]) => `${t}:${p}`) .join(","), }); setOpen(true); } async function openInspect(ch: SyncChannel, side: "remote" | "local" = "remote") { setInspectCh(ch); setInspectSide(side); setInspectOpen(true); setInspectData(null); setInspectLoading(true); try { const data = await inspectSyncChannel(session, ch.id, { side, include_sync_meta: true }); setInspectData(data); } catch (e: any) { message.error(e.message || String(e)); setInspectOpen(false); } finally { setInspectLoading(false); } } async function openPreview(table: string) { if (!inspectCh) return; setPreviewLoading(true); setPreviewOpen(true); setPreviewData(null); try { const data = await previewSyncTable(session, inspectCh.id, { table, side: inspectSide, limit: 50, }); setPreviewData(data); } catch (e: any) { message.error(e.message || String(e)); setPreviewOpen(false); } finally { setPreviewLoading(false); } } function formatCheckpointTime(iso?: string) { if (!iso) return "—"; const d = new Date(iso); if (Number.isNaN(d.getTime())) return iso; return d.toLocaleString(); } async function openRestore(r: SyncChannel) { setBusy(true); try { const meta = await getSyncCheckpointMeta(session, r.id); if (!meta.latest && !meta.previous) { message.warning("尚无成功同步快照。请先完成一次同步(含自动推送),约 30 秒后生成。"); return; } let which: "latest" | "previous" = meta.latest ? "latest" : "previous"; const options: { value: "latest" | "previous"; label: string }[] = []; if (meta.latest) { options.push({ value: "latest", label: `最近一次 · ${formatCheckpointTime(meta.latest.synced_at)}(${meta.latest.row_count ?? 0} 行 / ${meta.latest.table_count ?? 0} 表)`, }); } if (meta.previous) { options.push({ value: "previous", label: `上一代 · ${formatCheckpointTime(meta.previous.synced_at)}(${meta.previous.row_count ?? 0} 行 / ${meta.previous.table_count ?? 0} 表)`, }); } Modal.confirm({ title: `数据恢复 · ${r.name || r.id}`, width: 560, content: (

线上库恢复为所选同步快照(upsert 并删除快照中不存在的行)。本机数据需宇恒下次同步/pull 从线上补回。

若误删后又发生成功自动同步,最新快照可能已含删除后状态,请改选「上一代」。

恢复到 setFilterAgentId(v)} options={agents.map((a) => ({ value: a.agent_id, label: `${a.name}${a.channel_id ? " · 已绑通道" : ""}`, }))} /> ( {id} ), }, { title: "名称", dataIndex: "name" }, { title: "方向", dataIndex: "direction", render: (d: string) => {d}, }, { title: "本地", render: (_: unknown, r: SyncChannel) => `${r.local.driver}`, }, { title: "线上", render: (_: unknown, r: SyncChannel) => { const title = channelOnlineTitle(r, bindings); const ids = channelOnlineIds(r, bindings); const path = formatDsnHint(r.remote?.dsn || ""); return ( {title}
{r.remote?.driver} {ids ? ` · ${ids}` : ""} {path ? ` · ${path}` : ""}
); }, }, { title: "智能体/模块", render: (_: unknown, r: SyncChannel) => { const ag = agents.find((a) => a.agent_id === r.agent_id); const app = appOptions.find((o) => o.value === r.app_slug); const moduleLabel = r.app_slug ? app ? `${(app.label.split(" (")[0] || app.label).trim()}` : r.app_slug : "未绑模块"; return ( {ag ? ag.name : r.agent_id ? `#${r.agent_id}` : "—"}
{moduleLabel} {r.app_slug && app ? ( <>
{r.app_slug} ) : null}
); }, }, { title: "状态", render: (_: unknown, r: SyncChannel) => r.enabled ? 运行中 : 已停止, }, { title: "推送次数", render: (_: unknown, r: SyncChannel) => ( ↑{r.stats?.pushed_ok || 0} {r.stats?.pushed_skipped ? `(跳过 ${r.stats.pushed_skipped})` : ""}
push 成功次数,非表数
), }, { title: "操作", render: (_: unknown, r: SyncChannel) => ( {r.enabled ? ( ) : ( )} ), }, ]} /> 绑定码(发给终端开通) 关联本公司默认同步落点;终端用 POST /auth/bind-code/redeem{" "} 兑换。生产联调成员手机请用 13531041944,勿用超管号。
{c}, }, { title: "通道", dataIndex: "channel_id", render: (id: string) => id ? {id.slice(0, 8)}… : "—", }, { title: "次数", render: (_: unknown, r: BindCode) => `${r.used_count}/${r.max_uses}`, width: 80, }, { title: "过期", dataIndex: "expires_at", render: (t: string) => (t ? new Date(t).toLocaleString() : "—"), }, { title: "状态", width: 90, render: (_: unknown, r: BindCode) => { if (r.revoked) return 已撤销; if (r.used_count >= r.max_uses) return 已用尽; if (r.expires_at && new Date(r.expires_at).getTime() < Date.now()) return 已过期; return 可用; }, }, { title: "操作", width: 90, render: (_: unknown, r: BindCode) => ( ), }, ]} /> 库绑定(本地 ↔ 线上) local_database_id 与{" "} online_db_id{" "} 本就可以不同(映射关系,不是要求同名)。 可读名靠宇恒 ensure 时带 database_name /{" "} display_name ;通道「线上」列优先用本地可读名。默认隐藏联调烟雾 Binding( local_* / note=smoke)。 setShowSmokeBindings(e.target.checked)} > 显示联调烟雾 Binding {smokeHidden > 0 ? `(已隐藏 ${smokeHidden} 条)` : ""}
( {bindingLocalLabel(b)} {b.database_name ? ( <>
{b.local_database_id} ) : null}
), }, { title: "线上库名", render: (_: unknown, b: SyncBinding) => ( {bindingOnlineLabel(b)} {b.display_name ? ( <>
{b.online_db_id} ) : null}
), }, { title: "映射", render: (_: unknown, b: SyncBinding) => ( {bindingLocalLabel(b)} ↔ {bindingOnlineLabel(b)} ), }, { title: "通道", dataIndex: "channel_id", render: (id: string) => { if (!id) return "—"; const ch = items.find((c) => c.id === id); if (!ch) { return ( 通道已删除 · {id.slice(0, 8)}… ); } const tables = (ch.remote?.tables || []).join(", "); return ( {ch.name || id}
{tables ? `策略表(参考,非整库限制):${tables}` : "整库同步中(Binding 任意表可 push)"}
); }, }, { title: "备注", dataIndex: "note", ellipsis: true }, { title: "共享", width: 70, render: (_: unknown, b: SyncBinding) => b.shared ? 共享 : 个人, }, ]} /> 「同步修复」有最小间隔限流(默认 5 分钟)。 ↑ 推送次数 = agent push 成功次数,不是业务表数量。 本地有、线上没有的表(如尚未 push 的业务表)属正常,不是 Binding 映射错了。 Binding 标 shared=true 为公司共享库(仅管理员可设)。 setInspectOpen(false)} width={820} footer={ } > {inspectData?.dsn_hint ? ( {inspectData.driver} · {inspectData.dsn_hint} {inspectData.message ? ` · ${inspectData.message}` : ""} ) : null}
!String(t.name || "").startsWith("_ajz_"))} pagination={false} size="small" locale={{ emptyText: "暂无业务表(尚未 push 或连不上库)" }} columns={[ { title: "表名", dataIndex: "name" }, { title: "行数", dataIndex: "row_count", width: 90, render: (n: number) => {n}, }, { title: "字段数", dataIndex: "column_count", width: 90, }, { title: "字段", dataIndex: "columns", ellipsis: true, render: (cols: string[]) => (cols || []).join(", "), }, { title: "操作", width: 168, render: (_: unknown, t: { name: string; row_count: number }) => ( ), }, ]} /> {(inspectData?.tables || []).some((t) => String(t.name || "").startsWith("_ajz_")) ? ( <> 已隐藏同步系统表 _ajz_* (版本/outbox,不是业务模块表)。 ) : null} 本机有、这里没有的表 = 尚未 push/ensure。「删表」只删当前查看侧,不会同步到另一侧。 setPreviewOpen(false)} width={960} footer={ } >
({ ...row, __k: i }))} rowKey="__k" size="small" scroll={{ x: true }} pagination={false} columns={(previewData?.columns || []).map((c) => ({ title: c, dataIndex: c, ellipsis: true, render: (v: unknown) => v === null || v === undefined ? ( null ) : ( String(v) ), }))} /> setOpen(false)} width={720} footer={ } >
本地库 线上库 A(生产统一 Postgres)
); }