client_internal_test.go 2.4 KB
Newer Older
H
Helin Wang 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42
package master

import (
	"fmt"
	"net"
	"net/http"
	"net/rpc"
	"os"
	"strconv"
	"strings"
	"testing"
	"time"

	log "github.com/sirupsen/logrus"

	"github.com/PaddlePaddle/Paddle/go/connection"
	"github.com/PaddlePaddle/recordio"
)

const (
	totalTask    = 20
	chunkPerTask = 10
)

func init() {
	log.SetLevel(log.ErrorLevel)
}

func TestGetFinishTask(t *testing.T) {
	const path = "/tmp/master_client_test_0"

	l, err := net.Listen("tcp", ":0")
	if err != nil {
		panic(err)
	}

	ss := strings.Split(l.Addr().String(), ":")
	p, err := strconv.Atoi(ss[len(ss)-1])
	if err != nil {
		panic(err)
	}
	go func(l net.Listener) {
H
Helin Wang 已提交
43
		s, err := NewService(&InMemStore{}, chunkPerTask, time.Second, 1)
44 45 46 47
		if err != nil {
			panic(err)
		}

H
Helin Wang 已提交
48
		server := rpc.NewServer()
49
		err = server.Register(s)
H
Helin Wang 已提交
50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68
		if err != nil {
			panic(err)
		}

		mux := http.NewServeMux()
		mux.Handle(rpc.DefaultRPCPath, server)
		err = http.Serve(l, mux)
		if err != nil {
			panic(err)
		}
	}(l)

	f, err := os.Create(path)
	if err != nil {
		panic(err)
	}

	for i := 0; i < totalTask*chunkPerTask; i++ {
		w := recordio.NewWriter(f, -1, -1)
H
Helin Wang 已提交
69 70 71 72 73
		_, err = w.Write(nil)
		if err != nil {
			panic(err)
		}

H
Helin Wang 已提交
74
		// call Close to force RecordIO writing a chunk.
H
Helin Wang 已提交
75 76 77 78 79 80 81 82
		err = w.Close()
		if err != nil {
			panic(err)
		}
	}
	err = f.Close()
	if err != nil {
		panic(err)
H
Helin Wang 已提交
83 84
	}

H
Helin Wang 已提交
85
	// Manually intialize client to avoid calling c.getRecords()
H
Helin Wang 已提交
86 87
	c := &Client{}
	c.conn = connection.New()
88 89 90 91
	addr := fmt.Sprintf(":%d", p)
	ch := make(chan string, 1)
	ch <- addr
	go c.monitorMaster(ch)
H
Helin Wang 已提交
92 93 94 95 96
	err = c.SetDataset([]string{path})
	if err != nil {
		panic(err)
	}

H
Helin Wang 已提交
97 98 99 100 101
	checkOnePass := func(i int) {
		var tasks []Task
		for idx := 0; idx < totalTask; idx++ {
			task, err := c.getTask()
			if err != nil {
H
Helin Wang 已提交
102
				t.Fatalf("Error: %v, pass: %d\n", err, i)
H
Helin Wang 已提交
103 104 105 106 107 108
			}
			tasks = append(tasks, task)
		}

		_, err = c.getTask()
		if err == nil {
H
Helin Wang 已提交
109
			t.Fatalf("Should get error, pass: %d\n", i)
H
Helin Wang 已提交
110 111
		}

G
gongweibao 已提交
112
		err = c.taskFinished(tasks[0].Meta.ID)
H
Helin Wang 已提交
113
		if err != nil {
H
Helin Wang 已提交
114
			t.Fatalf("Error: %v, pass: %d\n", err, i)
H
Helin Wang 已提交
115
		}
G
gongweibao 已提交
116 117 118 119 120 121

		err = c.taskFailed(tasks[0].Meta)
		if err != nil {
			t.Fatalf("Error: %v, pass: %d\n", err, i)
		}

H
Helin Wang 已提交
122 123 124 125 126 127 128 129
		tasks = tasks[1:]
		task, err := c.getTask()
		if err != nil {
			t.Fatal(err)
		}
		tasks = append(tasks, task)

		for _, task := range tasks {
G
gongweibao 已提交
130
			err = c.taskFinished(task.Meta.ID)
H
Helin Wang 已提交
131
			if err != nil {
H
Helin Wang 已提交
132
				t.Fatalf("Error: %v, pass: %d\n", err, i)
H
Helin Wang 已提交
133 134 135 136 137 138 139 140
			}
		}
	}

	for i := 0; i < 10; i++ {
		checkOnePass(i)
	}
}