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
Expand Up @@ -271,17 +271,12 @@ class ActivityAggregatorActor @Inject()(implicit val cacheUtil: RedisCacheUtil)
}

private def updateActivityAggregates(courseAggregations: List[UserEnrolmentAgg], requestContext: RequestContext): Unit = {
logger.info(requestContext, s"updateActivityAggregates: Creating batch update queries for ${courseAggregations.size} aggregations")
val aggQueries = courseAggregations.map { agg =>
activityAggUtil.createActivityAggUpdateMap(agg.activityAgg)
}.asJava

if (!aggQueries.isEmpty) {
logger.info(requestContext, s"updateActivityAggregates: Executing batch update with ${aggQueries.size()} queries to ${activityAggDBInfo.getTableName}")
cassandraOperation.batchUpdate(activityAggDBInfo.getKeySpace, activityAggDBInfo.getTableName, aggQueries, requestContext)
logger.info(requestContext, s"updateActivityAggregates: Batch update completed successfully")
} else {
logger.warn(requestContext, s"updateActivityAggregates: No queries to execute")
cassandraOperation.batchUpdateWithPutAll(activityAggDBInfo.getKeySpace, activityAggDBInfo.getTableName, aggQueries, requestContext)
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,11 @@ class ActivityAggregatorActorTest
.returning(new Response())
.anyNumberOfTimes()

(cassandraOperation.batchUpdateWithPutAll(_: String, _: String, _: util.List[util.Map[String, util.Map[String, Object]]], _: RequestContext))
.expects(*, "user_activity_agg", *, *)
.returning(new Response())
.anyNumberOfTimes()

val actor = system.actorOf(Props(new TestableActivityAggregatorActor(cassandraOperation, redisUtil, deDupUtil, contentSearchUtil, certificateUtil)))

val request = createUpdateRequest(
Expand Down Expand Up @@ -128,6 +133,11 @@ class ActivityAggregatorActorTest
.returning(new Response())
.anyNumberOfTimes()

(cassandraOperation.batchUpdateWithPutAll(_: String, _: String, _: util.List[util.Map[String, util.Map[String, Object]]], _: RequestContext))
.expects(*, "user_activity_agg", *, *)
.returning(new Response())
.anyNumberOfTimes()

val actor = system.actorOf(Props(new TestableActivityAggregatorActor(cassandraOperation, redisUtil, deDupUtil, contentSearchUtil, certificateUtil)))

val request = createUpdateRequest(
Expand Down Expand Up @@ -194,6 +204,11 @@ class ActivityAggregatorActorTest
.returning(new Response())
.anyNumberOfTimes()

(cassandraOperation.batchUpdateWithPutAll(_: String, _: String, _: util.List[util.Map[String, util.Map[String, Object]]], _: RequestContext))
.expects(*, "user_activity_agg", *, *)
.returning(new Response())
.anyNumberOfTimes()

val actor = system.actorOf(Props(new TestableActivityAggregatorActor(cassandraOperation, redisUtil, deDupUtil, contentSearchUtil, certificateUtil)))

val contentsWithInvalid = new util.ArrayList[util.Map[String, AnyRef]]()
Expand Down Expand Up @@ -296,6 +311,11 @@ class ActivityAggregatorActorTest
.returning(new Response())
.anyNumberOfTimes()

(cassandraOperation.batchUpdateWithPutAll(_: String, _: String, _: util.List[util.Map[String, util.Map[String, Object]]], _: RequestContext))
.expects(*, "user_activity_agg", *, *)
.returning(new Response())
.anyNumberOfTimes()

val actor = system.actorOf(Props(new TestableActivityAggregatorActor(cassandraOperation, redisUtil, deDupUtil, contentSearchUtil, certificateUtil)))

val inputContents = new util.ArrayList[util.Map[String, AnyRef]]()
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
package org.sunbird.assessment.actor
import org.sunbird.actor.core.BaseActor
import javax.inject.Inject
import org.apache.pekko.actor.Props
import org.sunbird.common.exception.ProjectCommonException
import org.sunbird.common.request.{Request, RequestContext}
Expand All @@ -12,23 +11,35 @@ import org.sunbird.common.models.util.{JsonKey, LoggerUtil, ProjectUtil}
import scala.collection.JavaConverters._
import org.apache.commons.lang3.StringUtils

class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_contentService: Option[ContentService],_cassandraService: Option[CassandraService],_kafkaService: Option[KafkaService]) extends BaseActor {
class AssessmentAggregatorActor(
_cassandraService: Option[CassandraService],
_kafkaService: Option[KafkaService],
_redisService: Option[RedisService],
_contentService: Option[ContentService]
) extends BaseActor {

def this() = this(None, None, None, None)

private lazy val redisService = _redisService.getOrElse(new RedisService())
private lazy val contentService = _contentService.getOrElse(new ContentService())
private lazy val cassandraService = _cassandraService.getOrElse(AssessmentAggregatorActor.cassandraService)
private lazy val kafkaService = _kafkaService.getOrElse(AssessmentAggregatorActor.kafkaService)
private lazy val redisService = _redisService.getOrElse(AssessmentAggregatorActor.redisService)
private lazy val contentService = _contentService.getOrElse(AssessmentAggregatorActor.contentService)
private lazy val assessmentService = new AssessmentService(redisService, contentService)
private lazy val cassandraService = _cassandraService.getOrElse(new CassandraService())
private lazy val kafkaService = _kafkaService.getOrElse(new KafkaService())


override def onReceive(request: Request): Unit = {
request.getOperation match {
case "aggregateAssessment" => aggregateAssessment(request)
case _ => onReceiveUnsupportedOperation(request.getOperation)
}
}

private def aggregateAssessment(request: Request): Unit = {
val replyTo = sender()
try {
processAggregation(request, replyTo)
try {
processAggregation(request, replyTo)
} catch {
case ex: Exception =>
logger.error(request.getRequestContext, "Request failed", ex)
logger.error(request.getRequestContext, s"Assessment aggregation failed: ${ex.getMessage}", ex)
replyTo ! createErrorResponse("SERVER_ERROR", ex.getMessage, ResponseCode.SERVER_ERROR.getResponseCode)
}
}
Expand Down Expand Up @@ -64,7 +75,7 @@ class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_c
replyTo ! createSuccess(assessment.attemptId)
} catch {
case ex: Exception =>
logger.error(context, s"Assessment request failed. Reason: ${ex.getMessage} | Data: $body", ex)
logger.error(context, s"[ASSESSMENT_ACTOR] Request failed: ${ex.getMessage}", ex)
replyTo ! createErrorResponse("CLIENT_ERROR", ex.getMessage, ResponseCode.CLIENT_ERROR.getResponseCode)
}
}
Expand All @@ -89,7 +100,6 @@ class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_c
logger.warn(context, s"Sync Flow: No stored events found for userId=${request.userId}, contentId=${request.contentId}, attemptId=${request.attemptId}")
return List(request)
}
logger.info(context, s"Sync Flow: Recovered ${existing.size} attempt(s) for userId=${request.userId}, contentId=${request.contentId}")
existing.map(toSyncRequest(request, _))
}

Expand Down Expand Up @@ -125,16 +135,15 @@ class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_c
if (skipMissing) {
val totalQuestions = metadata.totalQuestions
if (totalQuestions > 0 && uniqueEvents.size > totalQuestions) {
logger.warn(context, s"Skipping assessment ${req.attemptId}: unique events (${uniqueEvents.size}) exceed total questions ($totalQuestions)")
logger.warn(context, s"[ASSESSMENT_ACTOR] SKIPPED: unique events (${uniqueEvents.size}) exceed total questions ($totalQuestions) for attemptId=${req.attemptId}")
return
}
}
val scoreMetrics = assessmentService.computeScoreMetrics(uniqueEvents)
val existing = cassandraService.getAssessment(req.attemptId, req.userId, req.courseId, req.batchId, req.contentId, context)
val existingTs = existing.map(_.lastAttemptedOn).getOrElse(0L)
logger.info(context, s"AssessmentAggregatorActor: Comparing timestamps for attemptId=${req.attemptId} | Incoming=${req.assessmentTimestamp} | Existing=$existingTs")
if (!req.ignoreTimestampValidation && existingTs > req.assessmentTimestamp) {
logger.info(context, s"Skipping stale assessment: ${req.attemptId}")
logger.warn(context, s"[ASSESSMENT_ACTOR] SKIPPED: Stale assessment attemptId=${req.attemptId}")
return
}
val result = AssessmentResult(req.attemptId, req.userId, req.courseId, req.batchId, req.contentId, scoreMetrics.totalScore, scoreMetrics.totalMaxScore, scoreMetrics.grandTotal, scoreMetrics.questions, existing.map(_.createdOn).getOrElse(System.currentTimeMillis()), req.assessmentTimestamp)
Expand All @@ -150,7 +159,10 @@ class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_c
val attemptId = assessmentService.getLatestAttemptId(agg)
if (ProjectUtil.getConfigValue("assessment_aggregator_publish_certificate") == "true") {
kafkaService.publishCertificateEvent(userId, courseId, batchId, attemptId)
logger.info(context, s"[ASSESSMENT_ACTOR] Published certificate event for attemptId=$attemptId")
}
} else {
logger.warn(context, s"[ASSESSMENT_ACTOR] No assessments found for userId=$userId, courseId=$courseId, batchId=$batchId, contentId=$contentId")
}
}

