server.go 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. package etcdserver
  2. import (
  3. "errors"
  4. "sync"
  5. "time"
  6. "code.google.com/p/go.net/context"
  7. pb "github.com/coreos/etcd/etcdserver2/etcdserverpb"
  8. "github.com/coreos/etcd/raft"
  9. "github.com/coreos/etcd/raft/raftpb"
  10. "github.com/coreos/etcd/store"
  11. "github.com/coreos/etcd/wait"
  12. )
  13. var (
  14. ErrUnknownMethod = errors.New("etcdserver: unknown method")
  15. ErrStopped = errors.New("etcdserver: server stopped")
  16. )
  17. type SendFunc func(m []raftpb.Message)
  18. type Response struct {
  19. // The last seen term raft was at when this request was built.
  20. Term int64
  21. // The last seen index raft was at when this request was built.
  22. Commit int64
  23. *store.Event
  24. *store.Watcher
  25. err error
  26. }
  27. type Server struct {
  28. once sync.Once
  29. w *wait.List
  30. done chan struct{}
  31. Node raft.Node
  32. Store store.Store
  33. // Send specifies the send function for sending msgs to peers. Send
  34. // MUST NOT block. It is okay to drop messages, since clients should
  35. // timeout and reissue their messages. If Send is nil, Server will
  36. // panic.
  37. Send SendFunc
  38. // Save specifies the save function for saving ents to stable storage.
  39. // Save MUST block until st and ents are on stable storage. If Send is
  40. // nil, Server will panic.
  41. Save func(st raftpb.State, ents []raftpb.Entry)
  42. }
  43. // Start prepares and starts server in a new goroutine.
  44. func Start(s *Server) {
  45. s.w = wait.New()
  46. s.done = make(chan struct{})
  47. go s.run()
  48. }
  49. func (s *Server) run() {
  50. for {
  51. select {
  52. case rd := <-s.Node.Ready():
  53. s.Save(rd.State, rd.Entries)
  54. s.Send(rd.Messages)
  55. // TODO(bmizerany): do this in the background, but take
  56. // care to apply entries in a single goroutine, and not
  57. // race them.
  58. for _, e := range rd.CommittedEntries {
  59. var r pb.Request
  60. if err := r.Unmarshal(e.Data); err != nil {
  61. panic("TODO: this is bad, what do we do about it?")
  62. }
  63. var resp Response
  64. resp.Event, resp.err = s.apply(context.TODO(), r)
  65. resp.Term = rd.Term
  66. resp.Commit = rd.Commit
  67. s.w.Trigger(r.Id, resp)
  68. }
  69. case <-s.done:
  70. return
  71. }
  72. }
  73. }
  74. // Stop stops the server, and shutsdown the running goroutine. Stop should be
  75. // called after a Start(s), otherwise it will block forever.
  76. func (s *Server) Stop() {
  77. s.done <- struct{}{}
  78. }
  79. // Do interprets r and performs an operation on s.Store according to r.Method
  80. // and other fields. If r.Method is "PUT", "PUT", or "DELETE, r will be sent
  81. // through consensus before performing its respective operation. Do will block
  82. // until an action is performed or there is an error.
  83. func (s *Server) Do(ctx context.Context, r pb.Request) (Response, error) {
  84. if r.Id == 0 {
  85. panic("r.Id cannot be 0")
  86. }
  87. switch r.Method {
  88. case "POST", "PUT", "DELETE":
  89. data, err := r.Marshal()
  90. if err != nil {
  91. return Response{}, err
  92. }
  93. ch := s.w.Register(r.Id)
  94. s.Node.Propose(ctx, data)
  95. select {
  96. case x := <-ch:
  97. resp := x.(Response)
  98. return resp, resp.err
  99. case <-ctx.Done():
  100. s.w.Trigger(r.Id, nil) // GC wait
  101. return Response{}, ctx.Err()
  102. case <-s.done:
  103. return Response{}, ErrStopped
  104. }
  105. case "GET":
  106. switch {
  107. case r.Wait:
  108. wc, err := s.Store.Watch(r.Path, r.Recursive, false, r.Since)
  109. if err != nil {
  110. return Response{}, err
  111. }
  112. return Response{Watcher: wc}, nil
  113. default:
  114. ev, err := s.Store.Get(r.Path, r.Recursive, r.Sorted)
  115. if err != nil {
  116. return Response{}, err
  117. }
  118. return Response{Event: ev}, nil
  119. }
  120. default:
  121. return Response{}, ErrUnknownMethod
  122. }
  123. }
  124. // apply interprets r as a call to store.X and returns an Response interpreted from store.Event
  125. func (s *Server) apply(ctx context.Context, r pb.Request) (*store.Event, error) {
  126. expr := time.Unix(0, r.Expiration)
  127. switch r.Method {
  128. case "POST":
  129. return s.Store.Create(r.Path, r.Dir, r.Val, true, expr)
  130. case "PUT":
  131. exists, existsSet := getBool(r.PrevExists)
  132. switch {
  133. case existsSet:
  134. if exists {
  135. return s.Store.Update(r.Path, r.Val, expr)
  136. } else {
  137. return s.Store.Create(r.Path, r.Dir, r.Val, false, expr)
  138. }
  139. case r.PrevIndex > 0 || r.PrevValue != "":
  140. return s.Store.CompareAndSwap(r.Path, r.PrevValue, r.PrevIndex, r.Val, expr)
  141. default:
  142. return s.Store.Set(r.Path, r.Dir, r.Val, expr)
  143. }
  144. case "DELETE":
  145. switch {
  146. case r.PrevIndex > 0 || r.PrevValue != "":
  147. return s.Store.CompareAndDelete(r.Path, r.PrevValue, r.PrevIndex)
  148. default:
  149. return s.Store.Delete(r.Path, r.Recursive, r.Dir)
  150. }
  151. default:
  152. // This should never be reached, but just in case:
  153. return nil, ErrUnknownMethod
  154. }
  155. }
  156. func getBool(v *bool) (vv bool, set bool) {
  157. if v == nil {
  158. return false, false
  159. }
  160. return *v, true
  161. }