|
| 1 | +package kv |
| 2 | + |
| 3 | +import ( |
| 4 | + "context" |
| 5 | + "testing" |
| 6 | + "time" |
| 7 | + |
| 8 | + "github.com/bootjp/elastickv/distribution" |
| 9 | + pb "github.com/bootjp/elastickv/proto" |
| 10 | + "github.com/bootjp/elastickv/store" |
| 11 | + "github.com/hashicorp/raft" |
| 12 | +) |
| 13 | + |
| 14 | +// helper to create single-node raft |
| 15 | +func newTestRaft(t *testing.T, id string, fsm raft.FSM) *raft.Raft { |
| 16 | + t.Helper() |
| 17 | + c := raft.DefaultConfig() |
| 18 | + c.LocalID = raft.ServerID(id) |
| 19 | + ldb := raft.NewInmemStore() |
| 20 | + sdb := raft.NewInmemStore() |
| 21 | + fss := raft.NewInmemSnapshotStore() |
| 22 | + addr, trans := raft.NewInmemTransport(raft.ServerAddress(id)) |
| 23 | + r, err := raft.NewRaft(c, fsm, ldb, sdb, fss, trans) |
| 24 | + if err != nil { |
| 25 | + t.Fatalf("new raft: %v", err) |
| 26 | + } |
| 27 | + cfg := raft.Configuration{Servers: []raft.Server{{ID: raft.ServerID(id), Address: addr}}} |
| 28 | + if err := r.BootstrapCluster(cfg).Error(); err != nil { |
| 29 | + t.Fatalf("bootstrap: %v", err) |
| 30 | + } |
| 31 | + |
| 32 | + // single node should eventually become leader |
| 33 | + for i := 0; i < 100; i++ { |
| 34 | + if r.State() == raft.Leader { |
| 35 | + break |
| 36 | + } |
| 37 | + time.Sleep(50 * time.Millisecond) |
| 38 | + } |
| 39 | + if r.State() != raft.Leader { |
| 40 | + t.Fatalf("node %s is not leader", id) |
| 41 | + } |
| 42 | + return r |
| 43 | +} |
| 44 | + |
| 45 | +func TestShardedTransactionManagerCommit(t *testing.T) { |
| 46 | + e := distribution.NewEngine() |
| 47 | + e.UpdateRoute([]byte("a"), []byte("m"), 1) |
| 48 | + e.UpdateRoute([]byte("m"), nil, 2) |
| 49 | + |
| 50 | + stm := NewShardedTransactionManager(e) |
| 51 | + |
| 52 | + // group 1 |
| 53 | + s1 := store.NewRbMemoryStore() |
| 54 | + l1 := store.NewRbMemoryStoreWithExpire(time.Minute) |
| 55 | + r1 := newTestRaft(t, "1", NewKvFSM(s1, l1)) |
| 56 | + defer r1.Shutdown() |
| 57 | + stm.Register(1, NewTransaction(r1)) |
| 58 | + |
| 59 | + // group 2 |
| 60 | + s2 := store.NewRbMemoryStore() |
| 61 | + l2 := store.NewRbMemoryStoreWithExpire(time.Minute) |
| 62 | + r2 := newTestRaft(t, "2", NewKvFSM(s2, l2)) |
| 63 | + defer r2.Shutdown() |
| 64 | + stm.Register(2, NewTransaction(r2)) |
| 65 | + |
| 66 | + reqs := []*pb.Request{ |
| 67 | + {IsTxn: false, Phase: pb.Phase_NONE, Mutations: []*pb.Mutation{{Op: pb.Op_PUT, Key: []byte("b"), Value: []byte("v1")}}}, |
| 68 | + {IsTxn: false, Phase: pb.Phase_NONE, Mutations: []*pb.Mutation{{Op: pb.Op_PUT, Key: []byte("x"), Value: []byte("v2")}}}, |
| 69 | + } |
| 70 | + |
| 71 | + _, err := stm.Commit(reqs) |
| 72 | + if err != nil { |
| 73 | + t.Fatalf("commit: %v", err) |
| 74 | + } |
| 75 | + |
| 76 | + v, err := s1.Get(context.Background(), []byte("b")) |
| 77 | + if err != nil || string(v) != "v1" { |
| 78 | + t.Fatalf("group1 value: %v %v", v, err) |
| 79 | + } |
| 80 | + v, err = s2.Get(context.Background(), []byte("x")) |
| 81 | + if err != nil || string(v) != "v2" { |
| 82 | + t.Fatalf("group2 value: %v %v", v, err) |
| 83 | + } |
| 84 | +} |
0 commit comments