diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 0e277038b7..30d5d3215e 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -09dfb247f6f59c61b7f2431932c4557fb26c22d8 +a70d7db187aa78a7fbfd3bfaf92372160cff320a \ No newline at end of file diff --git a/backend/migrations/20260309000000_kafka_offset_reset.down.sql b/backend/migrations/20260309000000_kafka_offset_reset.down.sql new file mode 100644 index 0000000000..d6df00484e --- /dev/null +++ b/backend/migrations/20260309000000_kafka_offset_reset.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE kafka_trigger DROP COLUMN auto_offset_reset; +ALTER TABLE kafka_trigger DROP COLUMN reset_offset; diff --git a/backend/migrations/20260309000000_kafka_offset_reset.up.sql b/backend/migrations/20260309000000_kafka_offset_reset.up.sql new file mode 100644 index 0000000000..7bcb23226d --- /dev/null +++ b/backend/migrations/20260309000000_kafka_offset_reset.up.sql @@ -0,0 +1,2 @@ +ALTER TABLE kafka_trigger ADD COLUMN auto_offset_reset VARCHAR(10) NOT NULL DEFAULT 'latest'; +ALTER TABLE kafka_trigger ADD COLUMN reset_offset BOOLEAN NOT NULL DEFAULT FALSE; diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 854514dfd8..d132f647db 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -12004,6 +12004,19 @@ paths: schema: type: string + /w/{workspace}/kafka_triggers/reset_offsets/{path}: + post: + summary: reset kafka trigger offsets to earliest + operationId: resetKafkaOffsets + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + responses: + "200": + description: kafka trigger offsets reset successfully + /w/{workspace}/nats_triggers/create: post: summary: create nats trigger @@ -22061,7 +22074,7 @@ components: description: Path to the Kafka resource containing connection configuration group_id: type: string - description: Kafka consumer group ID for this trigger + description: Kafka consumer group ID for this trigger topics: type: array items: @@ -22078,6 +22091,13 @@ components: required: - key - value + auto_offset_reset: + type: string + enum: + - latest + - earliest + default: latest + description: "Initial offset behavior when consumer group has no committed offset. 'latest' starts from new messages only, 'earliest' starts from the beginning." server_id: type: string description: ID of the server currently handling this trigger (internal) @@ -22138,6 +22158,13 @@ components: required: - key - value + auto_offset_reset: + type: string + enum: + - latest + - earliest + default: latest + description: "Initial offset behavior when consumer group has no committed offset." mode: $ref: "#/components/schemas/TriggerMode" error_handler_path: @@ -22190,6 +22217,13 @@ components: required: - key - value + auto_offset_reset: + type: string + enum: + - latest + - earliest + default: latest + description: "Initial offset behavior when consumer group has no committed offset." path: type: string description: The unique path identifier for this trigger diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte index 084652f431..ad5cce6f92 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte @@ -1,15 +1,21 @@ + (resetConfirmOpen = false)} +> + This will re-process all messages from the beginning of the topic. The consumer will restart + automatically. + + {#if mode === 'suspended'} + {#if edit && can_write} + + {/if} +
diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte index 1d3d62245c..9f1d34088f 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte @@ -3,6 +3,8 @@ import Section from '$lib/components/Section.svelte' import Subsection from '$lib/components/Subsection.svelte' import SchemaForm from '../../SchemaForm.svelte' + import Select from '$lib/components/select/Select.svelte' + import Label from '$lib/components/Label.svelte' import { workspaceStore } from '$lib/stores' import TestTriggerConnection from '../TestTriggerConnection.svelte' import TestingBadge from '../testingBadge.svelte' @@ -13,6 +15,7 @@ kafkaCfgValid?: boolean kafkaResourcePath?: string kafkaCfg?: Record + autoOffsetReset?: string can_write?: boolean showTestingBadge?: boolean } @@ -22,10 +25,16 @@ kafkaCfgValid = $bindable(false), kafkaResourcePath = $bindable(''), kafkaCfg = $bindable({}), + autoOffsetReset = $bindable('latest'), can_write = true, showTestingBadge = false }: Props = $props() + const offsetResetOptions = [ + { label: 'Latest (new messages only)', value: 'latest' }, + { label: 'Earliest (from beginning)', value: 'earliest' } + ] + const kafkaConfigSchema = { $schema: 'http://json-schema.org/draft-07/schema#', type: 'object', @@ -99,6 +108,21 @@ /> + +
+