* feat: add auto_commit option to Kafka triggers with manual commit API
Add ability to disable auto-commit on Kafka triggers so users can
manually commit offsets after processing messages. This prevents
message loss when processing fails.
Changes:
- Add `auto_commit` column to kafka_trigger table (default true)
- Add POST /kafka_triggers/commit_offsets/{path} endpoint using
BaseConsumer with manual assign() to avoid rebalance
- Enrich trigger_info payload with partition and offset fields
- Conditionally commit based on auto_commit setting
- Add auto-commit toggle to frontend Kafka trigger config
- Add commitKafkaOffsets helpers to Python and TypeScript SDKs
- Add integration tests for auto_commit DB defaults
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* feat: use DB-based pending commits for kafka manual offset commit
Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
* feat: pass trigger_path to all v2 preprocessors, secure commit_offsets endpoint, fix commit semantics
- Add trigger_path to v2 preprocessor event for all trigger types (kafka, nats, sqs, mqtt, gcp, postgres, websocket, http, email)
- Secure commit_offsets endpoint: infer trigger from job token (OptJobAuthed) instead of requiring trigger path parameter
- Fix auto_commit: only commit offset after successful job push
- Fix pending commits: commit offset+1 (Kafka semantics) and use CommitMode::Sync
- Update TS/Python clients and frontend preprocessor templates
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* feat: add advanced section badges and reorganize kafka trigger settings
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: remove dead wm_trigger assertions from kafka e2e test
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* sqlx
* refactor: remove unused advancedCollapsed state from all trigger editors
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* update ref
* chore: update ee-repo-ref to ed2c9d360e6fab866b9744cc79f50038d1fc7152
This commit updates the EE repository reference after PR #452 was merged in windmill-ee-private.
Previous ee-repo-ref: 5b31116a1d5a042c6a780732901cfd89584d1773
New ee-repo-ref: ed2c9d360e6fab866b9744cc79f50038d1fc7152
Automated by sync-ee-ref workflow.
* fix: use path-based auth for kafka commit_offsets endpoint
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* chore: update ee-repo-ref to fcd3ea52b0cc94fbe1159baf662a38da947456de
This commit updates the EE repository reference after PR #457 was merged in windmill-ee-private.
Previous ee-repo-ref: b3a5c33c92cb1b2caf7a65986d71da291ff72a35
New ee-repo-ref: fcd3ea52b0cc94fbe1159baf662a38da947456de
Automated by sync-ee-ref workflow.
---------
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
219 lines
5.7 KiB
Bash
Executable File
219 lines
5.7 KiB
Bash
Executable File
#!/usr/bin/env bash
|
|
set -eou pipefail
|
|
script_dirpath="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
|
|
|
rm -rf "${script_dirpath}/src"
|
|
|
|
npx --yes @hey-api/openapi-ts@0.43.0 --input "${script_dirpath}/../backend/windmill-api/openapi.yaml" --output "${script_dirpath}/src" --useOptions --schemas false
|
|
cat <<EOF - src/core/OpenAPI.ts > temp_file && mv temp_file src/core/OpenAPI.ts
|
|
const getEnv = (key: string) => {
|
|
if (typeof window === "undefined") {
|
|
if (typeof process !== "undefined") {
|
|
return process?.env?.[key];
|
|
}
|
|
// node
|
|
return globalThis?.process?.env?.[key];
|
|
}
|
|
// browser
|
|
return window?.process?.env?.[key];
|
|
};
|
|
|
|
const baseUrl = getEnv("BASE_INTERNAL_URL") ?? getEnv("BASE_URL") ?? "http://localhost:8000";
|
|
const baseUrlApi = (baseUrl ?? '') + "/api";
|
|
|
|
EOF
|
|
if [[ "$OSTYPE" == "darwin"* ]]; then
|
|
sed -i '' 's/WITH_CREDENTIALS: false/WITH_CREDENTIALS: true/g' src/core/OpenAPI.ts
|
|
sed -i '' 's/TOKEN: undefined/TOKEN: getEnv("WM_TOKEN")/g' src/core/OpenAPI.ts
|
|
sed -i '' "s/BASE: '\/api'/BASE: baseUrlApi/g" src/core/OpenAPI.ts
|
|
else
|
|
sed -i 's/WITH_CREDENTIALS: false/WITH_CREDENTIALS: true/g' src/core/OpenAPI.ts
|
|
sed -i 's/TOKEN: undefined/TOKEN: getEnv("WM_TOKEN")/g' src/core/OpenAPI.ts
|
|
sed -i "s/BASE: '\/api'/BASE: baseUrlApi/g" src/core/OpenAPI.ts
|
|
fi
|
|
|
|
|
|
|
|
cp "${script_dirpath}/client.ts" "${script_dirpath}/src/"
|
|
cp "${script_dirpath}/s3Types.ts" "${script_dirpath}/src/"
|
|
cp "${script_dirpath}/sqlUtils.ts" "${script_dirpath}/src/"
|
|
echo "" >> "${script_dirpath}/src/index.ts"
|
|
echo 'export type { DenoS3LightClientSettings } from "./s3Types";' >> "${script_dirpath}/src/index.ts"
|
|
echo "" >> "${script_dirpath}/src/index.ts"
|
|
echo 'export { type Base64, setClient, getVariable, setVariable, getResource, setResource, getResumeUrls, setState, setProgress, getProgress, getState, getIdToken, denoS3LightClientSettings, loadS3FileStream, loadS3File, writeS3File, signS3Objects, signS3Object, getPresignedS3PublicUrls, getPresignedS3PublicUrl, task, taskScript, taskFlow, workflow, step, sleep, parallel, waitForApproval, type TaskOptions, WorkflowCtx, _workflowCtx, setWorkflowCtx, StepSuspend, runScript, runScriptAsync, runScriptByPath, runScriptByHash, runScriptByPathAsync, runScriptByHashAsync, runFlow, runFlowAsync, waitJob, getRootJobId, setFlowUserState, getFlowUserState, usernameToEmail, requestInteractiveSlackApproval, type Sql, requestInteractiveTeamsApproval, appendToResultStream, streamResult, datatable, ducklake, type DatatableSqlTemplateFunction, type SqlTemplateFunction, type S3Object, type S3ObjectRecord, type S3ObjectURI, commitKafkaOffsets } from "./client";' >> "${script_dirpath}/src/index.ts"
|
|
|
|
# Build default export by combining client utilities + services
|
|
# This preserves backward compatibility for `import wmill from "windmill-client"`
|
|
# while enabling tree-shaking for named imports
|
|
cat >> "${script_dirpath}/src/index.ts" << 'INDEXEOF'
|
|
|
|
import {
|
|
setClient,
|
|
getVariable,
|
|
setVariable,
|
|
getResource,
|
|
setResource,
|
|
getResumeUrls,
|
|
setState,
|
|
setProgress,
|
|
getProgress,
|
|
getState,
|
|
getIdToken,
|
|
denoS3LightClientSettings,
|
|
loadS3FileStream,
|
|
loadS3File,
|
|
writeS3File,
|
|
signS3Objects,
|
|
signS3Object,
|
|
getPresignedS3PublicUrls,
|
|
getPresignedS3PublicUrl,
|
|
task,
|
|
taskScript,
|
|
taskFlow,
|
|
workflow,
|
|
step,
|
|
sleep,
|
|
parallel,
|
|
waitForApproval,
|
|
WorkflowCtx,
|
|
_workflowCtx,
|
|
setWorkflowCtx,
|
|
StepSuspend,
|
|
runScript,
|
|
runScriptAsync,
|
|
runScriptByPath,
|
|
runScriptByHash,
|
|
runScriptByPathAsync,
|
|
runScriptByHashAsync,
|
|
runFlow,
|
|
runFlowAsync,
|
|
waitJob,
|
|
getRootJobId,
|
|
setFlowUserState,
|
|
getFlowUserState,
|
|
usernameToEmail,
|
|
requestInteractiveSlackApproval,
|
|
requestInteractiveTeamsApproval,
|
|
appendToResultStream,
|
|
streamResult,
|
|
datatable,
|
|
ducklake,
|
|
SHARED_FOLDER,
|
|
getWorkspace,
|
|
getStatePath,
|
|
getInternalState,
|
|
setInternalState,
|
|
getResumeEndpoints,
|
|
getResult,
|
|
getResultMaybe,
|
|
resolveDefaultResource,
|
|
databaseUrlFromResource,
|
|
base64ToUint8Array,
|
|
uint8ArrayToBase64,
|
|
parseS3Object,
|
|
commitKafkaOffsets,
|
|
} from "./client";
|
|
|
|
import {
|
|
AdminService,
|
|
AuditService,
|
|
FlowService,
|
|
GranularAclService,
|
|
GroupService,
|
|
JobService,
|
|
ResourceService,
|
|
VariableService,
|
|
ScriptService,
|
|
ScheduleService,
|
|
SettingsService,
|
|
UserService,
|
|
WorkspaceService,
|
|
TeamsService,
|
|
} from "./services.gen";
|
|
|
|
const wmill = {
|
|
// Client utilities
|
|
setClient,
|
|
getVariable,
|
|
setVariable,
|
|
getResource,
|
|
setResource,
|
|
getResumeUrls,
|
|
setState,
|
|
setProgress,
|
|
getProgress,
|
|
getState,
|
|
getIdToken,
|
|
denoS3LightClientSettings,
|
|
loadS3FileStream,
|
|
loadS3File,
|
|
writeS3File,
|
|
signS3Objects,
|
|
signS3Object,
|
|
getPresignedS3PublicUrls,
|
|
getPresignedS3PublicUrl,
|
|
task,
|
|
taskScript,
|
|
taskFlow,
|
|
workflow,
|
|
step,
|
|
sleep,
|
|
parallel,
|
|
waitForApproval,
|
|
WorkflowCtx,
|
|
_workflowCtx,
|
|
setWorkflowCtx,
|
|
StepSuspend,
|
|
runScript,
|
|
runScriptAsync,
|
|
runScriptByPath,
|
|
runScriptByHash,
|
|
runScriptByPathAsync,
|
|
runScriptByHashAsync,
|
|
runFlow,
|
|
runFlowAsync,
|
|
waitJob,
|
|
getRootJobId,
|
|
setFlowUserState,
|
|
getFlowUserState,
|
|
usernameToEmail,
|
|
requestInteractiveSlackApproval,
|
|
requestInteractiveTeamsApproval,
|
|
appendToResultStream,
|
|
streamResult,
|
|
datatable,
|
|
ducklake,
|
|
SHARED_FOLDER,
|
|
getWorkspace,
|
|
getStatePath,
|
|
getInternalState,
|
|
setInternalState,
|
|
getResumeEndpoints,
|
|
getResult,
|
|
getResultMaybe,
|
|
resolveDefaultResource,
|
|
databaseUrlFromResource,
|
|
base64ToUint8Array,
|
|
uint8ArrayToBase64,
|
|
parseS3Object,
|
|
commitKafkaOffsets,
|
|
// Services
|
|
AdminService,
|
|
AuditService,
|
|
FlowService,
|
|
GranularAclService,
|
|
GroupService,
|
|
JobService,
|
|
ResourceService,
|
|
VariableService,
|
|
ScriptService,
|
|
ScheduleService,
|
|
SettingsService,
|
|
UserService,
|
|
WorkspaceService,
|
|
TeamsService,
|
|
};
|
|
|
|
export default wmill;
|
|
INDEXEOF
|