-
Notifications
You must be signed in to change notification settings - Fork 5
Expand file tree
/
Copy pathvcenter.go
More file actions
177 lines (143 loc) · 4.51 KB
/
Copy pathvcenter.go
File metadata and controls
177 lines (143 loc) · 4.51 KB
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
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
///////////////////////////////////////////////////////////////////////
// Copyright (c) 2017 VMware, Inc. All Rights Reserved.
// SPDX-License-Identifier: Apache-2.0
///////////////////////////////////////////////////////////////////////
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"reflect"
"strings"
"time"
"github.com/pkg/errors"
"github.com/satori/go.uuid"
"github.com/vmware/govmomi"
"github.com/vmware/govmomi/event"
"github.com/vmware/govmomi/vim25/soap"
"github.com/vmware/govmomi/vim25/types"
"github.com/vmware/dispatch/pkg/events"
"github.com/vmware/dispatch/pkg/utils"
)
// NO TESTS
const eventTypeVersion = "0.1"
type vCenterEvent struct {
Metadata interface{} `json:"metadata"`
Time time.Time `json:"time"`
Category string `json:"category"`
Message string `json:"message"`
}
// newDriver creates a new vCenter event driver
func newDriver(vcenterURL string, insecure bool) (*vCenterDriver, error) {
vClient, err := newVCenterClient(context.Background(), vcenterURL, insecure)
if err != nil {
return nil, err
}
manager := event.NewManager(vClient.Client)
return &vCenterDriver{
vcenterURL: vcenterURL,
insecure: insecure,
manager: manager,
client: vClient,
}, nil
}
type vCenterDriver struct {
vcenterURL string
insecure bool
manager *event.Manager
client *govmomi.Client
done func()
}
func (d *vCenterDriver) consume(topics []string) (<-chan *events.CloudEvent, error) {
ctx, cancel := context.WithCancel(context.Background())
d.done = cancel
eventsChan := make(chan *events.CloudEvent)
go func() {
err := d.manager.Events(
ctx, // context
// TODO: add support for filter customization
[]types.ManagedObjectReference{d.client.ServiceContent.RootFolder}, // object(s) to monitor
10, // maximum number of events per page passed to handler
true, // poll for events indefinitely
true, // ignore limit of monitored objects (10)
d.handler(eventsChan, false), // handler executed for each event page
)
if err != nil {
log.Printf("Error when reading events from vCenter: %+v", err)
}
close(eventsChan)
}()
return eventsChan, nil
}
func (d *vCenterDriver) topics() []string {
// TODO: generate it based on API WSDL
return nil
}
func (d *vCenterDriver) close() error {
d.done()
return nil
}
func (d *vCenterDriver) handler(events chan *events.CloudEvent, multiple bool) func(types.ManagedObjectReference, []types.BaseEvent) error {
return func(obj types.ManagedObjectReference, page []types.BaseEvent) error {
event.Sort(page) // sort by event time
for _, e := range page {
processedEvent, err := d.processEvent(e)
if err != nil {
return errors.Wrap(err, "error processing event")
}
events <- processedEvent
}
return nil
}
}
func (d *vCenterDriver) processEvent(e types.BaseEvent) (*events.CloudEvent, error) {
eventType := reflect.TypeOf(e).Elem().Name()
cat, err := d.manager.EventCategory(context.Background(), e)
if err != nil {
return nil, errors.Wrap(err, "error retrieving event category")
}
ve := &vCenterEvent{
Time: e.GetEvent().CreatedTime,
Category: cat,
Message: strings.TrimSpace(e.GetEvent().FullFormattedMessage),
}
// if this is a TaskEvent gather a little more information
if t, ok := e.(*types.TaskEvent); ok {
// some tasks won't have this information, so just use the event message
if t.Info.Entity != nil {
ve.Message = fmt.Sprintf("%s (target=%s %s)", ve.Message, t.Info.Entity.Type, t.Info.EntityName)
}
}
ve.Metadata = processEventMetadata(e)
topic := convertToTopic(eventType)
return d.dispatchEvent(topic, ve)
}
func (d *vCenterDriver) dispatchEvent(topic string, ve *vCenterEvent) (*events.CloudEvent, error) {
encoded, err := json.Marshal(*ve)
if err != nil {
return nil, err
}
event := events.CloudEvent{
CloudEventsVersion: "0.1",
EventType: topic,
EventTypeVersion: eventTypeVersion,
Source: "vcenter1", // TODO: make this configurable
EventID: uuid.NewV4().String(),
EventTime: time.Time{},
ContentType: "application/json",
Data: encoded,
}
return &event, nil
}
func newVCenterClient(ctx context.Context, vcenterURL string, insecure bool) (*govmomi.Client, error) {
url, err := soap.ParseURL(vcenterURL)
if err != nil {
return nil, err
}
return govmomi.NewClient(ctx, url, insecure)
}
func convertToTopic(eventType string) string {
eventType = strings.Replace(eventType, "Event", "", -1)
return utils.CamelCaseToLowerSeparated(eventType, ".")
}