migrate_command.go 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404
  1. // Copyright 2016 The etcd Authors
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package command
  15. import (
  16. "encoding/binary"
  17. "encoding/json"
  18. "fmt"
  19. "io"
  20. "os"
  21. "os/exec"
  22. "path/filepath"
  23. "time"
  24. "go.etcd.io/etcd/client"
  25. "go.etcd.io/etcd/etcdserver"
  26. "go.etcd.io/etcd/etcdserver/api"
  27. "go.etcd.io/etcd/etcdserver/api/membership"
  28. "go.etcd.io/etcd/etcdserver/api/snap"
  29. "go.etcd.io/etcd/etcdserver/api/v2error"
  30. "go.etcd.io/etcd/etcdserver/api/v2store"
  31. pb "go.etcd.io/etcd/etcdserver/etcdserverpb"
  32. "go.etcd.io/etcd/mvcc"
  33. "go.etcd.io/etcd/mvcc/backend"
  34. "go.etcd.io/etcd/mvcc/mvccpb"
  35. "go.etcd.io/etcd/pkg/pbutil"
  36. "go.etcd.io/etcd/pkg/types"
  37. "go.etcd.io/etcd/raft/raftpb"
  38. "go.etcd.io/etcd/wal"
  39. "go.etcd.io/etcd/wal/walpb"
  40. "github.com/gogo/protobuf/proto"
  41. "github.com/spf13/cobra"
  42. "go.uber.org/zap"
  43. )
  44. var (
  45. migrateExcludeTTLKey bool
  46. migrateDatadir string
  47. migrateWALdir string
  48. migrateTransformer string
  49. )
  50. // NewMigrateCommand returns the cobra command for "migrate".
  51. func NewMigrateCommand() *cobra.Command {
  52. mc := &cobra.Command{
  53. Use: "migrate",
  54. Short: "Migrates keys in a v2 store to a mvcc store",
  55. Run: migrateCommandFunc,
  56. }
  57. mc.Flags().BoolVar(&migrateExcludeTTLKey, "no-ttl", false, "Do not convert TTL keys")
  58. mc.Flags().StringVar(&migrateDatadir, "data-dir", "", "Path to the data directory")
  59. mc.Flags().StringVar(&migrateWALdir, "wal-dir", "", "Path to the WAL directory")
  60. mc.Flags().StringVar(&migrateTransformer, "transformer", "", "Path to the user-provided transformer program")
  61. return mc
  62. }
  63. func migrateCommandFunc(cmd *cobra.Command, args []string) {
  64. var (
  65. writer io.WriteCloser
  66. reader io.ReadCloser
  67. errc chan error
  68. )
  69. if migrateTransformer != "" {
  70. writer, reader, errc = startTransformer()
  71. } else {
  72. fmt.Println("using default transformer")
  73. writer, reader, errc = defaultTransformer()
  74. }
  75. st, index := rebuildStoreV2()
  76. be := prepareBackend()
  77. defer be.Close()
  78. go func() {
  79. writeStore(writer, st)
  80. writer.Close()
  81. }()
  82. readKeys(reader, be)
  83. mvcc.UpdateConsistentIndex(be, index)
  84. err := <-errc
  85. if err != nil {
  86. fmt.Println("failed to transform keys")
  87. ExitWithError(ExitError, err)
  88. }
  89. fmt.Println("finished transforming keys")
  90. }
  91. func prepareBackend() backend.Backend {
  92. var be backend.Backend
  93. bch := make(chan struct{})
  94. dbpath := filepath.Join(migrateDatadir, "member", "snap", "db")
  95. go func() {
  96. defer close(bch)
  97. be = backend.NewDefaultBackend(dbpath)
  98. }()
  99. select {
  100. case <-bch:
  101. case <-time.After(time.Second):
  102. fmt.Fprintf(os.Stderr, "waiting for etcd to close and release its lock on %q\n", dbpath)
  103. <-bch
  104. }
  105. tx := be.BatchTx()
  106. tx.Lock()
  107. tx.UnsafeCreateBucket([]byte("key"))
  108. tx.UnsafeCreateBucket([]byte("meta"))
  109. tx.Unlock()
  110. return be
  111. }
  112. func rebuildStoreV2() (v2store.Store, uint64) {
  113. var index uint64
  114. cl := membership.NewCluster(zap.NewExample(), "")
  115. waldir := migrateWALdir
  116. if len(waldir) == 0 {
  117. waldir = filepath.Join(migrateDatadir, "member", "wal")
  118. }
  119. snapdir := filepath.Join(migrateDatadir, "member", "snap")
  120. ss := snap.New(zap.NewExample(), snapdir)
  121. snapshot, err := ss.Load()
  122. if err != nil && err != snap.ErrNoSnapshot {
  123. ExitWithError(ExitError, err)
  124. }
  125. var walsnap walpb.Snapshot
  126. if snapshot != nil {
  127. walsnap.Index, walsnap.Term = snapshot.Metadata.Index, snapshot.Metadata.Term
  128. index = snapshot.Metadata.Index
  129. }
  130. w, err := wal.OpenForRead(zap.NewExample(), waldir, walsnap)
  131. if err != nil {
  132. ExitWithError(ExitError, err)
  133. }
  134. defer w.Close()
  135. _, _, ents, err := w.ReadAll()
  136. if err != nil {
  137. ExitWithError(ExitError, err)
  138. }
  139. st := v2store.New()
  140. if snapshot != nil {
  141. err := st.Recovery(snapshot.Data)
  142. if err != nil {
  143. ExitWithError(ExitError, err)
  144. }
  145. }
  146. cl.SetStore(st)
  147. cl.Recover(api.UpdateCapability)
  148. applier := etcdserver.NewApplierV2(zap.NewExample(), st, cl)
  149. for _, ent := range ents {
  150. if ent.Type == raftpb.EntryConfChange {
  151. var cc raftpb.ConfChange
  152. pbutil.MustUnmarshal(&cc, ent.Data)
  153. applyConf(cc, cl)
  154. continue
  155. }
  156. var raftReq pb.InternalRaftRequest
  157. if !pbutil.MaybeUnmarshal(&raftReq, ent.Data) { // backward compatible
  158. var r pb.Request
  159. pbutil.MustUnmarshal(&r, ent.Data)
  160. applyRequest(&r, applier)
  161. } else {
  162. if raftReq.V2 != nil {
  163. req := raftReq.V2
  164. applyRequest(req, applier)
  165. }
  166. }
  167. if ent.Index > index {
  168. index = ent.Index
  169. }
  170. }
  171. return st, index
  172. }
  173. func applyConf(cc raftpb.ConfChange, cl *membership.RaftCluster) {
  174. if err := cl.ValidateConfigurationChange(cc); err != nil {
  175. return
  176. }
  177. switch cc.Type {
  178. case raftpb.ConfChangeAddNode:
  179. m := new(membership.Member)
  180. if err := json.Unmarshal(cc.Context, m); err != nil {
  181. panic(err)
  182. }
  183. cl.AddMember(m)
  184. case raftpb.ConfChangeRemoveNode:
  185. cl.RemoveMember(types.ID(cc.NodeID))
  186. case raftpb.ConfChangeUpdateNode:
  187. m := new(membership.Member)
  188. if err := json.Unmarshal(cc.Context, m); err != nil {
  189. panic(err)
  190. }
  191. cl.UpdateRaftAttributes(m.ID, m.RaftAttributes)
  192. }
  193. }
  194. func applyRequest(req *pb.Request, applyV2 etcdserver.ApplierV2) {
  195. r := (*etcdserver.RequestV2)(req)
  196. r.TTLOptions()
  197. switch r.Method {
  198. case "POST":
  199. applyV2.Post(r)
  200. case "PUT":
  201. applyV2.Put(r)
  202. case "DELETE":
  203. applyV2.Delete(r)
  204. case "QGET":
  205. applyV2.QGet(r)
  206. case "SYNC":
  207. applyV2.Sync(r)
  208. default:
  209. panic("unknown command")
  210. }
  211. }
  212. func writeStore(w io.Writer, st v2store.Store) uint64 {
  213. all, err := st.Get("/1", true, true)
  214. if err != nil {
  215. if eerr, ok := err.(*v2error.Error); ok && eerr.ErrorCode == v2error.EcodeKeyNotFound {
  216. fmt.Println("no v2 keys to migrate")
  217. os.Exit(0)
  218. }
  219. ExitWithError(ExitError, err)
  220. }
  221. return writeKeys(w, all.Node)
  222. }
  223. func writeKeys(w io.Writer, n *v2store.NodeExtern) uint64 {
  224. maxIndex := n.ModifiedIndex
  225. nodes := n.Nodes
  226. // remove store v2 bucket prefix
  227. n.Key = n.Key[2:]
  228. if n.Key == "" {
  229. n.Key = "/"
  230. }
  231. if n.Dir {
  232. n.Nodes = nil
  233. }
  234. if !migrateExcludeTTLKey || n.TTL == 0 {
  235. b, err := json.Marshal(n)
  236. if err != nil {
  237. ExitWithError(ExitError, err)
  238. }
  239. fmt.Fprint(w, string(b))
  240. }
  241. for _, nn := range nodes {
  242. max := writeKeys(w, nn)
  243. if max > maxIndex {
  244. maxIndex = max
  245. }
  246. }
  247. return maxIndex
  248. }
  249. func readKeys(r io.Reader, be backend.Backend) error {
  250. for {
  251. length64, err := readInt64(r)
  252. if err != nil {
  253. if err == io.EOF {
  254. return nil
  255. }
  256. return err
  257. }
  258. buf := make([]byte, int(length64))
  259. if _, err = io.ReadFull(r, buf); err != nil {
  260. return err
  261. }
  262. var kv mvccpb.KeyValue
  263. err = proto.Unmarshal(buf, &kv)
  264. if err != nil {
  265. return err
  266. }
  267. mvcc.WriteKV(be, kv)
  268. }
  269. }
  270. func readInt64(r io.Reader) (int64, error) {
  271. var n int64
  272. err := binary.Read(r, binary.LittleEndian, &n)
  273. return n, err
  274. }
  275. func startTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  276. cmd := exec.Command(migrateTransformer)
  277. cmd.Stderr = os.Stderr
  278. writer, err := cmd.StdinPipe()
  279. if err != nil {
  280. ExitWithError(ExitError, err)
  281. }
  282. reader, rerr := cmd.StdoutPipe()
  283. if rerr != nil {
  284. ExitWithError(ExitError, rerr)
  285. }
  286. if err := cmd.Start(); err != nil {
  287. ExitWithError(ExitError, err)
  288. }
  289. errc := make(chan error, 1)
  290. go func() {
  291. errc <- cmd.Wait()
  292. }()
  293. return writer, reader, errc
  294. }
  295. func defaultTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  296. // transformer decodes v2 keys from sr
  297. sr, sw := io.Pipe()
  298. // transformer encodes v3 keys into dw
  299. dr, dw := io.Pipe()
  300. decoder := json.NewDecoder(sr)
  301. errc := make(chan error, 1)
  302. go func() {
  303. defer func() {
  304. sr.Close()
  305. dw.Close()
  306. }()
  307. for decoder.More() {
  308. node := &client.Node{}
  309. if err := decoder.Decode(node); err != nil {
  310. errc <- err
  311. return
  312. }
  313. kv := transform(node)
  314. if kv == nil {
  315. continue
  316. }
  317. data, err := proto.Marshal(kv)
  318. if err != nil {
  319. errc <- err
  320. return
  321. }
  322. buf := make([]byte, 8)
  323. binary.LittleEndian.PutUint64(buf, uint64(len(data)))
  324. if _, err := dw.Write(buf); err != nil {
  325. errc <- err
  326. return
  327. }
  328. if _, err := dw.Write(data); err != nil {
  329. errc <- err
  330. return
  331. }
  332. }
  333. errc <- nil
  334. }()
  335. return sw, dr, errc
  336. }
  337. func transform(n *client.Node) *mvccpb.KeyValue {
  338. const unKnownVersion = 1
  339. if n.Dir {
  340. return nil
  341. }
  342. kv := &mvccpb.KeyValue{
  343. Key: []byte(n.Key),
  344. Value: []byte(n.Value),
  345. CreateRevision: int64(n.CreatedIndex),
  346. ModRevision: int64(n.ModifiedIndex),
  347. Version: unKnownVersion,
  348. }
  349. return kv
  350. }