Skip to content

fix(inkless:retention): bound retention enforcement boundary scan to O(cap) - #722

Merged
viktorsomogyi merged 1 commit into
mainfrom
jeqo/retention-enforcer-bound-exec
Jul 28, 2026
Merged

fix(inkless:retention): bound retention enforcement boundary scan to O(cap)#722
viktorsomogyi merged 1 commit into
mainfrom
jeqo/retention-enforcer-bound-exec

Conversation

@jeqo

@jeqo jeqo commented Jul 27, 2026

Copy link
Copy Markdown
Contributor

The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every deletable batch, O(partition depth). PR #705 moved it out of the log lock so it no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain (retention lowered so almost the whole partition is deletable) a single scan on a multi-million-batch partition runs for seconds-to-minutes and exceeds the enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the still-running backend so the partition never drains.

V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The boundary can never be deeper than the cap because the delete is clamped there regardless; the +1 distinguishes an in-window boundary (delete up to it) from a deeper one that falls through to high_watermark and is re-clamped by the existing under-lock cap-probe. Each pass is O(cap), index-only on batches_by_last_offset_covering_idx, and a deep partition drains over successive enforcement cycles. max_batches_per_request = 0 keeps the original unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes oldest-first, so a snapshot-stale boundary can only under-delete).

Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the bound is on out of the box; it must stay above the per-interval expiry rate or retention falls behind.

Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via -Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms (flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is the sweet spot: short passes at 10x the drain throughput of 100.

Diff from previous version without comments:

                     WHERE topic_id = l_request.topic_id
                         AND partition = l_request.partition
                 ),
