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 @@ -20,6 +20,8 @@
package org.apache.iotdb.db.i18n;

public final class DataNodeQueryMessages {
public static final String MESSAGE_SCHEMA_REGION_IS_UNAVAILABLE_PLEASE_RETRY_LATER_D642510A =
"Schema region is unavailable. Please retry later.";

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Adds the retryable missing-region message to the English message catalog so the new validation result participates in the existing compile-time localization mechanism.


public static final String EXCEPTION_FRAGMENT_INSTANCE_ARG_IS_ALREADY_ARG_B44984B4 =
"Fragment instance %s is already %s";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
package org.apache.iotdb.db.i18n;

public final class DataNodeQueryMessages {
public static final String MESSAGE_SCHEMA_REGION_IS_UNAVAILABLE_PLEASE_RETRY_LATER_D642510A =
"SchemaRegion 不可用,请稍后重试。";

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Adds the same message key in Chinese to preserve locale parity. DataNode and ConfigNode compiled successfully with the Chinese profile; the full reactor stopped later at distribution because SNAPSHOT ZIP dependencies were unavailable.


public static final String EXCEPTION_FRAGMENT_INSTANCE_ARG_IS_ALREADY_ARG_B44984B4 =
"Fragment instance %s 已处于 %s 状态";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -842,16 +842,21 @@ private RegionExecutionResult executeInternalCreateMultiTimeSeries(
}

/**
* Check the quota before creating time series.
* Check region availability, data types and quota before creating time series.
*
* @return null if the quota is not exceeded, otherwise return the execution result.
* @return null if all checks pass, otherwise return the execution result.
*/
private RegionExecutionResult checkQuotaAndTypeBeforeCreatingTimeSeries(
final ISchemaRegion schemaRegion,
final PartialPath path,
final int size,
final List<String> measurements,
final List<TSDataType> dataTypes) {
// Shutdown may clear SchemaEngine while the region's validation lock still exists.
// Reject the request as retryable before dereferencing the missing region.
if (schemaRegion == null) {
return createSchemaRegionUnavailableResult();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This reuses the centralized retryable result for internal timeseries creation after the original guard. Keeping one result builder ensures shutdown races consistently return NO_AVAILABLE_REGION_GROUP (906) with the localized message across all schema-write prechecks; live-region quota and type validation remain unchanged.

}
for (int i = 0; i < measurements.size(); ++i) {
if (dataTypes.get(i) == TSDataType.OBJECT) {
final String errorStr =
Expand Down Expand Up @@ -933,6 +938,9 @@ private RegionExecutionResult executeAlterTimeSeries(
AlterTimeSeriesNode node, WritePlanNodeExecutionContext context, boolean receivedFromPipe) {
ISchemaRegion schemaRegion =
schemaEngine.getSchemaRegion((SchemaRegionId) context.getRegionId());
if (schemaRegion == null) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AlterTimeSeries fetches its measurement before consensus, so a cleared SchemaEngine entry could otherwise be dereferenced while the validation lock still exists. Returning the same retryable 906 result here preserves the shutdown behavior established for timeseries creation; the available-region path remains unchanged.

return createSchemaRegionUnavailableResult();
}
try {
MeasurementPath measurementPath = schemaRegion.fetchMeasurementPath(node.getPath());
if (node.isAlterView() && !measurementPath.getMeasurementSchema().isLogicalView()) {
Expand Down Expand Up @@ -1123,6 +1131,9 @@ private RegionExecutionResult executeCreateLogicalView(
final ISchemaRegion schemaRegion =
schemaEngine.getSchemaRegion((SchemaRegionId) context.getRegionId());
if (CONFIG.getSchemaRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)) {
if (schemaRegion == null) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ratis validates logical-view targets locally before writing consensus state. Guarding the missing region prevents a shutdown-time null dereference and makes the request retryable; SimpleConsensus continues to use its existing direct write path.

return createSchemaRegionUnavailableResult();
}
context.getRegionWriteValidationRWLock().writeLock().lock();
try {
// step 1. make sure all target paths do NOT exist.
Expand Down Expand Up @@ -1180,6 +1191,9 @@ private RegionExecutionResult executeCreateOrUpdateTableDevice(
final boolean receivedFromPipe) {
final ISchemaRegion schemaRegion =
schemaEngine.getSchemaRegion((SchemaRegionId) context.getRegionId());
if (schemaRegion == null) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Table-device quota validation also reads the schema region before consensus. This guard covers the same shutdown window as the timeseries path and returns 906 before checkSchemaQuota is called, while preserving normal validation for an available region.

return createSchemaRegionUnavailableResult();
}
try {
schemaRegion.checkSchemaQuota(node.getTableName(), node.getDeviceIdList());
} catch (final SchemaQuotaExceededException e) {
Expand All @@ -1192,6 +1206,15 @@ private RegionExecutionResult executeCreateOrUpdateTableDevice(
: PlanVisitor.super.visitCreateOrUpdateTableDevice(node, context);
}

private RegionExecutionResult createSchemaRegionUnavailableResult() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Centralizing the unavailable-region result keeps all guarded schema-write paths aligned on the retryable status and localized message. The helper also documents why a still-registered validation lock must not be treated as proof that the region is usable.

// A validation lock can outlive its region during shutdown, so all local schema checks
// must report a missing region as retryable even when the lock is still registered.
final String message =
DataNodeQueryMessages.MESSAGE_SCHEMA_REGION_IS_UNAVAILABLE_PLEASE_RETRY_LATER_D642510A;
return RegionExecutionResult.create(
false, message, RpcUtils.getStatus(TSStatusCode.NO_AVAILABLE_REGION_GROUP, message));
}

@Override
public RegionExecutionResult visitPipeEnrichedWritePlanNode(
final PipeEnrichedWritePlanNode node, final WritePlanNodeExecutionContext context) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,32 +22,248 @@
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.consensus.DataRegionId;
import org.apache.iotdb.commons.consensus.SchemaRegionId;
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.path.MeasurementPath;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId;
import org.apache.iotdb.commons.schema.view.viewExpression.leaf.TimeSeriesViewOperand;
import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.consensus.IConsensus;
import org.apache.iotdb.consensus.exception.ConsensusException;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.protocol.thrift.impl.DataNodeRegionManager;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.WritePlanNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.AlterTimeSeriesNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.InternalCreateMultiTimeSeriesNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.InternalCreateTimeSeriesNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.MeasurementGroup;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.view.CreateLogicalViewNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedWritePlanNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
import org.apache.iotdb.db.queryengine.plan.relational.planner.node.schema.CreateOrUpdateTableDeviceNode;
import org.apache.iotdb.db.queryengine.plan.statement.metadata.AlterTimeSeriesStatement.AlterType;
import org.apache.iotdb.db.schemaengine.SchemaEngine;
import org.apache.iotdb.db.schemaengine.schemaregion.ISchemaRegion;
import org.apache.iotdb.db.schemaengine.template.ClusterTemplateManager;
import org.apache.iotdb.db.trigger.executor.TriggerFireResult;
import org.apache.iotdb.db.trigger.executor.TriggerFireVisitor;
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.trigger.api.enums.TriggerEvent;

import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.enums.CompressionType;
import org.apache.tsfile.file.metadata.enums.TSEncoding;
import org.apache.tsfile.utils.Pair;
import org.junit.Test;
import org.mockito.Mockito;

import java.util.Collections;
import java.util.concurrent.locks.ReentrantReadWriteLock;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;

public class RegionWriteExecutorTest {

@Test
public void testAlterTimeSeriesRegionAvailability() throws Exception {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This regression test exercises AlterTimeSeries with the validation lock retained while the schema region is absent. It verifies the new retryable 906 result and the unchanged successful path for ordinary and Pipe requests under both consensus protocols.

// Alter validation fetches the measurement before consensus; a missing region must be
// retryable.
checkSchemaWriteWithRegionAvailability(
new AlterTimeSeriesNode(
new PlanNodeId("alter"),
new MeasurementPath("root.sg.d.s"),
AlterType.ADD_TAGS,
Collections.singletonMap("tag", "value"),
null,
null,
null,
false));
}

@Test
public void testCreateLogicalViewRegionAvailability() throws Exception {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test covers both logical-view consensus modes: Ratis performs local target validation and must reject a missing region with 906, while SimpleConsensus keeps its direct write behavior. The matrix also checks available regions, Pipe wrapping, and lock cleanup.

// Ratis checks target existence locally. SimpleConsensus must retain its direct write path.
checkSchemaWriteWithRegionAvailability(
new CreateLogicalViewNode(
new PlanNodeId("view"),
Collections.singletonMap(
new MeasurementPath("root.sg.d.v"), new TimeSeriesViewOperand("root.sg.d.s"))));
}

@Test
public void testCreateOrUpdateTableDeviceRegionAvailability() throws Exception {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This regression test covers table-device quota validation during the same region-cleared window. It verifies retryable failure without consensus writes, plus successful quota validation and writes when the region is available.

// Table-device quota validation has the same missing-region window as tree auto-creation.
checkSchemaWriteWithRegionAvailability(
new CreateOrUpdateTableDeviceNode(
new PlanNodeId("table"),
"db",
"table1",
Collections.singletonList(new Object[] {"device1"}),
Collections.emptyList(),
Collections.singletonList(new Object[0])));
}

private void checkSchemaWriteWithRegionAvailability(final WritePlanNode node) throws Exception {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The shared test matrix keeps the similar-path checks consistent across missing and available regions, both consensus protocols, and ordinary/Pipe requests. It also verifies the validation lock is released and that consensus is skipped only for the unavailable-region case.

final String originalProtocol =
IoTDBDescriptor.getInstance().getConfig().getSchemaRegionConsensusProtocolClass();
try {
for (String protocol :
new String[] {ConsensusFactory.RATIS_CONSENSUS, ConsensusFactory.SIMPLE_CONSENSUS}) {
IoTDBDescriptor.getInstance().getConfig().setSchemaRegionConsensusProtocolClass(protocol);
for (boolean available : new boolean[] {false, true}) {
for (boolean fromPipe : new boolean[] {false, true}) {
final IConsensus consensus = Mockito.mock(IConsensus.class);
final DataNodeRegionManager regionManager = Mockito.mock(DataNodeRegionManager.class);
final SchemaEngine schemaEngine = Mockito.mock(SchemaEngine.class);
final SchemaRegionId regionId = new SchemaRegionId(1);
final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
// Retain the lock even when the region has already been cleared during shutdown.
Mockito.when(regionManager.getRegionLock(regionId)).thenReturn(lock);
final ISchemaRegion region = Mockito.mock(ISchemaRegion.class);
Mockito.when(schemaEngine.getSchemaRegion(regionId))
.thenReturn(available ? region : null);
if (node instanceof AlterTimeSeriesNode alterNode) {
Mockito.when(region.fetchMeasurementPath(alterNode.getPath()))
.thenReturn(new MeasurementPath("root.sg.d.s", TSDataType.INT64));
}
final RegionWriteExecutor executor =
new RegionWriteExecutor(
Mockito.mock(IConsensus.class),
consensus,
regionManager,
schemaEngine,
Mockito.mock(ClusterTemplateManager.class),
Mockito.mock(TriggerFireVisitor.class));
final WritePlanNode request = fromPipe ? new PipeEnrichedWritePlanNode(node) : node;
Mockito.when(consensus.write(regionId, request))
.thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));

final boolean validatesRegion =
!(node instanceof CreateLogicalViewNode)
|| ConsensusFactory.RATIS_CONSENSUS.equals(protocol);
final boolean accepted = available || !validatesRegion;
final RegionExecutionResult result = executor.execute(regionId, request);
assertEquals(accepted, result.isAccepted());
assertEquals(
(accepted ? TSStatusCode.SUCCESS_STATUS : TSStatusCode.NO_AVAILABLE_REGION_GROUP)
.getStatusCode(),
result.getStatus().getCode());
assertFalse(lock.isWriteLocked());
assertEquals(0, lock.getReadLockCount());
if (accepted) {
Mockito.verify(consensus).write(regionId, request);
if (node instanceof AlterTimeSeriesNode alterNode) {
Mockito.verify(region).fetchMeasurementPath(alterNode.getPath());
} else if (node instanceof CreateOrUpdateTableDeviceNode tableNode) {
Mockito.verify(region)
.checkSchemaQuota(tableNode.getTableName(), tableNode.getDeviceIdList());
} else if (validatesRegion) {
Mockito.verify(region)
.checkMeasurementExistence(
new PartialPath("root.sg.d"), Collections.singletonList("v"), null);
}
} else {
Mockito.verifyZeroInteractions(consensus);
assertEquals(result.getMessage(), result.getStatus().getMessage());
}
}
}
}
} finally {
IoTDBDescriptor.getInstance()
.getConfig()
.setSchemaRegionConsensusProtocolClass(originalProtocol);
}
}

@Test
public void testInternalCreateTimeSeriesAfterSchemaRegionCleared() throws Exception {
// A shutdown can clear the region while its validation lock remains registered.
checkInternalCreateWithRegionAvailability(false, false);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reproduces the single-device internal creation failure by retaining the validation lock while returning a missing region. Before the fix this test received status 305 instead of the retryable 906; it now passes.

}

@Test
public void testInternalCreateMultiTimeSeriesAfterSchemaRegionCleared() throws Exception {
// Batch and Pipe auto-creation must return a retryable status for the same shutdown window.
checkInternalCreateWithRegionAvailability(true, false);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Covers batch internal creation against the same missing-region state so both auto-creation entry points have regression coverage. The helper also exercises Pipe-wrapped requests.

}

@Test
public void testInternalCreateWithAvailableSchemaRegion() throws Exception {
// The guard must preserve quota checks and consensus writes for both live-region requests.
checkInternalCreateWithRegionAvailability(false, true);
checkInternalCreateWithRegionAvailability(true, true);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Checks both single-device and batch requests with an available region to ensure that the early rejection does not suppress valid quota checks or consensus writes.

}

private void checkInternalCreateWithRegionAvailability(boolean batch, boolean available)
throws Exception {
final String originalProtocol =
IoTDBDescriptor.getInstance().getConfig().getSchemaRegionConsensusProtocolClass();
try {
for (String protocol :
new String[] {ConsensusFactory.SIMPLE_CONSENSUS, ConsensusFactory.RATIS_CONSENSUS}) {
IoTDBDescriptor.getInstance().getConfig().setSchemaRegionConsensusProtocolClass(protocol);
for (boolean fromPipe : new boolean[] {false, true}) {
final IConsensus consensus = Mockito.mock(IConsensus.class);
final DataNodeRegionManager regionManager = Mockito.mock(DataNodeRegionManager.class);
final SchemaEngine schemaEngine = Mockito.mock(SchemaEngine.class);
final SchemaRegionId regionId = new SchemaRegionId(1);
final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
Mockito.when(regionManager.getRegionLock(regionId)).thenReturn(lock);
final ISchemaRegion region = Mockito.mock(ISchemaRegion.class);
Mockito.when(schemaEngine.getSchemaRegion(regionId))
.thenReturn(available ? region : null);
final RegionWriteExecutor executor =
new RegionWriteExecutor(
Mockito.mock(IConsensus.class),
consensus,
regionManager,
schemaEngine,
Mockito.mock(ClusterTemplateManager.class),
Mockito.mock(TriggerFireVisitor.class));
final PartialPath device = new PartialPath("root.sg.d");
final MeasurementGroup measurements = new MeasurementGroup();
measurements.addMeasurement(
"s", TSDataType.INT64, TSEncoding.PLAIN, CompressionType.SNAPPY);
final WritePlanNode node =
batch
? new InternalCreateMultiTimeSeriesNode(
new PlanNodeId("test"),
Collections.singletonMap(device, new Pair<>(false, measurements)))
: new InternalCreateTimeSeriesNode(
new PlanNodeId("test"), device, measurements, false);
final WritePlanNode request = fromPipe ? new PipeEnrichedWritePlanNode(node) : node;
Mockito.when(consensus.write(regionId, request))
.thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));

final RegionExecutionResult result = executor.execute(regionId, request);
assertEquals(available, result.isAccepted());
assertEquals(
(available ? TSStatusCode.SUCCESS_STATUS : TSStatusCode.NO_AVAILABLE_REGION_GROUP)
.getStatusCode(),
result.getStatus().getCode());
assertFalse(lock.isWriteLocked());
assertEquals(0, lock.getReadLockCount());
if (available) {
Mockito.verify(region).checkSchemaQuota(device, 1);
Mockito.verify(consensus).write(regionId, request);
} else {
Mockito.verifyZeroInteractions(consensus);
assertEquals(1, measurements.size());
}
}
}
} finally {
IoTDBDescriptor.getInstance()
.getConfig()
.setSchemaRegionConsensusProtocolClass(originalProtocol);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The shared fixture runs SimpleConsensus and Ratis with ordinary and Pipe requests, using a real validation lock and mocked region/consensus services. Assertions check the returned status, lock release, and either the expected write or the absence of consensus interaction; configuration is restored in finally. All four tests in this class passed. This is deterministic executor-level coverage, not a full cluster shutdown stress test. The added imports support this fixture and its request variants.

}
}

@Test
public void testInsertRowNode() throws ConsensusException {

Expand Down
Loading