forked from Terry-Mao/gopush-cluster
-
Notifications
You must be signed in to change notification settings - Fork 0
/
zk.go
90 lines (83 loc) · 2.81 KB
/
zk.go
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
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
// Copyright © 2014 Terry Mao, LiuDing All rights reserved.
// This file is part of gopush-cluster.
// gopush-cluster is free software: you can redistribute it and/or modify
// it under the terms of the GNU General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
// gopush-cluster is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU General Public License for more details.
// You should have received a copy of the GNU General Public License
// along with gopush-cluster. If not, see <http://www.gnu.org/licenses/>.
// github.com/samuel/go-zookeeper
// Copyright (c) 2013, Samuel Stauffer <samuel@descolada.com>
// All rights reserved.
package main
import (
"fmt"
"github.com/samuel/go-zookeeper/zk"
"strings"
)
func InitZK() (*zk.Conn, error) {
// connect to zookeeper, get event from chan in goroutine(log)
conn, session, err := zk.Connect(Conf.ZookeeperAddr, Conf.ZookeeperTimeout)
if err != nil {
Log.Error("zk.Connect(\"%v\", %d) error(%v)", Conf.ZookeeperAddr, Conf.ZookeeperTimeout, err)
return nil, err
}
go func() {
for {
event := <-session
Log.Info("zookeeper get a event: %s", event.State.String())
}
}()
// create zk root path
tpath := ""
for _, str := range strings.Split(Conf.ZookeeperPath, "/")[1:] {
tpath += "/" + str
Log.Debug("create zookeeper path:%s", tpath)
_, err = conn.Create(tpath, []byte(""), 0, zk.WorldACL(zk.PermAll))
if err != nil {
if err == zk.ErrNodeExists {
Log.Warn("zk.create(\"%s\") exists", tpath)
} else {
Log.Error("zk.create(\"%s\") error(%v)", tpath, err)
return nil, err
}
}
}
// tcp, websocket and rpc bind address store in the zk
data := ""
for _, addr := range Conf.Addr {
data += fmt.Sprintf("tcp://%s,", addr)
}
data = strings.TrimRight(data, ",")
tpath, err = conn.Create(Conf.ZookeeperPath+"/", []byte(data), zk.FlagEphemeral|zk.FlagSequence, zk.WorldACL(zk.PermAll))
if err != nil {
Log.Error("conn.Create(\"%s\", \"%s\", zk.FlagEphemeral|zk.FlagSequence) error(%v)", Conf.ZookeeperPath, data, err)
return nil, err
}
Log.Debug("create a zookeeper node:%s", tpath)
// watch self
go func() {
for {
Log.Info("zk path: \"%s\" set a watch", tpath)
exist, _, watch, err := conn.ExistsW(tpath)
if err != nil {
Log.Error("zk.ExistsW(\"%s\") error(%v)", tpath, err)
Log.Warn("zk path: \"%s\" set watch failed, message kill itself", tpath)
KillSelf()
return
}
if !exist {
Log.Warn("zk path: \"%s\" not exist, message kill itself", tpath)
KillSelf()
return
}
event := <-watch
Log.Info("zk path: \"%s\" receive a event %v", tpath, event)
}
}()
return conn, nil
}