+                limited_batches AS (
+                    SELECT b.topic_id, b.partition, b.last_offset, b.base_offset, b.byte_size,
+                        batch_timestamp(b.timestamp_type, b.batch_max_timestamp, b.log_append_timestamp) AS effective_timestamp
+                    FROM batches b
+                    WHERE b.topic_id = l_request.topic_id
+                        AND b.partition = l_request.partition
+                    ORDER BY b.topic_id, b.partition, b.last_offset
+                    LIMIT (CASE WHEN max_batches_per_request > 0 THEN max_batches_per_request::bigint + 1 ELSE NULL END)
+                ),
                 augmented_batches AS (
-                    SELECT b.topic_id, b.partition, b.last_offset, b.base_offset,
+                    SELECT topic_id, partition, last_offset, base_offset,
                         (SELECT byte_size FROM selected_log)
-                            - SUM(b.byte_size) OVER (ORDER BY b.topic_id, b.partition, b.last_offset)
-                            + b.byte_size
+                            - SUM(byte_size) OVER (ORDER BY topic_id, partition, last_offset)
+                            + byte_size
                             AS reverse_agg_byte_size,
-                        batch_timestamp(b.timestamp_type, b.batch_max_timestamp, b.log_append_timestamp) AS effective_timestamp
-                    FROM batches b
-                    WHERE topic_id = l_request.topic_id
-                        AND partition = l_request.partition
-                    ORDER BY topic_id, partition, last_offset
+                        effective_timestamp
+                   FROM limited_batches

KC-354

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Pull request overview

This PR improves Inkless retention enforcement scalability by bounding the retention “boundary scan” to O(max_batches_per_request) rather than O(partition depth), preventing long-running scans from exceeding the control plane JDBC socket timeout on very deep partitions. It also flips the default retention.enforcement.max.batches.per.request from 0 (unbounded) to 1000 to enable the bounded behavior by default.

Changes:

  • Update enforce_retention_v2 to scan only the oldest (max_batches_per_request + 1) batches (or unbounded when the cap is 0), keeping per-pass work bounded.
  • Change the default retention enforcement cap to 1000 and update docs/tests accordingly.
  • Refresh benchmark guidance and regenerate jOOQ artifacts for schema version 22.

Reviewed changes

Copilot reviewed 6 out of 124 changed files in this pull request and generated no comments.

Show a summary per file
File Description
storage/inkless/src/test/java/io/aiven/inkless/control_plane/postgres/EnforceRetentionCommitBlockingBenchmarkTest.java Updates benchmark documentation and allows overriding the per-pass cap via a system property.
storage/inkless/src/test/java/io/aiven/inkless/control_plane/AbstractControlPlaneTest.java Adds a test to lock in correctness when the retention boundary falls within the bounded scan window.
storage/inkless/src/test/java/io/aiven/inkless/config/InklessConfigTest.java Updates expected default for maxBatchesPerEnforcementRequest() to 1000.
storage/inkless/src/main/resources/db/migration/V22__Retention_enforcement_bounded_scan.sql Implements the bounded boundary scan behavior for enforce_retention_v2.
storage/inkless/src/main/jooq/org/jooq/generated/UDTs.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/RepairDisklessLogResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/RepairDisklessLogRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/RepairDisklessLogResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/RepairDisklessLogRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/PruneBatchesBelowHighestTieredOffsetResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/PruneBatchesBelowHighestTieredOffsetRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogProducerStateV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchMetadataV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchInfoV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/AdvanceCrossTierLogStartResponseV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/AdvanceCrossTierLogStartRequestV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/PruneBatchesBelowHighestTieredOffsetResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/PruneBatchesBelowHighestTieredOffsetRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/RepairDisklessLogResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/RepairDisklessLogRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/PruneBatchesBelowHighestTieredOffsetResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/PruneBatchesBelowHighestTieredOffsetRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogProducerStateV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchMetadataV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchInfoV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/AdvanceCrossTierLogStartResponseV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/AdvanceCrossTierLogStartRequestV1Path.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogProducerStateV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchMetadataV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchInfoV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/AdvanceCrossTierLogStartResponseV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/udt/AdvanceCrossTierLogStartRequestV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/RepairDisklessLogV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/RepairDisklessLogV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/PruneBatchesBelowHighestTieredOffsetV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ProducerStateRecord.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/LogsRecord.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ListOffsetsV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/InitDisklessLogV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV2Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FilesRecord.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV2Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/DeleteRecordsV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/CommitFileV2Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/CommitFileV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/BatchesRecord.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/AdvanceCrossTierLogStartV1Record.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/PruneBatchesBelowHighestTieredOffsetV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/ProducerState.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Logs.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/ListOffsetsV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/InitDisklessLogV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV2.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Files.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV2.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/DeleteRecordsV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/CommitFileV2.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/CommitFileV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Batches.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/tables/AdvanceCrossTierLogStartV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/Tables.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/MarkFileToDeleteV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/FlushCommitRunV2.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteTopicV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteFilesV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteBatchV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/routines/BatchTimestamp.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/Routines.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/Keys.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/Indexes.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/PruneBatchesBelowHighestTieredOffsetErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/ListOffsetsResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/InitDisklessLogErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FindBatchesResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FileStateT.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FileReasonT.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/EnforceRetentionResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/DeleteRecordsResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/CommitBatchResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/enums/AdvanceCrossTierLogStartResponseErrorV1.java Regenerated jOOQ artifact for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/Domains.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/jooq/org/jooq/generated/DefaultSchema.java Regenerated jOOQ metadata for schema version 22.
storage/inkless/src/main/java/io/aiven/inkless/config/InklessConfig.java Updates config documentation and changes default retention enforcement cap to 1000.
docs/inkless/configs.rst Updates published configuration docs to reflect the new default and behavior.

…O(cap)

The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every
deletable batch, O(partition depth). PR #705 moved it out of the log lock so it
no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain
(retention lowered so almost the whole partition is deletable) a single scan on
a multi-million-batch partition runs for seconds-to-minutes and exceeds the
enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the
still-running backend so the partition never drains.

V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The
boundary can never be deeper than the cap because the delete is clamped there
regardless; the +1 distinguishes an in-window boundary (delete up to it) from a
deeper one that falls through to high_watermark and is re-clamped by the
existing under-lock cap-probe. Each pass is O(cap), index-only on
batches_by_last_offset_covering_idx, and a deep partition drains over
successive enforcement cycles. max_batches_per_request = 0 keeps the original
unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes
oldest-first, so a snapshot-stale boundary can only under-delete).

Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the
bound is on out of the box; it must stay above the per-interval expiry rate or
retention falls behind.

Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via
-Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k
batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms
(flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is
the sweet spot: short passes at 10x the drain throughput of 100.

[KC-354]
@jeqo
jeqo force-pushed the jeqo/retention-enforcer-bound-exec branch from 9fcab41 to 508d7fa Compare July 28, 2026 08:39
@jeqo
jeqo marked this pull request as ready for review July 28, 2026 08:42
@viktorsomogyi
viktorsomogyi self-requested a review July 28, 2026 08:58
@viktorsomogyi
viktorsomogyi merged commit 6ed0732 into main Jul 28, 2026
7 checks passed
@viktorsomogyi
viktorsomogyi deleted the jeqo/retention-enforcer-bound-exec branch July 28, 2026 10:22
giuseppelillo pushed a commit that referenced this pull request Jul 29, 2026
…O(cap) (#722)

The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every
deletable batch, O(partition depth). PR #705 moved it out of the log lock so it
no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain
(retention lowered so almost the whole partition is deletable) a single scan on
a multi-million-batch partition runs for seconds-to-minutes and exceeds the
enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the
still-running backend so the partition never drains.

V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The
boundary can never be deeper than the cap because the delete is clamped there
regardless; the +1 distinguishes an in-window boundary (delete up to it) from a
deeper one that falls through to high_watermark and is re-clamped by the
existing under-lock cap-probe. Each pass is O(cap), index-only on
batches_by_last_offset_covering_idx, and a deep partition drains over
successive enforcement cycles. max_batches_per_request = 0 keeps the original
unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes
oldest-first, so a snapshot-stale boundary can only under-delete).

Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the
bound is on out of the box; it must stay above the per-interval expiry rate or
retention falls behind.

Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via
-Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k
batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms
(flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is
the sweet spot: short passes at 10x the drain throughput of 100.

[KC-354]
giuseppelillo pushed a commit that referenced this pull request Jul 30, 2026
…O(cap) (#722)

The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every
deletable batch, O(partition depth). PR #705 moved it out of the log lock so it
no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain
(retention lowered so almost the whole partition is deletable) a single scan on
a multi-million-batch partition runs for seconds-to-minutes and exceeds the
enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the
still-running backend so the partition never drains.

V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The
boundary can never be deeper than the cap because the delete is clamped there
regardless; the +1 distinguishes an in-window boundary (delete up to it) from a
deeper one that falls through to high_watermark and is re-clamped by the
existing under-lock cap-probe. Each pass is O(cap), index-only on
batches_by_last_offset_covering_idx, and a deep partition drains over
successive enforcement cycles. max_batches_per_request = 0 keeps the original
unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes
oldest-first, so a snapshot-stale boundary can only under-delete).

Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the
bound is on out of the box; it must stay above the per-interval expiry rate or
retention falls behind.

Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via
-Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k
batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms
(flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is
the sweet spot: short passes at 10x the drain throughput of 100.

[KC-354]
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants