"use client"; import { useEffect, useMemo, useState } from "react"; import { useRouter } from "next/navigation"; import { toast } from "sonner"; import { api, type Connection, type Pipeline, type StreamConfig, type StreamInfo, } from "@/components/api"; import { keyValueLinesFromJson, parseKeyValueLines } from "@/components/kv"; import { Button } from "@/components/ui/button"; import { Card, CardContent, CardDescription, CardHeader, CardTitle, } from "@/components/ui/card"; import { Checkbox } from "@/components/ui/checkbox"; import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue, } from "@/components/ui/select"; import { Switch } from "@/components/ui/switch"; import { Textarea } from "@/components/ui/textarea"; const MODES = [ { value: "full-refresh", label: "full-refresh(全量覆盖,默认)" }, { value: "truncate", label: "truncate(清空后写入)" }, { value: "incremental", label: "incremental(增量)" }, { value: "snapshot", label: "snapshot(快照)" }, ]; export function PipelineForm({ initial }: { initial?: Pipeline }) { const router = useRouter(); const isEdit = !!initial; const [connections, setConnections] = useState([]); const [connError, setConnError] = useState(null); const [name, setName] = useState(initial?.name ?? ""); const [sourceId, setSourceId] = useState( initial ? String(initial.source_conn_id) : "" ); const [targetId, setTargetId] = useState( initial ? String(initial.target_conn_id) : "" ); const [mode, setMode] = useState(initial?.mode ?? "full-refresh"); const [schemaSync, setSchemaSync] = useState( initial ? !!initial.schema_sync : true ); const [schemaScope, setSchemaScope] = useState( initial?.schema_scope ?? "selected" ); const [includeFk, setIncludeFk] = useState(initial ? !!initial.include_fk : false); const [envText, setEnvText] = useState(() => initial ? keyValueLinesFromJson(initial.env) : "" ); const [streamList, setStreamList] = useState(null); const [streamsLoading, setStreamsLoading] = useState(false); const [selected, setSelected] = useState>(() => { if (!initial) return {}; try { const streams = JSON.parse(initial.streams || "[]") as StreamConfig[]; return Object.fromEntries(streams.map((s) => [s.name, s])); } catch { return {}; } }); const [saving, setSaving] = useState(false); useEffect(() => { api("/api/connections") .then(setConnections) .catch((e) => setConnError(e instanceof Error ? e.message : "加载失败")); }, []); // 编辑模式:加载现有流水线后自动拉取表列表,回填勾选状态 useEffect(() => { if (isEdit && sourceId) void loadStreams(); // eslint-disable-next-line react-hooks/exhaustive-deps }, [isEdit, sourceId]); const selectedList = useMemo(() => Object.values(selected), [selected]); async function loadStreams() { if (!sourceId) return; setStreamsLoading(true); try { const list = await api( `/api/connections/${sourceId}/streams` ); setStreamList(list); if (list.length === 0) toast.info("该连接没有可同步的表"); } catch (e) { toast.error(e instanceof Error ? e.message : "加载表失败"); } finally { setStreamsLoading(false); } } function toggleStream(s: StreamInfo, checked: boolean) { setSelected((m) => { const next = { ...m }; if (checked) next[s.name] = m[s.name] ?? { name: s.name }; else delete next[s.name]; return next; }); } function selectAll(on: boolean) { if (!streamList) return; setSelected( on ? Object.fromEntries( streamList.map((s) => [s.name, selected[s.name] ?? { name: s.name }]) ) : {} ); } function setStreamField( name: string, field: "primary_key" | "update_key", value: string ) { setSelected((m) => ({ ...m, [name]: { ...m[name], name, [field]: value || undefined }, })); } async function submit() { if (!name.trim()) return toast.error("请填写流水线名称"); if (!sourceId || !targetId) return toast.error("请选择源连接和目标连接"); if (sourceId === targetId) return toast.error("源连接和目标连接不能相同"); if (selectedList.length === 0) return toast.error("请至少选择一张表"); setSaving(true); const payload = { name: name.trim(), source_conn_id: Number(sourceId), target_conn_id: Number(targetId), schema_sync: schemaSync ? 1 : 0, schema_scope: schemaScope, include_fk: includeFk ? 1 : 0, mode, streams: JSON.stringify(selectedList), env: parseKeyValueLines(envText), }; try { const p = isEdit ? await api(`/api/pipelines/${initial.id}`, { method: "PUT", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), }) : await api("/api/pipelines", { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), }); toast.success(isEdit ? "流水线已保存" : "流水线已创建"); router.push(`/pipelines/${p.id}`); } catch (e) { toast.error(e instanceof Error ? e.message : isEdit ? "保存失败" : "创建失败"); setSaving(false); } } const cancelHref = isEdit ? `/pipelines/${initial.id}` : "/pipelines"; return (

{isEdit ? `编辑流水线:${initial.name}` : "新建流水线"}

{connError ? (
连接列表加载失败:{connError}
) : null} 基本信息
setName(e.target.value)} placeholder="例如:测试库 → 本地 dev" />
先用 Atlas 同步 schema(表/索引/约束),再迁移数据
setSchemaSync(!!c)} />
整库模式会同步源库所有表,且目标库独有的表可能被删除
关闭可避免选中表引用了未迁移表导致的失败;排除父表时引用它的外键会自动跳过
setIncludeFk(!!c)} disabled={!schemaSync} />