Skip to content

Commit 8341ec3

Browse files
authored
Merge pull request #1729 from guardian/live/add-channel-repo
Add DynamoDB LiveActivityChannelRepository
2 parents 44b90e4 + 16f2e8c commit 8341ec3

7 files changed

Lines changed: 293 additions & 0 deletions

File tree

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ project/plugins/project/
1717
dynamodb-local/
1818
dynamodb-local-common/
1919
dynamodb-local-schedule-lambda/
20+
dynamodb-local-live-activities/
2021

2122
# Scala-IDE specific
2223
.scala_dependencies

build.sbt

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -461,6 +461,7 @@ lazy val reportExtractor = lambda("reportextractor", "reportextractor", Some("co
461461
lazy val liveactivities = lambda("liveactivities", "liveactivities", Some("com.gu.liveactivities.LambdaLocalRun"))
462462
.dependsOn(common)
463463
.dependsOn(apiModels % "test->test", apiModels % "compile->compile")
464+
.settings(LocalDynamoDBLiveActivities.settings)
464465
.settings(
465466
libraryDependencies ++= Seq(
466467
"com.turo" % "pushy" % "0.13.10",
@@ -472,4 +473,10 @@ lazy val liveactivities = lambda("liveactivities", "liveactivities", Some("com.g
472473
// Hopefully this workaround can be removed once play-json-extensions either updates to Play 3.0 or is merged into play-json
473474
ExclusionRule(organization = "com.typesafe.play")
474475
),
476+
fork := true,
477+
startDynamoDBLocal := startDynamoDBLocal.dependsOn(Test / compile).value,
478+
Test / test := (Test / test).dependsOn(startDynamoDBLocal).value,
479+
Test / testOnly := (Test / testOnly).dependsOn(startDynamoDBLocal).evaluated,
480+
Test / testQuick := (Test / testQuick).dependsOn(startDynamoDBLocal).evaluated,
481+
Test / testOptions += dynamoDBLocalTestCleanup.value,
475482
)

common/src/main/scala/aws/AsyncDynamo.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,4 +41,5 @@ class AsyncDynamo(val client: AmazonDynamoDBAsync) {
4141
def query(request: QueryRequest): Future[QueryResult] = wrapAsyncMethod(client.queryAsync, request)
4242
def get(request: GetItemRequest): Future[GetItemResult] = wrapAsyncMethod(client.getItemAsync, request)
4343
def updateItem(request: UpdateItemRequest): Future[UpdateItemResult] = wrapAsyncMethod(client.updateItemAsync, request)
44+
def deleteItem(request: DeleteItemRequest): Future[DeleteItemResult] = wrapAsyncMethod(client.deleteItemAsync, request)
4445
}
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
package com.gu.liveactivities
2+
3+
import aws.AsyncDynamo
4+
import aws.DynamoJsonConversions.{fromAttributeMap, toAttributeMap}
5+
import cats.syntax.all._
6+
import com.amazonaws.services.dynamodbv2.model._
7+
import org.slf4j.{Logger, LoggerFactory}
8+
import play.api.libs.json._
9+
import tracking.Repository.RepositoryResult
10+
import tracking.RepositoryError
11+
12+
import scala.concurrent.{ExecutionContext, Future}
13+
import scala.jdk.CollectionConverters._
14+
15+
// MODELS /////////////////////////////////////////////
16+
17+
sealed trait LiveActivityData
18+
case class FootballLiveActivity(
19+
homeTeam: String,
20+
awayTeam: String,
21+
articleUrl: String
22+
) extends LiveActivityData
23+
24+
object FootballLiveActivity {
25+
implicit val format: OFormat[FootballLiveActivity] =
26+
Json.format[FootballLiveActivity]
27+
}
28+
29+
object LiveActivityData {
30+
implicit val format: OFormat[LiveActivityData] =
31+
new OFormat[LiveActivityData] {
32+
def writes(data: LiveActivityData): JsObject = data match {
33+
case f: FootballLiveActivity =>
34+
FootballLiveActivity.format.writes(f) + ("type" -> JsString(
35+
"football"
36+
))
37+
}
38+
def reads(json: JsValue): JsResult[LiveActivityData] =
39+
(json \ "type").validate[String].flatMap {
40+
case "football" => FootballLiveActivity.format.reads(json)
41+
case other => JsError(s"Unknown LiveActivityData type: $other")
42+
}
43+
}
44+
}
45+
46+
case class LiveActivityMapping(
47+
liveActivityId: String,
48+
channelId: String,
49+
data: Option[LiveActivityData]
50+
)
51+
object LiveActivityMapping {
52+
implicit val format: OFormat[LiveActivityMapping] =
53+
Json.format[LiveActivityMapping]
54+
}
55+
56+
// REPOSITORY /////////////////////////////////////////////
57+
58+
trait ChannelMappingsRepository {
59+
def saveMapping(mapping: LiveActivityMapping): Future[RepositoryResult[Unit]]
60+
def deleteMappingByActivityId(id: String): Future[RepositoryResult[Unit]]
61+
def getMappingByActivityId(
62+
id: String
63+
): Future[RepositoryResult[LiveActivityMapping]]
64+
}
65+
66+
class LiveActivityChannelRepository(client: AsyncDynamo, tableName: String)(
67+
implicit ec: ExecutionContext
68+
) extends ChannelMappingsRepository {
69+
70+
private val logger: Logger = LoggerFactory.getLogger(this.getClass)
71+
private val IdField = "liveActivityId"
72+
73+
override def saveMapping(
74+
mapping: LiveActivityMapping
75+
): Future[RepositoryResult[Unit]] = {
76+
val putItemRequest =
77+
new PutItemRequest()
78+
.withTableName(tableName)
79+
.withItem(toAttributeMap(mapping).asJava)
80+
.withConditionExpression(s"attribute_not_exists($IdField)")
81+
82+
client
83+
.putItem(putItemRequest)
84+
.map(_ => Right(()): RepositoryResult[Unit])
85+
.recover { case ex: Exception =>
86+
println("Error saving live activity mapping: " + ex.getMessage)
87+
val errorClass = ex.getClass.getName
88+
Left(RepositoryError(errorClass.toString))
89+
}
90+
}
91+
92+
override def getMappingByActivityId(
93+
id: String
94+
): Future[RepositoryResult[LiveActivityMapping]] = {
95+
val getItemRequest = new GetItemRequest()
96+
.withTableName(tableName)
97+
.withKey(Map(IdField -> new AttributeValue().withS(id)).asJava)
98+
.withConsistentRead(true)
99+
100+
client.get(getItemRequest) map { result =>
101+
for {
102+
item <- Either.fromOption(
103+
Option(result.getItem),
104+
RepositoryError("Live Activity not found")
105+
)
106+
parsed <- Either.fromOption(
107+
fromAttributeMap[LiveActivityMapping](item.asScala.toMap).asOpt,
108+
RepositoryError("Unable to parse live activity mapping")
109+
)
110+
} yield parsed
111+
}
112+
}
113+
114+
override def deleteMappingByActivityId(
115+
id: String
116+
): Future[RepositoryResult[Unit]] = {
117+
val deleteItemRequest = new DeleteItemRequest()
118+
.withTableName(tableName)
119+
.withKey(Map(IdField -> new AttributeValue().withS(id)).asJava)
120+
client.deleteItem(deleteItemRequest) map { _ => Right(()) }
121+
}
122+
123+
}
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
package com.gu.liveactivities
2+
3+
import aws.AsyncDynamo
4+
import com.amazonaws.auth.{
5+
AWSCredentials,
6+
AWSCredentialsProvider,
7+
AWSCredentialsProviderChain
8+
}
9+
import com.amazonaws.client.builder.AwsClientBuilder.EndpointConfiguration
10+
import com.amazonaws.regions.Regions
11+
import com.amazonaws.services.dynamodbv2.AmazonDynamoDBAsyncClientBuilder
12+
import com.amazonaws.services.dynamodbv2.model.{
13+
CreateTableRequest,
14+
DeleteTableRequest
15+
}
16+
import org.specs2.mutable.Specification
17+
import org.specs2.specification.{BeforeAfterAll, Scope}
18+
19+
trait DynamodbSpecification extends Specification with BeforeAfterAll {
20+
21+
sequential
22+
23+
val TableName: String
24+
25+
def createTableRequest: CreateTableRequest
26+
27+
val TestEndpoint = "http://localhost:8002"
28+
29+
override def beforeAll(): Unit = {
30+
awsClient.createTable(createTableRequest)
31+
}
32+
33+
override def afterAll(): Unit = {
34+
awsClient.deleteTable(new DeleteTableRequest(TableName))
35+
}
36+
37+
private def awsClient = {
38+
val chain = new AWSCredentialsProviderChain(new AWSCredentialsProvider {
39+
override def refresh(): Unit = {}
40+
41+
override def getCredentials: AWSCredentials = new AWSCredentials {
42+
override def getAWSAccessKeyId: String = "testid"
43+
override def getAWSSecretKey: String = "testkey"
44+
}
45+
})
46+
47+
val client = AmazonDynamoDBAsyncClientBuilder.standard()
48+
.withCredentials(chain)
49+
.withEndpointConfiguration( new EndpointConfiguration(TestEndpoint, Regions.EU_WEST_1.getName) )
50+
.build
51+
52+
client
53+
}
54+
55+
trait AsyncDynamoScope extends Scope {
56+
val asyncClient = new AsyncDynamo(awsClient)
57+
}
58+
}
Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
package com.gu.liveactivities
2+
3+
import com.amazonaws.services.dynamodbv2.model._
4+
import org.specs2.concurrent.ExecutionEnv
5+
import org.specs2.mock.Mockito
6+
import tracking.Repository.RepositoryResult
7+
import tracking.RepositoryError
8+
import scala.concurrent.Future
9+
import scala.jdk.CollectionConverters._
10+
11+
class LiveActivityChannelRepositoryTest(implicit ev: ExecutionEnv)
12+
extends DynamodbSpecification
13+
with Mockito {
14+
15+
override val TableName = "test-table"
16+
17+
"LiveActivityChannelRepository" should {
18+
"connect to local DynamoDB and create table" in new RepositoryScope {
19+
// This test will pass if the table is created successfully in beforeAll()
20+
1 must beEqualTo(1)
21+
}
22+
23+
"save a new channel mapping for an id if it does not exist" in new RepositoryScope
24+
with ExampleData {
25+
repository.saveMapping(footballMapping).flatMap { _ =>
26+
repository.getMappingByActivityId(footballMapping.liveActivityId)
27+
} must beEqualTo(Right(footballMapping)).await
28+
}
29+
30+
"Error if trying to save a new channel mapping for an id that already exists" in new RepositoryScope
31+
with ExampleData {
32+
repository.saveMapping(footballMapping).flatMap { _ =>
33+
repository.saveMapping(footballMapping)
34+
} must beLike[RepositoryResult[Unit]] {
35+
case Left(RepositoryError(msg))
36+
if msg.contains("ConditionalCheckFailed") =>
37+
ok
38+
}.await
39+
}
40+
41+
"delete a channel mapping for an activity id if it exists" in new RepositoryScope
42+
with ExampleData {
43+
44+
repository.saveMapping(footballMapping).flatMap { _ =>
45+
repository.deleteMappingByActivityId(footballMapping.liveActivityId)
46+
} must beEqualTo(Right(())).await
47+
}
48+
49+
"get a channel mapping for an activity id if it exists" in new RepositoryScope
50+
with ExampleData {
51+
repository.saveMapping(footballMapping).flatMap { _ =>
52+
repository.getMappingByActivityId(footballMapping.liveActivityId)
53+
} must beEqualTo(Right(footballMapping)).await
54+
}
55+
}
56+
57+
trait RepositoryScope extends AsyncDynamoScope {
58+
val repository = new LiveActivityChannelRepository(asyncClient, TableName)
59+
}
60+
61+
trait ExampleData {
62+
val footballData = FootballLiveActivity(
63+
homeTeam = "HomeTeamName",
64+
awayTeam = "AwayTeamName",
65+
articleUrl = "https://www.theguardian.com/football/test-article-id"
66+
)
67+
68+
val footballMapping = LiveActivityMapping(
69+
liveActivityId = "football-1234567",
70+
channelId = "test-channel-id",
71+
data = Some(footballData)
72+
)
73+
}
74+
75+
override def createTableRequest: CreateTableRequest = {
76+
val IdField = "liveActivityId"
77+
78+
new CreateTableRequest(
79+
TableName,
80+
List(new KeySchemaElement(IdField, KeyType.HASH)).asJava
81+
)
82+
.withAttributeDefinitions(
83+
List(
84+
new AttributeDefinition(IdField, ScalarAttributeType.S)
85+
).asJava
86+
)
87+
.withProvisionedThroughput(new ProvisionedThroughput(5L, 5L))
88+
}
89+
90+
}
Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,13 @@
1+
import com.localytics.sbt.dynamodb.DynamoDBLocalKeys
2+
import com.localytics.sbt.dynamodb.DynamoDBLocalKeys._
3+
import sbt._
4+
5+
object LocalDynamoDBLiveActivities {
6+
7+
val settings: Seq[Setting[_]] = DynamoDBLocalKeys.baseDynamoDBSettings ++ Seq(
8+
dynamoDBLocalDownloadDir := file("dynamodb-local-live-activities"),
9+
dynamoDBLocalInMemory := true,
10+
dynamoDBLocalVersion := "2018-04-11",
11+
dynamoDBLocalPort := 8002
12+
)
13+
}

0 commit comments

Comments
 (0)