Files
windmill/backend/windmill-worker/Cargo.toml
Ruben Fiszel 0dbd9c1231 perf: eliminate dual-connection DB pool contention across worker, queue, and api (#9798)
* perf: eliminate dual-connection DB pool contention across worker, queue, and api

Reuse the held transaction (or move pool reads before begin()) instead of
checking out a second pool connection while a tx is open, extending the
fix from #9789/#7861. Targets the per-worker pool (max 5) hot paths plus
several server-pool API handlers.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: pass owned pool to get_email_from_permissioned_as in http trigger handler

The generified signature takes impl PgExecutor; the http trigger handler
passed &db where db is already &DB, yielding &&Pool which does not impl
PgExecutor (only surfaced under the full feature set in CI).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: keep RLS-exposed reads on the non-RLS pool and isolate flow-eval reads in a savepoint

Addresses review of the dual-connection sweep:

- worker_flow: wrap the stop_after_all_iters_if reads in a SAVEPOINT. The
  caller swallows the error and keeps using tx, so a DB read failure must
  not leave the outer transaction aborted (it would fail the later commit).
  Matches the previous pool-read semantics.

- Revert reads that were moved onto an RLS (user_db) transaction back to the
  non-RLS pool, since RLS row-visibility/role context can change results:
  push_scheduled_job (email/tag/settings lookups; reachable with a user_db
  tx from api-schedule/api-flows), push_inner native-retry dedicated_worker
  routing (RLS isolation variants), resources.rs app-namespace folder
  auto-create (non-admins must not be blocked), and the script archive/delete
  UPDATEs. Non-RLS db.begin() reuse and move-before-begin are kept.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* test: failpoint proving the stop_after_all_iters_if savepoint isolates an aborted read

Adds a worker-crate failpoints feature and a data-driven hook: when the
stop_after_all_iters_if expr is the magic sentinel, the in-evaluation read runs
SELECT 1/0 to abort its (savepoint) transaction. The test asserts the flow still
completes (iteration marked failed) — which only holds if the savepoint keeps the
outer status-update transaction committable. Without the savepoint the abort would
poison the outer tx and the job would never complete.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-25 21:33:17 +00:00

159 lines
6.1 KiB
TOML

[package]
name = "windmill-worker"
version.workspace = true
authors.workspace = true
edition.workspace = true
[lib]
name = "windmill_worker"
path = "src/lib.rs"
[features]
default = []
private = ["windmill-worker-volumes/private", "windmill-queue/private", "windmill-common/private", "windmill-dep-map/private", "windmill-runtime-nativets?/private"]
mcp = ["windmill-ai/mcp", "dep:windmill-mcp"]
prometheus = ["dep:prometheus", "windmill-common/prometheus"]
enterprise = ["windmill-queue/enterprise", "windmill-git-sync/enterprise", "windmill-common/enterprise", "windmill-worker-volumes/enterprise", "windmill-runtime-nativets?/enterprise", "dep:pem", "dep:rsa", "dep:tokio-util", "dep:opentelemetry-proto", "dep:prost", "dep:hudsucker", "dep:rcgen", "dep:hyper-http-proxy", "dep:hyper-tls", "dep:hyper-util"]
mssql = ["dep:tiberius"]
mssql-kerberos = ["mssql", "tiberius/integrated-auth-gssapi"] # Linux/Unix integrated auth
mssql-winauth = ["mssql", "tiberius/winauth"] # Windows integrated auth
bigquery = ["dep:gcp_auth"]
benchmark = ["windmill-queue/benchmark", "windmill-common/benchmark"]
parquet = ["windmill-common/parquet", "windmill-object-store/parquet"]
flow_testing = []
failpoints = []
cloud = []
sqlx = []
deno_core = ["dep:windmill-runtime-nativets"]
libffi_mac = ["dep:libffi-sys"]
otel = ["windmill-common/otel", "dep:opentelemetry", "dep:tracing-opentelemetry"]
dind = ["dep:bollard"]
php = ["dep:windmill-parser-php"]
mysql = ["dep:mysql_async"]
oracledb = ["dep:oracle"]
python = ["dep:windmill-parser-py", "dep:windmill-parser-py-imports", "windmill-dep-map/python"]
csharp = ["dep:windmill-parser-csharp"]
rust = ["dep:windmill-parser-rust"]
nu = ["dep:windmill-parser-nu"]
java = ["dep:windmill-parser-java"]
ruby = ["dep:windmill-parser-ruby"]
rlang = ["dep:windmill-parser-r"]
duckdb = ["dep:libloading"]
quickjs = ["windmill-jseval/quickjs", "windmill-queue/quickjs"]
bedrock = ["windmill-ai/bedrock"]
[dependencies]
windmill-ai = { workspace = true, default-features = false }
windmill-queue.workspace = true
windmill-dep-map.workspace = true
windmill-audit.workspace = true # there isn't really a reason for audit-worth actions to happen in the worker.
windmill-common = { workspace = true, default-features = false }
windmill-types.workspace = true
windmill-object-store.workspace = true
windmill-worker-volumes.workspace = true
windmill-jseval.workspace = true
windmill-runtime-nativets = { workspace = true, optional = true }
windmill-mcp = { workspace = true, optional = true }
windmill-macros.workspace = true
windmill-parser.workspace = true
windmill-parser-ts.workspace = true
windmill-parser-go.workspace = true
windmill-parser-rust = { workspace = true, optional = true }
windmill-parser-csharp = { workspace = true, optional = true }
windmill-parser-nu = { workspace = true, optional = true }
windmill-parser-java = { workspace = true, optional = true }
windmill-parser-ruby = { workspace = true, optional = true }
windmill-parser-r = { workspace = true, optional = true }
windmill-parser-py = { workspace = true, optional = true }
windmill-parser-yaml.workspace = true
windmill-parser-py-imports = { workspace = true, optional = true }
windmill-parser-bash.workspace = true
windmill-parser-sql.workspace = true
windmill-parser-graphql.workspace = true
windmill-parser-php = { workspace = true, optional = true }
windmill-git-sync.workspace = true
flume.workspace = true
sqlx.workspace = true
uuid.workspace = true
ulid.workspace = true
tracing.workspace = true
tokio.workspace = true
tokio-stream.workspace = true
serde.workspace = true
serde_json.workspace = true
futures.workspace = true
async-recursion.workspace = true
async-trait.workspace = true
anyhow.workspace = true
derive_more.workspace = true
itertools.workspace = true
regex.workspace = true
prometheus = { workspace = true, optional = true }
lazy_static.workspace = true
quick_cache.workspace = true
chrono.workspace = true
dotenv.workspace = true
rand.workspace = true # TODO: Remove. only used by token creation hack.
const_format.workspace = true
mappable-rc.workspace = true
git-version.workspace = true
once_cell.workspace = true
tokio-postgres.workspace = true
bit-vec.workspace = true
url.workspace = true
async-stream.workspace = true
postgres-native-tls.workspace = true
native-tls.workspace = true
mysql_async = { workspace = true, optional = true }
base64.workspace = true
gcp_auth = { workspace = true, optional = true }
rust_decimal.workspace = true
jsonwebtoken.workspace = true
sha2.workspace = true
hmac.workspace = true
pem = { workspace = true, optional = true }
rsa = { workspace = true, optional = true }
urlencoding.workspace = true
# `fs` adds flock(2) for the cross-process Python install lock (shared cache mounts)
nix = { workspace = true, features = ["fs"] }
bytes.workspace = true
reqwest.workspace = true
reqwest-middleware.workspace = true
eventsource-stream.workspace = true
mime_guess.workspace = true
hex.workspace = true
tiberius = { workspace = true, optional = true }
tokio-util = { workspace = true, optional = true }
tar.workspace = true
convert_case.workspace = true
yaml-rust.workspace = true
backon.workspace = true
pep440_rs.workspace = true
process-wrap.workspace = true
async-once-cell.workspace = true
libloading = { workspace = true, optional = true }
opentelemetry-proto = { workspace = true, optional = true }
opentelemetry = { workspace = true, optional = true }
tracing-opentelemetry = { workspace = true, optional = true }
prost = { workspace = true, optional = true }
axum.workspace = true
bollard = { workspace = true, optional = true }
oracle = { workspace = true, optional = true }
hudsucker = { workspace = true, optional = true }
hyper-http-proxy = { workspace = true, optional = true }
hyper-tls = { workspace = true, optional = true }
hyper-util = { workspace = true, optional = true }
rcgen = { workspace = true, optional = true }
[target.'cfg(windows)'.dependencies]
windows = { version = "0.61", features = ["Win32_System_JobObjects", "Win32_System_Threading", "Win32_System_Diagnostics_ToolHelp"] }
[dev-dependencies]
tempfile.workspace = true
x509-parser.workspace = true
[build-dependencies]
libffi-sys = { workspace = true, optional = true }