Expand Down Expand Up @@ -225,5 +237,15 @@ class AssessmentAggregatorActor @Inject()(_redisService: Option[RedisService],_c
}

object AssessmentAggregatorActor {
def props(): Props = Props(new AssessmentAggregatorActor())
lazy val cassandraService = new CassandraService()
lazy val kafkaService = new KafkaService()
lazy val redisService = new RedisService()
lazy val contentService = new ContentService()

def props(): Props = Props(new AssessmentAggregatorActor(
Some(cassandraService),
Some(kafkaService),
Some(redisService),
Some(contentService)
))
}
Original file line number Diff line number Diff line change
Expand Up @@ -48,16 +48,18 @@ class CassandraService(optionalDao: Option[CassandraOperation] = None) {
try {
if (agg.aggregates.nonEmpty || agg.aggregateDetails.nonEmpty) {
val lastUpdated = agg.aggregates.map { case (k, _) => k -> new java.util.Date() }
val data = Map(
val compositeKey = Map(
"activity_id" -> cid,
"activity_type" -> "Course",
"context_id" -> s"cb:$bid",
"user_id" -> uid,
"user_id" -> uid
).asJava.asInstanceOf[java.util.Map[String, AnyRef]]
val updateAttributes = Map(
"aggregates" -> agg.aggregates.asJava,
"agg_details" -> agg.aggregateDetails.map(_.toJson).asJava,
"agg_last_updated" -> lastUpdated.asJava
).asJava.asInstanceOf[java.util.Map[String, AnyRef]]
dao.upsertRecord(keyspace, activityTable, data, ctx)
dao.updateRecordWithPutAll(keyspace, activityTable, updateAttributes, compositeKey, ctx)
}
} catch { case e: Exception => logger.error(s"Activity update failed for $uid", e); throw e }
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import org.sunbird.common.models.response.Response
import java.util.HashMap
import scala.collection.JavaConverters._
import scala.concurrent.duration._
import org.sunbird.assessment.service.{CassandraService, ContentMetadata, ContentService, KafkaService, RedisService}
import org.sunbird.assessment.service.{AssessmentService, CassandraService, ContentMetadata, ContentService, KafkaService, RedisService}
import org.sunbird.common.responsecode.ResponseCode

class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggregatorActorSpec"))
Expand All @@ -30,7 +30,9 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
TestKit.shutdownActorSystem(system)
}

def getActorRef = TestActorRef(new AssessmentAggregatorActor(Some(mRedis), Some(mContent), Some(mCassandra), Some(mKafka)))
def getActorRef = {
TestActorRef(new AssessmentAggregatorActor(Some(mCassandra), Some(mKafka), Some(mRedis), Some(mContent)))
}

"AssessmentAggregatorActor" should "silently ignore unknown message types (standard BaseActor behavior)" in {
val actorRef = getActorRef
Expand All @@ -48,6 +50,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put(JsonKey.USER_ID, "u1")
Expand Down Expand Up @@ -83,6 +86,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put(JsonKey.USER_ID, "u1")
Expand All @@ -109,6 +113,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()

Expand Down Expand Up @@ -140,6 +145,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
reset(mRedis, mContent, mCassandra, mKafka)
val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()

Expand Down Expand Up @@ -178,6 +184,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put(JsonKey.USER_ID, "u1")
Expand All @@ -196,6 +203,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
reset(mRedis, mContent, mCassandra, mKafka)
val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put(JsonKey.COURSE_ID, "c1")
Expand All @@ -211,13 +219,16 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
val existing = ExistingAssessment("a1", "cont1", System.currentTimeMillis(), System.currentTimeMillis(), 10.0, 10.0, List.empty)
when(mRedis.isValidContent(anyString, anyString)).thenReturn(true)
when(mRedis.getTotalQuestionsCount(anyString)).thenReturn(Some(10))
when(mCassandra.getAssessment(anyString, anyString, anyString, anyString, anyString, any[RequestContext])).thenReturn(Some(existing))
when(mCassandra.getUserAssessments(anyString, anyString, anyString, anyString, any[RequestContext])).thenReturn(List(existing))
when(mCassandra.getAssessment(anyString, anyString, anyString, anyString, anyString, any[RequestContext]))
.thenReturn(Some(existing))
when(mCassandra.getUserAssessments(anyString, anyString, anyString, anyString, any[RequestContext]))
.thenReturn(List(existing))

PropertiesCache.getInstance().saveConfigProperty("assessment_aggregator_publish_certificate", "true")

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put(JsonKey.USER_ID, "u1")
Expand All @@ -239,6 +250,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
reset(mRedis, mContent, mCassandra, mKafka)
val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put("userId", "u1")
Expand All @@ -255,21 +267,39 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
body.put("events", events)
request.setRequest(body)

when(mRedis.isValidContent(anyString, anyString)).thenReturn(true)
when(mRedis.getTotalQuestionsCount(anyString)).thenReturn(Some(10))
when(mCassandra.getAssessment(anyString, anyString, anyString, anyString, anyString, any[RequestContext]))
.thenReturn(Some(ExistingAssessment("att1", "cont1", 2000L, 1000L, 5.0, 10.0, List.empty)))
org.mockito.Mockito.doReturn(true).when(mRedis).isValidContent(org.mockito.ArgumentMatchers.anyString, org.mockito.ArgumentMatchers.anyString)
org.mockito.Mockito.doReturn(Some(10)).when(mRedis).getTotalQuestionsCount(org.mockito.ArgumentMatchers.anyString)

org.mockito.Mockito.doReturn(Some(ExistingAssessment("att1", "cont1", 2000L, 1000L, 5.0, 10.0, List.empty)))
.when(mCassandra).getAssessment(
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.any(classOf[RequestContext])
)

org.mockito.Mockito.doReturn(List(ExistingAssessment("att1", "cont1", 2000L, 1000L, 5.0, 10.0, List.empty)))
.when(mCassandra).getUserAssessments(
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.anyString,
org.mockito.ArgumentMatchers.any(classOf[RequestContext])
)

actorRef ! request
expectMsgType[Response]
verify(mCassandra, never).saveAssessment(any[AssessmentResult], any[RequestContext])
verify(mCassandra, never).saveAssessment(org.mockito.ArgumentMatchers.any(classOf[AssessmentResult]), org.mockito.ArgumentMatchers.any(classOf[RequestContext]))
}

it should "throw exception when content validation fails" in {
reset(mRedis, mContent, mCassandra, mKafka)
PropertiesCache.getInstance().saveConfigProperty("assessment_enable_content_validation", "true")
val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put("userId", "u1")
Expand All @@ -291,6 +321,7 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
// Passing a message that might cause an internal exception (e.g., if a mandatory field is missing in Request)
val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequest(null) // This should cause a NPE in processAggregation
actorRef ! request
expectMsgType[ProjectCommonException].getResponseCode should be (500)
Expand All @@ -300,11 +331,14 @@ class AssessmentAggregatorActorSpec extends TestKit(ActorSystem("AssessmentAggre
reset(mRedis, mContent, mCassandra, mKafka)
when(mRedis.isValidContent(any[String], any[String])).thenReturn(true)
when(mRedis.getTotalQuestionsCount(any[String])).thenReturn(Some(1))
when(mCassandra.getAssessment(any[String], any[String], any[String], any[String], any[String], any[RequestContext])).thenReturn(None)
when(mCassandra.getUserAssessments(any[String], any[String], any[String], any[String], any[RequestContext])).thenReturn(List.empty)
when(mCassandra.getAssessment(any[String], any[String], any[String], any[String], any[String], any[RequestContext]))
.thenReturn(None)
when(mCassandra.getUserAssessments(any[String], any[String], any[String], any[String], any[RequestContext]))
.thenReturn(List.empty)

val actorRef = getActorRef
val request = new Request()
request.setOperation("aggregateAssessment")
request.setRequestContext(mock[RequestContext])
val body = new HashMap[String, AnyRef]()
body.put("userId", "u1")
Expand Down
Loading
Loading