migrate_command.go 8.6 KB

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