|
|
|
@ -76,7 +76,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
private final Map<CalculatedFieldId, CalculatedFieldCtx> calculatedFields = new HashMap<>(); |
|
|
|
private final Map<EntityId, List<CalculatedFieldCtx>> entityIdCalculatedFields = new HashMap<>(); |
|
|
|
private final Map<EntityId, List<CalculatedFieldLink>> entityIdCalculatedFieldLinks = new HashMap<>(); |
|
|
|
private final Map<CalculatedFieldId, ScheduledFuture<?>> cfInvalidationScheduledTasks = new ConcurrentHashMap<>(); |
|
|
|
private final Map<CalculatedFieldId, ScheduledFuture<?>> cfDynamicArgumentsRefreshTasks = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private final CalculatedFieldProcessingService cfExecService; |
|
|
|
private final CalculatedFieldStateService cfStateService; |
|
|
|
@ -115,8 +115,8 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
calculatedFields.clear(); |
|
|
|
entityIdCalculatedFields.clear(); |
|
|
|
entityIdCalculatedFieldLinks.clear(); |
|
|
|
cfInvalidationScheduledTasks.values().forEach(future -> future.cancel(true)); |
|
|
|
cfInvalidationScheduledTasks.clear(); |
|
|
|
cfDynamicArgumentsRefreshTasks.values().forEach(future -> future.cancel(true)); |
|
|
|
cfDynamicArgumentsRefreshTasks.clear(); |
|
|
|
ctx.stop(ctx.getSelf()); |
|
|
|
} |
|
|
|
|
|
|
|
@ -146,7 +146,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
// We use copy on write lists to safely pass the reference to another actor for the iteration.
|
|
|
|
// Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead)
|
|
|
|
entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cfCtx); |
|
|
|
scheduleCalculatedFieldInvalidationMsgIfNeeded(cfCtx); |
|
|
|
scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(cfCtx); |
|
|
|
msg.getCallback().onSuccess(); |
|
|
|
} |
|
|
|
|
|
|
|
@ -339,7 +339,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
|
|
|
|
boolean hasSchedulingConfigChanges = newCfCtx.hasSchedulingConfigChanges(oldCfCtx); |
|
|
|
if (hasSchedulingConfigChanges) { |
|
|
|
cancelCfScheduledInvalidationTaskIfExists(cfId, false); |
|
|
|
cancelCfDynamicArgumentsRefreshTaskIfExists(cfId, false); |
|
|
|
} |
|
|
|
|
|
|
|
List<CalculatedFieldCtx> newCfList = new CopyOnWriteArrayList<>(); |
|
|
|
@ -382,7 +382,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
entityIdCalculatedFields.get(cfCtx.getEntityId()).remove(cfCtx); |
|
|
|
deleteLinks(cfCtx); |
|
|
|
|
|
|
|
cancelCfScheduledInvalidationTaskIfExists(cfId, true); |
|
|
|
cancelCfDynamicArgumentsRefreshTaskIfExists(cfId, true); |
|
|
|
|
|
|
|
EntityId entityId = cfCtx.getEntityId(); |
|
|
|
EntityType entityType = cfCtx.getEntityId().getEntityType(); |
|
|
|
@ -407,12 +407,12 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void cancelCfScheduledInvalidationTaskIfExists(CalculatedFieldId cfId, boolean cfDeleted) { |
|
|
|
var existingTask = cfInvalidationScheduledTasks.remove(cfId); |
|
|
|
private void cancelCfDynamicArgumentsRefreshTaskIfExists(CalculatedFieldId cfId, boolean cfDeleted) { |
|
|
|
var existingTask = cfDynamicArgumentsRefreshTasks.remove(cfId); |
|
|
|
if (existingTask != null) { |
|
|
|
existingTask.cancel(false); |
|
|
|
String reason = cfDeleted ? "deletion" : "update"; |
|
|
|
log.debug("[{}][{}] Cancelled scheduled invalidation task due to CF " + reason + "!", tenantId, cfId); |
|
|
|
log.debug("[{}][{}] Cancelled dynamic arguments refresh task due to CF " + reason + "!", tenantId, cfId); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -455,7 +455,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
var targetEntityId = link.entityId(); |
|
|
|
var targetEntityType = targetEntityId.getEntityType(); |
|
|
|
var cf = calculatedFields.get(link.cfId()); |
|
|
|
if (EntityType.DEVICE_PROFILE.equals(targetEntityType) || EntityType.ASSET_PROFILE.equals(targetEntityType)) { |
|
|
|
if (isProfileEntity(targetEntityType)) { |
|
|
|
// iterate over all entities that belong to profile and push the message for corresponding CF
|
|
|
|
var entityIds = entityProfileCache.getEntityIdsByProfileId(targetEntityId); |
|
|
|
if (!entityIds.isEmpty()) { |
|
|
|
@ -518,7 +518,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
private void initCf(CalculatedFieldCtx cfCtx, TbCallback callback, boolean forceStateReinit) { |
|
|
|
EntityId entityId = cfCtx.getEntityId(); |
|
|
|
EntityType entityType = cfCtx.getEntityId().getEntityType(); |
|
|
|
scheduleCalculatedFieldInvalidationMsgIfNeeded(cfCtx); |
|
|
|
scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(cfCtx); |
|
|
|
if (isProfileEntity(entityType)) { |
|
|
|
var entityIds = entityProfileCache.getEntityIdsByProfileId(entityId); |
|
|
|
if (!entityIds.isEmpty()) { |
|
|
|
@ -536,31 +536,31 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void scheduleCalculatedFieldInvalidationMsgIfNeeded(CalculatedFieldCtx cfCtx) { |
|
|
|
private void scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(CalculatedFieldCtx cfCtx) { |
|
|
|
CalculatedField cf = cfCtx.getCalculatedField(); |
|
|
|
CalculatedFieldConfiguration cfConfig = cf.getConfiguration(); |
|
|
|
if (!cfConfig.isScheduledUpdateEnabled()) { |
|
|
|
return; |
|
|
|
} |
|
|
|
if (cfInvalidationScheduledTasks.containsKey(cf.getId())) { |
|
|
|
log.debug("[{}][{}] Scheduled invalidation task for CF already exists!", tenantId, cf.getId()); |
|
|
|
if (cfDynamicArgumentsRefreshTasks.containsKey(cf.getId())) { |
|
|
|
log.debug("[{}][{}] Dynamic arguments refresh task for CF already exists!", tenantId, cf.getId()); |
|
|
|
return; |
|
|
|
} |
|
|
|
long refreshDynamicSourceInterval = TimeUnit.SECONDS.toMillis(cfConfig.getScheduledUpdateIntervalSec()); |
|
|
|
var scheduledMsg = new CalculatedFieldScheduledInvalidationMsg(tenantId, cfCtx.getCfId()); |
|
|
|
var scheduledMsg = new CalculatedFieldDynamicArgumentsRefreshMsg(tenantId, cfCtx.getCfId()); |
|
|
|
|
|
|
|
ScheduledFuture<?> scheduledFuture = systemContext |
|
|
|
.schedulePeriodicMsgWithDelay(ctx, scheduledMsg, refreshDynamicSourceInterval, refreshDynamicSourceInterval); |
|
|
|
cfInvalidationScheduledTasks.put(cf.getId(), scheduledFuture); |
|
|
|
log.debug("[{}][{}] Scheduled invalidation task for CF!", tenantId, cf.getId()); |
|
|
|
cfDynamicArgumentsRefreshTasks.put(cf.getId(), scheduledFuture); |
|
|
|
log.debug("[{}][{}] Scheduled dynamic arguments refresh task for CF!", tenantId, cf.getId()); |
|
|
|
} |
|
|
|
|
|
|
|
public void onScheduledInvalidationMsg(CalculatedFieldScheduledInvalidationMsg msg) { |
|
|
|
log.debug("[{}] [{}] Processing CF scheduled invalidation msg.", tenantId, msg.getCfId()); |
|
|
|
public void onDynamicArgumentsRefreshMsg(CalculatedFieldDynamicArgumentsRefreshMsg msg) { |
|
|
|
log.debug("[{}] [{}] Processing CF dynamic arguments refresh task.", tenantId, msg.getCfId()); |
|
|
|
CalculatedFieldCtx cfCtx = calculatedFields.get(msg.getCfId()); |
|
|
|
if (cfCtx == null) { |
|
|
|
log.debug("[{}][{}] Failed to find CF context, going to stop scheduled invalidations for CF.", tenantId, msg.getCfId()); |
|
|
|
cancelCfScheduledInvalidationTaskIfExists(msg.getCfId(), true); |
|
|
|
log.debug("[{}][{}] Failed to find CF context, going to stop dynamic arguments refresh task for CF.", tenantId, msg.getCfId()); |
|
|
|
cancelCfDynamicArgumentsRefreshTaskIfExists(msg.getCfId(), true); |
|
|
|
return; |
|
|
|
} |
|
|
|
EntityId entityId = cfCtx.getEntityId(); |
|
|
|
@ -571,7 +571,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
var multiCallback = new MultipleTbCallback(entityIds.size(), msg.getCallback()); |
|
|
|
entityIds.forEach(id -> { |
|
|
|
if (isMyPartition(id, multiCallback)) { |
|
|
|
InitCfInvalidationForEntity(id, msg.getCfId(), multiCallback); |
|
|
|
dynamicArgumentsRefreshForEntity(id, msg.getCfId(), multiCallback); |
|
|
|
} |
|
|
|
}); |
|
|
|
} else { |
|
|
|
@ -579,14 +579,14 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
} |
|
|
|
} else { |
|
|
|
if (isMyPartition(entityId, msg.getCallback())) { |
|
|
|
InitCfInvalidationForEntity(entityId, msg.getCfId(), msg.getCallback()); |
|
|
|
dynamicArgumentsRefreshForEntity(entityId, msg.getCfId(), msg.getCallback()); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void InitCfInvalidationForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) { |
|
|
|
log.debug("Pushing entity CF invalidation msg to specific actor [{}]", entityId); |
|
|
|
getOrCreateActor(entityId).tell(new EntityCalculatedFieldMarkStateDirtyMsg(tenantId, cfId, callback)); |
|
|
|
private void dynamicArgumentsRefreshForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) { |
|
|
|
log.debug("Pushing CF dynamic arguments refresh msg to specific actor [{}]", entityId); |
|
|
|
getOrCreateActor(entityId).tell(new EntityCalculatedFieldDynamicArgumentsRefreshMsg(tenantId, cfId, callback)); |
|
|
|
} |
|
|
|
|
|
|
|
private void deleteCfForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) { |
|
|
|
@ -652,10 +652,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware |
|
|
|
log.error("Failed to process calculated field record: {}", cf, e); |
|
|
|
} |
|
|
|
}); |
|
|
|
// TODO: why we need to do this loop if we do this inside the onFieldInitMsg?
|
|
|
|
calculatedFields.values().forEach(cf -> { |
|
|
|
entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cf); |
|
|
|
}); |
|
|
|
PageDataIterable<CalculatedFieldLink> cfls = new PageDataIterable<>(pageLink -> cfDaoService.findAllCalculatedFieldLinksByTenantId(tenantId, pageLink), cfSettings.getInitTenantFetchPackSize()); |
|
|
|
cfls.forEach(link -> { |
|
|
|
onLinkInitMsg(new CalculatedFieldLinkInitMsg(link.getTenantId(), link)); |
|
|
|
|