Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
@@ -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);
}
}

Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);

Expand All @@ -224,7 +231,7 @@ public String getAWSSecretKey() {
configuration.getDebug(),
firmwareUpdateStore,
groupFlipper,
configuration.getSenseUploadConfiguration(),
senseUploadDynamicConfig,
configuration.getOTAConfiguration(),
respCommandsDAODynamoDB,
configuration.getRingDuration()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -26,6 +27,7 @@ protected void run(Bootstrap<SuripuConfiguration> bootstrap, Namespace namespace

final AWSCredentialsProvider awsCredentialsProvider = new DefaultAWSCredentialsProviderChain();
createRingTimeHistoryTable(configuration, awsCredentialsProvider);
createServiceConfigurationTable(configuration, awsCredentialsProvider);
}


Expand All @@ -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<DynamoDBTableName, String> tableNames = configuration.dynamoDBConfiguration().tables();
final ImmutableMap<DynamoDBTableName, String> 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());
}
}

}
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
package com.hello.suripu.service.configuration;


import com.fasterxml.jackson.annotation.JsonFilter;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.yammer.dropwizard.config.Configuration;

import javax.validation.Valid;
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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;

Expand All @@ -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) {
Expand All @@ -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);
Expand Down Expand Up @@ -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())){
Expand Down