diff --git a/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamicConfiguration.java b/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamicConfiguration.java new file mode 100644 index 000000000..33892d7ee --- /dev/null +++ b/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamicConfiguration.java @@ -0,0 +1,68 @@ +package com.hello.suripu.core.configuration; + +import com.hello.suripu.core.db.ConfigurationDAODynamoDB; +import com.yammer.dropwizard.config.Configuration; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Created by jnorgan on 6/15/15. + */ +public class DynamicConfiguration { + + private static final Logger LOGGER = LoggerFactory.getLogger(DynamicConfiguration.class); + private ScheduledFuture scheduledFuture; + private final Integer pollingIntervalInSeconds; + private final ConfigurationDAODynamoDB configDAO; + private Configuration configuration; + + final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(5); + + public DynamicConfiguration(ConfigurationDAODynamoDB configDAO, final Integer pollingIntervalInSeconds) { + this.pollingIntervalInSeconds = pollingIntervalInSeconds; + this.configDAO = configDAO; + start(); + } + + private void startPolling() { + scheduledFuture = executorService.scheduleAtFixedRate(new Runnable() { + @Override + public void run() { + configuration = getData(); + } + } , pollingIntervalInSeconds, pollingIntervalInSeconds, TimeUnit.SECONDS); + } + + public void start() { + LOGGER.info("Starting polling config: {}", configDAO.configuration.getClass().getSimpleName()); + configuration = getData(); + startPolling(); + } + + public void stop() { + scheduledFuture.cancel(true); + LOGGER.info("Stopped polling config: {}", configDAO.configuration.getClass().getSimpleName()); + executorService.shutdown(); + LOGGER.info("ThreadPool shutdown"); + } + + public Configuration getConfiguration() { + return configuration; + } + + synchronized private Configuration getData() { + LOGGER.debug("Polling dynamic config: {}", configDAO.configuration.getClass().getSimpleName()); + + try { + final Configuration config = configDAO.getData(); + return config; + } catch (Exception e) { + LOGGER.error("Dynamic Config DAO method 'getData()' failed."); + } + return configuration; + } +} diff --git a/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamoDBTableName.java b/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamoDBTableName.java index 3a43bcd3c..4abd3fae0 100644 --- a/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamoDBTableName.java +++ b/suripu-core/src/main/java/com/hello/suripu/core/configuration/DynamoDBTableName.java @@ -29,7 +29,9 @@ public enum DynamoDBTableName { ALARM_LOG("alarm_log"), SMART_ALARM_LOG("smart_alarm_log"), BAYESNET_PRIORS("hmm_bayesnet_priors"), - BAYESNET_MODEL("hmm_bayesnet_models"); + BAYESNET_MODEL("hmm_bayesnet_models"), + CONFIGURATIONS("configurations"); + private String value; private DynamoDBTableName(String value) { diff --git a/suripu-core/src/main/java/com/hello/suripu/core/db/ConfigurationDAODynamoDB.java b/suripu-core/src/main/java/com/hello/suripu/core/db/ConfigurationDAODynamoDB.java new file mode 100644 index 000000000..84ba777dd --- /dev/null +++ b/suripu-core/src/main/java/com/hello/suripu/core/db/ConfigurationDAODynamoDB.java @@ -0,0 +1,119 @@ +package com.hello.suripu.core.db; + +import com.amazonaws.services.dynamodbv2.AmazonDynamoDB; +import com.amazonaws.services.dynamodbv2.AmazonDynamoDBClient; +import com.amazonaws.services.dynamodbv2.document.DynamoDB; +import com.amazonaws.services.dynamodbv2.document.Item; +import com.amazonaws.services.dynamodbv2.document.PutItemOutcome; +import com.amazonaws.services.dynamodbv2.document.Table; +import com.amazonaws.services.dynamodbv2.document.spec.GetItemSpec; +import com.amazonaws.services.dynamodbv2.model.AttributeDefinition; +import com.amazonaws.services.dynamodbv2.model.CreateTableRequest; +import com.amazonaws.services.dynamodbv2.model.CreateTableResult; +import com.amazonaws.services.dynamodbv2.model.KeySchemaElement; +import com.amazonaws.services.dynamodbv2.model.KeyType; +import com.amazonaws.services.dynamodbv2.model.ProvisionedThroughput; +import com.amazonaws.services.dynamodbv2.model.ScalarAttributeType; +import com.fasterxml.jackson.core.JsonGenerationException; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.ObjectReader; +import com.fasterxml.jackson.databind.ObjectWriter; +import com.fasterxml.jackson.databind.ser.FilterProvider; +import com.fasterxml.jackson.databind.ser.impl.SimpleBeanPropertyFilter; +import com.fasterxml.jackson.databind.ser.impl.SimpleFilterProvider; +import com.yammer.dropwizard.config.Configuration; +import java.io.IOException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Created by jnorgan on 6/15/15. + */ +public class ConfigurationDAODynamoDB { + private static final Logger LOGGER = LoggerFactory.getLogger(ConfigurationDAODynamoDB.class); + private static final String NAME_ATTRIBUTE_NAME = "name"; + private static final String NAMESPACE_ATTRIBUTE_NAME = "ns"; + private static final String VALUE_ATTRIBUTE_NAME = "config_values"; + private static final String TIMESTAMP_ATTRIBUTE_NAME = "timestamp"; + public final static String DATETIME_FORMAT = "yyyy-MM-dd HH:mm:ssZ"; + + public final Configuration configuration; + private final String namespace; + private final Table table; + + public ConfigurationDAODynamoDB(final Configuration configuration, final AmazonDynamoDB amazonDynamoDB, final String tableName, final String namespace) { + this.configuration = configuration; + this.table = new DynamoDB(amazonDynamoDB).getTable(tableName); + this.namespace = namespace; + } + + public Configuration getData() { + LOGGER.trace("Calling getData"); + + final GetItemSpec spec = new GetItemSpec() + .withPrimaryKey("ns", namespace, "name", configuration.getClass().getSimpleName()); + spec.withProjectionExpression(VALUE_ATTRIBUTE_NAME); + + final ObjectMapper mapper = new ObjectMapper(); + try + { + final String storedConfigString = table.getItem(spec).getJSON(VALUE_ATTRIBUTE_NAME); + final ObjectReader reader = mapper.reader(configuration.getClass()); + return reader.readValue(storedConfigString); + } catch (JsonGenerationException e) { + LOGGER.error("Failed to get data."); + } catch (IOException ioe) { + LOGGER.error("Failed to get data."); + } catch (NullPointerException npe) { + LOGGER.error("getData() failed. NPE."); + } + + //No stored config exists, persist the loaded config to DynamoDB + LOGGER.info("Configuration read from Dynamo failed. Attempting to persist loaded configuration values."); + + final String[] ignorableFieldNames = { "httpConfiguration", "loggingConfiguration", "http", "logging" }; + FilterProvider filters = new SimpleFilterProvider() + .addFilter("filter properties by name", + SimpleBeanPropertyFilter.serializeAllExcept( + ignorableFieldNames)); + final ObjectWriter objWriter = mapper.writer(filters); + try { + final String configJSON = objWriter.writeValueAsString(configuration); + put(configJSON); + } catch (JsonProcessingException jpe) { + LOGGER.error("Failed to serialize configuration."); + } + + return configuration; + } + + public void put(final String jsonString) { + Item item = new Item() + .withPrimaryKey("ns", namespace, "name", configuration.getClass().getSimpleName()) + .withJSON(VALUE_ATTRIBUTE_NAME, jsonString); + + PutItemOutcome outcome = table.putItem(item); + } + + public static CreateTableResult createTable(final String tableName, final AmazonDynamoDBClient dynamoDBClient){ + final CreateTableRequest request = new CreateTableRequest().withTableName(tableName); + + request.withKeySchema( + new KeySchemaElement().withAttributeName(NAMESPACE_ATTRIBUTE_NAME).withKeyType(KeyType.HASH), + new KeySchemaElement().withAttributeName(NAME_ATTRIBUTE_NAME).withKeyType(KeyType.RANGE) + ); + + request.withAttributeDefinitions( + new AttributeDefinition().withAttributeName(NAMESPACE_ATTRIBUTE_NAME).withAttributeType(ScalarAttributeType.S), + new AttributeDefinition().withAttributeName(NAME_ATTRIBUTE_NAME).withAttributeType(ScalarAttributeType.S) + ); + + request.setProvisionedThroughput(new ProvisionedThroughput() + .withReadCapacityUnits(1L) + .withWriteCapacityUnits(1L)); + + return dynamoDBClient.createTable(request); + } +} + diff --git a/suripu-service/src/main/java/com/hello/suripu/service/SuripuService.java b/suripu-service/src/main/java/com/hello/suripu/service/SuripuService.java index 56ade8d59..ba49173c3 100644 --- a/suripu-service/src/main/java/com/hello/suripu/service/SuripuService.java +++ b/suripu-service/src/main/java/com/hello/suripu/service/SuripuService.java @@ -14,10 +14,12 @@ import com.hello.dropwizard.mikkusu.resources.VersionResource; import com.hello.suripu.core.ObjectGraphRoot; import com.hello.suripu.core.clients.AmazonDynamoDBClientFactory; +import com.hello.suripu.core.configuration.DynamicConfiguration; import com.hello.suripu.core.configuration.DynamoDBTableName; import com.hello.suripu.core.configuration.QueueName; import com.hello.suripu.core.db.AccessTokenDAO; import com.hello.suripu.core.db.ApplicationsDAO; +import com.hello.suripu.core.db.ConfigurationDAODynamoDB; import com.hello.suripu.core.db.DeviceDAO; import com.hello.suripu.core.db.FeatureStore; import com.hello.suripu.core.db.FirmwareUpgradePathDAO; @@ -213,6 +215,11 @@ public String getAWSSecretKey() { final AmazonDynamoDB featuresDynamoDBClient = dynamoDBFactory.getForTable(DynamoDBTableName.FEATURES); final FeatureStore featureStore = new FeatureStore(featuresDynamoDBClient, tableNames.get(DynamoDBTableName.FEATURES), namespace); + + final AmazonDynamoDB configDynamoDBClient = dynamoDBFactory.getForTable(DynamoDBTableName.CONFIGURATIONS); + final ConfigurationDAODynamoDB senseUploadConfigDAO = new ConfigurationDAODynamoDB(configuration.getSenseUploadConfiguration(), configDynamoDBClient, tableNames.get(DynamoDBTableName.CONFIGURATIONS), "service_" + namespace); + final DynamicConfiguration senseUploadDynamicConfig = new DynamicConfiguration(senseUploadConfigDAO, 30); + final RolloutModule module = new RolloutModule(featureStore, 30); ObjectGraphRoot.getInstance().init(module); @@ -224,7 +231,7 @@ public String getAWSSecretKey() { configuration.getDebug(), firmwareUpdateStore, groupFlipper, - configuration.getSenseUploadConfiguration(), + senseUploadDynamicConfig, configuration.getOTAConfiguration(), respCommandsDAODynamoDB, configuration.getRingDuration() diff --git a/suripu-service/src/main/java/com/hello/suripu/service/cli/CreateDynamoDBTables.java b/suripu-service/src/main/java/com/hello/suripu/service/cli/CreateDynamoDBTables.java index 2fc1f12a3..a06900417 100644 --- a/suripu-service/src/main/java/com/hello/suripu/service/cli/CreateDynamoDBTables.java +++ b/suripu-service/src/main/java/com/hello/suripu/service/cli/CreateDynamoDBTables.java @@ -9,6 +9,7 @@ import com.google.common.collect.ImmutableMap; import com.hello.suripu.core.configuration.DynamoDBTableName; import com.hello.suripu.core.configuration.NewDynamoDBConfiguration; +import com.hello.suripu.core.db.ConfigurationDAODynamoDB; import com.hello.suripu.core.db.RingTimeHistoryDAODynamoDB; import com.hello.suripu.service.configuration.SuripuConfiguration; import com.yammer.dropwizard.cli.ConfiguredCommand; @@ -26,6 +27,7 @@ protected void run(Bootstrap bootstrap, Namespace namespace final AWSCredentialsProvider awsCredentialsProvider = new DefaultAWSCredentialsProviderChain(); createRingTimeHistoryTable(configuration, awsCredentialsProvider); + createServiceConfigurationTable(configuration, awsCredentialsProvider); } @@ -49,4 +51,25 @@ private void createRingTimeHistoryTable(final SuripuConfiguration configuration, } } + private void createServiceConfigurationTable(final SuripuConfiguration configuration, final AWSCredentialsProvider awsCredentialsProvider){ + final NewDynamoDBConfiguration config = configuration.dynamoDBConfiguration(); + final AmazonDynamoDBClient client = new AmazonDynamoDBClient(awsCredentialsProvider); + final ImmutableMap tableNames = configuration.dynamoDBConfiguration().tables(); + final ImmutableMap endpoints = configuration.dynamoDBConfiguration().endpoints(); + + final String tableName = tableNames.get(DynamoDBTableName.CONFIGURATIONS); + final String endpoint = endpoints.get(DynamoDBTableName.CONFIGURATIONS); + + client.setEndpoint(endpoint); + try { + client.describeTable(tableName); + System.out.println(String.format("%s already exists.", tableName)); + + } catch (AmazonServiceException exception) { + final CreateTableResult result = ConfigurationDAODynamoDB.createTable(tableName, client); + final TableDescription description = result.getTableDescription(); + System.out.println(tableName + ": " + description.getTableStatus()); + } + } + } diff --git a/suripu-service/src/main/java/com/hello/suripu/service/configuration/SenseUploadConfiguration.java b/suripu-service/src/main/java/com/hello/suripu/service/configuration/SenseUploadConfiguration.java index d535da408..f2d1edbd1 100644 --- a/suripu-service/src/main/java/com/hello/suripu/service/configuration/SenseUploadConfiguration.java +++ b/suripu-service/src/main/java/com/hello/suripu/service/configuration/SenseUploadConfiguration.java @@ -1,6 +1,7 @@ package com.hello.suripu.service.configuration; +import com.fasterxml.jackson.annotation.JsonFilter; import com.fasterxml.jackson.annotation.JsonProperty; import com.yammer.dropwizard.config.Configuration; @@ -8,6 +9,7 @@ import javax.validation.constraints.Max; import javax.validation.constraints.Min; +@JsonFilter("filter properties by name") public class SenseUploadConfiguration extends Configuration { private static final Integer DEFAULT_NON_PEAK_HOUR_LOWER_BOUND = 11; // non peak periods start at 11:00:00 private static final Integer DEFAULT_NON_PEAK_HOUR_UPPER_BOUND = 22; // non peak periods end at 22:59:59 diff --git a/suripu-service/src/main/java/com/hello/suripu/service/resources/ReceiveResource.java b/suripu-service/src/main/java/com/hello/suripu/service/resources/ReceiveResource.java index 90163ea86..e5d2fda98 100644 --- a/suripu-service/src/main/java/com/hello/suripu/service/resources/ReceiveResource.java +++ b/suripu-service/src/main/java/com/hello/suripu/service/resources/ReceiveResource.java @@ -7,6 +7,7 @@ import com.hello.suripu.api.ble.SenseCommandProtos; import com.hello.suripu.api.input.DataInputProtos; import com.hello.suripu.api.output.OutputProtos; +import com.hello.suripu.core.configuration.DynamicConfiguration; import com.hello.suripu.core.configuration.QueueName; import com.hello.suripu.core.db.KeyStore; import com.hello.suripu.core.db.KeyStoreDynamoDB; @@ -82,7 +83,7 @@ public class ReceiveResource extends BaseResource { private final FirmwareUpdateStore firmwareUpdateStore; private final GroupFlipper groupFlipper; - private final SenseUploadConfiguration senseUploadConfiguration; + private final DynamicConfiguration senseUploadDynamicConfiguration; private final OTAConfiguration otaConfiguration; private final ResponseCommandsDAODynamoDB responseCommandsDAODynamoDB; @@ -99,7 +100,7 @@ public ReceiveResource(final KeyStore keyStore, final Boolean debug, final FirmwareUpdateStore firmwareUpdateStore, final GroupFlipper groupFlipper, - final SenseUploadConfiguration senseUploadConfiguration, + final DynamicConfiguration senseUploadDynamicConfiguration, final OTAConfiguration otaConfiguration, final ResponseCommandsDAODynamoDB responseCommandsDAODynamoDB, final int ringDurationSec) { @@ -114,7 +115,7 @@ public ReceiveResource(final KeyStore keyStore, this.firmwareUpdateStore = firmwareUpdateStore; this.groupFlipper = groupFlipper; - this.senseUploadConfiguration = senseUploadConfiguration; + this.senseUploadDynamicConfiguration = senseUploadDynamicConfiguration; this.otaConfiguration = otaConfiguration; this.responseCommandsDAODynamoDB = responseCommandsDAODynamoDB; this.senseClockOutOfSync = Metrics.newMeter(ReceiveResource.class, "sense-clock-out-sync", "clock-out-of-sync", TimeUnit.SECONDS); @@ -366,7 +367,7 @@ private byte[] generateSyncResponse(final String deviceName, } final Boolean isReducedInterval = featureFlipper.deviceFeatureActive(FeatureFlipper.REDUCE_BATCH_UPLOAD_INTERVAL, deviceName, groups); - final int uploadCycle = computeNextUploadInterval(nextRingTime, now, senseUploadConfiguration, isReducedInterval); + final int uploadCycle = computeNextUploadInterval(nextRingTime, now, (SenseUploadConfiguration) senseUploadDynamicConfiguration.getConfiguration(), isReducedInterval); responseBuilder.setBatchSize(uploadCycle); if(shouldWriteRingTimeHistory(now, nextRingTime, responseBuilder.getBatchSize())){