From 1876d69210515db6ddc16b7aba2201b34b7ff9b8 Mon Sep 17 00:00:00 2001 From: Nathaniel Allred Date: Tue, 14 Jun 2022 20:07:33 -0500 Subject: [PATCH] add file storage telemetry (#20408) * cloud instances send file storage telemetry * Add core disabled plugins (`disabled_default_plugins`) per analytics team request --- services/telemetry/telemetry.go | 24 +- services/telemetry/telemetry_test.go | 324 +++++++++++++++++---------- 2 files changed, 222 insertions(+), 126 deletions(-) diff --git a/services/telemetry/telemetry.go b/services/telemetry/telemetry.go index 7088abf004..6ba851aac1 100644 --- a/services/telemetry/telemetry.go +++ b/services/telemetry/telemetry.go @@ -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, diff --git a/services/telemetry/telemetry_test.go b/services/telemetry/telemetry_test.go index 97947e14f1..1aa2ed5ef0 100644 --- a/services/telemetry/telemetry_test.go +++ b/services/telemetry/telemetry_test.go @@ -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"}))