Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
92 commits
Select commit Hold shift + click to select a range
93c1490
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 20, 2023
9b49a3b
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 20, 2023
ca23068
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 20, 2023
d76f191
104436:Changes related to publishing event to Kafka after ceritificat…
ravisaurav-tarento Sep 21, 2023
43a629b
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 21, 2023
75ca406
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 21, 2023
ec2301b
Updated changes
karthik-tarento Sep 22, 2023
19cb933
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 22, 2023
cf735fa
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 22, 2023
8e7fc61
4.8.6.1 post publish flink (#27)
karthik-tarento Sep 25, 2023
9571f71
Dev #104439 - Create new Flink job to validate the program certificat…
Sahil-tarento Sep 25, 2023
dd7ad93
Updated config and added logs
karthik-tarento Sep 26, 2023
12f85d7
Added code for metrics list
karthik-tarento Sep 26, 2023
349cac2
Adding New Fields in the Email Template - CourseProviderName, CourseP…
tarentomaheshvakkund Sep 28, 2023
92d511f
New Email Template - endpoint added for email notification (#31)
tarentomaheshvakkund Sep 30, 2023
f2c254d
Fix for NPE
karthik-tarento Oct 2, 2023
a175d16
Fix for type cast issue
karthik-tarento Oct 2, 2023
1ce7d10
Fix for cast when matching string
karthik-tarento Oct 2, 2023
ddebd30
Fixed for type casting
Sahil-tarento Oct 2, 2023
fa043c6
Added logs
karthik-tarento Oct 3, 2023
1526299
Handling null value for parentCollections
karthik-tarento Oct 3, 2023
0350b7e
Added debug logs
karthik-tarento Oct 3, 2023
8359a45
Handling null for versionKey
karthik-tarento Oct 3, 2023
f8b531b
Added versionKey
karthik-tarento Oct 3, 2023
f0b4fe2
Added logs
karthik-tarento Oct 3, 2023
cc0fc2f
Adding proper identifier in parentCollections
karthik-tarento Oct 3, 2023
6fc810d
Added logs
karthik-tarento Oct 3, 2023
ca5ac4e
Dev #106998-Need to enhance activity aggregate flink job to skip even…
tarentomaheshvakkund Oct 3, 2023
285b3c6
Added logs
karthik-tarento Oct 3, 2023
0b7b15c
Merge branch '4.8.6.1-certGenProgram' of https://github.com/sunbird-c…
karthik-tarento Oct 3, 2023
07a797b
Updated code to read cache values
karthik-tarento Oct 3, 2023
5b904c2
Updated key when fetching from redis
karthik-tarento Oct 3, 2023
b9ef989
Updated key when fetching from redis
karthik-tarento Oct 3, 2023
8c31a74
Fixed type cast issue
karthik-tarento Oct 3, 2023
3423a58
Dev #106996 - Need to create new Flink job to process activity aggreg…
Sahil-tarento Oct 3, 2023
97c24d6
Type cast issue fix
karthik-tarento Oct 3, 2023
36ccbc6
Merge branch '4.8.6.1-certGenProgram' of https://github.com/sunbird-c…
karthik-tarento Oct 3, 2023
68bede0
Added logs
karthik-tarento Oct 3, 2023
5fb91d5
Added the imports
karthik-tarento Oct 3, 2023
e1d6276
Fix for reading data from getCourseInfo
karthik-tarento Oct 4, 2023
441f4a2
Moved generating prog cert event before generating notification
karthik-tarento Oct 4, 2023
3bfced4
Added logs
karthik-tarento Oct 4, 2023
aaccd32
Added the loggers for IssueCertificateHelper.scala
Sahil-tarento Oct 4, 2023
efb558b
Added the changes for issueCertificate logic
Sahil-tarento Oct 4, 2023
f6987bd
Added the changes for program-activity-aggregator
Sahil-tarento Oct 4, 2023
9aa03e3
Enhanced course completion cert request (#32)
tarentomaheshvakkund Oct 4, 2023
b71e4f8
Added the changes for program-activity-aggregator
Sahil-tarento Oct 5, 2023
00d2b91
Added the changes for program-activity-aggregator Config fix
Sahil-tarento Oct 5, 2023
1d2a962
Added the config changes for program-activity-aggregator
Sahil-tarento Oct 5, 2023
b183194
Updated program activity aggregate updater flink job
karthik-tarento Oct 6, 2023
982ac6e
Fixed build issues
karthik-tarento Oct 6, 2023
f976b18
Updated config names
karthik-tarento Oct 6, 2023
41d3262
Added logs
karthik-tarento Oct 6, 2023
c8a3bb0
Dev #107756
Sahil-tarento Oct 6, 2023
dde3734
Added type casting
karthik-tarento Oct 6, 2023
4a260d0
Fix for type casting
karthik-tarento Oct 6, 2023
7be2af0
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 7, 2023
f2aa715
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 9, 2023
aa96b07
Added the changes for program-cert-pre-aggregator Add metrics parama
Sahil-tarento Oct 9, 2023
51d9d52
Added the changes for program-pre-processor code changes for complete…
Sahil-tarento Oct 9, 2023
abd7c28
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 10, 2023
1228b3e
Added the changes for program-activity-aggregator for new edge case
Sahil-tarento Oct 12, 2023
7e0f593
Added the changes for program-activity-aggregator code fix config fil…
Sahil-tarento Oct 12, 2023
cccccb1
Added the changes for program-activity-aggregator type casting fix
Sahil-tarento Oct 12, 2023
6302c6c
Added the changes for program-activity-aggregator code typecasting
Sahil-tarento Oct 12, 2023
277595a
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 12, 2023
8fcfd73
Added the changes for program-activity-aggregator event fix
Sahil-tarento Oct 12, 2023
910d651
Added the changes for program-activity-aggregator logger
Sahil-tarento Oct 12, 2023
01b1aea
Added the changes for program-activity-aggregator Stream FIx
Sahil-tarento Oct 12, 2023
2016d41
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 13, 2023
fabc764
Added the changes for program-activity-aggregator code changes for pr…
Sahil-tarento Oct 13, 2023
ab79a06
Added the loggers for the activity aggregrator
Sahil-tarento Oct 17, 2023
363e9c2
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 17, 2023
bbcb635
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 17, 2023
2f1a0f5
Added the changes for program-activity-aggregator DeDup fix
Sahil-tarento Oct 17, 2023
353b12a
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 17, 2023
f359fc9
Added the changes for program-activity-aggregator code fix typeCasting
Sahil-tarento Oct 18, 2023
ccc6bd5
Added the changes for program-activity-aggregator code fix typecasting
Sahil-tarento Oct 18, 2023
c00ea2a
Added the changes for program-activity-aggregator typecasting
Sahil-tarento Oct 18, 2023
7e2584a
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 18, 2023
05caa3c
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 18, 2023
23e7a7c
Added the changes for program-activity-aggregator code fix type casting
Sahil-tarento Oct 18, 2023
f803cfb
Added the logging
Sahil-tarento Oct 18, 2023
b2dc98a
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 18, 2023
96c5d3e
Added the changes for program-activity-aggregator code typecasting
Sahil-tarento Oct 18, 2023
855c5ad
Added the changes for program-activity-aggregator code fix
Sahil-tarento Oct 18, 2023
ffe4022
Optimized code to reduce cassandra queries
karthik-tarento Oct 19, 2023
4ca39c6
Merge branch '4.8.6.1-certGenProgram' of https://github.com/sunbird-c…
karthik-tarento Oct 19, 2023
661f142
Added the changes for new function for Activity Aggregrator
Sahil-tarento Oct 19, 2023
a0353e6
Added the changes for program-cert-pre-processor code fix
Sahil-tarento Oct 19, 2023
3a43796
Disabled user feed notification
karthik-tarento Oct 24, 2023
860e290
Merge branch 'cbrelease-4.8.7' into 4.8.6.1-certGenProgram
karthik-tarento Oct 26, 2023
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
Expand Up @@ -2,7 +2,6 @@ package org.sunbird.job.aggregate.functions

import java.lang.reflect.Type
import java.util.concurrent.TimeUnit

import com.datastax.driver.core.Row
import com.datastax.driver.core.querybuilder.{QueryBuilder, Select, Update}
import com.google.gson.Gson
Expand All @@ -19,10 +18,11 @@ import org.slf4j.LoggerFactory
import org.sunbird.job.cache.{DataCache, RedisConnect}
import org.sunbird.job.aggregate.domain.{UserContentConsumption, _}
import org.sunbird.job.aggregate.task.ActivityAggregateUpdaterConfig
import org.sunbird.job.util.{CassandraUtil, HttpUtil}
import org.sunbird.job.util.{CassandraUtil, HttpUtil, JSONUtil, ScalaJsonUtil}
import org.sunbird.job.{Metrics, WindowBaseProcessFunction}

import scala.collection.JavaConverters._
import scala.language.postfixOps

class ActivityAggregatesFunction(config: ActivityAggregateUpdaterConfig, httpUtil: HttpUtil, @transient var cassandraUtil: CassandraUtil = null)
(implicit val stringTypeInfo: TypeInformation[String])
Expand Down Expand Up @@ -57,8 +57,9 @@ class ActivityAggregatesFunction(config: ActivityAggregateUpdaterConfig, httpUti
events: Iterable[Map[String, AnyRef]],
metrics: Metrics): Unit = {

logger.debug("Input Events Size: " + events.toList.size)
logger.info("Input Events : " + JSONUtil.serialize(events.toList))
val inputUserConsumptionList: List[UserContentConsumption] = events
.filter(event=> verifyPrimaryCategory(event.getOrElse(config.courseId, "").asInstanceOf[String])(metrics, config, httpUtil, cache))
.groupBy(key => (key.get(config.courseId), key.get(config.batchId), key.get(config.userId)))
.values.map(value => {
metrics.incCounter(config.processedEnrolmentCount)
Expand All @@ -84,7 +85,7 @@ class ActivityAggregatesFunction(config: ActivityAggregateUpdaterConfig, httpUti
val userConsumptionQueries = finalUserConsumptionList.flatMap(userConsumption => getContentConsumptionQueries(userConsumption))
updateDB(config.thresholdBatchWriteSize, userConsumptionQueries)(metrics)


logger.info("The value for userConsumptionQueries:" + userConsumptionQueries)
val courseAggregations = finalUserConsumptionList.flatMap(userConsumption => {

// Course Level Agg using the merged data of ContentConsumption per user, course and batch.
Expand All @@ -106,13 +107,13 @@ class ActivityAggregatesFunction(config: ActivityAggregateUpdaterConfig, httpUti

// Saving enrolment completion data.
val collectionProgressList = courseAggregations.filter(agg => agg.collectionProgress.nonEmpty).map(agg => agg.collectionProgress.get)

logger.info("The value for collectionProgressList:" + collectionProgressList)
val collectionProgressUpdateList = collectionProgressList.filter(progress => !progress.completed)
context.output(config.collectionUpdateOutputTag, collectionProgressUpdateList)

logger.info("The value for collectionProgressUpdateList:" + collectionProgressUpdateList)
val collectionProgressCompleteList = collectionProgressList.filter(progress => progress.completed)
context.output(config.collectionCompleteOutputTag, collectionProgressCompleteList)

logger.info("The value for collectionProgressCompleteList:" + collectionProgressCompleteList)
// Content AUDIT Event generation and pushing to output tag.
finalUserConsumptionList.flatMap(userConsumption => contentAuditEvents(userConsumption)).foreach(event => context.output(config.auditEventOutputTag, gson.toJson(event)))
}
Expand Down Expand Up @@ -444,5 +445,133 @@ class ActivityAggregatesFunction(config: ActivityAggregateUpdaterConfig, httpUti
dbStatus
} else cacheStatus
}
}

def verifyPrimaryCategory(identifier: String)(
metrics: Metrics,
config: ActivityAggregateUpdaterConfig,
httpUtil: HttpUtil,
cache: DataCache
): Boolean = {
logger.info(
"Verify Program post-publish required for content: " + identifier
)
// Get the primary Categories for the courses here
var isValidProgram = false
val contentObj: java.util.Map[String, AnyRef] =
getCourseInfo(identifier)(metrics, config, cache, httpUtil)
if (!contentObj.isEmpty) {
val primaryCategory = contentObj.get("primaryCategory")
if (primaryCategory != null &&
(primaryCategory != "Program"
|| primaryCategory != "Curated Program"
|| primaryCategory != "Blended Program")) {
isValidProgram = true
}
logger.info("PrimaryCategory value is :" + primaryCategory + ", for Id: " + identifier)
} else {
logger.error("Failed to read content details for Id: " + identifier)
}
logger.info("is activity aggregator is skipping this event ? " + isValidProgram)
isValidProgram
}

def getCourseInfo(courseId: String)(
metrics: Metrics,
config: ActivityAggregateUpdaterConfig,
cache: DataCache,
httpUtil: HttpUtil
): java.util.Map[String, AnyRef] = {
val courseMetadata = cache.getWithRetry(courseId)
if (null == courseMetadata || courseMetadata.isEmpty) {
val url =
config.contentReadURL + "/" + courseId + "?fields=identifier,name,versionKey,parentCollections,primaryCategory"
val response = getAPICall(url, "content")(config, httpUtil, metrics)
val courseName = StringContext
.processEscapes(
response.getOrElse(config.name, "").asInstanceOf[String]
)
.filter(_ >= ' ')
val primaryCategory = StringContext
.processEscapes(
response.getOrElse(config.primaryCategory, "").asInstanceOf[String]
)
.filter(_ >= ' ')
val versionKey = StringContext
.processEscapes(
response.getOrElse(config.versionKey, "").asInstanceOf[String]
)
.filter(_ >= ' ')
val parentCollections = response
.getOrElse("parentCollections", List.empty[String])
.asInstanceOf[List[String]]
val courseInfoMap: java.util.Map[String, AnyRef] =
new java.util.HashMap[String, AnyRef]()
courseInfoMap.put("courseId", courseId)
courseInfoMap.put("courseName", courseName)
courseInfoMap.put("parentCollections", parentCollections)
courseInfoMap.put("primaryCategory", primaryCategory)
courseInfoMap.put("versionKey", versionKey)
courseInfoMap
} else {
val courseName = StringContext
.processEscapes(
courseMetadata.getOrElse(config.name, "").asInstanceOf[String]
)
.filter(_ >= ' ')
val primaryCategory = StringContext
.processEscapes(
courseMetadata
.getOrElse(config.primaryCategory, "")
.asInstanceOf[String]
)
.filter(_ >= ' ')
val versionKey = StringContext
.processEscapes(
courseMetadata.getOrElse(config.versionKey, "").asInstanceOf[String]
)
.filter(_ >= ' ')
val parentCollections = courseMetadata
.getOrElse("parentCollections", List.empty[String])
.asInstanceOf[List[String]]
val courseInfoMap: java.util.Map[String, AnyRef] =
new java.util.HashMap[String, AnyRef]()
courseInfoMap.put("courseId", courseId)
courseInfoMap.put("courseName", courseName)
courseInfoMap.put("parentCollections", parentCollections)
courseInfoMap.put("primaryCategory", primaryCategory)
courseInfoMap.put("versionKey", versionKey)
courseInfoMap
}

}

def getAPICall(url: String, responseParam: String)(
config: ActivityAggregateUpdaterConfig,
httpUtil: HttpUtil,
metrics: Metrics
): Map[String, AnyRef] = {
val response = httpUtil.get(url, config.defaultHeaders)
if (200 == response.status) {
ScalaJsonUtil
.deserialize[Map[String, AnyRef]](response.body)
.getOrElse("result", Map[String, AnyRef]())
.asInstanceOf[Map[String, AnyRef]]
.getOrElse(responseParam, Map[String, AnyRef]())
.asInstanceOf[Map[String, AnyRef]]
} else if (
400 == response.status && response.body.contains(
config.userAccBlockedErrCode
)
) {
metrics.incCounter(config.skippedEventCount)
logger.error(
s"Error while fetching user details for ${url}: " + response.status + " :: " + response.body
)
Map[String, AnyRef]()
} else {
throw new Exception(
s"Error from get API : ${url}, with response: ${response}"
)
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,11 @@ class CollectionProgressUpdateFunction(config: ActivityAggregateUpdaterConfig)(i
val row = getEnrolment(p.userId, p.courseId, p.batchId)(metrics)
(row != null && row.getInt("status") != 2)
} else events
logger.info("The event at progress: "+ events)
logger.info("Pending Enrollment: " + pendingEnrolments)
val enrolmentQueries = pendingEnrolments.map(collectionProgress => getEnrolmentUpdateQuery(collectionProgress))
updateDB(config.thresholdBatchWriteSize, enrolmentQueries)(metrics)
logger.info("enrolmentQueries Enrolement: " + enrolmentQueries)
// Create and update the checksum to DeDup store for the input events.
if (config.dedupEnabled) {
events.map(cp => cp.inputContents.map(c => DeDupHelper.getMessageId(cp.courseId, cp.batchId, cp.userId, c, 2)))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -132,4 +132,13 @@ class ActivityAggregateUpdaterConfig(override val config: Config) extends BaseJo
val searchServiceBasePath: String = config.getString("service.search.basePath")
val searchAPIURL = searchServiceBasePath + "/v3/search"

val contentServiceBase: String = config.getString("service.content.basePath")
val contentReadURL = contentServiceBase + "/content/v3/read/"
val name: String = "name"
val primaryCategory: String = "primaryCategory"
val versionKey: String = "versionKey"
val defaultHeaders = Map[String, String] ("Content-Type" -> "application/json")
val userAccBlockedErrCode = "UOS_USRRED0006"
val skippedEventCount = "skipped-events-count"

}
Original file line number Diff line number Diff line change
Expand Up @@ -3,33 +3,37 @@ package org.sunbird.job.collectioncert.functions
import java.text.SimpleDateFormat
import com.datastax.driver.core.querybuilder.QueryBuilder
import com.datastax.driver.core.{Row, TypeTokens}
import org.apache.commons.collections.CollectionUtils
import org.apache.commons.lang3.StringUtils
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper
import org.slf4j.LoggerFactory
import org.sunbird.job.Metrics
import org.sunbird.job.cache.DataCache
import org.sunbird.job.collectioncert.domain.{AssessedUser, AssessmentUserAttempt, BEJobRequestEvent, EnrolledUser, Event, EventObject}
import org.sunbird.job.collectioncert.task.CollectionCertPreProcessorConfig
import org.sunbird.job.util.{CassandraUtil, HttpUtil, ScalaJsonUtil}
import org.sunbird.job.util.{CassandraUtil, HttpUtil, ScalaJsonUtil, JSONUtil}

import scala.collection.JavaConverters._

trait IssueCertificateHelper {
private[this] val logger = LoggerFactory.getLogger(classOf[CollectionCertPreProcessorFn])


def issueCertificate(event:Event, template: Map[String, String])(cassandraUtil: CassandraUtil, cache:DataCache, contentCache: DataCache, metrics: Metrics, config: CollectionCertPreProcessorConfig, httpUtil: HttpUtil): String = {
//validCriteria
logger.info("issueCertificate i/p event =>"+event)
val criteria = validateTemplate(template, event.batchId)(config)
//validateEnrolmentCriteria
val certName = template.getOrElse(config.name, "")
logger.info("CertName" + certName)
val additionalProps: Map[String, List[String]] = ScalaJsonUtil.deserialize[Map[String, List[String]]](template.getOrElse("additionalProps", "{}"))
val enrolledUser: EnrolledUser = validateEnrolmentCriteria(event, criteria.getOrElse(config.enrollment, Map[String, AnyRef]()).asInstanceOf[Map[String, AnyRef]], certName, additionalProps)(metrics, cassandraUtil, config)
logger.info("enrolledUser" + enrolledUser)
//validateAssessmentCriteria
val assessedUser = validateAssessmentCriteria(event, criteria.getOrElse(config.assessment, Map[String, AnyRef]()).asInstanceOf[Map[String, AnyRef]], enrolledUser.userId, additionalProps)(metrics, cassandraUtil, contentCache, config)
logger.info("assessedUser" + assessedUser)
//validateUserCriteria
val userDetails = validateUser(assessedUser.userId, criteria.getOrElse(config.user, Map[String, AnyRef]()).asInstanceOf[Map[String, AnyRef]], additionalProps)(metrics, config, httpUtil)

logger.info("userDetails" + userDetails)
//generateCertificateEvent
if(userDetails.nonEmpty) {
generateCertificateEvent(event, template, userDetails, enrolledUser, assessedUser, additionalProps, certName)(metrics, config, cache, httpUtil)
Expand Down Expand Up @@ -210,10 +214,17 @@ trait IssueCertificateHelper {
val lastName = Option(userDetails.getOrElse("lastName", "").asInstanceOf[String]).getOrElse("")
def nullStringCheck(name:String):String = {if(StringUtils.equalsIgnoreCase("null", name)) "" else name}
val recipientName = nullStringCheck(firstName).concat(" ").concat(nullStringCheck(lastName)).trim
val courseName = getCourseName(event.courseId)(metrics, config, cache, httpUtil)
val courseInfo: java.util.Map[String, AnyRef] = getCourseInfo(event.courseId)(metrics, config, cache, httpUtil)
val courseName = courseInfo.getOrDefault("courseName", "").asInstanceOf[String]
val dateFormatter = new SimpleDateFormat("yyyy-MM-dd")
val related = getRelatedData(event, enrolledUser, assessedUser, userDetails, additionalProps, certName, courseName)(config)
val providerName = getCourseOrganisation(event.courseId)(metrics, config, cache, httpUtil)
val parentCollections: List[String] = Option(courseInfo.get(config.parentCollections))
.collect {
case list: java.util.List[_] =>
list.asInstanceOf[java.util.List[String]].asScala.toList
}
.getOrElse(List.empty)

val eData = Map[String, AnyRef] (
"issuedDate" -> dateFormatter.format(enrolledUser.issuedOn),
"data" -> List(Map[String, AnyRef]("recipientName" -> recipientName, "recipientId" -> event.userId)),
Expand All @@ -229,10 +240,13 @@ trait IssueCertificateHelper {
"basePath" -> config.certBasePath,
"related" -> related,
"name" -> certName,
"providerName" -> providerName,
"tag" -> event.batchId
"providerName" -> courseInfo.getOrDefault("providerName", "").asInstanceOf[String],
"tag" -> event.batchId,
"primaryCategory" -> courseInfo.getOrDefault("primaryCategory", "").asInstanceOf[String],
"parentCollections" -> parentCollections,
"coursePosterImage" -> courseInfo.getOrDefault("coursePosterImage", "").asInstanceOf[String],
)

logger.info("Constructured eData from preProcessor : " + JSONUtil.serialize(eData))
ScalaJsonUtil.serialize(BEJobRequestEvent(edata = eData, `object` = EventObject(id= event.userId)))
}

Expand All @@ -252,4 +266,43 @@ trait IssueCertificateHelper {
Map[String, Any]("batchId" -> event.batchId, "courseId" -> event.courseId, "type" -> certName) ++
locationProps ++ enrolledUser.additionalProps ++ assessedUser.additionalProps ++ userAdditionalProps ++ courseAdditionalProps
}

def getCourseInfo(courseId: String)(metrics: Metrics, config: CollectionCertPreProcessorConfig, cache: DataCache, httpUtil: HttpUtil): java.util.Map[String, AnyRef] = {
val courseMetadata = cache.getWithRetry(courseId)
if (null == courseMetadata || courseMetadata.isEmpty) {
val url = config.contentBasePath + config.contentReadApi + "/" + courseId + "?fields=name,parentCollections,primaryCategory,posterImage,organisation"
val response = getAPICall(url, "content")(config, httpUtil, metrics)
val courseName = StringContext.processEscapes(response.getOrElse(config.name, "").asInstanceOf[String]).filter(_ >= ' ')
val primaryCategory = StringContext.processEscapes(response.getOrElse(config.primaryCategory, "").asInstanceOf[String]).filter(_ >= ' ')
val posterImage: String = StringContext.processEscapes(response.getOrElse(config.posterImage, "").asInstanceOf[String]).filter(_ >= ' ')
val parentCollections = response.getOrElse("parentCollections", List.empty[String]).asInstanceOf[List[String]]
val orgData = response.get("organisation").toArray
val pm = orgData(0).toString
val providerName = pm.substring(1, pm.length - 1)
val courseInfoMap: java.util.Map[String, AnyRef] = new java.util.HashMap[String, AnyRef]()
courseInfoMap.put("courseId", courseId)
courseInfoMap.put("courseName", courseName)
courseInfoMap.put("parentCollections", parentCollections)
courseInfoMap.put("primaryCategory", primaryCategory)
courseInfoMap.put("coursePosterImage", posterImage)
courseInfoMap.put("providerName", providerName)
courseInfoMap
} else {
val courseName = StringContext.processEscapes(courseMetadata.getOrElse(config.name, "").asInstanceOf[String]).filter(_ >= ' ')
val primaryCategory = StringContext.processEscapes(courseMetadata.getOrElse("primarycategory", "").asInstanceOf[String]).filter(_ >= ' ')
val parentCollections = courseMetadata.getOrElse("parentcollections", new java.util.ArrayList()).asInstanceOf[java.util.ArrayList[String]]
val posterImage: String = StringContext.processEscapes(courseMetadata.getOrElse("posterimage", "").asInstanceOf[String]).filter(_ >= ' ')
val orgData = courseMetadata.get("organisation").toArray
val pm = orgData(0).toString
val providerName = pm.substring(1, pm.length - 1)
val courseInfoMap: java.util.Map[String, AnyRef] = new java.util.HashMap[String, AnyRef]()
courseInfoMap.put("courseId", courseId)
courseInfoMap.put("courseName", courseName)
courseInfoMap.put("parentCollections", parentCollections)
courseInfoMap.put("primaryCategory", primaryCategory)
courseInfoMap.put("coursePosterImage", posterImage)
courseInfoMap.put("providerName", providerName)
courseInfoMap
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,11 @@ class CollectionCertPreProcessorConfig(override val config: Config) extends Base
val assessmentContentTypes = if(config.hasPath("assessment.metrics.supported.contenttype")) config.getStringList("assessment.metrics.supported.contenttype") else util.Arrays.asList("SelfAssess")
val userAccBlockedErrCode = "UOS_USRRED0006"
val enableSuppressException: Boolean = if(config.hasPath("enable.suppress.exception")) config.getBoolean("enable.suppress.exception") else false

val contentServiceBase: String = config.getString("service.content.basePath")
val contentReadURL = contentServiceBase+ "/content/v3/read/"
val primaryCategory: String = "primaryCategory"
val versionKey: String = "versionKey"
val posterImage: String = "posterImage"
val parentCollections: String = "parentCollections"

}
Original file line number Diff line number Diff line change
Expand Up @@ -60,4 +60,9 @@ class Event(eventMap: java.util.Map[String, Any], partition: Int, offset: Long)

def providerName: String = readOrDefault[String]("edata.providerName", "")

def primaryCategory: String = readOrDefault[String]("edata.primaryCategory", "")

def parentCollections: List[String] = readOrDefault[List[String]]("edata.parentCollections", List.empty[String])

def coursePosterImage: String = readOrDefault[String]("edata.coursePosterImage", "")
}
Loading