-
Notifications
You must be signed in to change notification settings - Fork 48
/
info.go
99 lines (81 loc) · 2.63 KB
/
info.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
91
92
93
94
95
96
97
98
99
package messages
import (
"encoding/json"
"time"
)
const (
JobRunningResolution = "JOB_RUNNING_RESOLUTION"
JobDetectedLag = "JOB_DETECTED_LAG"
JobInitialMaxDelay = "JOB_INITIAL_MAX_DELAY"
FindLimitedResultSet = "FIND_LIMITED_RESULT_SET"
)
type MessageBlock struct {
TimestampedMessage
Code string `json:"messageCode"`
Level string `json:"messageLevel"`
NumInputTimeseries int `json:"numInputTimeSeries"`
// If the messageCode field in the message is known, this will be an
// instance that has more specific methods to access the known fields. You
// can always access the original content by treating this value as a
// map[string]interface{}.
Contents interface{} `json:"-"`
ContentsRaw map[string]interface{} `json:"contents"`
}
type InfoMessage struct {
BaseJSONChannelMessage
LogicalTimestampMillis uint64 `json:"logicalTimestampMs"`
MessageBlock `json:"message"`
}
func (im *InfoMessage) UnmarshalJSON(raw []byte) error {
type IM InfoMessage
if err := json.Unmarshal(raw, (*IM)(im)); err != nil {
return err
}
mb := &im.MessageBlock
switch mb.Code {
case JobRunningResolution:
mb.Contents = JobRunningResolutionContents(mb.ContentsRaw)
case JobDetectedLag:
mb.Contents = JobDetectedLagContents(mb.ContentsRaw)
case JobInitialMaxDelay:
mb.Contents = JobInitialMaxDelayContents(mb.ContentsRaw)
case FindLimitedResultSet:
mb.Contents = FindLimitedResultSetContents(mb.ContentsRaw)
default:
mb.Contents = mb.ContentsRaw
}
return nil
}
func (im *InfoMessage) LogicalTimestamp() time.Time {
return time.Unix(0, int64(im.LogicalTimestampMillis*uint64(time.Millisecond)))
}
type JobRunningResolutionContents map[string]interface{}
func (jm JobRunningResolutionContents) ResolutionMS() int {
field, _ := jm["resolutionMs"].(float64)
return int(field)
}
type JobDetectedLagContents map[string]interface{}
func (jm JobDetectedLagContents) LagMS() int {
field, _ := jm["lagMs"].(float64)
return int(field)
}
type JobInitialMaxDelayContents map[string]interface{}
func (jm JobInitialMaxDelayContents) MaxDelayMS() int {
field, _ := jm["maxDelayMs"].(float64)
return int(field)
}
type FindLimitedResultSetContents map[string]interface{}
func (jm FindLimitedResultSetContents) MatchedSize() int {
field, _ := jm["matchedSize"].(float64)
return int(field)
}
func (jm FindLimitedResultSetContents) LimitSize() int {
field, _ := jm["limitSize"].(float64)
return int(field)
}
// ExpiredTSIDMessage is received when a timeseries has expired and is no
// longer relvant to a computation.
type ExpiredTSIDMessage struct {
BaseJSONChannelMessage
TSID string `json:"tsId"`
}