diff --git a/api/src/main/java/org/apache/flink/agents/api/configuration/AgentConfigOptions.java b/api/src/main/java/org/apache/flink/agents/api/configuration/AgentConfigOptions.java index 782b4e4862..c1ec83e50d 100644 --- a/api/src/main/java/org/apache/flink/agents/api/configuration/AgentConfigOptions.java +++ b/api/src/main/java/org/apache/flink/agents/api/configuration/AgentConfigOptions.java @@ -94,6 +94,17 @@ public enum ConditionEvaluationFailureStrategy { public static final ConfigOption KAFKA_ACTION_STATE_TOMBSTONE_ENABLED = new ConfigOption<>("kafkaActionStateTombstoneEnabled", Boolean.class, false); + /** + * The separate, single-partition Kafka topic that stores committed checkpoint-aligned cleanup + * boundaries. It must use {@code cleanup.policy=compact} without delete retention. Setting this + * option enables boundary enforcement during recovery. It must not be the action-state data + * topic, and it must be dedicated to one job's recovery history. The action-state data topic + * must use {@code cleanup.policy=compact,delete}, {@code retention.ms=-1}, and {@code + * retention.bytes=-1} so only reviewed cleanup plans advance its prefix. + */ + public static final ConfigOption KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC = + new ConfigOption<>("kafkaActionStateCleanupControlTopic", String.class, null); + /** The config parameter specifies the Fluss bootstrap servers. */ public static final ConfigOption FLUSS_BOOTSTRAP_SERVERS = new ConfigOption<>("flussBootstrapServers", String.class, "localhost:9123"); diff --git a/docs/content/docs/operations/configuration.md b/docs/content/docs/operations/configuration.md index 4cc35929eb..a1d205829e 100644 --- a/docs/content/docs/operations/configuration.md +++ b/docs/content/docs/operations/configuration.md @@ -183,6 +183,30 @@ Here are the configuration options for Kafka-based Action State Store. | `kafkaActionStateTopicNumPartitions`| 64 | Integer | The config parameter specifies the number of partitions for the Kafka action state topic. | | `kafkaActionStateTopicReplicationFactor` | 1 | Integer | The config parameter specifies the replication factor for the Kafka action state topic. | | `kafkaActionStateTombstoneEnabled` | false | Boolean | Whether pruning sends tombstone records so log compaction can reclaim pruned keys on a compacted action-state topic. Off by default: pruning does not invalidate older restore points, but the topic continues to grow. When enabled, the checkpoint whose completion triggers pruning remains usable, but restoring an earlier checkpoint or savepoint may replay later tombstones and re-execute already completed actions. Enable only if the job never restores from earlier checkpoints or savepoints, or if re-executing actions is acceptable. | +| `kafkaActionStateCleanupControlTopic` | (none) | String | Separate, single-partition topic containing committed checkpoint-aligned cleanup boundaries. It must use `cleanup.policy=compact` without delete retention, differ from the action-state topic, and be dedicated to the same job recovery history. Setting it enables boundary enforcement during recovery. It cannot be combined with `kafkaActionStateTombstoneEnabled`. | + +##### Checkpoint-aligned Kafka cleanup + +Checkpoint-aligned cleanup is an explicit administrative operation. `KafkaActionStateCleanupTool` provides the same command-line workflow for Java and Python jobs. Run it with the Flink Agents runtime JAR and the matching Flink State Processor API JAR on the classpath: + +```text +plan --checkpoint PATH (--operator-uid UID | --operator-uid-hash HASH) --output FILE +apply --plan FILE --bootstrap-servers SERVERS --control-topic TOPIC [--replication-factor N] +``` + +The `plan` command reads all recovery markers from the selected checkpoint or savepoint through Flink's State Processor API, creates a deterministic plan using the earliest required offset per partition, and refuses to overwrite an existing output file. Legacy map markers can still be restored by a job, but cannot authorize deletion because they do not identify the physical Kafka topic. The `apply` command verifies the content-derived plan ID before contacting Kafka. Plan files must contain exactly one JSON document; trailing content is rejected. The optional `--replication-factor` controls replication when creating the control topic, defaults to `1`, and must be between `1` and `32767`. + +Review and retain the plan's deterministic JSON before applying it. The coordinator verifies that the selected offsets are still available, writes `COMMITTED` before calling Kafka `deleteRecords`, verifies every resulting beginning offset, and then writes `APPLIED`. Reapplying the same plan retries an interrupted committed operation without changing its boundary. Set `kafkaActionStateCleanupControlTopic` on the job only after the first plan is committed; recovery requires the configured topic to exist and contain a committed boundary, and fails closed if the topic is missing or empty. + +Physical prefix cleanup requires the action-state data topic to use `cleanup.policy=compact,delete`, `retention.ms=-1`, and `retention.bytes=-1`. The store uses those settings when it creates a new topic. Before applying cleanup to an existing topic, alter it to those values. The `delete` policy enables Kafka's `deleteRecords` API, while both retention limits remain disabled so Kafka cannot independently retire records that a supported checkpoint still needs. Apply validates these effective settings before writing `COMMITTED` and rechecks them immediately before and after deletion. + +For the first cleanup on an existing job, stop the job, create and apply the plan from the recovery point that will become the oldest supported one, and restart from that same point with `kafkaActionStateCleanupControlTopic` configured. This avoids a failover window in which the old running job does not yet know the committed boundary. Once the job is running with boundary enforcement enabled, later forward-only plans may be applied while it runs. Run only one `apply` operation at a time; concurrent incomparable plans fail closed and require operator intervention. Do not recreate the action-state topic or change its partitions while an apply operation is running: the coordinator checks the topic ID and partition set before and after deletion, but Kafka addresses `deleteRecords` by topic name and cannot atomically fence topic lifecycle changes. + +The action-state topic and control topic must be dedicated to one job's recovery history. Once any boundary has been committed, every future run and restore of that history must retain the same control-topic configuration; omitting or changing it removes logical boundary enforcement. The boundary can move only forward. A checkpoint whose marker is below the committed boundary is rejected even if Kafka has not finished physical deletion. + +Per-key tombstones and checkpoint-aligned cleanup are mutually exclusive. Tombstones are not tied to the selected recovery boundary and could invalidate a checkpoint that checkpoint-aligned cleanup promises to retain. + +To migrate a job that previously emitted tombstones, first stop it cleanly so its Kafka producer closes and all accepted tombstone sends finish. Restart from the newest recovery point that is valid under the tombstone mode's existing recovery trade-off, with `kafkaActionStateTombstoneEnabled=false` and no cleanup control topic configured. After that attempt completes a new checkpoint or savepoint, stop it again, create and apply the first cleanup plan from that new recovery point, and restart from the same point with `kafkaActionStateCleanupControlTopic` configured. The new marker is after every old tombstone; selecting an older recovery point cannot provide the same guarantee because replay may still encounter those tombstones. #### Fluss-based Action State Store diff --git a/docs/content/docs/operations/deployment.md b/docs/content/docs/operations/deployment.md index fd4d55d3d6..981e153dbd 100644 --- a/docs/content/docs/operations/deployment.md +++ b/docs/content/docs/operations/deployment.md @@ -120,6 +120,10 @@ See [Action State Store Configuration]({{< ref "docs/operations/configuration#ac **Note**: Enabling Kafka action-state tombstones can invalidate checkpoints or savepoints older than the prune and cause completed actions to execute again. See [Action State Store Configuration]({{< ref "docs/operations/configuration#action-state-store" >}}) for the recovery trade-off. {{< /hint >}} +{{< hint warning >}} +**Note**: A committed checkpoint-aligned Kafka cleanup boundary permanently makes older recovery points unsupported. Every subsequent run and restore must use the same cleanup control topic so the runtime can reject those recovery points before replay. See [Checkpoint-aligned Kafka cleanup]({{< ref "docs/operations/configuration#checkpoint-aligned-kafka-cleanup" >}}). +{{< /hint >}} + {{< hint info >}} **Note**: Exactly-once action consistency is guaranteed only if, after recovering from the same checkpoint, inputs for each key arrive in the same order as before recovery. If this ordering requirement is not met, the system falls back to exactly-once output consistency. {{< /hint >}} diff --git a/python/flink_agents/api/core_options.py b/python/flink_agents/api/core_options.py index 6921be7eb4..95b85a9b16 100644 --- a/python/flink_agents/api/core_options.py +++ b/python/flink_agents/api/core_options.py @@ -146,6 +146,12 @@ class AgentConfigOptions: default=False, ) + KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC = ConfigOption( + key="kafkaActionStateCleanupControlTopic", + config_type=str, + default=None, + ) + FLUSS_BOOTSTRAP_SERVERS = ConfigOption( key="flussBootstrapServers", config_type=str, diff --git a/python/flink_agents/api/tests/test_core_options.py b/python/flink_agents/api/tests/test_core_options.py index 3d1d05c3ef..a38c828f9a 100644 --- a/python/flink_agents/api/tests/test_core_options.py +++ b/python/flink_agents/api/tests/test_core_options.py @@ -108,6 +108,9 @@ def test_agent_config_options_are_explicitly_declared() -> None: options = _collect_config_options(AgentConfigOptions) assert options["BASE_LOG_DIR"].get_key() == "baseLogDir" assert options["KAFKA_BOOTSTRAP_SERVERS"].get_default_value() == "localhost:9092" + cleanup_control_topic = options["KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC"] + assert cleanup_control_topic.get_key() == "kafkaActionStateCleanupControlTopic" + assert cleanup_control_topic.get_default_value() is None assert options["EVENT_LOG_LEVEL"].get_default_value() is EventLogLevel.STANDARD assert options["EVENT_LOG_TRACE_ENABLED"].get_default_value() is False condition_failure = options["CONDITION_EVALUATION_FAILURE_STRATEGY"] diff --git a/runtime/pom.xml b/runtime/pom.xml index 5166d10344..e9e3521f2a 100644 --- a/runtime/pom.xml +++ b/runtime/pom.xml @@ -80,6 +80,12 @@ under the License. ${flink.version} provided + + org.apache.flink + flink-state-processor-api + ${flink.version} + provided + org.apache.flink flink-table-api-java @@ -122,6 +128,18 @@ under the License. kafka-clients ${kafka.version} + + org.testcontainers + kafka + 1.21.4 + test + + + commons-codec + commons-codec + 1.16.1 + test + org.apache.fluss @@ -261,4 +279,4 @@ under the License. - \ No newline at end of file + diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinator.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinator.java new file mode 100644 index 0000000000..c367761518 --- /dev/null +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinator.java @@ -0,0 +1,803 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import org.apache.flink.agents.plan.AgentConfiguration; +import org.apache.flink.annotation.Internal; +import org.apache.flink.util.ExceptionUtils; +import org.apache.flink.util.Preconditions; +import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.Config; +import org.apache.kafka.clients.admin.ConfigEntry; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.admin.OffsetSpec; +import org.apache.kafka.clients.admin.RecordsToDelete; +import org.apache.kafka.clients.admin.TopicDescription; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ConfigResource; +import org.apache.kafka.common.errors.TopicExistsException; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.TreeMap; +import java.util.UUID; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; + +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOMBSTONE_ENABLED; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_BOOTSTRAP_SERVERS; +import static org.apache.kafka.clients.CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG; + +/** + * Commits and applies explicit Kafka action-state prefix cleanup plans. + * + *

The control topic is the durable source of truth. A {@link Status#COMMITTED} record is written + * before any data is deleted, making that boundary the logical point of no return. Reapplying the + * same plan is idempotent: an unfinished committed operation retries deletion and advances to + * {@link Status#APPLIED} only after Kafka reports beginning offsets at or beyond the boundary. + */ +@Internal +public final class KafkaActionStateCleanupCoordinator implements AutoCloseable { + + private static final Duration OPERATION_TIMEOUT = Duration.ofSeconds(30); + + /** Durable lifecycle of one immutable cleanup plan. */ + public enum Status { + COMMITTED, + APPLIED + } + + private final Transport transport; + private final boolean committedBoundaryRequired; + private final String expectedDataTopic; + + /** + * Creates an administrative coordinator using the Kafka settings in the agent configuration. + * The control topic is created when absent, and the returned coordinator can commit and apply + * cleanup plans. + */ + public static KafkaActionStateCleanupCoordinator create(AgentConfiguration configuration) { + return create(configuration, true, false); + } + + /** Creates a read-only-boundary coordinator for action-state recovery. */ + static KafkaActionStateCleanupCoordinator createForRecovery(AgentConfiguration configuration) { + return create(configuration, false, true); + } + + private static KafkaActionStateCleanupCoordinator create( + AgentConfiguration configuration, + boolean createControlTopic, + boolean committedBoundaryRequired) { + Preconditions.checkNotNull(configuration, "Agent configuration must not be null"); + String dataTopic = + requireText( + configuration.get(KAFKA_ACTION_STATE_TOPIC), "Kafka action-state topic"); + String controlTopic = + requireText( + configuration.get(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC), + "Kafka action-state cleanup control topic"); + Preconditions.checkArgument( + !dataTopic.equals(controlTopic), + "Kafka action-state cleanup control topic must differ from the data topic"); + Preconditions.checkArgument( + !configuration.get(KAFKA_ACTION_STATE_TOMBSTONE_ENABLED), + "Per-key Kafka tombstones cannot be enabled with checkpoint-aligned cleanup"); + int replicationFactor = configuration.get(KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR); + if (createControlTopic) { + Preconditions.checkArgument( + replicationFactor > 0 && replicationFactor <= Short.MAX_VALUE, + "Kafka cleanup control topic replication factor must be between 1 and %s, but was %s", + Short.MAX_VALUE, + replicationFactor); + } + return new KafkaActionStateCleanupCoordinator( + new KafkaTransport( + configuration.get(KAFKA_BOOTSTRAP_SERVERS), + controlTopic, + replicationFactor, + createControlTopic), + committedBoundaryRequired, + dataTopic); + } + + KafkaActionStateCleanupCoordinator(Transport transport) { + this(transport, false, null); + } + + KafkaActionStateCleanupCoordinator(Transport transport, boolean committedBoundaryRequired) { + this(transport, committedBoundaryRequired, null); + } + + KafkaActionStateCleanupCoordinator( + Transport transport, boolean committedBoundaryRequired, String expectedDataTopic) { + this.transport = + Preconditions.checkNotNull(transport, "Cleanup transport must not be null"); + this.committedBoundaryRequired = committedBoundaryRequired; + this.expectedDataTopic = expectedDataTopic; + } + + /** + * Commits and physically applies a reviewed plan. If the plan was already committed, this + * resumes it without requiring another logical boundary change. + */ + public synchronized Status apply(KafkaActionStateCleanupPlan plan) throws Exception { + Preconditions.checkNotNull(plan, "Cleanup plan must not be null"); + validateExpectedDataTopic(plan.getTopic()); + TopicMetadata topicMetadata = transport.describeTopic(plan.getTopic()); + validatePlanTopic(plan, topicMetadata); + transport.validateDataTopicConfiguration(plan.getTopic()); + + Map operations = transport.readOperations(); + validateOperationSet(operations.values(), plan); + Operation existing = operations.get(plan.getPlanId()); + if (existing != null) { + Preconditions.checkState( + existing.getPlan().equals(plan), + "Cleanup control record %s does not match the supplied plan", + plan.getPlanId()); + if (existing.getStatus() == Status.APPLIED) { + return Status.APPLIED; + } + } else { + Map currentBoundary = effectiveBoundary(operations.values()); + Preconditions.checkArgument( + plan.advancesOrEquals(currentBoundary), + "Cleanup plan %s does not advance the committed boundary %s", + plan.getPlanId(), + currentBoundary); + transport.validateOffsetsAvailable( + plan.getTopic(), plan.getTopicId(), plan.getOffsets()); + transport.validateDataTopicConfiguration(plan.getTopic()); + transport.append(Operation.committed(plan)); + } + + operations = transport.readOperations(); + validateOperationSet(operations.values(), plan); + Map targetBoundary = effectiveBoundary(operations.values()); + validatePlanTopic(plan, transport.describeTopic(plan.getTopic())); + transport.validateDataTopicConfiguration(plan.getTopic()); + transport.deleteBefore(plan.getTopic(), plan.getTopicId(), targetBoundary); + validatePlanTopic(plan, transport.describeTopic(plan.getTopic())); + transport.validateDataTopicConfiguration(plan.getTopic()); + + for (Operation operation : operations.values()) { + if (operation.getStatus() == Status.COMMITTED) { + transport.append(operation.withStatus(Status.APPLIED)); + } + } + return Status.APPLIED; + } + + /** Returns the effective logical boundary committed for the supplied physical topic. */ + public synchronized Map getCommittedBoundary( + String topic, String topicId, Set partitions) throws Exception { + validateExpectedDataTopic(topic); + List operations = new ArrayList<>(transport.readOperations().values()); + if (operations.isEmpty()) { + Preconditions.checkState( + !committedBoundaryRequired, + "Kafka cleanup control topic contains no committed boundary; apply a cleanup plan before enabling checkpoint-aligned cleanup during recovery"); + return Collections.emptyMap(); + } + for (Operation operation : operations) { + KafkaActionStateCleanupPlan plan = operation.getPlan(); + Preconditions.checkState( + topic.equals(plan.getTopic()) && topicId.equals(plan.getTopicId()), + "Cleanup control topic contains a boundary for a different Kafka topic"); + Preconditions.checkState( + partitions.equals(plan.getOffsets().keySet()), + "Cleanup control topic contains a different Kafka partition set"); + } + validateComparable(operations); + return effectiveBoundary(operations); + } + + /** Rejects a restore point older than the durable logical cleanup boundary. */ + public synchronized void validateRecoveryOffsets( + String topic, String topicId, Map recoveryOffsets) throws Exception { + Map boundary = + getCommittedBoundary(topic, topicId, recoveryOffsets.keySet()); + boundary.forEach( + (partition, committedOffset) -> { + long requestedOffset = recoveryOffsets.get(partition); + Preconditions.checkState( + requestedOffset >= committedOffset, + "Cannot restore Kafka action state for %s-%s from offset %s because the committed cleanup boundary is %s", + topic, + partition, + requestedOffset, + committedOffset); + }); + } + + private static void validatePlanTopic( + KafkaActionStateCleanupPlan plan, TopicMetadata metadata) { + Preconditions.checkArgument( + plan.getTopicId().equals(metadata.getTopicId()), + "Kafka action-state topic %s has ID %s, but cleanup plan %s expects %s", + plan.getTopic(), + metadata.getTopicId(), + plan.getPlanId(), + plan.getTopicId()); + Preconditions.checkArgument( + plan.getOffsets().keySet().equals(metadata.getPartitions()), + "Kafka action-state topic %s has partitions %s, but cleanup plan %s expects %s", + plan.getTopic(), + metadata.getPartitions(), + plan.getPlanId(), + plan.getOffsets().keySet()); + } + + private static void validateOperationSet( + Iterable operations, KafkaActionStateCleanupPlan expectedPlan) { + for (Operation operation : operations) { + KafkaActionStateCleanupPlan plan = operation.getPlan(); + Preconditions.checkState( + expectedPlan.getTopic().equals(plan.getTopic()) + && expectedPlan.getTopicId().equals(plan.getTopicId()), + "Cleanup control topic is not dedicated to one Kafka action-state history"); + Preconditions.checkState( + expectedPlan.getOffsets().keySet().equals(plan.getOffsets().keySet()), + "Cleanup control topic contains incompatible Kafka partition sets"); + } + validateComparable(operations); + } + + private static void validateComparable(Iterable operations) { + List plans = new ArrayList<>(); + operations.forEach(operation -> plans.add(operation.getPlan())); + for (int left = 0; left < plans.size(); left++) { + for (int right = left + 1; right < plans.size(); right++) { + KafkaActionStateCleanupPlan leftPlan = plans.get(left); + KafkaActionStateCleanupPlan rightPlan = plans.get(right); + Preconditions.checkState( + leftPlan.advancesOrEquals(rightPlan.getOffsets()) + || rightPlan.advancesOrEquals(leftPlan.getOffsets()), + "Cleanup control topic contains incomparable boundaries %s and %s", + leftPlan.getPlanId(), + rightPlan.getPlanId()); + } + } + } + + private static Map effectiveBoundary(Iterable operations) { + Map boundary = new HashMap<>(); + for (Operation operation : operations) { + operation + .getPlan() + .getOffsets() + .forEach((partition, offset) -> boundary.merge(partition, offset, Math::max)); + } + return Collections.unmodifiableMap(new TreeMap<>(boundary)); + } + + private static String requireText(String value, String name) { + Preconditions.checkArgument( + value != null && !value.trim().isEmpty(), "%s must not be blank", name); + return value; + } + + private void validateExpectedDataTopic(String topic) { + Preconditions.checkArgument( + expectedDataTopic == null || expectedDataTopic.equals(topic), + "Kafka cleanup coordinator is configured for data topic %s, but received %s", + expectedDataTopic, + topic); + } + + static void validateControlTopicCleanupPolicy(String topic, String cleanupPolicy) { + Preconditions.checkState( + "compact".equals(cleanupPolicy), + "Kafka cleanup control topic %s must use cleanup.policy=compact without delete retention", + topic); + } + + static void validateDataTopicConfiguration( + String topic, String cleanupPolicy, String retentionMs, String retentionBytes) { + Set policies = new HashSet<>(); + if (cleanupPolicy != null) { + for (String policy : cleanupPolicy.split(",")) { + policies.add(policy.trim()); + } + } + Preconditions.checkState( + policies.equals(Set.of("compact", "delete")) + && "-1".equals(retentionMs) + && "-1".equals(retentionBytes), + "Kafka action-state topic %s must use cleanup.policy=compact,delete with retention.ms=-1 and retention.bytes=-1 before checkpoint-aligned cleanup can be committed", + topic); + } + + @Override + public void close() throws Exception { + transport.close(); + } + + interface Transport extends AutoCloseable { + TopicMetadata describeTopic(String topic) throws Exception; + + Map readOperations() throws Exception; + + void append(Operation operation) throws Exception; + + void validateDataTopicConfiguration(String topic) throws Exception; + + void validateOffsetsAvailable(String topic, String topicId, Map offsets) + throws Exception; + + void deleteBefore(String topic, String topicId, Map offsets) + throws Exception; + } + + static final class TopicMetadata { + private final String topicId; + private final Set partitions; + + TopicMetadata(String topicId, Set partitions) { + this.topicId = requireText(topicId, "Kafka topic ID"); + this.partitions = + Collections.unmodifiableSet( + new HashSet<>( + Preconditions.checkNotNull( + partitions, "Kafka partitions must not be null"))); + } + + String getTopicId() { + return topicId; + } + + Set getPartitions() { + return partitions; + } + } + + static final class Operation { + private static final int CURRENT_SCHEMA_VERSION = 1; + private static final ObjectMapper MAPPER = + new ObjectMapper() + .enable(JsonParser.Feature.STRICT_DUPLICATE_DETECTION) + .enable(DeserializationFeature.FAIL_ON_TRAILING_TOKENS); + private static final Set JSON_FIELDS = Set.of("schemaVersion", "status", "plan"); + + private final KafkaActionStateCleanupPlan plan; + private final Status status; + + private Operation(KafkaActionStateCleanupPlan plan, Status status) { + this.plan = Preconditions.checkNotNull(plan, "Cleanup plan must not be null"); + this.status = Preconditions.checkNotNull(status, "Cleanup status must not be null"); + } + + static Operation committed(KafkaActionStateCleanupPlan plan) { + return new Operation(plan, Status.COMMITTED); + } + + Operation withStatus(Status newStatus) { + Preconditions.checkArgument( + status != Status.APPLIED || newStatus == Status.APPLIED, + "An applied cleanup operation cannot return to committed"); + return new Operation(plan, newStatus); + } + + KafkaActionStateCleanupPlan getPlan() { + return plan; + } + + Status getStatus() { + return status; + } + + String toJson() throws Exception { + ObjectNode root = MAPPER.createObjectNode(); + root.put("schemaVersion", CURRENT_SCHEMA_VERSION); + root.put("status", status.name()); + root.set("plan", MAPPER.readTree(plan.toJson())); + return MAPPER.writeValueAsString(root); + } + + static Operation fromJson(String json) throws Exception { + JsonNode root = MAPPER.readTree(json); + Preconditions.checkArgument( + root != null && root.isObject(), "Cleanup control record must be JSON"); + Set fieldNames = new HashSet<>(); + root.fieldNames().forEachRemaining(fieldNames::add); + Preconditions.checkArgument( + fieldNames.equals(JSON_FIELDS), + "Cleanup control record fields must be exactly %s, but were %s", + JSON_FIELDS, + fieldNames); + JsonNode schemaVersion = root.get("schemaVersion"); + Preconditions.checkArgument( + schemaVersion.isIntegralNumber() + && schemaVersion.canConvertToInt() + && schemaVersion.intValue() == CURRENT_SCHEMA_VERSION, + "Unsupported cleanup control record schema %s", + schemaVersion); + JsonNode statusNode = root.get("status"); + Preconditions.checkArgument( + statusNode.isTextual(), "Cleanup control record status must be a string"); + Status status = Status.valueOf(statusNode.textValue()); + JsonNode planNode = root.get("plan"); + Preconditions.checkArgument(planNode != null, "Cleanup control record has no plan"); + return new Operation(KafkaActionStateCleanupPlan.fromJson(planNode.toString()), status); + } + } + + static void applyControlRecord( + Map operations, ConsumerRecord record) + throws Exception { + Preconditions.checkState( + record.key() != null, "Kafka cleanup control record must have a plan ID key"); + Preconditions.checkState( + record.value() != null, + "Kafka cleanup control topic contains a tombstone for plan %s; cleanup boundaries cannot be deleted", + record.key()); + Operation operation = Operation.fromJson(record.value()); + Preconditions.checkState( + record.key().equals(operation.getPlan().getPlanId()), + "Cleanup control record key does not match its plan ID"); + Operation previous = operations.get(record.key()); + Preconditions.checkState( + previous == null || previous.getPlan().equals(operation.getPlan()), + "Cleanup control record changed immutable plan %s", + record.key()); + Preconditions.checkState( + previous == null + || previous.getStatus() != Status.APPLIED + || operation.getStatus() == Status.APPLIED, + "Cleanup control record regressed applied plan %s", + record.key()); + operations.put(record.key(), operation); + } + + private static final class KafkaTransport implements Transport { + private final String controlTopic; + private final AdminClient adminClient; + private final Producer producer; + private final Consumer consumer; + + private KafkaTransport( + String bootstrapServers, + String controlTopic, + int replicationFactor, + boolean createControlTopic) { + this.controlTopic = requireText(controlTopic, "Kafka cleanup control topic"); + Properties common = new Properties(); + common.put( + BOOTSTRAP_SERVERS_CONFIG, + requireText(bootstrapServers, "Kafka bootstrap servers")); + this.adminClient = AdminClient.create(common); + Producer createdProducer = null; + Consumer createdConsumer = null; + try { + ensureControlTopic(replicationFactor, createControlTopic); + + if (createControlTopic) { + Properties producerProperties = new Properties(); + producerProperties.putAll(common); + producerProperties.put( + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put( + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put(ProducerConfig.ACKS_CONFIG, "all"); + producerProperties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); + createdProducer = new KafkaProducer<>(producerProperties); + } + + Properties consumerProperties = new Properties(); + consumerProperties.putAll(common); + consumerProperties.put( + ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put( + ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put( + ConsumerConfig.GROUP_ID_CONFIG, + "action-state-cleanup-control-" + UUID.randomUUID()); + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); + consumerProperties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); + consumerProperties.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); + createdConsumer = new KafkaConsumer<>(consumerProperties); + } catch (RuntimeException | Error failure) { + if (createdProducer != null) { + try { + createdProducer.close(); + } catch (Throwable closeFailure) { + failure.addSuppressed(closeFailure); + } + } + if (createdConsumer != null) { + try { + createdConsumer.close(); + } catch (Throwable closeFailure) { + failure.addSuppressed(closeFailure); + } + } + try { + adminClient.close(); + } catch (Throwable closeFailure) { + failure.addSuppressed(closeFailure); + } + throw failure; + } + this.producer = createdProducer; + this.consumer = createdConsumer; + } + + @Override + public TopicMetadata describeTopic(String topic) throws Exception { + TopicDescription description = + adminClient + .describeTopics(List.of(topic)) + .allTopicNames() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS) + .get(topic); + Preconditions.checkState(description != null, "Kafka topic does not exist: %s", topic); + Set partitions = new HashSet<>(); + description.partitions().forEach(value -> partitions.add(value.partition())); + return new TopicMetadata(description.topicId().toString(), partitions); + } + + @Override + public Map readOperations() throws Exception { + TopicPartition partition = new TopicPartition(controlTopic, 0); + consumer.assign(List.of(partition)); + long beginning = consumer.beginningOffsets(List.of(partition)).get(partition); + long end = consumer.endOffsets(List.of(partition)).get(partition); + consumer.seek(partition, beginning); + + Map operations = new LinkedHashMap<>(); + long deadline = System.nanoTime() + OPERATION_TIMEOUT.toNanos(); + while (consumer.position(partition) < end) { + ConsumerRecords records = consumer.poll(Duration.ofMillis(200)); + for (ConsumerRecord record : records) { + if (record.offset() >= end) { + continue; + } + applyControlRecord(operations, record); + } + if (consumer.position(partition) < end) { + Preconditions.checkState( + System.nanoTime() < deadline, + "Timed out reading Kafka cleanup control topic %s", + controlTopic); + } + } + return operations; + } + + @Override + public void append(Operation operation) throws Exception { + Preconditions.checkState( + producer != null, + "Recovery-only Kafka cleanup coordinator cannot append control records"); + producer.send( + new ProducerRecord<>( + controlTopic, + operation.getPlan().getPlanId(), + operation.toJson())) + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); + } + + @Override + public void validateDataTopicConfiguration(String topic) throws Exception { + Config config = describeTopicConfiguration(topic); + ConfigEntry cleanupPolicy = config.get("cleanup.policy"); + ConfigEntry retentionMs = config.get("retention.ms"); + ConfigEntry retentionBytes = config.get("retention.bytes"); + KafkaActionStateCleanupCoordinator.validateDataTopicConfiguration( + topic, + cleanupPolicy == null ? null : cleanupPolicy.value(), + retentionMs == null ? null : retentionMs.value(), + retentionBytes == null ? null : retentionBytes.value()); + } + + @Override + public void validateOffsetsAvailable( + String topic, String expectedTopicId, Map offsets) throws Exception { + validateTopicIdentity(topic, expectedTopicId, offsets.keySet()); + Map beginningOffsets = + listOffsets(topic, offsets.keySet(), OffsetSpec.earliest()); + Map endOffsets = + listOffsets(topic, offsets.keySet(), OffsetSpec.latest()); + validateTopicIdentity(topic, expectedTopicId, offsets.keySet()); + offsets.forEach( + (partition, requestedOffset) -> { + TopicPartition topicPartition = new TopicPartition(topic, partition); + Long beginningOffset = beginningOffsets.get(topicPartition); + Long endOffset = endOffsets.get(topicPartition); + Preconditions.checkArgument( + beginningOffset != null + && endOffset != null + && requestedOffset >= beginningOffset + && requestedOffset <= endOffset, + "Cannot commit Kafka cleanup boundary for %s at offset %s because the available range is [%s, %s]", + topicPartition, + requestedOffset, + beginningOffset, + endOffset); + }); + } + + @Override + public void deleteBefore(String topic, String expectedTopicId, Map offsets) + throws Exception { + validateTopicIdentity(topic, expectedTopicId, offsets.keySet()); + Map recordsToDelete = new HashMap<>(); + offsets.forEach( + (partition, offset) -> + recordsToDelete.put( + new TopicPartition(topic, partition), + RecordsToDelete.beforeOffset(offset))); + adminClient + .deleteRecords(recordsToDelete) + .all() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); + + Map beginningOffsets = + listOffsets(topic, offsets.keySet(), OffsetSpec.earliest()); + offsets.forEach( + (partition, target) -> { + TopicPartition topicPartition = new TopicPartition(topic, partition); + Long beginning = beginningOffsets.get(topicPartition); + Preconditions.checkState( + beginning != null && beginning >= target, + "Kafka cleanup for %s requested offset %s but beginning offset is %s", + topicPartition, + target, + beginning); + }); + validateTopicIdentity(topic, expectedTopicId, offsets.keySet()); + } + + private Map listOffsets( + String topic, Set partitions, OffsetSpec offsetSpec) throws Exception { + Map requests = new HashMap<>(); + partitions.forEach( + partition -> requests.put(new TopicPartition(topic, partition), offsetSpec)); + Map offsets = new HashMap<>(); + adminClient + .listOffsets(requests) + .all() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS) + .forEach((partition, info) -> offsets.put(partition, info.offset())); + return offsets; + } + + private void validateTopicIdentity( + String topic, String expectedTopicId, Set expectedPartitions) + throws Exception { + TopicMetadata metadata = describeTopic(topic); + Preconditions.checkState( + expectedTopicId.equals(metadata.getTopicId()), + "Kafka action-state topic %s changed during cleanup; expected ID %s but found %s", + topic, + expectedTopicId, + metadata.getTopicId()); + Preconditions.checkState( + expectedPartitions.equals(metadata.getPartitions()), + "Kafka action-state topic %s changed partitions during cleanup; expected %s but found %s", + topic, + expectedPartitions, + metadata.getPartitions()); + } + + private void ensureControlTopic(int replicationFactor, boolean createControlTopic) { + try { + boolean exists = + adminClient + .listTopics() + .names() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS) + .contains(controlTopic); + if (!exists) { + Preconditions.checkState( + createControlTopic, + "Kafka cleanup control topic %s does not exist; apply a cleanup plan before enabling checkpoint-aligned cleanup during recovery", + controlTopic); + NewTopic topic = new NewTopic(controlTopic, 1, (short) replicationFactor); + topic.configs(Map.of("cleanup.policy", "compact")); + try { + adminClient + .createTopics(List.of(topic)) + .all() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + if (!(e.getCause() instanceof TopicExistsException)) { + throw e; + } + } + } + + TopicMetadata metadata = describeTopic(controlTopic); + Preconditions.checkState( + metadata.getPartitions().equals(Set.of(0)), + "Kafka cleanup control topic %s must have exactly one partition", + controlTopic); + Config config = describeTopicConfiguration(controlTopic); + ConfigEntry cleanupPolicy = config.get("cleanup.policy"); + validateControlTopicCleanupPolicy( + controlTopic, cleanupPolicy == null ? null : cleanupPolicy.value()); + } catch (Exception e) { + throw new IllegalStateException( + "Failed to create or validate Kafka cleanup control topic " + controlTopic, + e); + } + } + + private Config describeTopicConfiguration(String topic) throws Exception { + ConfigResource resource = new ConfigResource(ConfigResource.Type.TOPIC, topic); + return adminClient + .describeConfigs(List.of(resource)) + .all() + .get(OPERATION_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS) + .get(resource); + } + + @Override + public void close() throws Exception { + Throwable firstFailure = null; + if (producer != null) { + try { + producer.close(); + } catch (Throwable failure) { + firstFailure = failure; + } + } + try { + consumer.close(); + } catch (Throwable failure) { + firstFailure = ExceptionUtils.firstOrSuppressed(failure, firstFailure); + } + try { + adminClient.close(); + } catch (Throwable failure) { + firstFailure = ExceptionUtils.firstOrSuppressed(failure, firstFailure); + } + if (firstFailure != null) { + ExceptionUtils.rethrowException(firstFailure); + } + } + } +} diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlan.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlan.java new file mode 100644 index 0000000000..2d6be91ecb --- /dev/null +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlan.java @@ -0,0 +1,332 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.apache.flink.annotation.Internal; +import org.apache.flink.util.Preconditions; +import org.apache.kafka.common.Uuid; + +import java.io.IOException; +import java.io.Serializable; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Iterator; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.TreeMap; + +/** Immutable, reviewable Kafka prefix-cleanup boundary derived from one recovery point. */ +@Internal +public final class KafkaActionStateCleanupPlan implements Serializable { + + private static final long serialVersionUID = 1L; + private static final int CURRENT_SCHEMA_VERSION = 1; + private static final ObjectMapper MAPPER = + new ObjectMapper() + .enable(JsonParser.Feature.STRICT_DUPLICATE_DETECTION) + .enable(DeserializationFeature.FAIL_ON_TRAILING_TOKENS); + private static final Set JSON_FIELDS = + Set.of("schemaVersion", "planId", "sourceRecoveryPoint", "topic", "topicId", "offsets"); + + private final int schemaVersion; + private final String planId; + private final String sourceRecoveryPoint; + private final String topic; + private final String topicId; + private final Map offsets; + + private KafkaActionStateCleanupPlan( + int schemaVersion, + String planId, + String sourceRecoveryPoint, + String topic, + String topicId, + Map offsets) { + Preconditions.checkArgument( + schemaVersion == CURRENT_SCHEMA_VERSION, + "Unsupported Kafka action-state cleanup plan schema %s", + schemaVersion); + this.sourceRecoveryPoint = requireText(sourceRecoveryPoint, "source recovery point"); + this.topic = requireText(topic, "Kafka action-state topic"); + this.topicId = requireText(topicId, "Kafka action-state topic ID"); + Preconditions.checkArgument( + !Uuid.ZERO_UUID.toString().equals(this.topicId), + "Checkpoint-aligned cleanup requires a broker-provided Kafka topic ID"); + this.offsets = immutableOffsets(offsets); + this.schemaVersion = schemaVersion; + String expectedPlanId = calculatePlanId(sourceRecoveryPoint, topic, topicId, this.offsets); + Preconditions.checkArgument( + planId == null || expectedPlanId.equals(planId), + "Kafka action-state cleanup plan ID does not match its contents"); + this.planId = expectedPlanId; + } + + /** + * Builds a plan from every union recovery marker in one selected checkpoint or savepoint. + * Legacy map markers are intentionally rejected because they do not identify the physical Kafka + * topic and therefore cannot authorize irreversible deletion safely. + */ + public static KafkaActionStateCleanupPlan fromRecoveryMarkers( + String sourceRecoveryPoint, List recoveryMarkers) { + Preconditions.checkNotNull(recoveryMarkers, "Recovery markers must not be null"); + Preconditions.checkArgument( + !recoveryMarkers.isEmpty(), "Recovery markers must not be empty"); + + String topic = null; + String topicId = null; + Map mergedOffsets = new HashMap<>(); + for (Object value : recoveryMarkers) { + Preconditions.checkArgument( + value instanceof KafkaActionStateRecoveryMarker, + "Checkpoint-aligned cleanup requires versioned Kafka recovery markers; legacy or unsupported marker found: %s", + value == null ? "null" : value.getClass().getName()); + KafkaActionStateRecoveryMarker marker = (KafkaActionStateRecoveryMarker) value; + Preconditions.checkArgument( + marker.getSchemaVersion() + == KafkaActionStateRecoveryMarker.CURRENT_SCHEMA_VERSION, + "Unsupported Kafka recovery marker schema %s", + marker.getSchemaVersion()); + if (topic == null) { + topic = marker.getTopic(); + topicId = marker.getTopicId(); + } else { + Preconditions.checkArgument( + topic.equals(marker.getTopic()) && topicId.equals(marker.getTopicId()), + "Recovery markers do not reference one physical Kafka topic"); + Preconditions.checkArgument( + mergedOffsets.keySet().equals(marker.getOffsets().keySet()), + "Recovery markers contain different Kafka partition sets"); + } + marker.getOffsets() + .forEach( + (partition, offset) -> + mergedOffsets.merge(partition, offset, Math::min)); + } + + return new KafkaActionStateCleanupPlan( + CURRENT_SCHEMA_VERSION, null, sourceRecoveryPoint, topic, topicId, mergedOffsets); + } + + /** Parses a previously reviewed plan and verifies its content-derived identifier. */ + public static KafkaActionStateCleanupPlan fromJson(String json) throws IOException { + JsonNode root = MAPPER.readTree(json); + Preconditions.checkArgument(root != null && root.isObject(), "Cleanup plan must be JSON"); + Set fieldNames = new HashSet<>(); + root.fieldNames().forEachRemaining(fieldNames::add); + Preconditions.checkArgument( + fieldNames.equals(JSON_FIELDS), + "Cleanup plan fields must be exactly %s, but were %s", + JSON_FIELDS, + fieldNames); + Map offsets = new HashMap<>(); + JsonNode offsetsNode = required(root, "offsets"); + Preconditions.checkArgument( + offsetsNode.isObject(), "Cleanup plan offsets must be an object"); + Iterator> fields = offsetsNode.fields(); + while (fields.hasNext()) { + Map.Entry field = fields.next(); + try { + Preconditions.checkArgument( + field.getValue().isIntegralNumber() && field.getValue().canConvertToLong(), + "Cleanup offset for partition %s must be an integer", + field.getKey()); + int partition = Integer.parseInt(field.getKey()); + Preconditions.checkArgument( + field.getKey().equals(Integer.toString(partition)), + "Cleanup plan partition %s must use canonical decimal form", + field.getKey()); + Long previous = offsets.put(partition, field.getValue().longValue()); + Preconditions.checkArgument( + previous == null, + "Cleanup plan contains duplicate Kafka partition %s", + partition); + } catch (NumberFormatException e) { + throw new IllegalArgumentException( + "Cleanup plan contains a non-integer Kafka partition: " + field.getKey(), + e); + } + } + JsonNode schemaVersion = required(root, "schemaVersion"); + Preconditions.checkArgument( + schemaVersion.isIntegralNumber() && schemaVersion.canConvertToInt(), + "Cleanup plan schemaVersion must be an integer"); + return new KafkaActionStateCleanupPlan( + schemaVersion.intValue(), + requiredText(root, "planId"), + requiredText(root, "sourceRecoveryPoint"), + requiredText(root, "topic"), + requiredText(root, "topicId"), + offsets); + } + + /** Returns deterministic JSON suitable for review and later application. */ + public String toJson() { + Map value = new LinkedHashMap<>(); + value.put("schemaVersion", schemaVersion); + value.put("planId", planId); + value.put("sourceRecoveryPoint", sourceRecoveryPoint); + value.put("topic", topic); + value.put("topicId", topicId); + value.put("offsets", new TreeMap<>(offsets)); + try { + return MAPPER.writerWithDefaultPrettyPrinter().writeValueAsString(value); + } catch (JsonProcessingException e) { + throw new IllegalStateException( + "Failed to serialize Kafka action-state cleanup plan", e); + } + } + + public int getSchemaVersion() { + return schemaVersion; + } + + public String getPlanId() { + return planId; + } + + public String getSourceRecoveryPoint() { + return sourceRecoveryPoint; + } + + public String getTopic() { + return topic; + } + + public String getTopicId() { + return topicId; + } + + public Map getOffsets() { + return offsets; + } + + boolean advancesOrEquals(Map boundary) { + if (boundary.isEmpty()) { + return true; + } + if (!offsets.keySet().equals(boundary.keySet())) { + return false; + } + return offsets.entrySet().stream() + .allMatch(entry -> entry.getValue() >= boundary.get(entry.getKey())); + } + + private static JsonNode required(JsonNode root, String name) { + JsonNode value = root.get(name); + Preconditions.checkArgument( + value != null && !value.isNull(), "Cleanup plan is missing %s", name); + return value; + } + + private static String requiredText(JsonNode root, String name) { + JsonNode value = required(root, name); + Preconditions.checkArgument(value.isTextual(), "Cleanup plan %s must be a string", name); + return value.textValue(); + } + + private static String requireText(String value, String name) { + Preconditions.checkArgument( + value != null && !value.trim().isEmpty(), "%s must not be blank", name); + return value; + } + + private static Map immutableOffsets(Map offsets) { + Preconditions.checkNotNull(offsets, "Cleanup offsets must not be null"); + Preconditions.checkArgument(!offsets.isEmpty(), "Cleanup offsets must not be empty"); + Map copy = new HashMap<>(); + offsets.forEach( + (partition, offset) -> { + Preconditions.checkNotNull(partition, "Kafka partition must not be null"); + Preconditions.checkNotNull(offset, "Kafka cleanup offset must not be null"); + Preconditions.checkArgument( + partition >= 0, "Kafka partition must be non-negative"); + Preconditions.checkArgument( + offset >= 0, "Kafka cleanup offset must be non-negative"); + copy.put(partition, offset); + }); + return Collections.unmodifiableMap(copy); + } + + private static String calculatePlanId( + String sourceRecoveryPoint, String topic, String topicId, Map offsets) { + StringBuilder canonical = + new StringBuilder() + .append(CURRENT_SCHEMA_VERSION) + .append('\n') + .append(sourceRecoveryPoint) + .append('\n') + .append(topic) + .append('\n') + .append(topicId) + .append('\n'); + new TreeMap<>(offsets) + .forEach( + (partition, offset) -> + canonical + .append(partition) + .append('=') + .append(offset) + .append('\n')); + try { + byte[] digest = + MessageDigest.getInstance("SHA-256") + .digest(canonical.toString().getBytes(StandardCharsets.UTF_8)); + StringBuilder result = new StringBuilder(digest.length * 2); + for (byte value : digest) { + result.append(String.format("%02x", value & 0xff)); + } + return result.toString(); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("SHA-256 is unavailable", e); + } + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof KafkaActionStateCleanupPlan)) { + return false; + } + KafkaActionStateCleanupPlan that = (KafkaActionStateCleanupPlan) object; + return schemaVersion == that.schemaVersion + && planId.equals(that.planId) + && sourceRecoveryPoint.equals(that.sourceRecoveryPoint) + && topic.equals(that.topic) + && topicId.equals(that.topicId) + && offsets.equals(that.offsets); + } + + @Override + public int hashCode() { + return Objects.hash(schemaVersion, planId, sourceRecoveryPoint, topic, topicId, offsets); + } +} diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupTool.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupTool.java new file mode 100644 index 0000000000..51cb4c5a06 --- /dev/null +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupTool.java @@ -0,0 +1,168 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.agents.plan.AgentConfiguration; +import org.apache.flink.api.common.RuntimeExecutionMode; +import org.apache.flink.api.common.typeinfo.TypeInformation; +import org.apache.flink.state.api.OperatorIdentifier; +import org.apache.flink.state.api.SavepointReader; +import org.apache.flink.streaming.api.datastream.DataStream; +import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; +import org.apache.flink.util.CloseableIterator; +import org.apache.flink.util.Preconditions; + +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardOpenOption; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; + +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_BOOTSTRAP_SERVERS; + +/** Command-line entry point for planning and applying checkpoint-aligned Kafka cleanup. */ +public final class KafkaActionStateCleanupTool { + + private KafkaActionStateCleanupTool() {} + + public static void main(String[] args) throws Exception { + Preconditions.checkArgument(args.length > 0, usage()); + if (args.length == 1 && ("--help".equals(args[0]) || "-h".equals(args[0]))) { + System.out.println(usage()); + return; + } + String command = args[0]; + Map options = parseOptions(args); + if ("plan".equals(command)) { + createPlan(options); + } else if ("apply".equals(command)) { + applyPlan(options); + } else { + throw new IllegalArgumentException("Unknown command: " + command + "\n" + usage()); + } + } + + private static void createPlan(Map options) throws Exception { + requireOnly(options, Set.of("checkpoint", "output", "operator-uid", "operator-uid-hash")); + String checkpoint = required(options, "checkpoint"); + String output = required(options, "output"); + String uid = options.get("operator-uid"); + String uidHash = options.get("operator-uid-hash"); + Preconditions.checkArgument( + (uid == null) != (uidHash == null), + "Exactly one of --operator-uid or --operator-uid-hash is required"); + + StreamExecutionEnvironment environment = + StreamExecutionEnvironment.getExecutionEnvironment(); + environment.setRuntimeMode(RuntimeExecutionMode.BATCH); + SavepointReader reader = SavepointReader.read(environment, checkpoint); + OperatorIdentifier identifier = + uid == null + ? OperatorIdentifier.forUidHash(uidHash) + : OperatorIdentifier.forUid(uid); + DataStream markerStream = + reader.readUnionState( + identifier, + KafkaActionStateRecoveryMarker.UNION_STATE_NAME, + TypeInformation.of(Object.class)); + List markers = new ArrayList<>(); + try (CloseableIterator iterator = markerStream.executeAndCollect()) { + iterator.forEachRemaining(markers::add); + } + + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers(checkpoint, markers); + Files.writeString( + Path.of(output), + plan.toJson() + System.lineSeparator(), + StandardCharsets.UTF_8, + StandardOpenOption.CREATE_NEW, + StandardOpenOption.WRITE); + System.out.println( + "Created cleanup plan " + + plan.getPlanId() + + " at " + + Path.of(output).toAbsolutePath()); + } + + private static void applyPlan(Map options) throws Exception { + requireOnly( + options, + Set.of("plan", "bootstrap-servers", "control-topic", "replication-factor")); + Path planPath = Path.of(required(options, "plan")); + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromJson( + Files.readString(planPath, StandardCharsets.UTF_8)); + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(KAFKA_BOOTSTRAP_SERVERS, required(options, "bootstrap-servers")); + configuration.set( + KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, required(options, "control-topic")); + configuration.set(KAFKA_ACTION_STATE_TOPIC, plan.getTopic()); + if (options.containsKey("replication-factor")) { + int replicationFactor = Integer.parseInt(options.get("replication-factor")); + configuration.set(KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR, replicationFactor); + } + + try (KafkaActionStateCleanupCoordinator coordinator = + KafkaActionStateCleanupCoordinator.create(configuration)) { + KafkaActionStateCleanupCoordinator.Status status = coordinator.apply(plan); + System.out.println("Cleanup plan " + plan.getPlanId() + " is " + status); + } + } + + private static Map parseOptions(String[] args) { + Preconditions.checkArgument( + (args.length - 1) % 2 == 0, "Options must be provided as --name value pairs"); + Map options = new LinkedHashMap<>(); + for (int index = 1; index < args.length; index += 2) { + String name = args[index]; + Preconditions.checkArgument(name.startsWith("--"), "Invalid option: %s", name); + String previous = options.put(name.substring(2), args[index + 1]); + Preconditions.checkArgument(previous == null, "Duplicate option: %s", name); + } + return options; + } + + private static String required(Map options, String name) { + String value = options.get(name); + Preconditions.checkArgument( + value != null && !value.trim().isEmpty(), "Missing required option --%s", name); + return value; + } + + private static void requireOnly(Map options, Set allowed) { + options.keySet() + .forEach( + option -> + Preconditions.checkArgument( + allowed.contains(option), "Unknown option: --%s", option)); + } + + private static String usage() { + return "Usage:\n" + + " plan --checkpoint PATH (--operator-uid UID | --operator-uid-hash HASH) --output FILE\n" + + " apply --plan FILE --bootstrap-servers SERVERS --control-topic TOPIC [--replication-factor N]"; + } +} diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarker.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarker.java new file mode 100644 index 0000000000..0beead4c66 --- /dev/null +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarker.java @@ -0,0 +1,119 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.annotation.Internal; + +import java.io.Serializable; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; + +import static org.apache.flink.util.Preconditions.checkArgument; +import static org.apache.flink.util.Preconditions.checkNotNull; + +/** Versioned Kafka topic identity and per-partition offsets used to rebuild action state. */ +@Internal +public final class KafkaActionStateRecoveryMarker implements Serializable { + + private static final long serialVersionUID = 1L; + + public static final int CURRENT_SCHEMA_VERSION = 1; + public static final String UNION_STATE_NAME = "recoveryMarker"; + + private final int schemaVersion; + private final String topic; + private final String topicId; + private final Map offsets; + + KafkaActionStateRecoveryMarker(String topic, String topicId, Map offsets) { + this(CURRENT_SCHEMA_VERSION, topic, topicId, offsets); + } + + KafkaActionStateRecoveryMarker( + int schemaVersion, String topic, String topicId, Map offsets) { + this.schemaVersion = schemaVersion; + this.topic = checkNotNull(topic, "Kafka topic must not be null"); + this.topicId = checkNotNull(topicId, "Kafka topic ID must not be null"); + checkArgument(!topic.trim().isEmpty(), "Kafka topic must not be blank"); + checkArgument(!topicId.trim().isEmpty(), "Kafka topic ID must not be blank"); + checkNotNull(offsets, "Kafka recovery offsets must not be null"); + checkArgument(!offsets.isEmpty(), "Kafka recovery offsets must not be empty"); + offsets.forEach( + (partition, offset) -> { + checkNotNull(partition, "Kafka partition must not be null"); + checkNotNull(offset, "Kafka recovery offset must not be null"); + checkArgument(partition >= 0, "Kafka partition must be non-negative"); + checkArgument(offset >= 0, "Kafka recovery offset must be non-negative"); + }); + this.offsets = new HashMap<>(offsets); + } + + public int getSchemaVersion() { + return schemaVersion; + } + + public String getTopic() { + return topic; + } + + public String getTopicId() { + return topicId; + } + + public Map getOffsets() { + return Collections.unmodifiableMap(offsets); + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof KafkaActionStateRecoveryMarker)) { + return false; + } + KafkaActionStateRecoveryMarker that = (KafkaActionStateRecoveryMarker) object; + return schemaVersion == that.schemaVersion + && topic.equals(that.topic) + && topicId.equals(that.topicId) + && offsets.equals(that.offsets); + } + + @Override + public int hashCode() { + return Objects.hash(schemaVersion, topic, topicId, offsets); + } + + @Override + public String toString() { + return "KafkaActionStateRecoveryMarker{" + + "schemaVersion=" + + schemaVersion + + ", topic='" + + topic + + '\'' + + ", topicId='" + + topicId + + '\'' + + ", offsets=" + + offsets + + '}'; + } +} diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java index 1bbfbe1593..60a72a80f3 100644 --- a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java @@ -25,8 +25,10 @@ import org.apache.flink.util.ExceptionUtils; import org.apache.flink.util.Preconditions; import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.DescribeTopicsResult; import org.apache.kafka.clients.admin.ListTopicsResult; import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.admin.TopicDescription; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -45,14 +47,18 @@ import java.time.Duration; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Properties; +import java.util.Set; import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.function.IntPredicate; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC; import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOMBSTONE_ENABLED; import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC; import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC_NUM_PARTITIONS; @@ -60,6 +66,7 @@ import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_BOOTSTRAP_SERVERS; import static org.apache.kafka.clients.CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG; import static org.apache.kafka.clients.consumer.ConsumerConfig.CLIENT_ID_CONFIG; +import static org.apache.kafka.clients.consumer.ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG; import static org.apache.kafka.clients.consumer.ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG; import static org.apache.kafka.clients.consumer.ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG; @@ -75,6 +82,7 @@ public class KafkaActionStateStore implements ActionStateStore { private static final Duration CONSUMER_POLL_TIMEOUT = Duration.ofMillis(1000); + private static final Duration REBUILD_PROGRESS_TIMEOUT = Duration.ofSeconds(30); private static final Logger LOG = LoggerFactory.getLogger(KafkaActionStateStore.class); // A cold AdminClient's first metadata round-trip routinely exceeds 100ms even against a // local broker; 100ms made store initialization fail nondeterministically at startup. @@ -105,6 +113,13 @@ public class KafkaActionStateStore implements ActionStateStore { private final ActionStateKeyEncoder keyEncoder; + // Kafka topic identity and partition set captured when this store instance starts + private final KafkaTopicMetadata topicMetadata; + private final TopicMetadataLoader topicMetadataLoader; + + // Present only when checkpoint-aligned cleanup boundary enforcement is configured. + private final KafkaActionStateCleanupCoordinator cleanupCoordinator; + @VisibleForTesting KafkaActionStateStore( Map actionStates, @@ -113,14 +128,32 @@ public class KafkaActionStateStore implements ActionStateStore { Consumer consumer, String topic, ActionStateKeyEncoder keyEncoder) { + this(actionStates, agentConfiguration, producer, consumer, topic, keyEncoder, null); + } + + @VisibleForTesting + KafkaActionStateStore( + Map actionStates, + AgentConfiguration agentConfiguration, + Producer producer, + Consumer consumer, + String topic, + ActionStateKeyEncoder keyEncoder, + KafkaActionStateCleanupCoordinator cleanupCoordinator) { this.actionStates = actionStates; this.producer = producer; this.consumer = consumer; this.topic = topic; + this.topicMetadataLoader = () -> testTopicMetadata(consumer, topic); + this.topicMetadata = testTopicMetadata(consumer, topic); this.latestKeySeqNum = new HashMap<>(); this.agentConfiguration = agentConfiguration; this.tombstoneEnabled = agentConfiguration.get(KAFKA_ACTION_STATE_TOMBSTONE_ENABLED); this.keyEncoder = Preconditions.checkNotNull(keyEncoder, "keyEncoder cannot be null"); + Preconditions.checkArgument( + cleanupCoordinator == null || !this.tombstoneEnabled, + "Per-key Kafka tombstones cannot be enabled with checkpoint-aligned cleanup"); + this.cleanupCoordinator = cleanupCoordinator; } /** @@ -136,16 +169,46 @@ public KafkaActionStateStore( this.latestKeySeqNum = new HashMap<>(); this.agentConfiguration = agentConfiguration; this.tombstoneEnabled = agentConfiguration.get(KAFKA_ACTION_STATE_TOMBSTONE_ENABLED); + String cleanupControlTopic = + agentConfiguration.get(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC); + Preconditions.checkArgument( + cleanupControlTopic == null || !cleanupControlTopic.trim().isEmpty(), + "Kafka action-state cleanup control topic must not be blank"); + Preconditions.checkArgument( + cleanupControlTopic == null || !this.tombstoneEnabled, + "Per-key Kafka tombstones cannot be enabled with checkpoint-aligned cleanup"); this.topic = Preconditions.checkNotNull( agentConfiguration.get(KAFKA_ACTION_STATE_TOPIC), "Kafka action state topic must be configured"); // create the topic if not exists maybeCreateTopic(); - Properties producerProp = createProducerProp(); - this.producer = new KafkaProducer<>(producerProp); - Properties consumerProp = createConsumerProp(); - this.consumer = new KafkaConsumer<>(consumerProp); + this.topicMetadataLoader = this::loadTopicMetadata; + try { + this.topicMetadata = topicMetadataLoader.load(); + } catch (Exception e) { + throw new RuntimeException("Failed to load Kafka topic metadata for " + topic, e); + } + KafkaActionStateCleanupCoordinator createdCoordinator = null; + Producer createdProducer = null; + Consumer createdConsumer = null; + try { + createdCoordinator = + cleanupControlTopic == null + ? null + : KafkaActionStateCleanupCoordinator.createForRecovery( + agentConfiguration); + createdProducer = new KafkaProducer<>(createProducerProp()); + createdConsumer = new KafkaConsumer<>(createConsumerProp()); + } catch (RuntimeException | Error failure) { + closeAfterInitializationFailure(createdConsumer, failure); + closeAfterInitializationFailure(createdProducer, failure); + closeAfterInitializationFailure(createdCoordinator, failure); + throw failure; + } + this.cleanupCoordinator = createdCoordinator; + this.producer = createdProducer; + this.consumer = createdConsumer; LOG.info("Initialized KafkaActionStateStore with topic: {}", topic); } @@ -221,59 +284,64 @@ private boolean checkDivergence(String businessKeyIdentity, long seqNum) { @Override public void rebuildState(List recoveryMarkers) { + Preconditions.checkNotNull(recoveryMarkers, "Recovery markers must not be null"); LOG.info("Rebuilding state from {} recovery markers", recoveryMarkers.size()); - if (recoveryMarkers.isEmpty()) { - LOG.info("No recovery markers, skipping state rebuild"); - return; - } - try { - Map partitionMap = new HashMap<>(); - // Process recovery markers to get the smallest offsets for each partition - for (Object marker : recoveryMarkers) { - if (marker instanceof Map) { - @SuppressWarnings("unchecked") - Map markerMap = (Map) marker; - for (Map.Entry entry : markerMap.entrySet()) { - Long offset = - partitionMap.computeIfPresent( - entry.getKey(), - (key, value) -> Math.min(value, entry.getValue())); - partitionMap.put( - entry.getKey(), offset == null ? entry.getValue() : offset); - } + KafkaTopicMetadata currentTopicMetadata = loadVerifiedTopicMetadata(); + if (recoveryMarkers.isEmpty()) { + if (cleanupCoordinator != null) { + Map boundary = + cleanupCoordinator.getCommittedBoundary( + topic, + currentTopicMetadata.getTopicId(), + currentTopicMetadata.getPartitions()); + Preconditions.checkState( + boundary.isEmpty(), + "Cannot initialize Kafka action state without a recovery marker because cleanup boundary %s is committed", + boundary); } + LOG.info("No recovery markers, skipping state rebuild"); + return; } - // Build list of TopicPartitions to assign + Map partitionMap = + mergeAndValidateRecoveryMarkers(recoveryMarkers, currentTopicMetadata); + + if (cleanupCoordinator != null) { + cleanupCoordinator.validateRecoveryOffsets( + topic, currentTopicMetadata.getTopicId(), partitionMap); + } + loadVerifiedTopicMetadata(); + List partitionsToAssign = new ArrayList<>(); for (Integer partition : partitionMap.keySet()) { partitionsToAssign.add(new TopicPartition(topic, partition)); } - // Assign partitions - consumer.assign(partitionsToAssign); + Map beginningOffsets = + consumer.beginningOffsets(partitionsToAssign); + Map replayEndOffsets = consumer.endOffsets(partitionsToAssign); + validateRecoveryOffsets(partitionMap, beginningOffsets, replayEndOffsets); - // Seek to marker offsets - if (!partitionMap.isEmpty()) { - // Seek to marker offsets - partitionMap.forEach( - (partition, offset) -> - consumer.seek(new TopicPartition(topic, partition), offset)); - } + consumer.assign(partitionsToAssign); + partitionMap.forEach( + (partition, offset) -> + consumer.seek(new TopicPartition(topic, partition), offset)); - // Poll for records and rebuild state until the latest offsets - while (true) { + Map lastPositions = consumerPositions(replayEndOffsets); + long lastProgressNanos = System.nanoTime(); + while (!hasReachedReplayEnd(lastPositions, replayEndOffsets)) { ConsumerRecords records = consumer.poll(CONSUMER_POLL_TIMEOUT); - if (records.isEmpty()) { - // reaches to the end of the topic - break; - } - // Deserialization failures throw from poll() itself and are handled by the // outer catch, so records here are always fully deserialized. for (ConsumerRecord record : records) { + TopicPartition recordPartition = + new TopicPartition(record.topic(), record.partition()); + Long replayEndOffset = replayEndOffsets.get(recordPartition); + if (replayEndOffset == null || record.offset() >= replayEndOffset) { + continue; + } if (!keyEncoder.isKeyRetained(ownershipFilter, record.key())) { continue; } @@ -285,15 +353,167 @@ public void rebuildState(List recoveryMarkers) { } } - // Commit offsets manually - consumer.commitSync(); + Map currentPositions = consumerPositions(replayEndOffsets); + if (!currentPositions.equals(lastPositions)) { + lastPositions = currentPositions; + lastProgressNanos = System.nanoTime(); + } else if (System.nanoTime() - lastProgressNanos + >= REBUILD_PROGRESS_TIMEOUT.toNanos()) { + throw new IllegalStateException( + String.format( + "Kafka action-state replay made no progress for %s; current positions are %s and target end offsets are %s", + REBUILD_PROGRESS_TIMEOUT, currentPositions, replayEndOffsets)); + } } + loadVerifiedTopicMetadata(); LOG.info("Completed rebuilding state, recovered {} states", actionStates.size()); } catch (Exception e) { throw new RuntimeException("Failed to rebuild state from Kafka", e); } } + private Map mergeAndValidateRecoveryMarkers( + List recoveryMarkers, KafkaTopicMetadata topicMetadata) { + Map mergedOffsets = new HashMap<>(); + boolean foundVersionedMarker = false; + boolean foundLegacyMarker = false; + + for (Object marker : recoveryMarkers) { + if (marker == null) { + throw new IllegalArgumentException( + "Kafka action-state recovery marker must not be null"); + } + Map offsets; + if (marker instanceof KafkaActionStateRecoveryMarker) { + foundVersionedMarker = true; + KafkaActionStateRecoveryMarker versionedMarker = + (KafkaActionStateRecoveryMarker) marker; + validateVersionedMarker(versionedMarker, topicMetadata); + offsets = versionedMarker.getOffsets(); + } else if (marker instanceof Map) { + foundLegacyMarker = true; + offsets = readLegacyOffsets((Map) marker); + validatePartitionSet(offsets.keySet(), topicMetadata.getPartitions()); + } else { + throw new IllegalArgumentException( + "Unsupported Kafka action-state recovery marker: " + + marker.getClass().getName()); + } + + offsets.forEach( + (partition, offset) -> mergedOffsets.merge(partition, offset, Math::min)); + } + + if (foundVersionedMarker && foundLegacyMarker) { + throw new IllegalArgumentException( + "Cannot restore from a mixture of versioned and legacy Kafka recovery markers"); + } + if (cleanupCoordinator != null && foundLegacyMarker) { + throw new IllegalArgumentException( + "Checkpoint-aligned cleanup requires versioned Kafka recovery markers"); + } + return mergedOffsets; + } + + private void validateVersionedMarker( + KafkaActionStateRecoveryMarker marker, KafkaTopicMetadata topicMetadata) { + if (marker.getSchemaVersion() != KafkaActionStateRecoveryMarker.CURRENT_SCHEMA_VERSION) { + throw new IllegalArgumentException( + String.format( + "Unsupported Kafka action-state recovery marker schema %d, expected %d", + marker.getSchemaVersion(), + KafkaActionStateRecoveryMarker.CURRENT_SCHEMA_VERSION)); + } + if (!topic.equals(marker.getTopic())) { + throw new IllegalStateException( + String.format( + "Kafka action-state recovery marker references topic %s, but the configured topic is %s", + marker.getTopic(), topic)); + } + if (!topicMetadata.getTopicId().equals(marker.getTopicId())) { + throw new IllegalStateException( + String.format( + "Kafka action-state topic %s has ID %s, but the recovery marker expects %s; the topic may have been recreated", + topic, topicMetadata.getTopicId(), marker.getTopicId())); + } + validatePartitionSet(marker.getOffsets().keySet(), topicMetadata.getPartitions()); + } + + private Map readLegacyOffsets(Map marker) { + Map offsets = new HashMap<>(); + for (Map.Entry entry : marker.entrySet()) { + if (!(entry.getKey() instanceof Integer) || !(entry.getValue() instanceof Long)) { + throw new IllegalArgumentException( + "Legacy Kafka action-state recovery markers must map Integer partitions to Long offsets"); + } + int partition = (Integer) entry.getKey(); + long offset = (Long) entry.getValue(); + if (partition < 0 || offset < 0) { + throw new IllegalArgumentException( + "Legacy Kafka action-state recovery marker partitions and offsets must be non-negative"); + } + offsets.put(partition, offset); + } + return offsets; + } + + private void validatePartitionSet(Set markerPartitions, Set topicPartitions) { + if (!markerPartitions.equals(topicPartitions)) { + throw new IllegalStateException( + String.format( + "Kafka action-state recovery marker contains partitions %s, but topic %s currently has partitions %s", + markerPartitions, topic, topicPartitions)); + } + } + + private void validateRecoveryOffsets( + Map requestedOffsets, + Map beginningOffsets, + Map endOffsets) { + requestedOffsets.forEach( + (partition, requestedOffset) -> { + TopicPartition topicPartition = new TopicPartition(topic, partition); + Long beginningOffset = beginningOffsets.get(topicPartition); + Long endOffset = endOffsets.get(topicPartition); + if (beginningOffset == null || endOffset == null) { + throw new IllegalStateException( + String.format( + "Kafka did not return beginning and end offsets for %s", + topicPartition)); + } + if (requestedOffset < beginningOffset || requestedOffset > endOffset) { + throw new IllegalStateException( + String.format( + "Cannot rebuild Kafka action state for %s: requested offset %d is outside the available range [%d, %d]", + topicPartition, + requestedOffset, + beginningOffset, + endOffset)); + } + }); + } + + private Map consumerPositions( + Map replayEndOffsets) { + Map positions = new HashMap<>(); + replayEndOffsets.forEach( + (partition, replayEndOffset) -> + positions.put( + partition, + Math.min(consumer.position(partition), replayEndOffset))); + return positions; + } + + private boolean hasReachedReplayEnd( + Map positions, Map replayEndOffsets) { + for (Map.Entry entry : replayEndOffsets.entrySet()) { + if (positions.get(entry.getKey()) < entry.getValue()) { + return false; + } + } + return true; + } + @Override public void setOwnershipFilter(IntPredicate ownershipFilter) { this.ownershipFilter = ownershipFilter; @@ -352,29 +572,50 @@ public void pruneState(Object key, long seqNum) { LOG.debug("Pruned state for key: {} up to sequence number: {}", key, seqNum); } - /** - * In kafka's implementation, we always return the end offsets of each partitions as recovery - * markers. - */ + /** Returns a versioned marker containing the Kafka topic identity and current end offsets. */ @Override public Object getRecoveryMarker() { - Map recoveryMarker = new HashMap<>(); - try { + KafkaTopicMetadata currentTopicMetadata = loadVerifiedTopicMetadata(); List partitions = new ArrayList<>(); - for (PartitionInfo partitionInfo : consumer.partitionsFor(topic)) { - partitions.add(new TopicPartition(topic, partitionInfo.partition())); + for (Integer partition : currentTopicMetadata.getPartitions()) { + partitions.add(new TopicPartition(topic, partition)); } Map endOffsets = consumer.endOffsets(partitions); + Map recoveryOffsets = new HashMap<>(); for (Map.Entry entry : endOffsets.entrySet()) { - recoveryMarker.put(entry.getKey().partition(), entry.getValue()); + recoveryOffsets.put(entry.getKey().partition(), entry.getValue()); + } + if (!recoveryOffsets.keySet().equals(currentTopicMetadata.getPartitions())) { + throw new IllegalStateException( + String.format( + "Kafka topic %s returned end offsets for partitions %s, expected %s", + topic, + recoveryOffsets.keySet(), + currentTopicMetadata.getPartitions())); } + return new KafkaActionStateRecoveryMarker( + topic, currentTopicMetadata.getTopicId(), recoveryOffsets); } catch (Exception e) { LOG.error("Failed to verify Kafka topic: {}", topic, e); throw new RuntimeException("Failed to verify Kafka topic", e); } + } - return recoveryMarker; + private KafkaTopicMetadata loadVerifiedTopicMetadata() throws Exception { + KafkaTopicMetadata currentTopicMetadata = topicMetadataLoader.load(); + if (!topicMetadata.getTopicId().equals(currentTopicMetadata.getTopicId()) + || !topicMetadata.getPartitions().equals(currentTopicMetadata.getPartitions())) { + throw new IllegalStateException( + String.format( + "Kafka action-state topic %s changed while the job was running; expected ID %s and partitions %s but found ID %s and partitions %s", + topic, + topicMetadata.getTopicId(), + topicMetadata.getPartitions(), + currentTopicMetadata.getTopicId(), + currentTopicMetadata.getPartitions())); + } + return currentTopicMetadata; } @Override @@ -397,11 +638,44 @@ public void close() throws Exception { firstException = ExceptionUtils.firstOrSuppressed(t, firstException); } } + if (cleanupCoordinator != null) { + try { + cleanupCoordinator.close(); + } catch (Throwable t) { + firstException = ExceptionUtils.firstOrSuppressed(t, firstException); + } + } if (firstException != null) { ExceptionUtils.rethrowException(firstException); } } + private KafkaTopicMetadata loadTopicMetadata() throws Exception { + try (AdminClient adminClient = AdminClient.create(createCommonKafkaConfig())) { + DescribeTopicsResult result = adminClient.describeTopics(List.of(topic)); + TopicDescription description = + result.allTopicNames() + .get(DEFAULT_FUTURE_GET_TIMEOUT_MS, TimeUnit.MILLISECONDS) + .get(topic); + if (description == null) { + throw new IllegalStateException("Kafka topic does not exist: " + topic); + } + Set partitions = new HashSet<>(); + description.partitions().forEach(partition -> partitions.add(partition.partition())); + return new KafkaTopicMetadata(description.topicId().toString(), partitions); + } + } + + private static KafkaTopicMetadata testTopicMetadata( + Consumer consumer, String topic) { + Set partitions = new HashSet<>(); + List partitionInfos = consumer.partitionsFor(topic); + if (partitionInfos != null) { + partitionInfos.forEach(partition -> partitions.add(partition.partition())); + } + return new KafkaTopicMetadata("test-topic-id:" + topic, partitions); + } + private void maybeCreateTopic() { try (AdminClient adminClient = AdminClient.create(createCommonKafkaConfig())) { ListTopicsResult topics = adminClient.listTopics(); @@ -415,8 +689,13 @@ private void maybeCreateTopic() { agentConfiguration .get(KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR) .shortValue()); - // enable topic compaction - newTopic.configs(Map.of("cleanup.policy", "compact")); + // Keep compaction, enable explicit prefix deletion, and disable Kafka's independent + // time- and size-based deletion so retained recovery points remain available. + newTopic.configs( + Map.of( + "cleanup.policy", "compact,delete", + "retention.ms", "-1", + "retention.bytes", "-1")); adminClient.createTopics(List.of(newTopic)).all().get(); LOG.info("Created Kafka topic: {}", topic); } else { @@ -434,6 +713,18 @@ private Properties createCommonKafkaConfig() { return props; } + private static void closeAfterInitializationFailure( + AutoCloseable resource, Throwable initializationFailure) { + if (resource == null) { + return; + } + try { + resource.close(); + } catch (Throwable closeFailure) { + initializationFailure.addSuppressed(closeFailure); + } + } + private Properties createProducerProp() { Properties producerProps = new Properties(); producerProps.putAll(createCommonKafkaConfig()); @@ -447,7 +738,8 @@ private Properties createProducerProp() { return producerProps; } - private Properties createConsumerProp() { + @VisibleForTesting + Properties createConsumerProp() { Properties consumerProps = new Properties(); consumerProps.putAll(createCommonKafkaConfig()); @@ -456,8 +748,36 @@ private Properties createConsumerProp() { consumerProps.put(VALUE_DESERIALIZER_CLASS_CONFIG, ActionStateKafkaSeder.class.getName()); consumerProps.put( ConsumerConfig.GROUP_ID_CONFIG, "action-state-rebuild-" + UUID.randomUUID()); - consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); + consumerProps.put(ENABLE_AUTO_COMMIT_CONFIG, false); return consumerProps; } + + static final class KafkaTopicMetadata { + private final String topicId; + private final Set partitions; + + KafkaTopicMetadata(String topicId, Set partitions) { + this.topicId = Preconditions.checkNotNull(topicId, "Topic ID must not be null"); + this.partitions = + Collections.unmodifiableSet( + new HashSet<>( + Preconditions.checkNotNull( + partitions, "Partitions must not be null"))); + } + + String getTopicId() { + return topicId; + } + + Set getPartitions() { + return partitions; + } + } + + @FunctionalInterface + private interface TopicMetadataLoader { + KafkaTopicMetadata load() throws Exception; + } } diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/operator/DurableExecutionManager.java b/runtime/src/main/java/org/apache/flink/agents/runtime/operator/DurableExecutionManager.java index fdf76ab897..49f10995bb 100644 --- a/runtime/src/main/java/org/apache/flink/agents/runtime/operator/DurableExecutionManager.java +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/operator/DurableExecutionManager.java @@ -54,6 +54,7 @@ import static org.apache.flink.agents.api.configuration.AgentConfigOptions.ACTION_STATE_STORE_BACKEND; import static org.apache.flink.agents.runtime.actionstate.ActionStateStore.BackendType.FLUSS; import static org.apache.flink.agents.runtime.actionstate.ActionStateStore.BackendType.KAFKA; +import static org.apache.flink.agents.runtime.actionstate.KafkaActionStateRecoveryMarker.UNION_STATE_NAME; /** * Owns the durable-execution side of {@link ActionExecutionOperator}: the optional {@link @@ -85,7 +86,6 @@ class DurableExecutionManager implements ActionStatePersister, AutoCloseable { private static final Logger LOG = LoggerFactory.getLogger(DurableExecutionManager.class); - private static final String RECOVERY_MARKER_STATE_NAME = "recoveryMarker"; private static final String LAST_COMPLETED_SEQUENCE_NUMBER_STATE_NAME = "lastCompletedSequenceNumber"; @@ -147,7 +147,7 @@ void initRecoveryMarkerState(OperatorStateBackend operatorStateBackend) throws E recoveryMarkerOpState = operatorStateBackend.getUnionListState( new ListStateDescriptor<>( - RECOVERY_MARKER_STATE_NAME, TypeInformation.of(Object.class))); + UNION_STATE_NAME, TypeInformation.of(Object.class))); } } @@ -215,7 +215,7 @@ void handleRecovery( ListState markerState = operatorStateBackend.getUnionListState( new ListStateDescriptor<>( - RECOVERY_MARKER_STATE_NAME, TypeInformation.of(Object.class))); + UNION_STATE_NAME, TypeInformation.of(Object.class))); Iterable recoveryMarkers = markerState.get(); if (recoveryMarkers != null) { recoveryMarkers.forEach(markers::add); diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/ActionStateSerializerRestoreTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/ActionStateSerializerRestoreTest.java index f837b29e3f..9b26d6c834 100644 --- a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/ActionStateSerializerRestoreTest.java +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/ActionStateSerializerRestoreTest.java @@ -35,6 +35,7 @@ import org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.MockConsumer; +import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; @@ -96,6 +97,9 @@ void recoveryKeepsCompletedStateAfterPojoSubclassCacheChanges(boolean useSubclas Map cache = new HashMap<>(); MockConsumer consumer = new MockConsumer<>("earliest"); TopicPartition partition = new TopicPartition(TOPIC, 0); + consumer.updatePartitions( + TOPIC, List.of(new PartitionInfo(TOPIC, 0, null, null, null))); + consumer.updateEndOffsets(Map.of(partition, 1L)); consumer.assign(List.of(partition)); consumer.updateBeginningOffsets(Map.of(partition, 0L)); consumer.addRecord(new ConsumerRecord<>(TOPIC, 0, 0, stateKey, completed)); diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorIntegrationTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorIntegrationTest.java new file mode 100644 index 0000000000..114f2cfa64 --- /dev/null +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorIntegrationTest.java @@ -0,0 +1,148 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.agents.plan.AgentConfiguration; +import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.OffsetSpec; +import org.apache.kafka.clients.admin.TopicDescription; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.serialization.StringSerializer; +import org.junit.jupiter.api.Test; +import org.testcontainers.DockerClientFactory; +import org.testcontainers.kafka.KafkaContainer; +import org.testcontainers.utility.DockerImageName; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeUnit; + +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC_NUM_PARTITIONS; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_BOOTSTRAP_SERVERS; +import static org.apache.flink.agents.runtime.actionstate.ActionStateTestUtils.createKeyEncoder; +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +/** Kafka integration tests for checkpoint-aligned action-state cleanup. */ +class KafkaActionStateCleanupCoordinatorIntegrationTest { + + private static final long TIMEOUT_SECONDS = 30; + + @Test + void testCommitsDeletesVerifiesAndRecoversBoundary() throws Exception { + assumeTrue( + DockerClientFactory.instance().isDockerAvailable(), + "Docker is required for the Kafka cleanup integration test"); + + try (KafkaContainer kafka = + new KafkaContainer(DockerImageName.parse("apache/kafka-native:3.8.0"))) { + kafka.start(); + String dataTopic = "action-state"; + String controlTopic = "action-state-control"; + Properties common = new Properties(); + common.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, kafka.getBootstrapServers()); + + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(KAFKA_BOOTSTRAP_SERVERS, kafka.getBootstrapServers()); + configuration.set(KAFKA_ACTION_STATE_TOPIC, dataTopic); + configuration.set(KAFKA_ACTION_STATE_TOPIC_NUM_PARTITIONS, 2); + try (KafkaActionStateStore ignored = + new KafkaActionStateStore(configuration, createKeyEncoder(128))) { + // Exercise the exact topic configuration created by the production store. + } + + String topicId; + try (AdminClient admin = AdminClient.create(common)) { + TopicDescription description = + admin.describeTopics(List.of(dataTopic)) + .allTopicNames() + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS) + .get(dataTopic); + topicId = description.topicId().toString(); + } + + Properties producerProperties = new Properties(); + producerProperties.putAll(common); + producerProperties.put( + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put( + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + try (KafkaProducer producer = new KafkaProducer<>(producerProperties)) { + for (int offset = 0; offset < 3; offset++) { + producer.send( + new ProducerRecord<>( + dataTopic, 0, "p0-" + offset, "value-" + offset)) + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } + for (int offset = 0; offset < 2; offset++) { + producer.send( + new ProducerRecord<>( + dataTopic, 1, "p1-" + offset, "value-" + offset)) + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } + } + + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of( + new KafkaActionStateRecoveryMarker( + dataTopic, topicId, Map.of(0, 2L, 1, 1L)))); + configuration.set(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, controlTopic); + + try (KafkaActionStateCleanupCoordinator coordinator = + KafkaActionStateCleanupCoordinator.create(configuration)) { + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + } + + try (KafkaActionStateCleanupCoordinator recoveryCoordinator = + KafkaActionStateCleanupCoordinator.createForRecovery(configuration)) { + assertThat( + recoveryCoordinator.getCommittedBoundary( + dataTopic, topicId, Set.of(0, 1))) + .isEqualTo(Map.of(0, 2L, 1, 1L)); + } + + try (AdminClient admin = AdminClient.create(common)) { + Map requests = new HashMap<>(); + requests.put(new TopicPartition(dataTopic, 0), OffsetSpec.earliest()); + requests.put(new TopicPartition(dataTopic, 1), OffsetSpec.earliest()); + Map beginningOffsets = new HashMap<>(); + admin.listOffsets(requests) + .all() + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS) + .forEach( + (partition, result) -> + beginningOffsets.put(partition, result.offset())); + assertThat(beginningOffsets) + .containsEntry(new TopicPartition(dataTopic, 0), 2L) + .containsEntry(new TopicPartition(dataTopic, 1), 1L); + } + } + } +} diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorTest.java new file mode 100644 index 0000000000..51f8eec5ca --- /dev/null +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupCoordinatorTest.java @@ -0,0 +1,660 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.agents.plan.AgentConfiguration; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; + +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOMBSTONE_ENABLED; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR; +import static org.apache.flink.agents.api.configuration.AgentConfigOptions.KAFKA_BOOTSTRAP_SERVERS; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests for {@link KafkaActionStateCleanupCoordinator}. */ +class KafkaActionStateCleanupCoordinatorTest { + + @Test + void testCommitsBeforeDeleteAndMarksAppliedAfterVerification() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + + assertThat(transport.events) + .containsExactly("append:COMMITTED", "delete:{0=10, 1=20}", "append:APPLIED"); + assertThat(transport.deletedOffsets).containsExactlyInAnyOrderEntriesOf(plan.getOffsets()); + assertThat(coordinator.getCommittedBoundary("action-state", "topic-id", Set.of(0, 1))) + .isEqualTo(plan.getOffsets()); + } + + @Test + void testDeleteFailureLeavesCommittedPlanForRetry() throws Exception { + FakeTransport transport = new FakeTransport(); + transport.failNextDelete = true; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThatThrownBy(() -> coordinator.apply(plan)) + .isInstanceOf(IllegalStateException.class) + .hasMessage("simulated delete failure"); + assertThat(transport.operations.get(plan.getPlanId()).getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.COMMITTED); + + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + assertThat(transport.events) + .containsExactly( + "append:COMMITTED", + "delete:{0=10, 1=20}", + "delete:{0=10, 1=20}", + "append:APPLIED"); + } + + @Test + void testAppliedRecordFailureRetriesDeletionFromCommittedPlan() throws Exception { + FakeTransport transport = new FakeTransport(); + transport.failNextAppliedAppend = true; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThatThrownBy(() -> coordinator.apply(plan)) + .isInstanceOf(IllegalStateException.class) + .hasMessage("simulated applied append failure"); + assertThat(transport.operations.get(plan.getPlanId()).getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.COMMITTED); + + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + assertThat(transport.events) + .containsExactly( + "append:COMMITTED", + "delete:{0=10, 1=20}", + "delete:{0=10, 1=20}", + "append:APPLIED"); + } + + @Test + void testCommittedBoundaryRejectsOldRestoreBeforeDeleteSucceeds() throws Exception { + FakeTransport transport = new FakeTransport(); + transport.failNextDelete = true; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThatThrownBy(() -> coordinator.apply(plan)).isInstanceOf(IllegalStateException.class); + + assertThatThrownBy( + () -> + coordinator.validateRecoveryOffsets( + "action-state", "topic-id", Map.of(0, 9L, 1, 20L))) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("committed cleanup boundary is 10"); + coordinator.validateRecoveryOffsets("action-state", "topic-id", Map.of(0, 10L, 1, 21L)); + } + + @Test + void testRecoveryFailsClosedWhenControlHistoryIsEmpty() { + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(new FakeTransport(), true); + + assertThatThrownBy( + () -> + coordinator.getCommittedBoundary( + "action-state", "topic-id", Set.of(0, 1))) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("contains no committed boundary"); + } + + @Test + void testRejectsBackwardBoundary() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + coordinator.apply(plan("checkpoint-42", 10L, 20L)); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-41", 9L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("does not advance"); + } + + @Test + void testRejectsBoundaryBelowAvailableBeginningBeforeCommit() { + FakeTransport transport = new FakeTransport(); + transport.beginningOffsets.put(0, 11L); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("available range is [11"); + assertThat(transport.operations).isEmpty(); + assertThat(transport.events).isEmpty(); + } + + @Test + void testRejectsBoundaryAboveAvailableEndBeforeCommit() { + FakeTransport transport = new FakeTransport(); + transport.endOffsets.put(1, 19L); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("available range is [0, 19]"); + assertThat(transport.operations).isEmpty(); + assertThat(transport.events).isEmpty(); + } + + @Test + void testRejectsUnsafeDataTopicConfigurationBeforeCommit() { + FakeTransport transport = new FakeTransport(); + transport.failDataTopicConfigurationOnCheck = 1; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalStateException.class) + .hasMessage("unsafe data topic configuration"); + assertThat(transport.operations).isEmpty(); + assertThat(transport.events).isEmpty(); + } + + @Test + void testRevalidatesDataTopicConfigurationBeforeDelete() { + FakeTransport transport = new FakeTransport(); + transport.failDataTopicConfigurationOnCheck = 3; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThatThrownBy(() -> coordinator.apply(plan)) + .isInstanceOf(IllegalStateException.class) + .hasMessage("unsafe data topic configuration"); + assertThat(transport.operations.get(plan.getPlanId()).getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.COMMITTED); + assertThat(transport.events).containsExactly("append:COMMITTED"); + } + + @Test + void testRevalidatesDataTopicConfigurationImmediatelyBeforeCommit() { + FakeTransport transport = new FakeTransport(); + transport.failDataTopicConfigurationOnCheck = 2; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalStateException.class) + .hasMessage("unsafe data topic configuration"); + assertThat(transport.operations).isEmpty(); + assertThat(transport.events).isEmpty(); + } + + @Test + void testRevalidatesDataTopicConfigurationAfterDelete() { + FakeTransport transport = new FakeTransport(); + transport.failDataTopicConfigurationOnCheck = 4; + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + + assertThatThrownBy(() -> coordinator.apply(plan)) + .isInstanceOf(IllegalStateException.class) + .hasMessage("unsafe data topic configuration"); + assertThat(transport.operations.get(plan.getPlanId()).getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.COMMITTED); + assertThat(transport.events).containsExactly("append:COMMITTED", "delete:{0=10, 1=20}"); + } + + @Test + void testCommittedPlanRetriesAfterKafkaBeginningPassesBoundary() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + transport.operations.put( + plan.getPlanId(), KafkaActionStateCleanupCoordinator.Operation.committed(plan)); + transport.beginningOffsets.put(0, 11L); + transport.beginningOffsets.put(1, 21L); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + assertThat(transport.events).containsExactly("delete:{0=10, 1=20}", "append:APPLIED"); + } + + @Test + void testRejectsIncomparableCommittedBoundariesBeforeDeletion() { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupPlan left = plan("checkpoint-left", 10L, 30L); + KafkaActionStateCleanupPlan right = plan("checkpoint-right", 20L, 20L); + transport.operations.put( + left.getPlanId(), KafkaActionStateCleanupCoordinator.Operation.committed(left)); + transport.operations.put( + right.getPlanId(), KafkaActionStateCleanupCoordinator.Operation.committed(right)); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(right)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("incomparable boundaries"); + assertThat(transport.events).isEmpty(); + } + + @Test + void testAppliedPlanIsIdempotent() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + coordinator.apply(plan); + transport.events.clear(); + + assertThat(coordinator.apply(plan)) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + assertThat(transport.events).isEmpty(); + } + + @Test + void testCloseClosesTransport() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + coordinator.close(); + + assertThat(transport.closed).isTrue(); + } + + @Test + void testLaterPlanAdvancesEveryPartition() throws Exception { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + coordinator.apply(plan("checkpoint-41", 10L, 20L)); + transport.events.clear(); + + KafkaActionStateCleanupPlan later = plan("checkpoint-42", 15L, 25L); + coordinator.apply(later); + + assertThat(transport.events) + .containsExactly("append:COMMITTED", "delete:{0=15, 1=25}", "append:APPLIED"); + assertThat(coordinator.getCommittedBoundary("action-state", "topic-id", Set.of(0, 1))) + .isEqualTo(later.getOffsets()); + } + + @Test + void testControlRecordJsonRoundTrip() throws Exception { + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + KafkaActionStateCleanupCoordinator.Operation operation = + KafkaActionStateCleanupCoordinator.Operation.committed(plan); + + KafkaActionStateCleanupCoordinator.Operation restored = + KafkaActionStateCleanupCoordinator.Operation.fromJson(operation.toJson()); + + assertThat(restored.getPlan()).isEqualTo(plan); + assertThat(restored.getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.COMMITTED); + } + + @Test + void testRejectsControlTopicTombstoneWithoutRollingBackBoundary() { + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + Map operations = + new LinkedHashMap<>(); + operations.put( + plan.getPlanId(), KafkaActionStateCleanupCoordinator.Operation.committed(plan)); + ConsumerRecord tombstone = + new ConsumerRecord<>("action-state-control", 0, 1L, plan.getPlanId(), null); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.applyControlRecord( + operations, tombstone)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("cleanup boundaries cannot be deleted"); + assertThat(operations).containsKey(plan.getPlanId()); + } + + @Test + void testRejectsAppliedStatusRegressionWithoutChangingBoundary() throws Exception { + KafkaActionStateCleanupPlan plan = plan("checkpoint-42", 10L, 20L); + KafkaActionStateCleanupCoordinator.Operation applied = + KafkaActionStateCleanupCoordinator.Operation.committed(plan) + .withStatus(KafkaActionStateCleanupCoordinator.Status.APPLIED); + Map operations = + new LinkedHashMap<>(); + operations.put(plan.getPlanId(), applied); + ConsumerRecord regression = + new ConsumerRecord<>( + "action-state-control", + 0, + 1L, + plan.getPlanId(), + KafkaActionStateCleanupCoordinator.Operation.committed(plan).toJson()); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.applyControlRecord( + operations, regression)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("regressed applied plan"); + assertThat(operations.get(plan.getPlanId()).getStatus()) + .isEqualTo(KafkaActionStateCleanupCoordinator.Status.APPLIED); + } + + @ParameterizedTest + @ValueSource(strings = {"{}", "null", "garbage"}) + void testRejectsTrailingControlRecordContent(String trailingContent) throws Exception { + KafkaActionStateCleanupCoordinator.Operation operation = + KafkaActionStateCleanupCoordinator.Operation.committed( + plan("checkpoint-42", 10L, 20L)); + ConsumerRecord record = + new ConsumerRecord<>( + "action-state-control", + 0, + 0L, + operation.getPlan().getPlanId(), + operation.toJson() + trailingContent); + Map operations = + new LinkedHashMap<>(); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.applyControlRecord( + operations, record)) + .isInstanceOf(IOException.class); + assertThat(operations).isEmpty(); + } + + @ParameterizedTest + @ValueSource(ints = {-1, 0, 32768, 65537}) + void testRejectsInvalidReplicationFactorBeforeCreatingKafkaClients(int replicationFactor) { + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(KAFKA_ACTION_STATE_TOPIC, "action-state"); + configuration.set(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, "action-state-control"); + configuration.set(KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR, replicationFactor); + // Fail immediately if validation reaches client setup instead of rejecting the factor. + configuration.set(KAFKA_BOOTSTRAP_SERVERS, ""); + + assertThatThrownBy(() -> KafkaActionStateCleanupCoordinator.create(configuration)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("replication factor") + .hasMessageContaining("32767") + .hasMessageContaining(Integer.toString(replicationFactor)); + } + + @Test + void testRejectsFractionalControlRecordSchema() throws Exception { + KafkaActionStateCleanupCoordinator.Operation operation = + KafkaActionStateCleanupCoordinator.Operation.committed( + plan("checkpoint-42", 10L, 20L)); + String malformed = + operation.toJson().replace("\"schemaVersion\":1", "\"schemaVersion\":1.5"); + + assertThatThrownBy(() -> KafkaActionStateCleanupCoordinator.Operation.fromJson(malformed)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Unsupported cleanup control record schema"); + } + + @Test + void testRejectsRecreatedDataTopicBeforeCommit() { + FakeTransport transport = new FakeTransport(); + transport.metadata = + new KafkaActionStateCleanupCoordinator.TopicMetadata( + "recreated-topic-id", Set.of(0, 1)); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("expects topic-id"); + assertThat(transport.operations).isEmpty(); + } + + @Test + void testRevalidatesDataTopicImmediatelyBeforeDelete() { + FakeTransport transport = new FakeTransport(); + transport.metadataByDescribeCall.put( + 2, + new KafkaActionStateCleanupCoordinator.TopicMetadata( + "recreated-topic-id", Set.of(0, 1))); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("expects topic-id"); + assertThat(transport.events).containsExactly("append:COMMITTED"); + } + + @Test + void testRevalidatesDataTopicAfterDeleteBeforeMarkingApplied() { + FakeTransport transport = new FakeTransport(); + transport.metadataByDescribeCall.put( + 3, + new KafkaActionStateCleanupCoordinator.TopicMetadata( + "recreated-topic-id", Set.of(0, 1))); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("expects topic-id"); + assertThat(transport.events).containsExactly("append:COMMITTED", "delete:{0=10, 1=20}"); + assertThat(transport.operations.values()) + .allMatch( + operation -> + operation.getStatus() + == KafkaActionStateCleanupCoordinator.Status.COMMITTED); + } + + @Test + void testRejectsPlanForDifferentConfiguredDataTopicBeforeCommit() { + FakeTransport transport = new FakeTransport(); + KafkaActionStateCleanupCoordinator coordinator = + new KafkaActionStateCleanupCoordinator(transport, false, "other-action-state"); + + assertThatThrownBy(() -> coordinator.apply(plan("checkpoint-42", 10L, 20L))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("configured for data topic other-action-state"); + assertThat(transport.operations).isEmpty(); + } + + @Test + void testRejectsTombstonesWithCheckpointAlignedCleanupBeforeConnecting() { + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(KAFKA_ACTION_STATE_TOPIC, "action-state"); + configuration.set(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, "action-state-control"); + configuration.set(KAFKA_ACTION_STATE_TOMBSTONE_ENABLED, true); + + assertThatThrownBy(() -> KafkaActionStateCleanupCoordinator.create(configuration)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("tombstones cannot be enabled"); + } + + @Test + void testRejectsSharedDataAndControlTopicBeforeConnecting() { + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(KAFKA_ACTION_STATE_TOPIC, "action-state"); + configuration.set(KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, "action-state"); + + assertThatThrownBy(() -> KafkaActionStateCleanupCoordinator.create(configuration)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("must differ from the data topic"); + } + + @Test + void testControlTopicRequiresCompactionWithoutDeleteRetention() { + assertThatCode( + () -> + KafkaActionStateCleanupCoordinator + .validateControlTopicCleanupPolicy( + "action-state-control", "compact")) + .doesNotThrowAnyException(); + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator + .validateControlTopicCleanupPolicy( + "action-state-control", "compact,delete")) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("cleanup.policy=compact without delete retention"); + } + + @Test + void testDataTopicRequiresExplicitDeletionWithoutAutomaticRetention() { + assertThatCode( + () -> + KafkaActionStateCleanupCoordinator.validateDataTopicConfiguration( + "action-state", "delete,compact", "-1", "-1")) + .doesNotThrowAnyException(); + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.validateDataTopicConfiguration( + "action-state", "compact", "-1", "-1")) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("cleanup.policy=compact,delete"); + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.validateDataTopicConfiguration( + "action-state", "compact,delete", "604800000", "-1")) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("retention.ms=-1"); + assertThatThrownBy( + () -> + KafkaActionStateCleanupCoordinator.validateDataTopicConfiguration( + "action-state", "compact,delete", "-1", "1048576")) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("retention.bytes=-1"); + } + + private static KafkaActionStateCleanupPlan plan( + String recoveryPoint, long partitionZero, long partitionOne) { + return KafkaActionStateCleanupPlan.fromRecoveryMarkers( + recoveryPoint, + List.of( + new KafkaActionStateRecoveryMarker( + "action-state", + "topic-id", + Map.of(0, partitionZero, 1, partitionOne)))); + } + + private static final class FakeTransport + implements KafkaActionStateCleanupCoordinator.Transport { + private KafkaActionStateCleanupCoordinator.TopicMetadata metadata = + new KafkaActionStateCleanupCoordinator.TopicMetadata("topic-id", Set.of(0, 1)); + private final Map + metadataByDescribeCall = new HashMap<>(); + private final Map operations = + new LinkedHashMap<>(); + private final List events = new ArrayList<>(); + private final Map deletedOffsets = new HashMap<>(); + private final Map beginningOffsets = new HashMap<>(Map.of(0, 0L, 1, 0L)); + private final Map endOffsets = + new HashMap<>(Map.of(0, Long.MAX_VALUE, 1, Long.MAX_VALUE)); + private boolean failNextDelete; + private boolean failNextAppliedAppend; + private int dataTopicConfigurationChecks; + private int failDataTopicConfigurationOnCheck = -1; + private boolean closed; + private int describeCalls; + + @Override + public KafkaActionStateCleanupCoordinator.TopicMetadata describeTopic(String topic) { + describeCalls++; + return metadataByDescribeCall.getOrDefault(describeCalls, metadata); + } + + @Override + public Map readOperations() { + return new LinkedHashMap<>(operations); + } + + @Override + public void append(KafkaActionStateCleanupCoordinator.Operation operation) { + if (operation.getStatus() == KafkaActionStateCleanupCoordinator.Status.APPLIED + && failNextAppliedAppend) { + failNextAppliedAppend = false; + throw new IllegalStateException("simulated applied append failure"); + } + events.add("append:" + operation.getStatus()); + operations.put(operation.getPlan().getPlanId(), operation); + } + + @Override + public void validateDataTopicConfiguration(String topic) { + dataTopicConfigurationChecks++; + if (dataTopicConfigurationChecks == failDataTopicConfigurationOnCheck) { + throw new IllegalStateException("unsafe data topic configuration"); + } + } + + @Override + public void validateOffsetsAvailable( + String topic, String topicId, Map offsets) { + offsets.forEach( + (partition, requestedOffset) -> { + long beginningOffset = beginningOffsets.get(partition); + long endOffset = endOffsets.get(partition); + if (requestedOffset < beginningOffset || requestedOffset > endOffset) { + throw new IllegalArgumentException( + String.format( + "Cannot commit Kafka cleanup boundary for %s-%s at offset %s because the available range is [%s, %s]", + topic, + partition, + requestedOffset, + beginningOffset, + endOffset)); + } + }); + } + + @Override + public void deleteBefore(String topic, String topicId, Map offsets) { + events.add("delete:" + new java.util.TreeMap<>(offsets)); + if (failNextDelete) { + failNextDelete = false; + throw new IllegalStateException("simulated delete failure"); + } + deletedOffsets.putAll(offsets); + } + + @Override + public void close() { + closed = true; + } + } +} diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlanTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlanTest.java new file mode 100644 index 0000000000..0fd2556687 --- /dev/null +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupPlanTest.java @@ -0,0 +1,199 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.kafka.common.Uuid; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.io.IOException; +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.assertj.core.api.Assertions.entry; + +/** Tests for {@link KafkaActionStateCleanupPlan}. */ +class KafkaActionStateCleanupPlanTest { + + @Test + void testUsesPerPartitionMinimumAcrossUnionMarkers() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of(marker(Map.of(0, 15L, 1, 30L)), marker(Map.of(0, 10L, 1, 35L)))); + + assertThat(plan.getOffsets()).containsOnly(entry(0, 10L), entry(1, 30L)); + assertThat(plan.getSourceRecoveryPoint()).isEqualTo("checkpoint-42"); + assertThat(plan.getPlanId()).hasSize(64); + } + + @Test + void testJsonRoundTripPreservesReviewedPlan() throws Exception { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "s3://checkpoints/savepoint-1", List.of(marker(Map.of(0, 10L, 1, 20L)))); + + KafkaActionStateCleanupPlan restored = KafkaActionStateCleanupPlan.fromJson(plan.toJson()); + + assertThat(restored).isEqualTo(plan); + } + + @ParameterizedTest + @ValueSource(strings = {"{}", "null", "garbage"}) + void testRejectsTrailingPlanContent(String trailingContent) { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + + assertThatThrownBy( + () -> KafkaActionStateCleanupPlan.fromJson(plan.toJson() + trailingContent)) + .isInstanceOf(IOException.class); + } + + @Test + void testRejectsTamperedJson() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String tampered = plan.toJson().replace("\"0\" : 10", "\"0\" : 11"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(tampered)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("plan ID does not match"); + } + + @Test + void testRejectsNonStringPlanId() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String malformed = plan.toJson().replace('"' + plan.getPlanId() + '"', "123"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(malformed)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Cleanup plan planId must be a string"); + } + + @Test + void testRejectsUnknownPlanField() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String malformed = plan.toJson().replaceFirst("\\{", "{ \"unexpected\" : true,"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(malformed)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("fields must be exactly"); + } + + @Test + void testRejectsDuplicatePlanField() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String malformed = plan.toJson().replaceFirst("\\{", "{ \"schemaVersion\" : 1,"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(malformed)) + .isInstanceOf(IOException.class) + .hasMessageContaining("Duplicate field 'schemaVersion'"); + } + + @Test + void testRejectsNonCanonicalPartitionKeyAlias() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String malformed = plan.toJson().replace("\"0\" : 10", "\"+0\" : 999,\n \"0\" : 10"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(malformed)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("partition +0 must use canonical decimal form"); + } + + @Test + void testRejectsFractionalOffsets() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker(Map.of(0, 10L, 1, 20L)))); + String fractional = plan.toJson().replace("\"0\" : 10", "\"0\" : 10.5"); + + assertThatThrownBy(() -> KafkaActionStateCleanupPlan.fromJson(fractional)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("must be an integer"); + } + + @Test + void testRejectsLegacyMarkersForDestructiveCleanup() { + assertThatThrownBy( + () -> + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(Map.of(0, 10L)))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("requires versioned Kafka recovery markers"); + } + + @Test + void testRejectsUnavailableTopicIdentity() { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + "action-state", Uuid.ZERO_UUID.toString(), Map.of(0, 10L)); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", List.of(marker))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("broker-provided Kafka topic ID"); + } + + @Test + void testRejectsMarkersFromDifferentPhysicalTopics() { + KafkaActionStateRecoveryMarker otherTopic = + new KafkaActionStateRecoveryMarker( + "action-state", "other-topic-id", Map.of(0, 10L, 1, 20L)); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of(marker(Map.of(0, 10L, 1, 20L)), otherTopic))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("one physical Kafka topic"); + } + + @Test + void testRejectsMarkersWithDifferentPartitionSets() { + KafkaActionStateRecoveryMarker missingPartition = + new KafkaActionStateRecoveryMarker("action-state", "topic-id", Map.of(0, 10L)); + + assertThatThrownBy( + () -> + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of(marker(Map.of(0, 10L, 1, 20L)), missingPartition))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("different Kafka partition sets"); + } + + private static KafkaActionStateRecoveryMarker marker(Map offsets) { + return new KafkaActionStateRecoveryMarker("action-state", "topic-id", offsets); + } +} diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupToolTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupToolTest.java new file mode 100644 index 0000000000..09af0a877e --- /dev/null +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateCleanupToolTest.java @@ -0,0 +1,223 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.api.common.functions.RichMapFunction; +import org.apache.flink.api.common.state.ListState; +import org.apache.flink.api.common.state.ListStateDescriptor; +import org.apache.flink.api.common.typeinfo.TypeInformation; +import org.apache.flink.core.execution.JobClient; +import org.apache.flink.core.execution.SavepointFormatType; +import org.apache.flink.runtime.state.FunctionInitializationContext; +import org.apache.flink.runtime.state.FunctionSnapshotContext; +import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction; +import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; +import org.apache.flink.streaming.api.functions.sink.legacy.SinkFunction; +import org.apache.flink.streaming.api.functions.source.legacy.RichParallelSourceFunction; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; + +/** Argument-contract tests for {@link KafkaActionStateCleanupTool}. */ +class KafkaActionStateCleanupToolTest { + + @Test + void testRequiresCommand() { + assertThatThrownBy(() -> KafkaActionStateCleanupTool.main(new String[0])) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Usage:"); + } + + @Test + void testHelpDoesNotRequireKafkaOrFlink() { + assertDoesNotThrow(() -> KafkaActionStateCleanupTool.main(new String[] {"--help"})); + } + + @Test + void testRejectsUnknownApplyOptionBeforeReadingPlan() { + assertThatThrownBy( + () -> + KafkaActionStateCleanupTool.main( + new String[] { + "apply", + "--plan", + "missing.json", + "--bootstrap-servers", + "localhost:9092", + "--control-topic", + "control", + "--typo", + "value" + })) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Unknown option: --typo"); + } + + @Test + void testPlanRequiresExactlyOneOperatorIdentifier() { + assertThatThrownBy( + () -> + KafkaActionStateCleanupTool.main( + new String[] { + "plan", + "--checkpoint", + "checkpoint", + "--output", + "plan.json", + "--operator-uid", + "uid", + "--operator-uid-hash", + "hash" + })) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Exactly one"); + } + + @Test + void testPlanReadsRealUnionStateFromSavepoint(@TempDir Path tempDirectory) throws Exception { + String operatorUid = "cleanup-plan-test"; + Path readyPath = tempDirectory.resolve("ready"); + Path savepointDirectory = tempDirectory.resolve("savepoints"); + Path planPath = tempDirectory.resolve("cleanup-plan.json"); + Files.createDirectories(savepointDirectory); + StreamExecutionEnvironment environment = + StreamExecutionEnvironment.getExecutionEnvironment(); + environment.setParallelism(1); + environment.enableCheckpointing(1000); + environment + .addSource(new WaitingSource()) + .map( + new RecoveryMarkerStatefulMapFunction( + readyPath.toString(), + List.of( + marker(Map.of(0, 15L, 1, 20L)), + marker(Map.of(0, 10L, 1, 25L))))) + .uid(operatorUid) + .setParallelism(1) + .addSink(new DiscardingSink<>()) + .setParallelism(1); + + JobClient jobClient = environment.executeAsync("write Kafka cleanup test savepoint"); + String savepoint; + try { + waitUntilReady(readyPath); + savepoint = + jobClient + .triggerSavepoint( + savepointDirectory.toString(), SavepointFormatType.CANONICAL) + .get(30, TimeUnit.SECONDS); + } finally { + jobClient.cancel().get(30, TimeUnit.SECONDS); + } + + KafkaActionStateCleanupTool.main( + new String[] { + "plan", + "--checkpoint", + savepoint, + "--operator-uid", + operatorUid, + "--output", + planPath.toString() + }); + + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromJson( + Files.readString(planPath, StandardCharsets.UTF_8)); + assertThat(plan.getOffsets()).isEqualTo(Map.of(0, 10L, 1, 20L)); + assertThat(plan.getSourceRecoveryPoint()).isEqualTo(savepoint); + } + + private static KafkaActionStateRecoveryMarker marker(Map offsets) { + return new KafkaActionStateRecoveryMarker("action-state", "topic-id", offsets); + } + + private static void waitUntilReady(Path readyPath) throws Exception { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(30); + while (!Files.exists(readyPath) && System.nanoTime() < deadline) { + Thread.sleep(20); + } + assertThat(readyPath).exists(); + } + + private static final class WaitingSource extends RichParallelSourceFunction { + + private volatile boolean running = true; + + @Override + public void run(SourceContext context) throws Exception { + synchronized (context.getCheckpointLock()) { + context.collect(1); + } + while (running) { + Thread.sleep(20); + } + } + + @Override + public void cancel() { + running = false; + } + } + + private static final class RecoveryMarkerStatefulMapFunction + extends RichMapFunction implements CheckpointedFunction { + + private transient ListState markerState; + private final String readyPath; + private final List markers; + + private RecoveryMarkerStatefulMapFunction(String readyPath, List markers) { + this.readyPath = readyPath; + this.markers = List.copyOf(markers); + } + + @Override + public Integer map(Integer value) throws Exception { + Files.writeString(Path.of(readyPath), "ready", StandardCharsets.UTF_8); + return value; + } + + @Override + public void snapshotState(FunctionSnapshotContext context) throws Exception { + markerState.update(markers); + } + + @Override + public void initializeState(FunctionInitializationContext context) throws Exception { + markerState = + context.getOperatorStateStore() + .getUnionListState( + new ListStateDescriptor<>( + KafkaActionStateRecoveryMarker.UNION_STATE_NAME, + TypeInformation.of(Object.class))); + } + } + + private static final class DiscardingSink implements SinkFunction {} +} diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarkerTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarkerTest.java new file mode 100644 index 0000000000..e989999a90 --- /dev/null +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateRecoveryMarkerTest.java @@ -0,0 +1,65 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.flink.agents.runtime.actionstate; + +import org.apache.flink.api.common.serialization.SerializerConfigImpl; +import org.apache.flink.api.common.typeinfo.TypeInformation; +import org.apache.flink.api.common.typeutils.TypeSerializer; +import org.apache.flink.core.memory.DataInputDeserializer; +import org.apache.flink.core.memory.DataOutputSerializer; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.assertj.core.api.Assertions.entry; + +/** Tests for {@link KafkaActionStateRecoveryMarker}. */ +class KafkaActionStateRecoveryMarkerTest { + + @Test + void testFlinkStateSerializationRoundTrip() throws Exception { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + "action-state", "topic-id", Map.of(0, 10L, 1, 20L)); + TypeSerializer serializer = + TypeInformation.of(Object.class).createSerializer(new SerializerConfigImpl()); + DataOutputSerializer output = new DataOutputSerializer(256); + + serializer.serialize(marker, output); + Object restored = + serializer.deserialize(new DataInputDeserializer(output.getCopyOfBuffer())); + + assertThat(restored).isEqualTo(marker); + } + + @Test + void testOffsetsAreDefensivelyCopied() { + Map offsets = new HashMap<>(Map.of(0, 10L)); + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker("action-state", "topic-id", offsets); + + offsets.put(1, 20L); + + assertThat(marker.getOffsets()).containsExactly(entry(0, 10L)); + assertThatThrownBy(() -> marker.getOffsets().put(1, 20L)) + .isInstanceOf(UnsupportedOperationException.class); + } +} diff --git a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java index ec5b772f2c..28a6aebbfb 100644 --- a/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java +++ b/runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java @@ -23,6 +23,7 @@ import org.apache.flink.agents.plan.AgentConfiguration; import org.apache.flink.agents.plan.actions.Action; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.producer.Callback; @@ -51,15 +52,20 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Properties; +import java.util.Set; import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import static org.apache.flink.agents.runtime.actionstate.ActionStateTestUtils.KEY_SERIALIZER; import static org.apache.flink.agents.runtime.actionstate.ActionStateTestUtils.createKeyEncoder; import static org.apache.flink.agents.runtime.actionstate.ActionStateTestUtils.generateKey; import static org.apache.kafka.clients.consumer.internals.AutoOffsetResetStrategy.EARLIEST; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.catchThrowable; +import static org.assertj.core.api.Assertions.entry; import static org.junit.jupiter.api.Assertions.*; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; @@ -80,6 +86,34 @@ public class KafkaActionStateStoreTest { private ActionState testActionState; private Map actionStates; + @Test + void testRejectsBlankCleanupControlTopicBeforeConnecting() { + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set(AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, " "); + + assertThatThrownBy( + () -> + new KafkaActionStateStore( + configuration, createKeyEncoder(MAX_PARALLELISM))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Kafka action-state cleanup control topic must not be blank"); + } + + @Test + void testRejectsTombstonesWithCleanupControlTopicBeforeConnecting() { + AgentConfiguration configuration = new AgentConfiguration(); + configuration.set( + AgentConfigOptions.KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC, "control-topic"); + configuration.set(AgentConfigOptions.KAFKA_ACTION_STATE_TOMBSTONE_ENABLED, true); + + assertThatThrownBy( + () -> + new KafkaActionStateStore( + configuration, createKeyEncoder(MAX_PARALLELISM))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("tombstones cannot be enabled"); + } + @BeforeEach void setUp() throws Exception { mockProducer = @@ -88,7 +122,31 @@ void setUp() throws Exception { new ActionStateKeyPartitioner(), new StringSerializer(), new ActionStateKafkaSeder()); - mockConsumer = new MockConsumer<>(EARLIEST.name()); + mockConsumer = + new MockConsumer(EARLIEST.name()) { + @Override + public synchronized void commitSync() { + throw new AssertionError( + "Action-state replay must not commit consumer-group offsets"); + } + }; + mockConsumer.updatePartitions( + TEST_TOPIC, + List.of( + new PartitionInfo(TEST_TOPIC, 0, null, null, null), + new PartitionInfo(TEST_TOPIC, 1, null, null, null))); + mockConsumer.updateBeginningOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 0L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 0L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); mockConsumer.assign( List.of(new TopicPartition(TEST_TOPIC, 0), new TopicPartition(TEST_TOPIC, 1))); actionStates = new HashMap<>(); @@ -201,36 +259,35 @@ void testGetCleansFutureStateForKeyContainingUnderscore() throws Exception { void testRecoveryMarker() throws Exception { // Test getting initial recovery marker Object initialMarker = actionStateStore.getRecoveryMarker(); - assertNotNull(initialMarker); - assertTrue(initialMarker instanceof Map); - assertTrue(((Map) initialMarker).isEmpty()); + assertThat(initialMarker).isInstanceOf(KafkaActionStateRecoveryMarker.class); + KafkaActionStateRecoveryMarker initialRecoveryMarker = + (KafkaActionStateRecoveryMarker) initialMarker; + assertThat(initialRecoveryMarker.getSchemaVersion()) + .isEqualTo(KafkaActionStateRecoveryMarker.CURRENT_SCHEMA_VERSION); + assertThat(initialRecoveryMarker.getTopic()).isEqualTo(TEST_TOPIC); + assertThat(initialRecoveryMarker.getTopicId()).isEqualTo("test-topic-id:" + TEST_TOPIC); + assertThat(initialRecoveryMarker.getOffsets()).containsOnly(entry(0, 0L), entry(1, 0L)); - mockConsumer.updatePartitions( - TEST_TOPIC, - List.of( - new PartitionInfo(TEST_TOPIC, 0, null, null, null), - new PartitionInfo(TEST_TOPIC, 1, null, null, null))); mockConsumer.updateEndOffsets( Map.of( new TopicPartition(TEST_TOPIC, 0), 5L, new TopicPartition(TEST_TOPIC, 1), 3L)); - for (int i = 0; i < 5; i++) { - mockConsumer.addRecord( - new ConsumerRecord<>( - TEST_TOPIC, - 0, - i++, - "key", - new ActionState(null, null, null, null, null, false))); - } // Test getting recovery marker after putting state Object secondMarker = actionStateStore.getRecoveryMarker(); - assertTrue(initialMarker instanceof Map); - assertFalse(((Map) secondMarker).isEmpty()); - assertThat((Map) secondMarker).containsEntry(0, 5L); - assertThat((Map) secondMarker).containsEntry(1, 3L); + assertThat(secondMarker).isInstanceOf(KafkaActionStateRecoveryMarker.class); + assertThat(((KafkaActionStateRecoveryMarker) secondMarker).getOffsets()) + .containsOnly(entry(0, 5L), entry(1, 3L)); + } + + @Test + void testRebuildConsumerDisablesOffsetResetAndAutoCommit() { + Properties consumerProperties = actionStateStore.createConsumerProp(); + + assertThat(consumerProperties) + .containsEntry(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none") + .containsEntry(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); } @Test @@ -455,6 +512,13 @@ void testRebuildState() throws Exception { mockConsumer.addRecord( new ConsumerRecord<>(record.topic(), 0, i++, record.key(), record.value())); } + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + i, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + actionStates.clear(); actionStateStore.rebuildState(recoveryMarkers); @@ -473,24 +537,351 @@ void testRebuildState() throws Exception { .isEqualTo(thirdState); } + @Test + void testRebuildStateFromVersionedMarker() throws Exception { + String stateKey = generateKey(TEST_KEY, 1L, testAction, testEvent, MAX_PARALLELISM); + mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 0L, stateKey, testActionState)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 1L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 0L, 1, 0L)); + + actionStateStore.rebuildState(List.of(marker)); + + assertThat(actionStates).containsEntry(stateKey, testActionState); + } + + @Test + void testRebuildStateContinuesAfterEmptyPollAndStopsAtCapturedEnd() throws Exception { + String includedKey = generateKey(TEST_KEY, 1L, testAction, testEvent, MAX_PARALLELISM); + String laterKey = generateKey(TEST_KEY, 2L, testAction, testEvent, MAX_PARALLELISM); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 1L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + mockConsumer.schedulePollTask(() -> {}); + mockConsumer.schedulePollTask( + () -> { + mockConsumer.addRecord( + new ConsumerRecord<>(TEST_TOPIC, 0, 0L, includedKey, testActionState)); + mockConsumer.addRecord( + new ConsumerRecord<>(TEST_TOPIC, 0, 1L, laterKey, testActionState)); + }); + + actionStateStore.rebuildState(List.of(Map.of(0, 0L, 1, 0L))); + + assertThat(actionStates) + .containsEntry(includedKey, testActionState) + .doesNotContainKey(laterKey); + } + + @Test + void testRebuildStateRejectsUnavailableEarlierOffset() { + mockConsumer.updateBeginningOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 5L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 10L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(Map.of(0, 4L, 1, 0L)))); + + assertThat(error) + .hasRootCauseMessage( + "Cannot rebuild Kafka action state for test-action-state-0: requested offset 4 is outside the available range [5, 10]"); + } + + @Test + void testRebuildStateRejectsOffsetBeyondEnd() { + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 10L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(Map.of(0, 11L, 1, 0L)))); + + assertThat(error) + .hasRootCauseMessage( + "Cannot rebuild Kafka action state for test-action-state-0: requested offset 11 is outside the available range [0, 10]"); + } + + @Test + void testRebuildStateRejectsRecreatedTopic() { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "old-topic-id", Map.of(0, 0L, 1, 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(marker))); + + assertThat(error) + .hasRootCauseMessage( + "Kafka action-state topic test-action-state has ID test-topic-id:test-action-state, but the recovery marker expects old-topic-id; the topic may have been recreated"); + } + + @Test + void testRebuildStateRejectsUnsupportedMarkerSchema() { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + 99, TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 0L, 1, 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(marker))); + + assertThat(error) + .hasRootCauseMessage( + "Unsupported Kafka action-state recovery marker schema 99, expected 1"); + } + + @Test + void testRebuildStateRejectsMixedLegacyAndVersionedMarkers() { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 0L, 1, 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(marker, Map.of(0, 0L, 1, 0L)))); + + assertThat(error) + .hasRootCauseMessage( + "Cannot restore from a mixture of versioned and legacy Kafka recovery markers"); + } + + @Test + void testRebuildStateRejectsChangedPartitionSet() { + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(Map.of(0, 0L)))); + + assertThat(error) + .hasRootCauseMessage( + "Kafka action-state recovery marker contains partitions [0], but topic test-action-state currently has partitions [0, 1]"); + } + + @Test + void testRecoveryMarkerRejectsTopicPartitionChange() { + mockConsumer.updatePartitions( + TEST_TOPIC, + List.of( + new PartitionInfo(TEST_TOPIC, 0, null, null, null), + new PartitionInfo(TEST_TOPIC, 1, null, null, null), + new PartitionInfo(TEST_TOPIC, 2, null, null, null))); + + RuntimeException error = + assertThrows(RuntimeException.class, actionStateStore::getRecoveryMarker); + + assertThat(error) + .hasRootCauseMessage( + "Kafka action-state topic test-action-state changed while the job was running; expected ID test-topic-id:test-action-state and partitions [0, 1] but found ID test-topic-id:test-action-state and partitions [0, 1, 2]"); + } + + @Test + void testRebuildStateRefreshesTopicMetadataBeforeReplay() { + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 0L, 1, 0L)); + mockConsumer.updatePartitions( + TEST_TOPIC, + List.of( + new PartitionInfo(TEST_TOPIC, 0, null, null, null), + new PartitionInfo(TEST_TOPIC, 1, null, null, null), + new PartitionInfo(TEST_TOPIC, 2, null, null, null))); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> actionStateStore.rebuildState(List.of(marker))); + + assertThat(error) + .hasRootCauseMessage( + "Kafka action-state topic test-action-state changed while the job was running; expected ID test-topic-id:test-action-state and partitions [0, 1] but found ID test-topic-id:test-action-state and partitions [0, 1, 2]"); + } + + @Test + void testRebuildStateRechecksTopicMetadataImmediatelyBeforeReadingOffsets() { + AtomicInteger metadataLoads = new AtomicInteger(); + MockConsumer changingConsumer = + new MockConsumer(EARLIEST.name()) { + @Override + public synchronized List partitionsFor(String topic) { + if (metadataLoads.incrementAndGet() == 3) { + updatePartitions( + topic, + List.of( + new PartitionInfo(topic, 0, null, null, null), + new PartitionInfo(topic, 1, null, null, null), + new PartitionInfo(topic, 2, null, null, null))); + } + return super.partitionsFor(topic); + } + }; + changingConsumer.updatePartitions( + TEST_TOPIC, + List.of( + new PartitionInfo(TEST_TOPIC, 0, null, null, null), + new PartitionInfo(TEST_TOPIC, 1, null, null, null))); + KafkaActionStateStore store = + new KafkaActionStateStore( + new HashMap<>(), + new AgentConfiguration(), + mockProducer, + changingConsumer, + TEST_TOPIC, + createKeyEncoder(MAX_PARALLELISM)); + KafkaActionStateRecoveryMarker marker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 0L, 1, 0L)); + + RuntimeException error = + assertThrows(RuntimeException.class, () -> store.rebuildState(List.of(marker))); + + assertThat(error) + .hasRootCauseMessage( + "Kafka action-state topic test-action-state changed while the job was running; expected ID test-topic-id:test-action-state and partitions [0, 1] but found ID test-topic-id:test-action-state and partitions [0, 1, 2]"); + } + + @Test + void testRebuildStateRejectsMarkerOlderThanCommittedCleanupBoundary() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of( + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, + "test-topic-id:" + TEST_TOPIC, + Map.of(0, 5L, 1, 0L)))); + KafkaActionStateStore cleanupAwareStore = cleanupAwareStore(committedCoordinator(plan)); + KafkaActionStateRecoveryMarker oldMarker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 4L, 1, 0L)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> cleanupAwareStore.rebuildState(List.of(oldMarker))); + + assertThat(error) + .hasRootCauseMessage( + "Cannot restore Kafka action state for test-action-state-0 from offset 4 because the committed cleanup boundary is 5"); + } + + @Test + void testRebuildStateRejectsLegacyMarkerWhenCleanupBoundaryIsConfigured() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of( + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, + "test-topic-id:" + TEST_TOPIC, + Map.of(0, 0L, 1, 0L)))); + KafkaActionStateStore cleanupAwareStore = cleanupAwareStore(committedCoordinator(plan)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> cleanupAwareStore.rebuildState(List.of(Map.of(0, 0L, 1, 0L)))); + + assertThat(error) + .hasRootCauseMessage( + "Checkpoint-aligned cleanup requires versioned Kafka recovery markers"); + } + + @Test + void testRebuildStateRejectsMissingMarkerAfterCleanupWasCommitted() { + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "checkpoint-42", + List.of( + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, + "test-topic-id:" + TEST_TOPIC, + Map.of(0, 5L, 1, 0L)))); + KafkaActionStateStore cleanupAwareStore = cleanupAwareStore(committedCoordinator(plan)); + + RuntimeException error = + assertThrows( + RuntimeException.class, + () -> cleanupAwareStore.rebuildState(Collections.emptyList())); + + assertThat(error) + .hasRootCauseMessage( + "Cannot initialize Kafka action state without a recovery marker because cleanup boundary {0=5, 1=0} is committed"); + } + @Test void testRebuildStateRemovesTombstonedKeys() throws Exception { // Arrange - two state records followed by a tombstone for the first key - List recoveryMarkers = List.of(Map.of(0, 0L)); + List recoveryMarkers = List.of(Map.of(0, 0L, 1, 0L)); String key1 = generateKey(TEST_KEY, 1L, testAction, testEvent, MAX_PARALLELISM); String key2 = generateKey(TEST_KEY, 2L, testAction, testEvent, MAX_PARALLELISM); mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 0L, key1, testActionState)); mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 1L, key2, testActionState)); mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 2L, key1, null)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 3L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); - // Act actionStateStore.rebuildState(recoveryMarkers); - - // Assert - the tombstoned key is removed, the other key is restored assertThat(actionStates).doesNotContainKey(key1); assertThat(actionStates.get(key2)).isEqualTo(testActionState); } + @Test + void testCleanupMigrationMarkerAfterExistingTombstonesIsUsable() throws Exception { + String stateKey = generateKey(TEST_KEY, 1L, testAction, testEvent, MAX_PARALLELISM); + mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 0L, stateKey, testActionState)); + mockConsumer.addRecord(new ConsumerRecord<>(TEST_TOPIC, 0, 1L, stateKey, null)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 2L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + KafkaActionStateRecoveryMarker postTombstoneMarker = + new KafkaActionStateRecoveryMarker( + TEST_TOPIC, "test-topic-id:" + TEST_TOPIC, Map.of(0, 2L, 1, 0L)); + KafkaActionStateCleanupPlan plan = + KafkaActionStateCleanupPlan.fromRecoveryMarkers( + "post-tombstone-checkpoint", List.of(postTombstoneMarker)); + KafkaActionStateStore cleanupAwareStore = cleanupAwareStore(committedCoordinator(plan)); + + cleanupAwareStore.rebuildState(List.of(postTombstoneMarker)); + + assertThat(actionStates).isEmpty(); + } + private static class TestAppender extends AbstractAppender { private final List messages = Collections.synchronizedList(new ArrayList<>()); @@ -508,7 +899,6 @@ private List getMessages() { return messages; } } - /** * After recovery, only the keys accepted by the ownership filter should enter the in-memory * cache. Here key "A" is owned and "B" is foreign, so "B" must be skipped while "A" is kept. @@ -525,6 +915,12 @@ void testRebuildStateFiltersForeignKeys() throws Exception { new ConsumerRecord<>(TEST_TOPIC, 0, offset++, stateKeyA, testActionState)); mockConsumer.addRecord( new ConsumerRecord<>(TEST_TOPIC, 0, offset++, stateKeyB, testActionState)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + offset, + new TopicPartition(TEST_TOPIC, 1), + 0L)); List recoveryMarkers = List.of(Map.of(0, 0L, 1, 0L)); @@ -553,6 +949,12 @@ void testRebuildStateKeepsAllKeysWhenNoFilter() throws Exception { new ConsumerRecord<>(TEST_TOPIC, 0, offset++, stateKeyA, testActionState)); mockConsumer.addRecord( new ConsumerRecord<>(TEST_TOPIC, 0, offset++, stateKeyB, testActionState)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + offset, + new TopicPartition(TEST_TOPIC, 1), + 0L)); List recoveryMarkers = List.of(Map.of(0, 0L, 1, 0L)); @@ -651,6 +1053,13 @@ void testRebuildStateRejectsUnrecognizedFormatKeys(boolean tombstone) { new ConsumerRecord<>( TEST_TOPIC, 0, 0L, legacyKey, tombstone ? null : testActionState)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 1L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + List recoveryMarkers = List.of(Map.of(0, 0L, 1, 0L)); actionStateStore.setOwnershipFilter(kg -> true); @@ -672,6 +1081,13 @@ void testRebuildStateRejectsKeyWithUnparsableKeyGroup() throws Exception { mockConsumer.addRecord( new ConsumerRecord<>(TEST_TOPIC, 0, 0L, unparseableGroupKey, testActionState)); + mockConsumer.updateEndOffsets( + Map.of( + new TopicPartition(TEST_TOPIC, 0), + 1L, + new TopicPartition(TEST_TOPIC, 1), + 0L)); + List recoveryMarkers = List.of(Map.of(0, 0L, 1, 0L)); actionStateStore.setOwnershipFilter(kg -> true); @@ -685,6 +1101,56 @@ void testRebuildStateRejectsKeyWithUnparsableKeyGroup() throws Exception { .hasMessageContaining("Invalid key-group"); } + private KafkaActionStateStore cleanupAwareStore( + KafkaActionStateCleanupCoordinator coordinator) { + return new KafkaActionStateStore( + actionStates, + new AgentConfiguration(), + mockProducer, + mockConsumer, + TEST_TOPIC, + createKeyEncoder(MAX_PARALLELISM), + coordinator); + } + + private static KafkaActionStateCleanupCoordinator committedCoordinator( + KafkaActionStateCleanupPlan plan) { + KafkaActionStateCleanupCoordinator.Operation operation = + KafkaActionStateCleanupCoordinator.Operation.committed(plan); + return new KafkaActionStateCleanupCoordinator( + new KafkaActionStateCleanupCoordinator.Transport() { + @Override + public KafkaActionStateCleanupCoordinator.TopicMetadata describeTopic( + String topic) { + return new KafkaActionStateCleanupCoordinator.TopicMetadata( + plan.getTopicId(), Set.copyOf(plan.getOffsets().keySet())); + } + + @Override + public Map + readOperations() { + return Map.of(plan.getPlanId(), operation); + } + + @Override + public void append(KafkaActionStateCleanupCoordinator.Operation value) {} + + @Override + public void validateDataTopicConfiguration(String topic) {} + + @Override + public void validateOffsetsAvailable( + String topic, String topicId, Map offsets) {} + + @Override + public void deleteBefore( + String topic, String topicId, Map offsets) {} + + @Override + public void close() {} + }); + } + /** Contract: the consumer is closed even when closing the producer throws. */ @Test @SuppressWarnings("unchecked")