/usr/share/gocode/src/github.com/hashicorp/serf/cmd/serf/command/agent/ipc_event_stream.go is in golang-github-hashicorp-serf-dev 0.8.1+git20171021.c20a0b1~ds1-4.
This file is owned by root:root, with mode 0o644.
The actual contents of the file can be viewed below.
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 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 | package agent
import (
"fmt"
"github.com/hashicorp/serf/serf"
"log"
)
type streamClient interface {
Send(*responseHeader, interface{}) error
RegisterQuery(*serf.Query) uint64
}
// eventStream is used to stream events to a client over IPC
type eventStream struct {
client streamClient
eventCh chan serf.Event
filters []EventFilter
logger *log.Logger
seq uint64
}
func newEventStream(client streamClient, filters []EventFilter, seq uint64, logger *log.Logger) *eventStream {
es := &eventStream{
client: client,
eventCh: make(chan serf.Event, 512),
filters: filters,
logger: logger,
seq: seq,
}
go es.stream()
return es
}
func (es *eventStream) HandleEvent(e serf.Event) {
// Check the event
for _, f := range es.filters {
if f.Invoke(e) {
goto HANDLE
}
}
return
// Do a non-blocking send
HANDLE:
select {
case es.eventCh <- e:
default:
es.logger.Printf("[WARN] agent.ipc: Dropping event to %v", es.client)
}
}
func (es *eventStream) Stop() {
close(es.eventCh)
}
func (es *eventStream) stream() {
var err error
for event := range es.eventCh {
switch e := event.(type) {
case serf.MemberEvent:
err = es.sendMemberEvent(e)
case serf.UserEvent:
err = es.sendUserEvent(e)
case *serf.Query:
err = es.sendQuery(e)
default:
err = fmt.Errorf("Unknown event type: %s", event.EventType().String())
}
if err != nil {
es.logger.Printf("[ERR] agent.ipc: Failed to stream event to %v: %v",
es.client, err)
return
}
}
}
// sendMemberEvent is used to send a single member event
func (es *eventStream) sendMemberEvent(me serf.MemberEvent) error {
members := make([]Member, 0, len(me.Members))
for _, m := range me.Members {
sm := Member{
Name: m.Name,
Addr: m.Addr,
Port: m.Port,
Tags: m.Tags,
Status: m.Status.String(),
ProtocolMin: m.ProtocolMin,
ProtocolMax: m.ProtocolMax,
ProtocolCur: m.ProtocolCur,
DelegateMin: m.DelegateMin,
DelegateMax: m.DelegateMax,
DelegateCur: m.DelegateCur,
}
members = append(members, sm)
}
header := responseHeader{
Seq: es.seq,
Error: "",
}
rec := memberEventRecord{
Event: me.String(),
Members: members,
}
return es.client.Send(&header, &rec)
}
// sendUserEvent is used to send a single user event
func (es *eventStream) sendUserEvent(ue serf.UserEvent) error {
header := responseHeader{
Seq: es.seq,
Error: "",
}
rec := userEventRecord{
Event: ue.EventType().String(),
LTime: ue.LTime,
Name: ue.Name,
Payload: ue.Payload,
Coalesce: ue.Coalesce,
}
return es.client.Send(&header, &rec)
}
// sendQuery is used to send a single query event
func (es *eventStream) sendQuery(q *serf.Query) error {
id := es.client.RegisterQuery(q)
header := responseHeader{
Seq: es.seq,
Error: "",
}
rec := queryEventRecord{
Event: q.EventType().String(),
ID: id,
LTime: q.LTime,
Name: q.Name,
Payload: q.Payload,
}
return es.client.Send(&header, &rec)
}
|