|
|
@ -19,7 +19,8 @@ import com.datastax.driver.core.BoundStatement; |
|
|
import com.datastax.driver.core.PreparedStatement; |
|
|
import com.datastax.driver.core.PreparedStatement; |
|
|
import com.datastax.driver.core.ResultSet; |
|
|
import com.datastax.driver.core.ResultSet; |
|
|
import com.datastax.driver.core.ResultSetFuture; |
|
|
import com.datastax.driver.core.ResultSetFuture; |
|
|
import com.datastax.driver.core.utils.UUIDs; |
|
|
import com.datastax.driver.core.querybuilder.QueryBuilder; |
|
|
|
|
|
import com.datastax.driver.core.querybuilder.Select; |
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
@ -29,8 +30,9 @@ import org.springframework.beans.factory.annotation.Value; |
|
|
import org.springframework.core.env.Environment; |
|
|
import org.springframework.core.env.Environment; |
|
|
import org.springframework.stereotype.Component; |
|
|
import org.springframework.stereotype.Component; |
|
|
import org.thingsboard.server.common.data.audit.AuditLog; |
|
|
import org.thingsboard.server.common.data.audit.AuditLog; |
|
|
import org.thingsboard.server.common.data.id.AuditLogId; |
|
|
import org.thingsboard.server.common.data.id.CustomerId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
|
|
import org.thingsboard.server.common.data.id.UserId; |
|
|
import org.thingsboard.server.common.data.page.TimePageLink; |
|
|
import org.thingsboard.server.common.data.page.TimePageLink; |
|
|
import org.thingsboard.server.dao.DaoUtil; |
|
|
import org.thingsboard.server.dao.DaoUtil; |
|
|
import org.thingsboard.server.dao.model.ModelConstants; |
|
|
import org.thingsboard.server.dao.model.ModelConstants; |
|
|
@ -52,6 +54,7 @@ import java.util.Optional; |
|
|
import java.util.UUID; |
|
|
import java.util.UUID; |
|
|
import java.util.concurrent.ExecutorService; |
|
|
import java.util.concurrent.ExecutorService; |
|
|
import java.util.concurrent.Executors; |
|
|
import java.util.concurrent.Executors; |
|
|
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
|
import static com.datastax.driver.core.querybuilder.QueryBuilder.eq; |
|
|
import static com.datastax.driver.core.querybuilder.QueryBuilder.eq; |
|
|
import static org.thingsboard.server.dao.model.ModelConstants.*; |
|
|
import static org.thingsboard.server.dao.model.ModelConstants.*; |
|
|
@ -82,7 +85,11 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
private String partitioning; |
|
|
private String partitioning; |
|
|
private TsPartitionDate tsFormat; |
|
|
private TsPartitionDate tsFormat; |
|
|
|
|
|
|
|
|
private PreparedStatement[] saveStmts; |
|
|
private PreparedStatement partitionInsertStmt; |
|
|
|
|
|
private PreparedStatement saveByTenantStmt; |
|
|
|
|
|
private PreparedStatement saveByTenantIdAndUserIdStmt; |
|
|
|
|
|
private PreparedStatement saveByTenantIdAndEntityIdStmt; |
|
|
|
|
|
private PreparedStatement saveByTenantIdAndCustomerIdStmt; |
|
|
|
|
|
|
|
|
private boolean isInstall() { |
|
|
private boolean isInstall() { |
|
|
return environment.acceptsProfiles("install"); |
|
|
return environment.acceptsProfiles("install"); |
|
|
@ -123,11 +130,9 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
public ListenableFuture<Void> saveByTenantId(AuditLog auditLog) { |
|
|
public ListenableFuture<Void> saveByTenantId(AuditLog auditLog) { |
|
|
log.debug("Save saveByTenantId [{}] ", auditLog); |
|
|
log.debug("Save saveByTenantId [{}] ", auditLog); |
|
|
|
|
|
|
|
|
AuditLogId auditLogId = new AuditLogId(UUIDs.timeBased()); |
|
|
|
|
|
|
|
|
|
|
|
long partition = toPartitionTs(LocalDate.now().atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
long partition = toPartitionTs(LocalDate.now().atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
BoundStatement stmt = getSaveByTenantStmt().bind(); |
|
|
BoundStatement stmt = getSaveByTenantStmt().bind(); |
|
|
stmt.setUUID(0, auditLogId.getId()) |
|
|
stmt = stmt.setUUID(0, auditLog.getId().getId()) |
|
|
.setUUID(1, auditLog.getTenantId().getId()) |
|
|
.setUUID(1, auditLog.getTenantId().getId()) |
|
|
.setUUID(2, auditLog.getEntityId().getId()) |
|
|
.setUUID(2, auditLog.getEntityId().getId()) |
|
|
.setString(3, auditLog.getEntityId().getEntityType().name()) |
|
|
.setString(3, auditLog.getEntityId().getEntityType().name()) |
|
|
@ -140,17 +145,44 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
public ListenableFuture<Void> saveByTenantIdAndEntityId(AuditLog auditLog) { |
|
|
public ListenableFuture<Void> saveByTenantIdAndEntityId(AuditLog auditLog) { |
|
|
log.debug("Save saveByTenantIdAndEntityId [{}] ", auditLog); |
|
|
log.debug("Save saveByTenantIdAndEntityId [{}] ", auditLog); |
|
|
|
|
|
|
|
|
AuditLogId auditLogId = new AuditLogId(UUIDs.timeBased()); |
|
|
|
|
|
|
|
|
|
|
|
BoundStatement stmt = getSaveByTenantIdAndEntityIdStmt().bind(); |
|
|
BoundStatement stmt = getSaveByTenantIdAndEntityIdStmt().bind(); |
|
|
stmt.setUUID(0, auditLogId.getId()) |
|
|
stmt = setSaveStmtVariables(stmt, auditLog); |
|
|
.setUUID(1, auditLog.getTenantId().getId()) |
|
|
|
|
|
.setUUID(2, auditLog.getEntityId().getId()) |
|
|
|
|
|
.setString(3, auditLog.getEntityId().getEntityType().name()) |
|
|
|
|
|
.setString(4, auditLog.getActionType().name()); |
|
|
|
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public ListenableFuture<Void> saveByTenantIdAndCustomerId(AuditLog auditLog) { |
|
|
|
|
|
log.debug("Save saveByTenantIdAndCustomerId [{}] ", auditLog); |
|
|
|
|
|
|
|
|
|
|
|
BoundStatement stmt = getSaveByTenantIdAndCustomerIdStmt().bind(); |
|
|
|
|
|
stmt = setSaveStmtVariables(stmt, auditLog); |
|
|
|
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public ListenableFuture<Void> saveByTenantIdAndUserId(AuditLog auditLog) { |
|
|
|
|
|
log.debug("Save saveByTenantIdAndUserId [{}] ", auditLog); |
|
|
|
|
|
|
|
|
|
|
|
BoundStatement stmt = getSaveByTenantIdAndUserIdStmt().bind(); |
|
|
|
|
|
stmt = setSaveStmtVariables(stmt, auditLog); |
|
|
|
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private BoundStatement setSaveStmtVariables(BoundStatement stmt, AuditLog auditLog) { |
|
|
|
|
|
return stmt.setUUID(0, auditLog.getId().getId()) |
|
|
|
|
|
.setUUID(1, auditLog.getTenantId().getId()) |
|
|
|
|
|
.setUUID(2, auditLog.getCustomerId().getId()) |
|
|
|
|
|
.setUUID(3, auditLog.getEntityId().getId()) |
|
|
|
|
|
.setString(4, auditLog.getEntityId().getEntityType().name()) |
|
|
|
|
|
.setString(5, auditLog.getEntityName()) |
|
|
|
|
|
.setUUID(6, auditLog.getUserId().getId()) |
|
|
|
|
|
.setString(7, auditLog.getUserName()) |
|
|
|
|
|
.setString(8, auditLog.getActionType().name()) |
|
|
|
|
|
.setString(9, auditLog.getActionData() != null ? auditLog.getActionData().toString() : null) |
|
|
|
|
|
.setString(10, auditLog.getActionStatus().name()) |
|
|
|
|
|
.setString(11, auditLog.getActionFailureDetails()); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> savePartitionsByTenantId(AuditLog auditLog) { |
|
|
public ListenableFuture<Void> savePartitionsByTenantId(AuditLog auditLog) { |
|
|
log.debug("Save savePartitionsByTenantId [{}] ", auditLog); |
|
|
log.debug("Save savePartitionsByTenantId [{}] ", auditLog); |
|
|
@ -163,35 +195,66 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
return getFuture(executeAsyncWrite(stmt), rs -> null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private PreparedStatement getPartitionInsertStmt() { |
|
|
private PreparedStatement getSaveByTenantIdAndEntityIdStmt() { |
|
|
// TODO: ADD CACHE LOGIC
|
|
|
if (saveByTenantIdAndEntityIdStmt == null) { |
|
|
return getSession().prepare(INSERT_INTO + ModelConstants.AUDIT_LOG_BY_TENANT_ID_PARTITIONS_CF + |
|
|
saveByTenantIdAndEntityIdStmt = getSaveByTenantIdAndCFName(ModelConstants.AUDIT_LOG_BY_ENTITY_ID_CF); |
|
|
"(" + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
} |
|
|
"," + ModelConstants.AUDIT_LOG_PARTITION_PROPERTY + ")" + |
|
|
return saveByTenantIdAndEntityIdStmt; |
|
|
" VALUES(?, ?)"); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private PreparedStatement getSaveByTenantIdAndEntityIdStmt() { |
|
|
private PreparedStatement getSaveByTenantIdAndCustomerIdStmt() { |
|
|
// TODO: ADD CACHE LOGIC
|
|
|
if (saveByTenantIdAndCustomerIdStmt == null) { |
|
|
return getSession().prepare(INSERT_INTO + ModelConstants.AUDIT_LOG_BY_ENTITY_ID_CF + |
|
|
saveByTenantIdAndCustomerIdStmt = getSaveByTenantIdAndCFName(ModelConstants.AUDIT_LOG_BY_CUSTOMER_ID_CF); |
|
|
"(" + ModelConstants.AUDIT_LOG_ID_PROPERTY + |
|
|
} |
|
|
"," + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
return saveByTenantIdAndCustomerIdStmt; |
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_TYPE_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_TYPE_PROPERTY + ")" + |
|
|
|
|
|
" VALUES(?, ?, ?, ?, ?)"); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private PreparedStatement getSaveByTenantStmt() { |
|
|
private PreparedStatement getSaveByTenantIdAndUserIdStmt() { |
|
|
// TODO: ADD CACHE LOGIC
|
|
|
if (saveByTenantIdAndUserIdStmt == null) { |
|
|
return getSession().prepare(INSERT_INTO + ModelConstants.AUDIT_LOG_BY_TENANT_ID_CF + |
|
|
saveByTenantIdAndUserIdStmt = getSaveByTenantIdAndCFName(ModelConstants.AUDIT_LOG_BY_USER_ID_CF); |
|
|
|
|
|
} |
|
|
|
|
|
return saveByTenantIdAndUserIdStmt; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private PreparedStatement getSaveByTenantIdAndCFName(String cfName) { |
|
|
|
|
|
return getSession().prepare(INSERT_INTO + cfName + |
|
|
"(" + ModelConstants.AUDIT_LOG_ID_PROPERTY + |
|
|
"(" + ModelConstants.AUDIT_LOG_ID_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_CUSTOMER_ID_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_ID_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_ID_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_TYPE_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_TYPE_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_NAME_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_USER_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_USER_NAME_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_TYPE_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_TYPE_PROPERTY + |
|
|
"," + ModelConstants.AUDIT_LOG_PARTITION_PROPERTY + ")" + |
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_DATA_PROPERTY + |
|
|
" VALUES(?, ?, ?, ?, ?, ?)"); |
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_STATUS_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_FAILURE_DETAILS_PROPERTY + ")" + |
|
|
|
|
|
" VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private PreparedStatement getPartitionInsertStmt() { |
|
|
|
|
|
if (partitionInsertStmt == null) { |
|
|
|
|
|
partitionInsertStmt = getSession().prepare(INSERT_INTO + ModelConstants.AUDIT_LOG_BY_TENANT_ID_PARTITIONS_CF + |
|
|
|
|
|
"(" + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_PARTITION_PROPERTY + ")" + |
|
|
|
|
|
" VALUES(?, ?)"); |
|
|
|
|
|
} |
|
|
|
|
|
return partitionInsertStmt; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private PreparedStatement getSaveByTenantStmt() { |
|
|
|
|
|
if (saveByTenantStmt == null) { |
|
|
|
|
|
saveByTenantStmt = getSession().prepare(INSERT_INTO + ModelConstants.AUDIT_LOG_BY_TENANT_ID_CF + |
|
|
|
|
|
"(" + ModelConstants.AUDIT_LOG_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_ID_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ENTITY_TYPE_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_ACTION_TYPE_PROPERTY + |
|
|
|
|
|
"," + ModelConstants.AUDIT_LOG_PARTITION_PROPERTY + ")" + |
|
|
|
|
|
" VALUES(?, ?, ?, ?, ?, ?)"); |
|
|
|
|
|
} |
|
|
|
|
|
return saveByTenantStmt; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private long toPartitionTs(long ts) { |
|
|
private long toPartitionTs(long ts) { |
|
|
@ -199,7 +262,6 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); |
|
|
return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public List<AuditLog> findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, TimePageLink pageLink) { |
|
|
public List<AuditLog> findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, TimePageLink pageLink) { |
|
|
log.trace("Try to find audit logs by tenant [{}], entity [{}] and pageLink [{}]", tenantId, entityId, pageLink); |
|
|
log.trace("Try to find audit logs by tenant [{}], entity [{}] and pageLink [{}]", tenantId, entityId, pageLink); |
|
|
@ -212,31 +274,76 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao<AuditLo |
|
|
return DaoUtil.convertDataList(entities); |
|
|
return DaoUtil.convertDataList(entities); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public List<AuditLog> findAuditLogsByTenantIdAndCustomerId(UUID tenantId, CustomerId customerId, TimePageLink pageLink) { |
|
|
|
|
|
log.trace("Try to find audit logs by tenant [{}], customer [{}] and pageLink [{}]", tenantId, customerId, pageLink); |
|
|
|
|
|
List<AuditLogEntity> entities = findPageWithTimeSearch(AUDIT_LOG_BY_CUSTOMER_ID_CF, |
|
|
|
|
|
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId), |
|
|
|
|
|
eq(ModelConstants.AUDIT_LOG_CUSTOMER_ID_PROPERTY, customerId.getId())), |
|
|
|
|
|
pageLink); |
|
|
|
|
|
log.trace("Found audit logs by tenant [{}], customer [{}] and pageLink [{}]", tenantId, customerId, pageLink); |
|
|
|
|
|
return DaoUtil.convertDataList(entities); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public List<AuditLog> findAuditLogsByTenantIdAndUserId(UUID tenantId, UserId userId, TimePageLink pageLink) { |
|
|
|
|
|
log.trace("Try to find audit logs by tenant [{}], user [{}] and pageLink [{}]", tenantId, userId, pageLink); |
|
|
|
|
|
List<AuditLogEntity> entities = findPageWithTimeSearch(AUDIT_LOG_BY_USER_ID_CF, |
|
|
|
|
|
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId), |
|
|
|
|
|
eq(ModelConstants.AUDIT_LOG_USER_ID_PROPERTY, userId.getId())), |
|
|
|
|
|
pageLink); |
|
|
|
|
|
log.trace("Found audit logs by tenant [{}], user [{}] and pageLink [{}]", tenantId, userId, pageLink); |
|
|
|
|
|
return DaoUtil.convertDataList(entities); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public List<AuditLog> findAuditLogsByTenantId(UUID tenantId, TimePageLink pageLink) { |
|
|
public List<AuditLog> findAuditLogsByTenantId(UUID tenantId, TimePageLink pageLink) { |
|
|
log.trace("Try to find audit logs by tenant [{}] and pageLink [{}]", tenantId, pageLink); |
|
|
log.trace("Try to find audit logs by tenant [{}] and pageLink [{}]", tenantId, pageLink); |
|
|
|
|
|
|
|
|
// TODO: ADD AUDIT LOG PARTITION CURSOR LOGIC
|
|
|
|
|
|
|
|
|
|
|
|
long minPartition; |
|
|
long minPartition; |
|
|
long maxPartition; |
|
|
|
|
|
|
|
|
|
|
|
if (pageLink.getStartTime() != null && pageLink.getStartTime() != 0) { |
|
|
if (pageLink.getStartTime() != null && pageLink.getStartTime() != 0) { |
|
|
minPartition = toPartitionTs(pageLink.getStartTime()); |
|
|
minPartition = toPartitionTs(pageLink.getStartTime()); |
|
|
} else { |
|
|
} else { |
|
|
minPartition = toPartitionTs(LocalDate.now().minusMonths(1).atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
minPartition = toPartitionTs(LocalDate.now().minusMonths(1).atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
long maxPartition; |
|
|
if (pageLink.getEndTime() != null && pageLink.getEndTime() != 0) { |
|
|
if (pageLink.getEndTime() != null && pageLink.getEndTime() != 0) { |
|
|
maxPartition = toPartitionTs(pageLink.getEndTime()); |
|
|
maxPartition = toPartitionTs(pageLink.getEndTime()); |
|
|
} else { |
|
|
} else { |
|
|
maxPartition = toPartitionTs(LocalDate.now().atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
maxPartition = toPartitionTs(LocalDate.now().atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli()); |
|
|
} |
|
|
} |
|
|
List<AuditLogEntity> entities = findPageWithTimeSearch(AUDIT_LOG_BY_TENANT_ID_CF, |
|
|
|
|
|
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId), |
|
|
List<Long> partitions = fetchPartitions(tenantId, minPartition, maxPartition) |
|
|
eq(ModelConstants.AUDIT_LOG_PARTITION_PROPERTY, maxPartition)), |
|
|
.all() |
|
|
pageLink); |
|
|
.stream() |
|
|
|
|
|
.map(row -> row.getLong(ModelConstants.PARTITION_COLUMN)) |
|
|
|
|
|
.collect(Collectors.toList()); |
|
|
|
|
|
|
|
|
|
|
|
AuditLogQueryCursor cursor = new AuditLogQueryCursor(tenantId, pageLink, partitions); |
|
|
|
|
|
List<AuditLogEntity> entities = fetchSequentiallyWithLimit(cursor); |
|
|
log.trace("Found audit logs by tenant [{}] and pageLink [{}]", tenantId, pageLink); |
|
|
log.trace("Found audit logs by tenant [{}] and pageLink [{}]", tenantId, pageLink); |
|
|
return DaoUtil.convertDataList(entities); |
|
|
return DaoUtil.convertDataList(entities); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private List<AuditLogEntity> fetchSequentiallyWithLimit(AuditLogQueryCursor cursor) { |
|
|
|
|
|
if (cursor.isFull() || !cursor.hasNextPartition()) { |
|
|
|
|
|
return cursor.getData(); |
|
|
|
|
|
} else { |
|
|
|
|
|
cursor.addData(findPageWithTimeSearch(AUDIT_LOG_BY_TENANT_ID_CF, |
|
|
|
|
|
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, cursor.getTenantId()), |
|
|
|
|
|
eq(ModelConstants.AUDIT_LOG_PARTITION_PROPERTY, cursor.getNextPartition())), |
|
|
|
|
|
cursor.getPageLink())); |
|
|
|
|
|
return fetchSequentiallyWithLimit(cursor); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private ResultSet fetchPartitions(UUID tenantId, long minPartition, long maxPartition) { |
|
|
|
|
|
Select.Where select = QueryBuilder.select(ModelConstants.AUDIT_LOG_PARTITION_PROPERTY).from(ModelConstants.AUDIT_LOG_BY_TENANT_ID_PARTITIONS_CF) |
|
|
|
|
|
.where(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId)); |
|
|
|
|
|
select.and(QueryBuilder.gte(ModelConstants.PARTITION_COLUMN, minPartition)); |
|
|
|
|
|
select.and(QueryBuilder.lte(ModelConstants.PARTITION_COLUMN, maxPartition)); |
|
|
|
|
|
return getSession().execute(select); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|