* cloud instances send file storage telemetry
* Add core disabled plugins (`disabled_default_plugins`) per analytics team request
Этот коммит содержится в:
Nathaniel Allred
2022-06-14 20:07:33 -05:00
коммит произвёл GitHub
родитель 906fab07bd
Коммит 1876d69210
2 изменённых файлов: 222 добавлений и 126 удалений

Просмотреть файл

@@ -331,7 +331,7 @@ func (ts *TelemetryService) trackActivity() {
activeUsersMonthlyCount = r.Data.(int64)
}
ts.SendTelemetry(TrackActivity, map[string]interface{}{
activity := map[string]interface{}{
"registered_users": userCount,
"bot_accounts": botAccountsCount,
"guest_accounts": guestAccountsCount,
@@ -350,7 +350,17 @@ func (ts *TelemetryService) trackActivity() {
"slash_commands": slashCommandsCount,
"incoming_webhooks": incomingWebhooksCount,
"outgoing_webhooks": outgoingWebhooksCount,
})
}
if license := ts.srv.License(); license != nil && license.Features.Cloud != nil && *license.Features.Cloud {
var tmpStorage int64
if usage, err := ts.dbStore.FileInfo().GetStorageUsage(true, false); err == nil {
tmpStorage = usage
}
activity["storage_bytes"] = utils.RoundOffToZeroes(float64(tmpStorage))
}
ts.SendTelemetry(TrackActivity, activity)
}
func (ts *TelemetryService) trackConfig() {
@@ -855,6 +865,7 @@ func (ts *TelemetryService) trackPlugins() {
webappEnabledCount := 0
backendEnabledCount := 0
totalDisabledCount := 0
totalCoreDisabledCount := 0
webappDisabledCount := 0
backendDisabledCount := 0
brokenManifestCount := 0
@@ -886,14 +897,18 @@ func (ts *TelemetryService) trackPlugins() {
if plugin.Manifest.HasWebapp() {
webappDisabledCount += 1
}
if _, isCorePlugin := model.InstalledIntegrationsIgnoredPlugins[plugin.Manifest.Id]; isCorePlugin {
totalCoreDisabledCount += 1
}
}
if plugin.Manifest.SettingsSchema != nil {
settingsCount += 1
}
}
} else {
totalEnabledCount = -1 // -1 to indicate disabled or error
totalDisabledCount = -1 // -1 to indicate disabled or error
totalEnabledCount = -1 // -1 to indicate disabled or error
totalCoreDisabledCount = -1 // -1 to indicate disabled or error
totalDisabledCount = -1 // -1 to indicate disabled or error
}
ts.SendTelemetry(TrackPlugins, map[string]interface{}{
@@ -901,6 +916,7 @@ func (ts *TelemetryService) trackPlugins() {
"enabled_webapp_plugins": webappEnabledCount,
"enabled_backend_plugins": backendEnabledCount,
"disabled_plugins": totalDisabledCount,
"disabled_default_plugins": totalCoreDisabledCount,
"disabled_webapp_plugins": webappDisabledCount,
"disabled_backend_plugins": backendDisabledCount,
"plugins_with_settings": settingsCount,

Просмотреть файл

@@ -8,6 +8,7 @@ import (
"crypto/ecdsa"
"encoding/json"
"errors"
"fmt"
"io/ioutil"
"net/http"
"net/http/httptest"
@@ -35,12 +36,117 @@ type FakeConfigService struct {
cfg *model.Config
}
type testTelemetryPayload struct {
MessageId string
SentAt time.Time
Batch []struct {
MessageId string
UserId string
Event string
Timestamp time.Time
Properties map[string]interface{}
}
Context struct {
Library struct {
Name string
Version string
}
}
}
type testBatch struct {
MessageId string
UserId string
Event string
Timestamp time.Time
Properties map[string]interface{}
}
func assertPayload(t *testing.T, actual testTelemetryPayload, event string, properties map[string]interface{}) {
t.Helper()
assert.NotEmpty(t, actual.MessageId)
assert.False(t, actual.SentAt.IsZero())
if assert.Len(t, actual.Batch, 1) {
assert.NotEmpty(t, actual.Batch[0].MessageId, "message id should not be empty")
assert.Equal(t, testTelemetryID, actual.Batch[0].UserId)
if event != "" {
assert.Equal(t, event, actual.Batch[0].Event)
}
assert.False(t, actual.Batch[0].Timestamp.IsZero(), "batch timestamp should not be the zero value")
if properties != nil {
assert.Equal(t, properties, actual.Batch[0].Properties)
}
}
assert.Equal(t, "analytics-go", actual.Context.Library.Name)
assert.Equal(t, "3.3.0", actual.Context.Library.Version)
}
func collectBatches(t *testing.T, info *[]testBatch, pchan chan testTelemetryPayload) {
t.Helper()
for {
select {
case result := <-pchan:
assertPayload(t, result, "", nil)
*info = append(*info, result.Batch[0])
case <-time.After(time.Second * 1):
return
}
}
}
func makeTelemetryServiceAndReceiver(t *testing.T, cloudLicense bool) (*TelemetryService, chan testTelemetryPayload, *model.Config, func()) {
cfg := &model.Config{}
cfg.SetDefaults()
serverIfaceMock, storeMock, deferredAssertions, cleanUp := initializeMocks(cfg, cloudLicense)
testLogger, _ := mlog.NewLogger()
logCfg, _ := config.MloggerConfigFromLoggerConfig(&cfg.LogSettings, nil, config.GetLogFileLocation)
if errCfg := testLogger.ConfigureTargets(logCfg, nil); errCfg != nil {
panic("failed to configure test logger: " + errCfg.Error())
}
pchan := make(chan testTelemetryPayload, 100)
receiver := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, err := ioutil.ReadAll(r.Body)
require.NoError(t, err)
var p testTelemetryPayload
err = json.Unmarshal(body, &p)
require.NoError(t, err)
pchan <- p
}))
service := New(serverIfaceMock, storeMock, searchengine.NewBroker(cfg), testLogger)
service.TelemetryID = testTelemetryID
service.rudderClient = nil
service.initRudder(receiver.URL, RudderKey)
// initializing rudder send a client identify message
select {
case identifyMessage := <-pchan:
assertPayload(t, identifyMessage, "", nil)
case <-time.After(time.Second * 1):
require.Fail(t, "Did not receive ID message")
}
return service, pchan, cfg, func() {
receiver.Close()
testLogger.Shutdown()
cleanUp()
deferredAssertions(t)
}
}
const testTelemetryID = "test-telemetry-id-12345"
func (fcs *FakeConfigService) Config() *model.Config { return fcs.cfg }
func (fcs *FakeConfigService) AddConfigListener(f func(old, current *model.Config)) string { return "" }
func (fcs *FakeConfigService) RemoveConfigListener(key string) {}
func (fcs *FakeConfigService) AsymmetricSigningKey() *ecdsa.PrivateKey { return nil }
func initializeMocks(cfg *model.Config) (*mocks.ServerIface, *storeMocks.Store, func(t *testing.T), func()) {
func initializeMocks(cfg *model.Config, cloudLicense bool) (*mocks.ServerIface, *storeMocks.Store, func(t *testing.T), func()) {
serverIfaceMock := &mocks.ServerIface{}
logger, _ := mlog.NewLogger()
@@ -63,7 +169,11 @@ func initializeMocks(cfg *model.Config) (*mocks.ServerIface, *storeMocks.Store,
nil)
serverIfaceMock.On("GetPluginsEnvironment").Return(pluginEnv, nil)
serverIfaceMock.On("License").Return(model.NewTestLicense(), nil)
if cloudLicense {
serverIfaceMock.On("License").Return(model.NewTestLicense("cloud"), nil)
} else {
serverIfaceMock.On("License").Return(model.NewTestLicense(), nil)
}
serverIfaceMock.On("GetRoleByName", context.Background(), "system_admin").Return(&model.Role{Permissions: []string{"sa-test1", "sa-test2"}}, nil)
serverIfaceMock.On("GetRoleByName", context.Background(), "system_user").Return(&model.Role{Permissions: []string{"su-test1", "su-test2"}}, nil)
serverIfaceMock.On("GetRoleByName", context.Background(), "system_user_manager").Return(&model.Role{Permissions: []string{"sum-test1", "sum-test2"}}, nil)
@@ -265,6 +375,8 @@ func TestPluginActivated(t *testing.T) {
assert.False(t, pluginActivated(states, "none"))
}
const keyStorageBytes = "storage_bytes"
func TestPluginVersion(t *testing.T) {
plugins := []*model.BundleInfo{
{
@@ -290,44 +402,8 @@ func TestRudderTelemetry(t *testing.T) {
t.SkipNow()
}
type batch struct {
MessageId string
UserId string
Event string
Timestamp time.Time
Properties map[string]interface{}
}
type payload struct {
MessageId string
SentAt time.Time
Batch []struct {
MessageId string
UserId string
Event string
Timestamp time.Time
Properties map[string]interface{}
}
Context struct {
Library struct {
Name string
Version string
}
}
}
data := make(chan payload, 100)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, err := ioutil.ReadAll(r.Body)
require.NoError(t, err)
var p payload
err = json.Unmarshal(body, &p)
require.NoError(t, err)
data <- p
}))
defer server.Close()
service, pchan, cfg, teardown := makeTelemetryServiceAndReceiver(t, false)
defer teardown()
marketplaceServer := httptest.NewServer(http.HandlerFunc(func(res http.ResponseWriter, req *http.Request) {
res.WriteHeader(http.StatusOK)
@@ -342,52 +418,13 @@ func TestRudderTelemetry(t *testing.T) {
res.Write(json)
}))
defer func() { marketplaceServer.Close() }()
telemetryID := "test-telemetry-id-12345"
cfg := &model.Config{}
cfg.SetDefaults()
serverIfaceMock, storeMock, deferredAssertions, cleanUp := initializeMocks(cfg)
defer cleanUp()
defer deferredAssertions(t)
testLogger, _ := mlog.NewLogger()
logCfg, _ := config.MloggerConfigFromLoggerConfig(&cfg.LogSettings, nil, config.GetLogFileLocation)
if errCfg := testLogger.ConfigureTargets(logCfg, nil); errCfg != nil {
panic("failed to configure test logger: " + errCfg.Error())
}
defer testLogger.Shutdown()
telemetryService := New(serverIfaceMock, storeMock, searchengine.NewBroker(cfg), testLogger)
telemetryService.TelemetryID = telemetryID
telemetryService.rudderClient = nil
telemetryService.initRudder(server.URL, RudderKey)
assertPayload := func(t *testing.T, actual payload, event string, properties map[string]interface{}) {
t.Helper()
assert.NotEmpty(t, actual.MessageId)
assert.False(t, actual.SentAt.IsZero())
if assert.Len(t, actual.Batch, 1) {
assert.NotEmpty(t, actual.Batch[0].MessageId, "message id should not be empty")
assert.Equal(t, telemetryID, actual.Batch[0].UserId)
if event != "" {
assert.Equal(t, event, actual.Batch[0].Event)
}
assert.False(t, actual.Batch[0].Timestamp.IsZero(), "batch timestamp should not be the zero value")
if properties != nil {
assert.Equal(t, properties, actual.Batch[0].Properties)
}
}
assert.Equal(t, "analytics-go", actual.Context.Library.Name)
assert.Equal(t, "3.3.0", actual.Context.Library.Version)
}
defer marketplaceServer.Close()
collectInfo := func(info *[]string) {
t.Helper()
for {
select {
case result := <-data:
case result := <-pchan:
assertPayload(t, result, "", nil)
*info = append(*info, result.Batch[0].Event)
case <-time.After(time.Second * 1):
@@ -396,34 +433,13 @@ func TestRudderTelemetry(t *testing.T) {
}
}
collectBatches := func(info *[]batch) {
t.Helper()
for {
select {
case result := <-data:
assertPayload(t, result, "", nil)
*info = append(*info, result.Batch[0])
case <-time.After(time.Second * 1):
return
}
}
}
// Should send a client identify message
select {
case identifyMessage := <-data:
assertPayload(t, identifyMessage, "", nil)
case <-time.After(time.Second * 1):
require.Fail(t, "Did not receive ID message")
}
t.Run("Send", func(t *testing.T) {
testValue := "test-send-value-6789"
telemetryService.SendTelemetry("Testing Telemetry", map[string]interface{}{
service.SendTelemetry("Testing Telemetry", map[string]interface{}{
"hey": testValue,
})
select {
case result := <-data:
case result := <-pchan:
assertPayload(t, result, "Testing Telemetry", map[string]interface{}{
"hey": testValue,
})
@@ -434,7 +450,7 @@ func TestRudderTelemetry(t *testing.T) {
// Plugins remain disabled at this point
t.Run("SendDailyTelemetryPluginsDisabled", func(t *testing.T) {
telemetryService.sendDailyTelemetry(true)
service.sendDailyTelemetry(true)
var info []string
// Collect the info sent.
@@ -477,7 +493,7 @@ func TestRudderTelemetry(t *testing.T) {
// th.Server.UpdateConfig(func(cfg *model.Config) { *cfg.PluginSettings.Enable = true })
t.Run("SendDailyTelemetry", func(t *testing.T) {
telemetryService.sendDailyTelemetry(true)
service.sendDailyTelemetry(true)
var info []string
// Collect the info sent.
@@ -516,10 +532,10 @@ func TestRudderTelemetry(t *testing.T) {
}
})
t.Run("Telemetry for Marketplace plugins is returned", func(t *testing.T) {
telemetryService.trackPluginConfig(telemetryService.srv.Config(), marketplaceServer.URL)
service.trackPluginConfig(service.srv.Config(), marketplaceServer.URL)
var batches []batch
collectBatches(&batches)
var batches []testBatch
collectBatches(t, &batches, pchan)
for _, b := range batches {
if b.Event == TrackConfigPlugin {
@@ -534,10 +550,10 @@ func TestRudderTelemetry(t *testing.T) {
})
t.Run("Telemetry for known plugins is returned, if request to Marketplace fails", func(t *testing.T) {
telemetryService.trackPluginConfig(telemetryService.srv.Config(), "http://some.random.invalid.url")
service.trackPluginConfig(service.srv.Config(), "http://some.random.invalid.url")
var batches []batch
collectBatches(&batches)
var batches []testBatch
collectBatches(t, &batches, pchan)
for _, b := range batches {
if b.Event == TrackConfigPlugin {
@@ -555,16 +571,41 @@ func TestRudderTelemetry(t *testing.T) {
if !strings.Contains(RudderKey, "placeholder") {
t.Skipf("Skipping telemetry on production builds")
}
telemetryService.sendDailyTelemetry(false)
service.sendDailyTelemetry(false)
select {
case <-data:
case <-pchan:
require.Fail(t, "Should not send telemetry when the rudder key is not set")
case <-time.After(time.Second * 1):
// Did not receive telemetry
}
})
t.Run("SendDailyTelemetryNonCloud", func(t *testing.T) {
if !strings.Contains(RudderKey, "placeholder") {
t.Skipf("Skipping telemetry on production builds")
}
service.sendDailyTelemetry(true)
var batches []testBatch
collectBatches(t, &batches, pchan)
var activityEvent testBatch
var found bool
for _, testBatch := range batches {
if testBatch.Event == TrackActivity {
activityEvent = testBatch
found = true
break
}
}
require.True(t, found, fmt.Sprintf("Expected to receive %q event, but received %q", TrackActivity, activityEvent.Event))
_, ok := activityEvent.Properties[keyStorageBytes]
require.False(t, ok, fmt.Sprintf("Expected non-cloud payload not to contain %q, got %+v", keyStorageBytes, activityEvent.Properties))
})
t.Run("SendDailyTelemetryDisabled", func(t *testing.T) {
if !strings.Contains(RudderKey, "placeholder") {
t.Skipf("Skipping telemetry on production builds")
@@ -574,10 +615,10 @@ func TestRudderTelemetry(t *testing.T) {
*cfg.LogSettings.EnableDiagnostics = true
}()
telemetryService.sendDailyTelemetry(true)
service.sendDailyTelemetry(true)
select {
case <-data:
case <-pchan:
require.Fail(t, "Should not send telemetry when they are disabled")
case <-time.After(time.Second * 1):
// Did not receive telemetry
@@ -586,10 +627,10 @@ func TestRudderTelemetry(t *testing.T) {
t.Run("TestInstallationType", func(t *testing.T) {
os.Unsetenv(EnvVarInstallType)
telemetryService.sendDailyTelemetry(true)
service.sendDailyTelemetry(true)
var batches []batch
collectBatches(&batches)
var batches []testBatch
collectBatches(t, &batches, pchan)
for _, b := range batches {
if b.Event == TrackServer {
@@ -600,8 +641,8 @@ func TestRudderTelemetry(t *testing.T) {
os.Setenv(EnvVarInstallType, "docker")
defer os.Unsetenv(EnvVarInstallType)
batches = []batch{}
collectBatches(&batches)
batches = []testBatch{}
collectBatches(t, &batches, pchan)
for _, b := range batches {
if b.Event == TrackServer {
@@ -619,13 +660,52 @@ func TestRudderTelemetry(t *testing.T) {
defer os.Unsetenv("RudderKey")
defer os.Unsetenv("RudderDataplaneURL")
config := telemetryService.getRudderConfig()
config := service.getRudderConfig()
assert.Equal(t, "arudderstackplace", config.DataplaneURL)
assert.Equal(t, "abc123", config.RudderKey)
})
}
func TestRudderTelemetryCloud(t *testing.T) {
if !strings.Contains(RudderKey, "placeholder") {
t.Skipf("Skipping telemetry on production builds")
}
if testing.Short() {
t.SkipNow()
}
service, pchan, _, teardown := makeTelemetryServiceAndReceiver(t, true)
defer teardown()
fileInfoStore := storeMocks.FileInfoStore{}
mockBytes := int64(1000000000)
fileInfoStore.On("GetStorageUsage", true, false).Return(mockBytes, nil)
defer fileInfoStore.AssertExpectations(t)
service.dbStore.(*storeMocks.Store).On("FileInfo").Return(&fileInfoStore)
service.sendDailyTelemetry(true)
var batches []testBatch
collectBatches(t, &batches, pchan)
var activityEvent testBatch
var found bool
for _, batch := range batches {
if batch.Event == TrackActivity {
activityEvent = batch
found = true
break
}
}
require.True(t, found, fmt.Sprintf("Expected to receive %q event, but received %q: %+v", TrackActivity, activityEvent.Event, activityEvent))
storageBytes, ok := activityEvent.Properties[keyStorageBytes]
require.True(t, ok, fmt.Sprintf("Expected payload to contain %q", keyStorageBytes))
require.Equal(t, mockBytes, int64(storageBytes.(float64)), fmt.Sprintf("Expected storage usage of %d bytes", mockBytes))
}
func TestIsDefaultArray(t *testing.T) {
assert.True(t, isDefaultArray([]string{"one", "two"}, []string{"one", "two"}))
assert.False(t, isDefaultArray([]string{"one", "two"}, []string{"one", "two", "three"}))