diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 83858bb6cc..fb6fca772c 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -125,7 +125,7 @@ dependencies = [ "getrandom 0.3.4", "once_cell", "version_check", - "zerocopy", + "zerocopy 0.8.33", ] [[package]] @@ -1302,6 +1302,34 @@ dependencies = [ "tracing", ] +[[package]] +name = "axum" +version = "0.6.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3b829e4e32b91e643de6eafe82b1d90675f5874230191a4ffbc1b336dec4d6bf" +dependencies = [ + "async-trait", + "axum-core 0.3.4", + "bitflags 1.3.2", + "bytes", + "futures-util", + "http 0.2.12", + "http-body 0.4.6", + "hyper 0.14.32", + "itoa", + "matchit 0.7.3", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "rustversion", + "serde", + "sync_wrapper 0.1.2", + "tower 0.4.13", + "tower-layer", + "tower-service", +] + [[package]] name = "axum" version = "0.7.9" @@ -1330,7 +1358,7 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", - "sync_wrapper", + "sync_wrapper 1.0.2", "tokio", "tower 0.5.3", "tower-layer", @@ -1364,7 +1392,7 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", - "sync_wrapper", + "sync_wrapper 1.0.2", "tokio", "tower 0.5.3", "tower-layer", @@ -1372,6 +1400,23 @@ dependencies = [ "tracing", ] +[[package]] +name = "axum-core" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "759fa577a247914fd3f7f76d62972792636412fbfd634cd452f6a385a74d2d2c" +dependencies = [ + "async-trait", + "bytes", + "futures-util", + "http 0.2.12", + "http-body 0.4.6", + "mime", + "rustversion", + "tower-layer", + "tower-service", +] + [[package]] name = "axum-core" version = "0.4.5" @@ -1387,7 +1432,7 @@ dependencies = [ "mime", "pin-project-lite", "rustversion", - "sync_wrapper", + "sync_wrapper 1.0.2", "tower-layer", "tower-service", "tracing", @@ -1406,7 +1451,7 @@ dependencies = [ "http-body-util", "mime", "pin-project-lite", - "sync_wrapper", + "sync_wrapper 1.0.2", "tower-layer", "tower-service", "tracing", @@ -1541,6 +1586,29 @@ dependencies = [ "serde", ] +[[package]] +name = "bindgen" +version = "0.66.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2b84e06fc203107bfbad243f4aba2af864eb7db3b1cf46ea0a023b0b433d2a7" +dependencies = [ + "bitflags 2.9.4", + "cexpr", + "clang-sys", + "lazy_static", + "lazycell", + "log", + "peeking_take_while", + "prettyplease", + "proc-macro2", + "quote", + "regex", + "rustc-hash 1.1.0", + "shlex", + "syn 2.0.114", + "which 4.4.2", +] + [[package]] name = "bindgen" version = "0.69.5" @@ -3889,7 +3957,7 @@ dependencies = [ "tokio-socks", "tokio-util", "tower 0.5.3", - "tower-http", + "tower-http 0.6.8", "tower-service", ] @@ -4031,7 +4099,7 @@ dependencies = [ "http-body-util", "log", "num-bigint", - "prost", + "prost 0.13.5", "prost-build", "rand 0.8.5", "rusqlite", @@ -4654,7 +4722,7 @@ dependencies = [ "deno_error", "futures", "num-bigint", - "prost", + "prost 0.13.5", "serde", "uuid", ] @@ -4674,7 +4742,7 @@ dependencies = [ "futures", "http 1.4.0", "log", - "prost", + "prost 0.13.5", "rand 0.8.5", "serde", "serde_json", @@ -6340,7 +6408,7 @@ dependencies = [ "thiserror 1.0.69", "tokio", "tokio-retry2", - "tonic", + "tonic 0.12.3", "tower 0.4.13", "tracing", ] @@ -6351,9 +6419,9 @@ version = "0.16.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "886aa8ec755382a1fdf4651f6e6ec01f2f3bf49f2cb0f068b9a74cafd574a715" dependencies = [ - "prost", + "prost 0.13.5", "prost-types", - "tonic", + "tonic 0.12.3", ] [[package]] @@ -6516,7 +6584,7 @@ dependencies = [ "num-traits", "rand 0.9.0", "rand_distr 0.5.1", - "zerocopy", + "zerocopy 0.8.33", ] [[package]] @@ -6576,6 +6644,15 @@ dependencies = [ "syn 2.0.114", ] +[[package]] +name = "hashlink" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8094feaf31ff591f651a2664fb9cfd92bba7a60ce3197265e9482ebe753c8f7" +dependencies = [ + "hashbrown 0.14.5", +] + [[package]] name = "hashlink" version = "0.9.1" @@ -6858,6 +6935,12 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "http-range-header" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "add0ab9360ddbd88cfeb3bd9574a1d85cfdfa14db10b3e21d3700dbc4328758f" + [[package]] name = "httparse" version = "1.10.1" @@ -7006,6 +7089,24 @@ dependencies = [ "tokio-rustls 0.24.1", ] +[[package]] +name = "hyper-rustls" +version = "0.25.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "399c78f9338483cb7e630c8474b07268983c6bd5acee012e4211f9f7bb21b070" +dependencies = [ + "futures-util", + "http 0.2.12", + "hyper 0.14.32", + "log", + "rustls 0.22.4", + "rustls-native-certs 0.7.3", + "rustls-pki-types", + "tokio", + "tokio-rustls 0.25.0", + "webpki-roots 0.26.11", +] + [[package]] name = "hyper-rustls" version = "0.26.0" @@ -7044,6 +7145,18 @@ dependencies = [ "webpki-roots 1.0.5", ] +[[package]] +name = "hyper-timeout" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbb958482e8c7be4bc3cf272a766a2b0bf1a6755e7a6ae777f017a31d11b13b1" +dependencies = [ + "hyper 0.14.32", + "pin-project-lite", + "tokio", + "tokio-io-timeout", +] + [[package]] name = "hyper-timeout" version = "0.5.2" @@ -7777,7 +7890,7 @@ dependencies = [ "hyper 1.8.1", "hyper-http-proxy", "hyper-rustls 0.27.7", - "hyper-timeout", + "hyper-timeout 0.5.2", "hyper-util", "jsonpath-rust", "k8s-openapi", @@ -7792,7 +7905,7 @@ dependencies = [ "tokio", "tokio-util", "tower 0.5.3", - "tower-http", + "tower-http 0.6.8", "tracing", ] @@ -8036,6 +8149,142 @@ dependencies = [ "redox_syscall 0.7.0", ] +[[package]] +name = "libsql" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2329faffc510cc3c6b4f00169a39177cc7099d3ed7647fc92f7cf26e53a8d976" +dependencies = [ + "anyhow", + "async-stream", + "async-trait", + "base64 0.21.7", + "bincode", + "bitflags 2.9.4", + "bytes", + "chrono", + "crc32fast", + "fallible-iterator 0.3.0", + "futures", + "http 0.2.12", + "hyper 0.14.32", + "hyper-rustls 0.25.0", + "libsql-hrana", + "libsql-sqlite3-parser", + "libsql-sys", + "libsql_replication", + "parking_lot", + "serde", + "serde_json", + "thiserror 1.0.69", + "tokio", + "tokio-stream", + "tokio-util", + "tonic 0.11.0", + "tonic-web", + "tower 0.4.13", + "tower-http 0.4.4", + "tracing", + "uuid", + "zerocopy 0.7.35", +] + +[[package]] +name = "libsql-ffi" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6cd1c1662822495393327856774f6803be25d85bfdcd5b9d4af35458f5daaf75" +dependencies = [ + "bindgen 0.66.1", + "cc", + "cmake", + "glob", +] + +[[package]] +name = "libsql-hrana" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "646d0aa75e412769018422f0da798f72e93bd51964f0b2ddad4317aa779ae444" +dependencies = [ + "base64 0.21.7", + "bytes", + "prost 0.12.6", + "serde", +] + +[[package]] +name = "libsql-rusqlite" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a4ce3a78c6e3c2b23b02ab6272df8340e1c53380497979d456882254f348d5f" +dependencies = [ + "bitflags 2.9.4", + "fallible-iterator 0.2.0", + "fallible-streaming-iterator", + "hashlink 0.8.4", + "libsql-ffi", + "smallvec", +] + +[[package]] +name = "libsql-sqlite3-parser" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15a90128c708356af8f7d767c9ac2946692c9112b4f74f07b99a01a60680e413" +dependencies = [ + "bitflags 2.9.4", + "cc", + "fallible-iterator 0.3.0", + "indexmap 2.11.1", + "log", + "memchr", + "phf 0.11.3", + "phf_codegen", + "phf_shared 0.11.3", + "uncased", +] + +[[package]] +name = "libsql-sys" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2a3c326fcfc36fe7578238d5ee6b58c529f8c76372acd61ec50267529cdaff95" +dependencies = [ + "bytes", + "libsql-ffi", + "libsql-rusqlite", + "once_cell", + "tracing", + "zerocopy 0.7.35", +] + +[[package]] +name = "libsql_replication" +version = "0.9.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d9a2e469ac8400659bd31f81a745908bcc5cb6b40be2f2ff8de90b15bec5501" +dependencies = [ + "aes 0.8.3", + "async-stream", + "async-trait", + "bytes", + "cbc", + "libsql-rusqlite", + "libsql-sys", + "parking_lot", + "prost 0.12.6", + "serde", + "thiserror 1.0.69", + "tokio", + "tokio-stream", + "tokio-util", + "tonic 0.11.0", + "tracing", + "uuid", + "zerocopy 0.7.35", +] + [[package]] name = "libsqlite3-sys" version = "0.30.1" @@ -9554,11 +9803,11 @@ dependencies = [ "opentelemetry-http", "opentelemetry-proto 0.27.0", "opentelemetry_sdk 0.27.1", - "prost", + "prost 0.13.5", "serde_json", "thiserror 1.0.69", "tokio", - "tonic", + "tonic 0.12.3", "tracing", ] @@ -9571,9 +9820,9 @@ dependencies = [ "hex", "opentelemetry 0.27.1", "opentelemetry_sdk 0.27.1", - "prost", + "prost 0.13.5", "serde", - "tonic", + "tonic 0.12.3", ] [[package]] @@ -9586,9 +9835,9 @@ dependencies = [ "hex", "opentelemetry 0.29.1", "opentelemetry_sdk 0.29.0", - "prost", + "prost 0.13.5", "serde", - "tonic", + "tonic 0.12.3", "tracing", ] @@ -9874,6 +10123,12 @@ dependencies = [ "hmac", ] +[[package]] +name = "peeking_take_while" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "19b17cddbe7ec3f8bc800887bab5e717348c95ea2ca0b1bf0837fb964dc67099" + [[package]] name = "pem" version = "1.1.1" @@ -10041,6 +10296,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "67eabc2ef2a60eb7faa00097bd1ffdb5bd28e62bf39990626a582201b7a754e5" dependencies = [ "siphasher 1.0.1", + "uncased", ] [[package]] @@ -10277,7 +10533,7 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" dependencies = [ - "zerocopy", + "zerocopy 0.8.33", ] [[package]] @@ -10451,6 +10707,16 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "prost" +version = "0.12.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "deb1435c188b76130da55f17a466d252ff7b1418b2ad3e037d127b94e3411f29" +dependencies = [ + "bytes", + "prost-derive 0.12.6", +] + [[package]] name = "prost" version = "0.13.5" @@ -10458,7 +10724,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2796faa41db3ec313a31f7624d9286acf277b52de526150b7e69f3debf891ee5" dependencies = [ "bytes", - "prost-derive", + "prost-derive 0.13.5", ] [[package]] @@ -10474,13 +10740,26 @@ dependencies = [ "once_cell", "petgraph", "prettyplease", - "prost", + "prost 0.13.5", "prost-types", "regex", "syn 2.0.114", "tempfile", ] +[[package]] +name = "prost-derive" +version = "0.12.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "81bddcdb20abf9501610992b6759a4c888aef7d1a7247ef75e2404275ac24af1" +dependencies = [ + "anyhow", + "itertools 0.12.1", + "proc-macro2", + "quote", + "syn 2.0.114", +] + [[package]] name = "prost-derive" version = "0.13.5" @@ -10500,7 +10779,7 @@ version = "0.13.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52c2c1bf36ddb1a1c396b3601a3cec27c2462e45f07c386894ec3ccf5332bd16" dependencies = [ - "prost", + "prost 0.13.5", ] [[package]] @@ -10725,7 +11004,7 @@ checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" dependencies = [ "rand_chacha 0.9.0", "rand_core 0.9.5", - "zerocopy", + "zerocopy 0.8.33", ] [[package]] @@ -11061,13 +11340,13 @@ dependencies = [ "serde", "serde_json", "serde_urlencoded", - "sync_wrapper", + "sync_wrapper 1.0.2", "tokio", "tokio-native-tls", "tokio-rustls 0.26.4", "tokio-util", "tower 0.5.3", - "tower-http", + "tower-http 0.6.8", "tower-service", "url", "wasm-bindgen", @@ -11108,12 +11387,12 @@ dependencies = [ "serde", "serde_json", "serde_urlencoded", - "sync_wrapper", + "sync_wrapper 1.0.2", "tokio", "tokio-rustls 0.26.4", "tokio-util", "tower 0.5.3", - "tower-http", + "tower-http 0.6.8", "tower-service", "url", "wasm-bindgen", @@ -13426,6 +13705,12 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "sync_wrapper" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2047c6ded9c721764247e62cd3b03c09ffc529b2ba5b10ec482ae507a4a70160" + [[package]] name = "sync_wrapper" version = "1.0.2" @@ -14054,6 +14339,16 @@ dependencies = [ "tracing", ] +[[package]] +name = "tokio-io-timeout" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0bd86198d9ee903fedd2f9a2e72014287c0d9167e4ae43b5853007205dda1b76" +dependencies = [ + "pin-project-lite", + "tokio", +] + [[package]] name = "tokio-macros" version = "2.5.0" @@ -14202,6 +14497,17 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-test" +version = "0.4.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f6d24790a10a7af737693a3e8f1d03faef7e6ca0cc99aae5066f533766de545" +dependencies = [ + "futures-core", + "tokio", + "tokio-stream", +] + [[package]] name = "tokio-tungstenite" version = "0.21.0" @@ -14336,6 +14642,33 @@ dependencies = [ "winnow 0.7.14", ] +[[package]] +name = "tonic" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76c4eb7a4e9ef9d4763600161f12f5070b92a578e1b634db88a6887844c91a13" +dependencies = [ + "async-stream", + "async-trait", + "axum 0.6.20", + "base64 0.21.7", + "bytes", + "h2 0.3.27", + "http 0.2.12", + "http-body 0.4.6", + "hyper 0.14.32", + "hyper-timeout 0.4.1", + "percent-encoding", + "pin-project", + "prost 0.12.6", + "tokio", + "tokio-stream", + "tower 0.4.13", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tonic" version = "0.12.3" @@ -14353,11 +14686,11 @@ dependencies = [ "http-body 1.0.1", "http-body-util", "hyper 1.8.1", - "hyper-timeout", + "hyper-timeout 0.5.2", "hyper-util", "percent-encoding", "pin-project", - "prost", + "prost 0.13.5", "rustls-native-certs 0.8.3", "rustls-pemfile 2.2.0", "socket2 0.5.10", @@ -14371,6 +14704,26 @@ dependencies = [ "webpki-roots 0.26.11", ] +[[package]] +name = "tonic-web" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc3b0e1cedbf19fdfb78ef3d672cb9928e0a91a9cb4629cc0c916e8cff8aaaa1" +dependencies = [ + "base64 0.21.7", + "bytes", + "http 0.2.12", + "http-body 0.4.6", + "hyper 0.14.32", + "pin-project", + "tokio-stream", + "tonic 0.11.0", + "tower-http 0.4.4", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tower" version = "0.4.13" @@ -14400,7 +14753,7 @@ dependencies = [ "futures-core", "futures-util", "pin-project-lite", - "sync_wrapper", + "sync_wrapper 1.0.2", "tokio", "tokio-util", "tower-layer", @@ -14425,6 +14778,43 @@ dependencies = [ "tower-service", ] +[[package]] +name = "tower-http" +version = "0.4.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c5bb1d698276a2443e5ecfabc1008bf15a36c12e6a7176e7bf089ea9131140" +dependencies = [ + "bitflags 2.9.4", + "bytes", + "futures-core", + "futures-util", + "http 0.2.12", + "http-body 0.4.6", + "http-range-header", + "pin-project-lite", + "tower 0.4.13", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "tower-http" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e9cd434a998747dd2c4276bc96ee2e0c7a2eadf3cae88e52be55a05fa9053f5" +dependencies = [ + "bitflags 2.9.4", + "bytes", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "pin-project-lite", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tower-http" version = "0.6.8" @@ -14790,6 +15180,15 @@ dependencies = [ "web-time", ] +[[package]] +name = "uncased" +version = "0.9.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e1b88fcfe09e89d3866a5c11019378088af2d24c3fbd4f0543f96b479ec90697" +dependencies = [ + "version_check", +] + [[package]] name = "unic-char-property" version = "0.9.0" @@ -15718,10 +16117,10 @@ dependencies = [ "tokio-stream", "tokio-tungstenite 0.24.0", "tokio-util", - "tonic", + "tonic 0.12.3", "tower 0.5.3", "tower-cookies", - "tower-http", + "tower-http 0.6.8", "tracing", "tracing-subscriber", "ulid", @@ -15866,7 +16265,7 @@ dependencies = [ "tokio-postgres 0.7.13", "tokio-stream", "tokio-util", - "tonic", + "tonic 0.12.3", "tracing", "tracing-appender", "tracing-opentelemetry", @@ -15920,6 +16319,26 @@ dependencies = [ "windmill-common", ] +[[package]] +name = "windmill-local" +version = "0.1.0" +dependencies = [ + "anyhow", + "axum 0.7.9", + "chrono", + "libsql", + "serde", + "serde_json", + "thiserror 1.0.69", + "tokio", + "tokio-test", + "tower 0.4.13", + "tower-http 0.5.2", + "tracing", + "tracing-subscriber", + "uuid", +] + [[package]] name = "windmill-macros" version = "1.613.4" @@ -17128,13 +17547,34 @@ dependencies = [ "synstructure 0.13.2", ] +[[package]] +name = "zerocopy" +version = "0.7.35" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1b9b4fd18abc82b8136838da5d50bae7bdea537c574d8dc1a34ed098d6c166f0" +dependencies = [ + "byteorder", + "zerocopy-derive 0.7.35", +] + [[package]] name = "zerocopy" version = "0.8.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "668f5168d10b9ee831de31933dc111a459c97ec93225beb307aed970d1372dfd" dependencies = [ - "zerocopy-derive", + "zerocopy-derive 0.8.33", +] + +[[package]] +name = "zerocopy-derive" +version = "0.7.35" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fa4f8080344d4671fb4e831a13ad1e68092748387dfc4f55e356242fae12ce3e" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.114", ] [[package]] diff --git a/backend/Cargo.toml b/backend/Cargo.toml index b100e9ef0b..6b91bdb179 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -18,6 +18,7 @@ members = [ "./windmill-indexer", "./windmill-macros", "./windmill-oauth", + "./windmill-local", "./parsers/windmill-parser", "./parsers/windmill-parser-ts", "./parsers/windmill-parser-go", diff --git a/backend/windmill-local/Cargo.toml b/backend/windmill-local/Cargo.toml new file mode 100644 index 0000000000..b80ddc4418 --- /dev/null +++ b/backend/windmill-local/Cargo.toml @@ -0,0 +1,46 @@ +[package] +name = "windmill-local" +version = "0.1.0" +edition = "2021" +description = "Windmill local mode with libSQL/Turso support for preview execution" + +[dependencies] +# libSQL - Turso's SQLite fork (used for local and remote Turso connections) +# Note: libsql IS the Turso database driver - Turso is built on libSQL +libsql = "0.9" + +# Async runtime +tokio = { version = "1", features = ["full", "sync", "macros"] } + +# Serialization +serde = { version = "1", features = ["derive"] } +serde_json = "1" + +# Error handling +anyhow = "1" +thiserror = "1" + +# UUID for job IDs +uuid = { version = "1", features = ["v4", "serde"] } + +# Timestamps +chrono = { version = "0.4", features = ["serde"] } + +# Tracing +tracing = "0.1" + +# HTTP server +axum = "0.7" +tower = "0.4" +tower-http = { version = "0.5", features = ["cors", "trace"] } + +# For script execution (will integrate with windmill-worker later) +# windmill-common = { path = "../windmill-common" } + +[dev-dependencies] +tokio-test = "0.4" +tracing-subscriber = { version = "0.3", features = ["env-filter"] } + +[[example]] +name = "local_server" +path = "examples/local_server.rs" diff --git a/backend/windmill-local/examples/local_server.rs b/backend/windmill-local/examples/local_server.rs new file mode 100644 index 0000000000..29a718f99c --- /dev/null +++ b/backend/windmill-local/examples/local_server.rs @@ -0,0 +1,65 @@ +//! Example: Run the Windmill local server +//! +//! This demonstrates running a local Windmill server with libSQL/SQLite backend. +//! +//! Run with: +//! cargo run -p windmill-local --example local_server +//! +//! Test with: +//! # Health check +//! curl http://localhost:8000/health +//! +//! # Run a bash preview (async) +//! curl -X POST http://localhost:8000/api/w/local/jobs/run/preview \ +//! -H "Content-Type: application/json" \ +//! -d '{"content": "echo Hello World", "language": "bash"}' +//! +//! # Run a bash preview and wait for result +//! curl -X POST http://localhost:8000/api/w/local/jobs/run_wait_result/preview \ +//! -H "Content-Type: application/json" \ +//! -d '{"content": "echo 42", "language": "bash"}' +//! +//! # Run a Python preview +//! curl -X POST http://localhost:8000/api/w/local/jobs/run_wait_result/preview \ +//! -H "Content-Type: application/json" \ +//! -d '{"content": "def main(x=1): return x * 2", "language": "python3", "args": {"x": 21}}' +//! +//! # Run a flow preview (linear flow with two steps) +//! curl -X POST http://localhost:8000/api/w/local/jobs/run_wait_result/preview_flow \ +//! -H "Content-Type: application/json" \ +//! -d '{ +//! "value": { +//! "modules": [ +//! {"id": "step1", "value": {"type": "rawscript", "language": "bash", "content": "echo 10"}}, +//! {"id": "step2", "value": {"type": "identity"}} +//! ] +//! }, +//! "args": {} +//! }' + +use std::net::SocketAddr; +use windmill_local::LocalServer; + +#[tokio::main] +async fn main() { + // Initialize tracing + tracing_subscriber::fmt() + .with_env_filter( + tracing_subscriber::EnvFilter::from_default_env() + .add_directive("windmill_local=info".parse().unwrap()), + ) + .init(); + + let addr: SocketAddr = "0.0.0.0:8000".parse().unwrap(); + println!("Starting Windmill Local Server on {}", addr); + println!(); + println!("Test endpoints:"); + println!(" Health: curl http://localhost:8000/health"); + println!(" Preview: curl -X POST http://localhost:8000/api/w/local/jobs/run_wait_result/preview \\"); + println!(" -H 'Content-Type: application/json' \\"); + println!(" -d '{{\"content\": \"echo Hello\", \"language\": \"bash\"}}'"); + println!(); + + let server = LocalServer::new(addr).await.expect("Failed to create server"); + server.run().await.expect("Server error"); +} diff --git a/backend/windmill-local/src/db.rs b/backend/windmill-local/src/db.rs new file mode 100644 index 0000000000..99e9a0d2d0 --- /dev/null +++ b/backend/windmill-local/src/db.rs @@ -0,0 +1,161 @@ +//! Database connection and initialization for local mode + +use libsql::{Builder, Connection, Database}; +use std::sync::Arc; +use tokio::sync::Mutex; + +use crate::error::Result; +use crate::schema; + +/// Local database wrapper +/// +/// For local mode, we use a single connection with in-process coordination. +/// This simplifies the implementation since we don't need `FOR UPDATE SKIP LOCKED`. +pub struct LocalDb { + #[allow(dead_code)] + db: Database, + /// Single connection for all operations (simplifies transaction handling) + conn: Arc>, +} + +impl LocalDb { + /// Create an in-memory database (for testing/ephemeral use) + pub async fn in_memory() -> Result { + let db = Builder::new_local(":memory:").build().await?; + let conn = db.connect()?; + let local_db = Self { + db, + conn: Arc::new(Mutex::new(conn)), + }; + local_db.init_schema().await?; + Ok(local_db) + } + + /// Create a file-based database + pub async fn file(path: &str) -> Result { + let db = Builder::new_local(path).build().await?; + let conn = db.connect()?; + let local_db = Self { + db, + conn: Arc::new(Mutex::new(conn)), + }; + local_db.init_schema().await?; + Ok(local_db) + } + + /// Create a Turso remote database connection + /// This would be used for the multi-writer scenario + pub async fn turso_remote(url: &str, auth_token: &str) -> Result { + let db = Builder::new_remote(url.to_string(), auth_token.to_string()) + .build() + .await?; + let conn = db.connect()?; + let local_db = Self { + db, + conn: Arc::new(Mutex::new(conn)), + }; + local_db.init_schema().await?; + Ok(local_db) + } + + /// Initialize the schema + async fn init_schema(&self) -> Result<()> { + let conn = self.conn.lock().await; + // Execute schema as multiple statements + conn.execute_batch(schema::SCHEMA).await?; + Ok(()) + } + + /// Reset the database (drop and recreate all tables) + pub async fn reset(&self) -> Result<()> { + let conn = self.conn.lock().await; + conn.execute_batch(schema::DROP_SCHEMA).await?; + conn.execute_batch(schema::SCHEMA).await?; + Ok(()) + } + + /// Get a reference to the connection (locked) + pub async fn conn(&self) -> tokio::sync::MutexGuard<'_, Connection> { + self.conn.lock().await + } + + /// Execute a simple query that returns no rows + pub async fn execute(&self, sql: &str, params: impl libsql::params::IntoParams) -> Result { + let conn = self.conn.lock().await; + let rows_affected = conn.execute(sql, params).await?; + Ok(rows_affected) + } + + /// Execute a query and return all rows + pub async fn query( + &self, + sql: &str, + params: impl libsql::params::IntoParams, + ) -> Result { + let conn = self.conn.lock().await; + let rows = conn.query(sql, params).await?; + Ok(rows) + } + + /// Execute a batch of statements (for transactions) + pub async fn execute_batch(&self, sql: &str) -> Result<()> { + let conn = self.conn.lock().await; + conn.execute_batch(sql).await?; + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_in_memory_db() { + let db = LocalDb::in_memory().await.unwrap(); + + // Verify tables exist + let rows = db + .query( + "SELECT name FROM sqlite_master WHERE type='table' ORDER BY name", + (), + ) + .await + .unwrap(); + + let mut tables = Vec::new(); + let mut rows = rows; + while let Some(row) = rows.next().await.unwrap() { + let name: String = row.get(0).unwrap(); + tables.push(name); + } + + assert!(tables.contains(&"v2_job".to_string())); + assert!(tables.contains(&"v2_job_queue".to_string())); + assert!(tables.contains(&"v2_job_completed".to_string())); + } + + #[tokio::test] + async fn test_reset_db() { + let db = LocalDb::in_memory().await.unwrap(); + + // Insert a job + db.execute( + "INSERT INTO v2_job (id, kind) VALUES ('test-uuid', 'preview')", + (), + ) + .await + .unwrap(); + + // Reset + db.reset().await.unwrap(); + + // Verify job is gone + let mut rows = db + .query("SELECT COUNT(*) FROM v2_job", ()) + .await + .unwrap(); + let row = rows.next().await.unwrap().unwrap(); + let count: i64 = row.get(0).unwrap(); + assert_eq!(count, 0); + } +} diff --git a/backend/windmill-local/src/error.rs b/backend/windmill-local/src/error.rs new file mode 100644 index 0000000000..315c4ff65b --- /dev/null +++ b/backend/windmill-local/src/error.rs @@ -0,0 +1,29 @@ +//! Error types for local mode + +use thiserror::Error; + +#[derive(Error, Debug)] +pub enum LocalError { + #[error("Database error: {0}")] + Database(#[from] libsql::Error), + + #[error("Serialization error: {0}")] + Serialization(#[from] serde_json::Error), + + #[error("Job not found: {0}")] + JobNotFound(uuid::Uuid), + + #[error("Invalid job state: {0}")] + InvalidJobState(String), + + #[error("Queue is empty")] + QueueEmpty, + + #[error("Execution error: {0}")] + Execution(String), + + #[error("Timeout")] + Timeout, +} + +pub type Result = std::result::Result; diff --git a/backend/windmill-local/src/executor.rs b/backend/windmill-local/src/executor.rs new file mode 100644 index 0000000000..a62038ec5e --- /dev/null +++ b/backend/windmill-local/src/executor.rs @@ -0,0 +1,330 @@ +//! Simple script executor for local mode +//! +//! This is a minimal executor that supports a few languages for demonstration. +//! A full implementation would integrate with windmill-worker's execution logic. + +use std::process::Stdio; +use tokio::process::Command; +use tokio::io::AsyncReadExt; + +use crate::error::{LocalError, Result}; +use crate::jobs::ScriptLang; + +/// Result of script execution +#[derive(Debug)] +pub struct ExecutionResult { + pub success: bool, + pub result: serde_json::Value, + pub logs: String, +} + +/// Execute a script with the given language and arguments +pub async fn execute_script( + language: ScriptLang, + code: &str, + args: &serde_json::Value, +) -> Result { + match language { + ScriptLang::Bash => execute_bash(code, args).await, + ScriptLang::Python3 => execute_python(code, args).await, + ScriptLang::Deno => execute_deno(code, args).await, + ScriptLang::Bun => execute_bun(code, args).await, + _ => Ok(ExecutionResult { + success: false, + result: serde_json::json!({ + "error": format!("Language {:?} not supported in local mode yet", language) + }), + logs: format!("Language {:?} not supported in local mode", language), + }), + } +} + +/// Execute a bash script +async fn execute_bash(code: &str, args: &serde_json::Value) -> Result { + // Create environment variables from args + let mut env_vars = Vec::new(); + if let serde_json::Value::Object(map) = args { + for (key, value) in map { + let val_str = match value { + serde_json::Value::String(s) => s.clone(), + _ => value.to_string(), + }; + env_vars.push((key.to_uppercase(), val_str)); + } + } + + let mut child = Command::new("bash") + .arg("-c") + .arg(code) + .envs(env_vars) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| LocalError::Execution(format!("Failed to spawn bash: {}", e)))?; + + let status = child + .wait() + .await + .map_err(|e| LocalError::Execution(format!("Failed to wait for bash: {}", e)))?; + + let mut stdout = String::new(); + let mut stderr = String::new(); + + if let Some(mut out) = child.stdout.take() { + out.read_to_string(&mut stdout).await.ok(); + } + if let Some(mut err) = child.stderr.take() { + err.read_to_string(&mut stderr).await.ok(); + } + + let logs = format!("{}{}", stdout, stderr); + let success = status.success(); + + // Try to parse stdout as JSON, otherwise use as string + let result = if success { + let trimmed = stdout.trim(); + serde_json::from_str(trimmed).unwrap_or_else(|_| serde_json::json!(trimmed)) + } else { + serde_json::json!({ "error": stderr.trim() }) + }; + + Ok(ExecutionResult { + success, + result, + logs, + }) +} + +/// Execute a Python script +async fn execute_python(code: &str, args: &serde_json::Value) -> Result { + // Wrap the code to handle args and return JSON result + let wrapped_code = format!( + r#" +import json +import sys + +# Args passed as JSON +args = json.loads('''{}''') + +# User code +{} + +# Call main if it exists +if 'main' in dir(): + result = main(**args) + print(json.dumps(result)) +"#, + serde_json::to_string(args)?, + code + ); + + let mut child = Command::new("python3") + .arg("-c") + .arg(&wrapped_code) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| LocalError::Execution(format!("Failed to spawn python3: {}", e)))?; + + let status = child + .wait() + .await + .map_err(|e| LocalError::Execution(format!("Failed to wait for python3: {}", e)))?; + + let mut stdout = String::new(); + let mut stderr = String::new(); + + if let Some(mut out) = child.stdout.take() { + out.read_to_string(&mut stdout).await.ok(); + } + if let Some(mut err) = child.stderr.take() { + err.read_to_string(&mut stderr).await.ok(); + } + + let logs = format!("{}{}", stdout, stderr); + let success = status.success(); + + let result = if success { + let trimmed = stdout.trim(); + // Get the last line as result (in case there's debug output) + let last_line = trimmed.lines().last().unwrap_or(""); + serde_json::from_str(last_line).unwrap_or_else(|_| serde_json::json!(trimmed)) + } else { + serde_json::json!({ "error": stderr.trim() }) + }; + + Ok(ExecutionResult { + success, + result, + logs, + }) +} + +/// Execute a Deno/TypeScript script +async fn execute_deno(code: &str, args: &serde_json::Value) -> Result { + // Wrap the code to handle args and return JSON result + let wrapped_code = format!( + r#" +const args = {}; + +{} + +// Call main if it exists +if (typeof main === 'function') {{ + const result = await main(args); + console.log(JSON.stringify(result)); +}} +"#, + serde_json::to_string(args)?, + code + ); + + let mut child = Command::new("deno") + .arg("eval") + .arg("--unstable") + .arg(&wrapped_code) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| LocalError::Execution(format!("Failed to spawn deno: {}", e)))?; + + let status = child + .wait() + .await + .map_err(|e| LocalError::Execution(format!("Failed to wait for deno: {}", e)))?; + + let mut stdout = String::new(); + let mut stderr = String::new(); + + if let Some(mut out) = child.stdout.take() { + out.read_to_string(&mut stdout).await.ok(); + } + if let Some(mut err) = child.stderr.take() { + err.read_to_string(&mut stderr).await.ok(); + } + + let logs = format!("{}{}", stdout, stderr); + let success = status.success(); + + let result = if success { + let trimmed = stdout.trim(); + let last_line = trimmed.lines().last().unwrap_or(""); + serde_json::from_str(last_line).unwrap_or_else(|_| serde_json::json!(trimmed)) + } else { + serde_json::json!({ "error": stderr.trim() }) + }; + + Ok(ExecutionResult { + success, + result, + logs, + }) +} + +/// Execute a Bun/TypeScript script +async fn execute_bun(code: &str, args: &serde_json::Value) -> Result { + // Similar to Deno but using Bun + let wrapped_code = format!( + r#" +const args = {}; + +{} + +// Call main if it exists +if (typeof main === 'function') {{ + const result = await main(args); + console.log(JSON.stringify(result)); +}} +"#, + serde_json::to_string(args)?, + code + ); + + let mut child = Command::new("bun") + .arg("eval") + .arg(&wrapped_code) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| LocalError::Execution(format!("Failed to spawn bun: {}", e)))?; + + let status = child + .wait() + .await + .map_err(|e| LocalError::Execution(format!("Failed to wait for bun: {}", e)))?; + + let mut stdout = String::new(); + let mut stderr = String::new(); + + if let Some(mut out) = child.stdout.take() { + out.read_to_string(&mut stdout).await.ok(); + } + if let Some(mut err) = child.stderr.take() { + err.read_to_string(&mut stderr).await.ok(); + } + + let logs = format!("{}{}", stdout, stderr); + let success = status.success(); + + let result = if success { + let trimmed = stdout.trim(); + let last_line = trimmed.lines().last().unwrap_or(""); + serde_json::from_str(last_line).unwrap_or_else(|_| serde_json::json!(trimmed)) + } else { + serde_json::json!({ "error": stderr.trim() }) + }; + + Ok(ExecutionResult { + success, + result, + logs, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_bash_execution() { + let result = execute_script( + ScriptLang::Bash, + "echo 42", + &serde_json::json!({}), + ) + .await + .unwrap(); + + assert!(result.success); + // Output "42" is parsed as JSON number + assert_eq!(result.result, serde_json::json!(42)); + } + + #[tokio::test] + async fn test_bash_with_args() { + let result = execute_script( + ScriptLang::Bash, + "echo $NAME", + &serde_json::json!({"name": "world"}), + ) + .await + .unwrap(); + + assert!(result.success); + assert_eq!(result.result, serde_json::json!("world")); + } + + #[tokio::test] + async fn test_bash_json_output() { + let result = execute_script( + ScriptLang::Bash, + r#"echo '{"key": "value"}'"#, + &serde_json::json!({}), + ) + .await + .unwrap(); + + assert!(result.success); + assert_eq!(result.result, serde_json::json!({"key": "value"})); + } +} diff --git a/backend/windmill-local/src/jobs.rs b/backend/windmill-local/src/jobs.rs new file mode 100644 index 0000000000..713f2a54e3 --- /dev/null +++ b/backend/windmill-local/src/jobs.rs @@ -0,0 +1,495 @@ +//! Job types and operations for local mode + +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +use crate::db::LocalDb; +use crate::error::{LocalError, Result}; + +/// Job kind (mirrors windmill-common JobKind) +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum JobKind { + Script, + Preview, + Flow, + FlowPreview, + Dependencies, + FlowDependencies, + ScriptHub, + Identity, + Http, + Graphql, + Postgresql, + Noop, + AppDependencies, + DeploymentCallback, + SingleScriptFlow, + FlowScript, + FlowNode, + AppScript, +} + +impl JobKind { + pub fn as_str(&self) -> &'static str { + match self { + JobKind::Script => "script", + JobKind::Preview => "preview", + JobKind::Flow => "flow", + JobKind::FlowPreview => "flowpreview", + JobKind::Dependencies => "dependencies", + JobKind::FlowDependencies => "flowdependencies", + JobKind::ScriptHub => "script_hub", + JobKind::Identity => "identity", + JobKind::Http => "http", + JobKind::Graphql => "graphql", + JobKind::Postgresql => "postgresql", + JobKind::Noop => "noop", + JobKind::AppDependencies => "appdependencies", + JobKind::DeploymentCallback => "deploymentcallback", + JobKind::SingleScriptFlow => "singlescriptflow", + JobKind::FlowScript => "flowscript", + JobKind::FlowNode => "flownode", + JobKind::AppScript => "appscript", + } + } +} + +/// Script language +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum ScriptLang { + Python3, + Deno, + Go, + Bash, + Postgresql, + Nativets, + Bun, + Mysql, + Bigquery, + Snowflake, + Graphql, + Powershell, + Mssql, + Php, + Bunnative, + Rust, + Ansible, + Csharp, + Oracledb, + Nu, + Java, + Duckdb, +} + +impl ScriptLang { + pub fn as_str(&self) -> &'static str { + match self { + ScriptLang::Python3 => "python3", + ScriptLang::Deno => "deno", + ScriptLang::Go => "go", + ScriptLang::Bash => "bash", + ScriptLang::Postgresql => "postgresql", + ScriptLang::Nativets => "nativets", + ScriptLang::Bun => "bun", + ScriptLang::Mysql => "mysql", + ScriptLang::Bigquery => "bigquery", + ScriptLang::Snowflake => "snowflake", + ScriptLang::Graphql => "graphql", + ScriptLang::Powershell => "powershell", + ScriptLang::Mssql => "mssql", + ScriptLang::Php => "php", + ScriptLang::Bunnative => "bunnative", + ScriptLang::Rust => "rust", + ScriptLang::Ansible => "ansible", + ScriptLang::Csharp => "csharp", + ScriptLang::Oracledb => "oracledb", + ScriptLang::Nu => "nu", + ScriptLang::Java => "java", + ScriptLang::Duckdb => "duckdb", + } + } + + pub fn from_str(s: &str) -> Option { + match s { + "python3" => Some(ScriptLang::Python3), + "deno" => Some(ScriptLang::Deno), + "go" => Some(ScriptLang::Go), + "bash" => Some(ScriptLang::Bash), + "postgresql" => Some(ScriptLang::Postgresql), + "nativets" => Some(ScriptLang::Nativets), + "bun" => Some(ScriptLang::Bun), + "mysql" => Some(ScriptLang::Mysql), + "bigquery" => Some(ScriptLang::Bigquery), + "snowflake" => Some(ScriptLang::Snowflake), + "graphql" => Some(ScriptLang::Graphql), + "powershell" => Some(ScriptLang::Powershell), + "mssql" => Some(ScriptLang::Mssql), + "php" => Some(ScriptLang::Php), + "bunnative" => Some(ScriptLang::Bunnative), + "rust" => Some(ScriptLang::Rust), + "ansible" => Some(ScriptLang::Ansible), + "csharp" => Some(ScriptLang::Csharp), + "oracledb" => Some(ScriptLang::Oracledb), + "nu" => Some(ScriptLang::Nu), + "java" => Some(ScriptLang::Java), + "duckdb" => Some(ScriptLang::Duckdb), + _ => None, + } + } +} + +/// Job status for completed jobs +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum JobStatus { + Success, + Failure, + Canceled, + Skipped, +} + +impl JobStatus { + pub fn as_str(&self) -> &'static str { + match self { + JobStatus::Success => "success", + JobStatus::Failure => "failure", + JobStatus::Canceled => "canceled", + JobStatus::Skipped => "skipped", + } + } + + pub fn from_str(s: &str) -> Option { + match s { + "success" => Some(JobStatus::Success), + "failure" => Some(JobStatus::Failure), + "canceled" => Some(JobStatus::Canceled), + "skipped" => Some(JobStatus::Skipped), + _ => None, + } + } +} + +/// Preview job request (simplified from windmill-api Preview struct) +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PreviewRequest { + pub content: String, + pub language: ScriptLang, + #[serde(default)] + pub args: serde_json::Value, + pub lock: Option, + pub tag: Option, +} + +/// Flow preview request +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct FlowPreviewRequest { + pub value: serde_json::Value, // FlowValue as JSON + #[serde(default)] + pub args: serde_json::Value, + pub tag: Option, +} + +/// A queued job (combines v2_job and v2_job_queue) +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct QueuedJob { + pub id: Uuid, + pub workspace_id: String, + pub kind: JobKind, + pub script_lang: Option, + pub raw_code: Option, + pub raw_lock: Option, + pub raw_flow: Option, + pub args: serde_json::Value, + pub tag: String, + pub created_at: DateTime, + pub scheduled_for: DateTime, + pub running: bool, + pub parent_job: Option, + pub root_job: Option, + pub flow_step_id: Option, +} + +/// A completed job +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CompletedJob { + pub id: Uuid, + pub workspace_id: String, + pub status: JobStatus, + pub result: serde_json::Value, + pub started_at: Option>, + pub completed_at: DateTime, + pub duration_ms: Option, +} + +/// Push a preview job to the queue +pub async fn push_preview(db: &LocalDb, req: PreviewRequest) -> Result { + let id = Uuid::new_v4(); + let now = Utc::now().to_rfc3339(); + let args_json = serde_json::to_string(&req.args)?; + let tag = req.tag.as_deref().unwrap_or("deno"); + + // Insert into v2_job + db.execute( + r#" + INSERT INTO v2_job (id, kind, script_lang, raw_code, raw_lock, args, tag, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) + "#, + libsql::params![ + id.to_string(), + JobKind::Preview.as_str(), + req.language.as_str(), + req.content, + req.lock, + args_json, + tag, + now.clone(), + ], + ) + .await?; + + // Insert into v2_job_queue + db.execute( + r#" + INSERT INTO v2_job_queue (id, tag, created_at, scheduled_for, running) + VALUES (?1, ?2, ?3, ?4, 0) + "#, + libsql::params![id.to_string(), tag, now.clone(), now], + ) + .await?; + + // Insert into v2_job_runtime (for heartbeat tracking) + db.execute( + "INSERT INTO v2_job_runtime (id) VALUES (?1)", + libsql::params![id.to_string()], + ) + .await?; + + // Insert into job_perms (simplified) + db.execute( + "INSERT INTO job_perms (job_id) VALUES (?1)", + libsql::params![id.to_string()], + ) + .await?; + + tracing::info!("Pushed preview job: {}", id); + Ok(id) +} + +/// Push a flow preview job to the queue +pub async fn push_flow_preview(db: &LocalDb, req: FlowPreviewRequest) -> Result { + let id = Uuid::new_v4(); + let now = Utc::now().to_rfc3339(); + let args_json = serde_json::to_string(&req.args)?; + let flow_json = serde_json::to_string(&req.value)?; + let tag = req.tag.as_deref().unwrap_or("flow"); + + // Insert into v2_job + db.execute( + r#" + INSERT INTO v2_job (id, kind, raw_flow, args, tag, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6) + "#, + libsql::params![ + id.to_string(), + JobKind::FlowPreview.as_str(), + flow_json, + args_json, + tag, + now.clone(), + ], + ) + .await?; + + // Insert into v2_job_queue + db.execute( + r#" + INSERT INTO v2_job_queue (id, tag, created_at, scheduled_for, running) + VALUES (?1, ?2, ?3, ?4, 0) + "#, + libsql::params![id.to_string(), tag, now.clone(), now], + ) + .await?; + + // Insert into v2_job_runtime + db.execute( + "INSERT INTO v2_job_runtime (id) VALUES (?1)", + libsql::params![id.to_string()], + ) + .await?; + + // Insert into job_perms + db.execute( + "INSERT INTO job_perms (job_id) VALUES (?1)", + libsql::params![id.to_string()], + ) + .await?; + + // Insert initial flow status + db.execute( + "INSERT INTO v2_job_status (id, flow_status) VALUES (?1, '{}')", + libsql::params![id.to_string()], + ) + .await?; + + tracing::info!("Pushed flow preview job: {}", id); + Ok(id) +} + +/// Get a completed job result (for polling) +pub async fn get_completed_job(db: &LocalDb, id: Uuid) -> Result> { + let mut rows = db + .query( + r#" + SELECT id, workspace_id, status, result, started_at, completed_at, duration_ms + FROM v2_job_completed + WHERE id = ?1 + "#, + libsql::params![id.to_string()], + ) + .await?; + + if let Some(row) = rows.next().await? { + let status_str: String = row.get(2)?; + let status = JobStatus::from_str(&status_str) + .ok_or_else(|| LocalError::InvalidJobState(status_str))?; + + let result_str: Option = row.get(3)?; + let result: serde_json::Value = result_str + .map(|s| serde_json::from_str(&s)) + .transpose()? + .unwrap_or(serde_json::Value::Null); + + let started_at_str: Option = row.get(4)?; + let started_at = started_at_str + .map(|s| DateTime::parse_from_rfc3339(&s).map(|dt| dt.with_timezone(&Utc))) + .transpose() + .ok() + .flatten(); + + let completed_at_str: String = row.get(5)?; + let completed_at = DateTime::parse_from_rfc3339(&completed_at_str) + .map(|dt| dt.with_timezone(&Utc)) + .unwrap_or_else(|_| Utc::now()); + + let duration_ms: Option = row.get(6)?; + + Ok(Some(CompletedJob { + id, + workspace_id: row.get(1)?, + status, + result, + started_at, + completed_at, + duration_ms, + })) + } else { + Ok(None) + } +} + +/// Mark a job as completed with result +pub async fn complete_job( + db: &LocalDb, + id: Uuid, + status: JobStatus, + result: serde_json::Value, + started_at: DateTime, +) -> Result<()> { + let now = Utc::now(); + let duration_ms = (now - started_at).num_milliseconds(); + let result_json = serde_json::to_string(&result)?; + + db.execute( + r#" + INSERT INTO v2_job_completed (id, workspace_id, status, result, started_at, completed_at, duration_ms) + SELECT ?1, workspace_id, ?2, ?3, ?4, ?5, ?6 + FROM v2_job WHERE id = ?1 + "#, + libsql::params![ + id.to_string(), + status.as_str(), + result_json, + started_at.to_rfc3339(), + now.to_rfc3339(), + duration_ms, + ], + ) + .await?; + + // Remove from queue + db.execute( + "DELETE FROM v2_job_queue WHERE id = ?1", + libsql::params![id.to_string()], + ) + .await?; + + tracing::info!("Completed job {}: {:?}", id, status); + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_push_preview() { + let db = LocalDb::in_memory().await.unwrap(); + + let req = PreviewRequest { + content: "export function main() { return 42; }".to_string(), + language: ScriptLang::Deno, + args: serde_json::json!({}), + lock: None, + tag: None, + }; + + let id = push_preview(&db, req).await.unwrap(); + + // Verify job exists in queue + let mut rows = db + .query( + "SELECT running FROM v2_job_queue WHERE id = ?1", + libsql::params![id.to_string()], + ) + .await + .unwrap(); + + let row = rows.next().await.unwrap().unwrap(); + let running: i64 = row.get(0).unwrap(); + assert_eq!(running, 0); + } + + #[tokio::test] + async fn test_complete_job() { + let db = LocalDb::in_memory().await.unwrap(); + + let req = PreviewRequest { + content: "export function main() { return 42; }".to_string(), + language: ScriptLang::Deno, + args: serde_json::json!({}), + lock: None, + tag: None, + }; + + let id = push_preview(&db, req).await.unwrap(); + let started_at = Utc::now(); + + complete_job( + &db, + id, + JobStatus::Success, + serde_json::json!(42), + started_at, + ) + .await + .unwrap(); + + // Verify job is in completed + let completed = get_completed_job(&db, id).await.unwrap().unwrap(); + assert_eq!(completed.status, JobStatus::Success); + assert_eq!(completed.result, serde_json::json!(42)); + } +} diff --git a/backend/windmill-local/src/lib.rs b/backend/windmill-local/src/lib.rs new file mode 100644 index 0000000000..525187f794 --- /dev/null +++ b/backend/windmill-local/src/lib.rs @@ -0,0 +1,30 @@ +//! Windmill Local Mode +//! +//! This crate provides a minimal local mode for Windmill using libSQL (SQLite/Turso) +//! instead of PostgreSQL. The goal is to support preview execution end-to-end +//! with a lightweight, embedded database. +//! +//! ## Scope +//! - Script preview execution +//! - Flow preview execution +//! - In-memory or file-based SQLite storage +//! - Remote Turso database support for multi-writer scenarios +//! +//! ## Non-goals (for this experiment) +//! - Full feature parity with PostgreSQL mode +//! - Multi-worker support (single embedded worker) +//! - Persistence of scripts/flows (only jobs) + +pub mod db; +pub mod schema; +pub mod jobs; +pub mod queue; +pub mod executor; +pub mod worker; +pub mod server; +pub mod error; + +pub use db::LocalDb; +pub use error::LocalError; +pub use worker::Worker; +pub use server::LocalServer; diff --git a/backend/windmill-local/src/queue.rs b/backend/windmill-local/src/queue.rs new file mode 100644 index 0000000000..cec559db92 --- /dev/null +++ b/backend/windmill-local/src/queue.rs @@ -0,0 +1,296 @@ +//! Queue operations for local mode +//! +//! Since local mode uses a single worker, we don't need the complex +//! `FOR UPDATE SKIP LOCKED` mechanism. Instead, we use simple atomic +//! operations with the database lock. + +use chrono::{DateTime, Utc}; +use uuid::Uuid; + +use crate::db::LocalDb; +use crate::error::{LocalError, Result}; +use crate::jobs::{JobKind, QueuedJob, ScriptLang}; + +/// Pull the next job from the queue +/// +/// This is simplified from the PostgreSQL version since we have a single worker +/// and use the connection mutex for coordination. +pub async fn pull_job(db: &LocalDb) -> Result> { + let now = Utc::now().to_rfc3339(); + + // Get the next job (ordered by priority, then scheduled_for) + // We use a transaction-like approach: SELECT then UPDATE + let mut rows = db + .query( + r#" + SELECT q.id, j.workspace_id, j.kind, j.script_lang, j.raw_code, j.raw_lock, + j.raw_flow, j.args, q.tag, j.created_at, q.scheduled_for, + j.parent_job, j.root_job, j.flow_step_id + FROM v2_job_queue q + JOIN v2_job j ON q.id = j.id + WHERE q.running = 0 AND q.scheduled_for <= ?1 + ORDER BY q.priority DESC, q.scheduled_for ASC + LIMIT 1 + "#, + libsql::params![now], + ) + .await?; + + let Some(row) = rows.next().await? else { + return Ok(None); + }; + + let id_str: String = row.get(0)?; + let id = Uuid::parse_str(&id_str).map_err(|e| LocalError::InvalidJobState(e.to_string()))?; + + // Mark as running + let started_at = Utc::now().to_rfc3339(); + db.execute( + "UPDATE v2_job_queue SET running = 1, started_at = ?2 WHERE id = ?1", + libsql::params![id_str.clone(), started_at], + ) + .await?; + + // Parse the job fields + let kind_str: String = row.get(2)?; + let kind = match kind_str.as_str() { + "preview" => JobKind::Preview, + "flowpreview" => JobKind::FlowPreview, + "script" => JobKind::Script, + "flow" => JobKind::Flow, + "flowscript" => JobKind::FlowScript, + "flownode" => JobKind::FlowNode, + _ => JobKind::Preview, // Default + }; + + let lang_str: Option = row.get(3)?; + let script_lang = lang_str.and_then(|s| ScriptLang::from_str(&s)); + + let raw_code: Option = row.get(4)?; + let raw_lock: Option = row.get(5)?; + + let raw_flow_str: Option = row.get(6)?; + let raw_flow = raw_flow_str + .map(|s| serde_json::from_str(&s)) + .transpose()?; + + let args_str: Option = row.get(7)?; + let args: serde_json::Value = args_str + .map(|s| serde_json::from_str(&s)) + .transpose()? + .unwrap_or(serde_json::Value::Object(serde_json::Map::new())); + + let tag: String = row.get(8)?; + + let created_at_str: String = row.get(9)?; + let created_at = DateTime::parse_from_rfc3339(&created_at_str) + .map(|dt| dt.with_timezone(&Utc)) + .unwrap_or_else(|_| Utc::now()); + + let scheduled_for_str: String = row.get(10)?; + let scheduled_for = DateTime::parse_from_rfc3339(&scheduled_for_str) + .map(|dt| dt.with_timezone(&Utc)) + .unwrap_or_else(|_| Utc::now()); + + let parent_job_str: Option = row.get(11)?; + let parent_job = parent_job_str.and_then(|s| Uuid::parse_str(&s).ok()); + + let root_job_str: Option = row.get(12)?; + let root_job = root_job_str.and_then(|s| Uuid::parse_str(&s).ok()); + + let flow_step_id: Option = row.get(13)?; + + Ok(Some(QueuedJob { + id, + workspace_id: row.get(1)?, + kind, + script_lang, + raw_code, + raw_lock, + raw_flow, + args, + tag, + created_at, + scheduled_for, + running: true, + parent_job, + root_job, + flow_step_id, + })) +} + +/// Get queue statistics +pub async fn queue_stats(db: &LocalDb) -> Result { + let mut rows = db + .query( + r#" + SELECT + COUNT(*) as total, + SUM(CASE WHEN running = 1 THEN 1 ELSE 0 END) as running, + SUM(CASE WHEN running = 0 THEN 1 ELSE 0 END) as pending + FROM v2_job_queue + "#, + (), + ) + .await?; + + let row = rows.next().await?.ok_or(LocalError::QueueEmpty)?; + + Ok(QueueStats { + total: row.get::(0)? as u64, + running: row.get::(1).unwrap_or(0) as u64, + pending: row.get::(2).unwrap_or(0) as u64, + }) +} + +#[derive(Debug, Clone)] +pub struct QueueStats { + pub total: u64, + pub running: u64, + pub pending: u64, +} + +/// Update job heartbeat (ping) +pub async fn ping_job(db: &LocalDb, id: Uuid) -> Result<()> { + let now = Utc::now().to_rfc3339(); + db.execute( + "UPDATE v2_job_runtime SET ping = ?2 WHERE id = ?1", + libsql::params![id.to_string(), now], + ) + .await?; + Ok(()) +} + +/// Update flow status for a running flow job +pub async fn update_flow_status( + db: &LocalDb, + id: Uuid, + flow_status: &serde_json::Value, +) -> Result<()> { + let status_json = serde_json::to_string(flow_status)?; + db.execute( + "UPDATE v2_job_status SET flow_status = ?2 WHERE id = ?1", + libsql::params![id.to_string(), status_json], + ) + .await?; + Ok(()) +} + +/// Push a child job for flow execution +pub async fn push_flow_child_job( + db: &LocalDb, + parent_id: Uuid, + root_id: Uuid, + step_id: &str, + kind: JobKind, + script_lang: Option, + raw_code: Option<&str>, + args: &serde_json::Value, +) -> Result { + let id = Uuid::new_v4(); + let now = Utc::now().to_rfc3339(); + let args_json = serde_json::to_string(args)?; + + // Insert into v2_job + db.execute( + r#" + INSERT INTO v2_job (id, kind, script_lang, raw_code, args, parent_job, root_job, flow_step_id, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) + "#, + libsql::params![ + id.to_string(), + kind.as_str(), + script_lang.map(|l| l.as_str()), + raw_code, + args_json, + parent_id.to_string(), + root_id.to_string(), + step_id, + now.clone(), + ], + ) + .await?; + + // Insert into v2_job_queue + db.execute( + r#" + INSERT INTO v2_job_queue (id, tag, created_at, scheduled_for, running) + VALUES (?1, 'flow', ?2, ?3, 0) + "#, + libsql::params![id.to_string(), now.clone(), now], + ) + .await?; + + // Insert into v2_job_runtime + db.execute( + "INSERT INTO v2_job_runtime (id) VALUES (?1)", + libsql::params![id.to_string()], + ) + .await?; + + Ok(id) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::jobs::{push_preview, PreviewRequest}; + + #[tokio::test] + async fn test_pull_job() { + let db = LocalDb::in_memory().await.unwrap(); + + // Push a job + let req = PreviewRequest { + content: "export function main() { return 42; }".to_string(), + language: ScriptLang::Deno, + args: serde_json::json!({}), + lock: None, + tag: None, + }; + let pushed_id = push_preview(&db, req).await.unwrap(); + + // Pull it + let job = pull_job(&db).await.unwrap().unwrap(); + assert_eq!(job.id, pushed_id); + assert!(job.running); + assert_eq!(job.kind, JobKind::Preview); + + // Queue should now be empty (job is running) + let job2 = pull_job(&db).await.unwrap(); + assert!(job2.is_none()); + } + + #[tokio::test] + async fn test_queue_stats() { + let db = LocalDb::in_memory().await.unwrap(); + + // Initially empty + let stats = queue_stats(&db).await.unwrap(); + assert_eq!(stats.total, 0); + + // Push two jobs + let req = PreviewRequest { + content: "test".to_string(), + language: ScriptLang::Deno, + args: serde_json::json!({}), + lock: None, + tag: None, + }; + push_preview(&db, req.clone()).await.unwrap(); + push_preview(&db, req).await.unwrap(); + + let stats = queue_stats(&db).await.unwrap(); + assert_eq!(stats.total, 2); + assert_eq!(stats.pending, 2); + assert_eq!(stats.running, 0); + + // Pull one + pull_job(&db).await.unwrap(); + + let stats = queue_stats(&db).await.unwrap(); + assert_eq!(stats.total, 2); + assert_eq!(stats.pending, 1); + assert_eq!(stats.running, 1); + } +} diff --git a/backend/windmill-local/src/schema.rs b/backend/windmill-local/src/schema.rs new file mode 100644 index 0000000000..8d2a61ca9a --- /dev/null +++ b/backend/windmill-local/src/schema.rs @@ -0,0 +1,218 @@ +//! SQLite schema for local mode +//! +//! This is a minimal schema supporting preview job execution. +//! Key differences from PostgreSQL: +//! - ENUMs are TEXT with CHECK constraints +//! - JSONB is JSON (stored as TEXT in SQLite) +//! - Arrays are JSON arrays +//! - No FOR UPDATE SKIP LOCKED (single worker, in-process coordination) + +/// SQL to create the minimal schema for local mode preview execution +pub const SCHEMA: &str = r#" +-- Job kinds (equivalent to PostgreSQL ENUM) +-- Values: script, preview, flow, flowpreview, dependencies, flowdependencies, +-- script_hub, identity, http, graphql, postgresql, noop, appdependencies, +-- deploymentcallback, singlescriptflow, flowscript, flownode, appscript + +-- Job status (equivalent to PostgreSQL ENUM) +-- Values: success, failure, canceled, skipped + +-- Script languages (equivalent to PostgreSQL ENUM) +-- Values: python3, deno, go, bash, postgresql, nativets, bun, mysql, bigquery, +-- snowflake, graphql, powershell, mssql, php, bunnative, rust, ansible, +-- csharp, oracledb, nu, java, duckdb + +-- Main job table (minimal for preview) +CREATE TABLE IF NOT EXISTS v2_job ( + id TEXT PRIMARY KEY, -- UUID as TEXT + workspace_id TEXT NOT NULL DEFAULT 'local', + + -- Raw code for preview jobs + raw_code TEXT, + raw_lock TEXT, + raw_flow TEXT, -- JSON for flow definitions + + -- Job metadata + tag TEXT DEFAULT 'deno', + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + created_by TEXT NOT NULL DEFAULT 'local_user', + + -- Permission context (simplified for local mode) + permissioned_as TEXT NOT NULL DEFAULT 'u/local_user', + permissioned_as_email TEXT DEFAULT 'local@windmill.local', + + -- Job type info + kind TEXT NOT NULL DEFAULT 'preview' CHECK (kind IN ( + 'script', 'preview', 'flow', 'flowpreview', 'dependencies', + 'flowdependencies', 'script_hub', 'identity', 'http', 'graphql', + 'postgresql', 'noop', 'appdependencies', 'deploymentcallback', + 'singlescriptflow', 'flowscript', 'flownode', 'appscript' + )), + + -- Script execution details + script_lang TEXT CHECK (script_lang IN ( + 'python3', 'deno', 'go', 'bash', 'postgresql', 'nativets', 'bun', + 'mysql', 'bigquery', 'snowflake', 'graphql', 'powershell', 'mssql', + 'php', 'bunnative', 'rust', 'ansible', 'csharp', 'oracledb', 'nu', + 'java', 'duckdb' + )), + + -- Flow execution details + parent_job TEXT, -- UUID reference + root_job TEXT, -- UUID reference + flow_step INTEGER, + flow_step_id TEXT, + flow_innermost_root_job TEXT, + + -- Execution settings + timeout INTEGER, + priority INTEGER DEFAULT 0, + same_worker INTEGER DEFAULT 0, -- BOOLEAN as INTEGER + visible_to_owner INTEGER DEFAULT 1, + + -- Arguments (JSON) + args TEXT, -- JSON object + + -- Pre-run error if validation failed + pre_run_error TEXT +); + +-- Job queue table +CREATE TABLE IF NOT EXISTS v2_job_queue ( + id TEXT PRIMARY KEY, -- UUID, references v2_job.id + workspace_id TEXT NOT NULL DEFAULT 'local', + + -- Timestamps + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + started_at TEXT, + scheduled_for TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + + -- Queue state + running INTEGER NOT NULL DEFAULT 0, -- BOOLEAN + canceled_by TEXT, + canceled_reason TEXT, + + -- Suspend state (for approval flows) + suspend INTEGER DEFAULT 0, + suspend_until TEXT, + + -- Execution settings + tag TEXT DEFAULT 'deno', + priority INTEGER DEFAULT 0, + worker TEXT, + + FOREIGN KEY (id) REFERENCES v2_job(id) +); + +-- Index for queue ordering (simulates queue_sort_v2) +CREATE INDEX IF NOT EXISTS idx_queue_sort ON v2_job_queue ( + priority DESC, scheduled_for ASC, tag +) WHERE running = 0; + +-- Job runtime tracking (heartbeat/ping) +CREATE TABLE IF NOT EXISTS v2_job_runtime ( + id TEXT PRIMARY KEY, -- UUID, references v2_job.id + ping TEXT, -- Timestamp + memory_peak INTEGER, + + FOREIGN KEY (id) REFERENCES v2_job(id) +); + +-- Completed jobs with results +CREATE TABLE IF NOT EXISTS v2_job_completed ( + id TEXT PRIMARY KEY, -- UUID, references v2_job.id + workspace_id TEXT NOT NULL DEFAULT 'local', + + -- Timing + started_at TEXT, + completed_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + duration_ms INTEGER, + + -- Result + result TEXT, -- JSON + result_columns TEXT, -- JSON array of column names + + -- Status + status TEXT NOT NULL DEFAULT 'success' CHECK (status IN ( + 'success', 'failure', 'canceled', 'skipped' + )), + + -- Cancellation details + canceled_by TEXT, + canceled_reason TEXT, + + -- Flow status (for flow jobs) + flow_status TEXT, -- JSON + + -- Execution details + memory_peak INTEGER, + worker TEXT, + deleted INTEGER DEFAULT 0, -- BOOLEAN + + FOREIGN KEY (id) REFERENCES v2_job(id) +); + +-- Index for completed job lookup by workspace and time +CREATE INDEX IF NOT EXISTS idx_completed_workspace_time ON v2_job_completed ( + workspace_id, completed_at DESC +); + +-- Flow status tracking (separate from completed to allow updates during execution) +CREATE TABLE IF NOT EXISTS v2_job_status ( + id TEXT PRIMARY KEY, -- UUID, references v2_job.id + flow_status TEXT, -- JSON object tracking flow module execution + flow_leaf_jobs TEXT, -- JSON object + workflow_as_code_status TEXT, -- JSON object + + FOREIGN KEY (id) REFERENCES v2_job(id) +); + +-- Simplified job permissions (for local mode, mostly unused) +CREATE TABLE IF NOT EXISTS job_perms ( + job_id TEXT PRIMARY KEY, + email TEXT DEFAULT 'local@windmill.local', + username TEXT DEFAULT 'local_user', + is_admin INTEGER DEFAULT 1, -- BOOLEAN + is_operator INTEGER DEFAULT 0, -- BOOLEAN + workspace_id TEXT DEFAULT 'local', + groups TEXT, -- JSON array + folders TEXT, -- JSON array of objects + + FOREIGN KEY (job_id) REFERENCES v2_job(id) +); + +-- Simple audit log (optional for local mode) +CREATE TABLE IF NOT EXISTS audit ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + workspace_id TEXT DEFAULT 'local', + timestamp TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + username TEXT DEFAULT 'local_user', + operation TEXT NOT NULL, + action_kind TEXT CHECK (action_kind IN ('create', 'update', 'delete', 'execute')), + resource TEXT, + parameters TEXT -- JSON +); + +-- Job logs storage +CREATE TABLE IF NOT EXISTS job_logs ( + job_id TEXT PRIMARY KEY, + workspace_id TEXT DEFAULT 'local', + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + logs TEXT, + log_offset INTEGER DEFAULT 0, + + FOREIGN KEY (job_id) REFERENCES v2_job(id) +); +"#; + +/// SQL to drop all tables (for testing/reset) +pub const DROP_SCHEMA: &str = r#" +DROP TABLE IF EXISTS job_logs; +DROP TABLE IF EXISTS audit; +DROP TABLE IF EXISTS job_perms; +DROP TABLE IF EXISTS v2_job_status; +DROP TABLE IF EXISTS v2_job_completed; +DROP TABLE IF EXISTS v2_job_runtime; +DROP TABLE IF EXISTS v2_job_queue; +DROP TABLE IF EXISTS v2_job; +"#; diff --git a/backend/windmill-local/src/server.rs b/backend/windmill-local/src/server.rs new file mode 100644 index 0000000000..474a34757f --- /dev/null +++ b/backend/windmill-local/src/server.rs @@ -0,0 +1,419 @@ +//! HTTP server for local mode +//! +//! Provides a minimal API compatible with Windmill's preview endpoints. + +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; + +use axum::{ + extract::{Path, State}, + http::StatusCode, + response::IntoResponse, + routing::{get, post}, + Json, Router, +}; +use serde::{Deserialize, Serialize}; +use tokio::sync::watch; +use tower_http::cors::{Any, CorsLayer}; +use tower_http::trace::TraceLayer; +use uuid::Uuid; + +use crate::db::LocalDb; +use crate::error::Result; +use crate::jobs::{ + get_completed_job, push_flow_preview, push_preview, FlowPreviewRequest, + JobStatus, PreviewRequest, ScriptLang, +}; +use crate::worker::Worker; + +/// Application state shared across handlers +pub struct AppState { + pub db: Arc, +} + +/// Local server that runs the API and embedded worker +pub struct LocalServer { + db: Arc, + addr: SocketAddr, +} + +impl LocalServer { + /// Create a new local server + pub async fn new(addr: SocketAddr) -> Result { + let db = Arc::new(LocalDb::in_memory().await?); + Ok(Self { db, addr }) + } + + /// Create a local server with a file-based database + pub async fn with_file(addr: SocketAddr, db_path: &str) -> Result { + let db = Arc::new(LocalDb::file(db_path).await?); + Ok(Self { db, addr }) + } + + /// Create a local server connected to a remote Turso database + pub async fn with_turso(addr: SocketAddr, url: &str, auth_token: &str) -> Result { + let db = Arc::new(LocalDb::turso_remote(url, auth_token).await?); + Ok(Self { db, addr }) + } + + /// Run the server + pub async fn run(self) -> Result<()> { + let (shutdown_tx, shutdown_rx) = watch::channel(false); + + // Start the embedded worker + let worker_db = self.db.clone(); + let worker_handle = tokio::spawn(async move { + let mut worker = Worker::new(worker_db, shutdown_rx); + if let Err(e) = worker.run().await { + tracing::error!("Worker error: {}", e); + } + }); + + // Build the router + let state = Arc::new(AppState { db: self.db }); + let app = create_router(state); + + // Run the server + tracing::info!("Local server listening on {}", self.addr); + let listener = tokio::net::TcpListener::bind(self.addr).await.unwrap(); + + // Handle graceful shutdown + let shutdown_signal = async move { + tokio::signal::ctrl_c() + .await + .expect("Failed to install CTRL+C signal handler"); + tracing::info!("Shutdown signal received"); + shutdown_tx.send(true).ok(); + }; + + axum::serve(listener, app) + .with_graceful_shutdown(shutdown_signal) + .await + .unwrap(); + + // Wait for worker to finish + worker_handle.await.ok(); + + Ok(()) + } +} + +/// Create the API router +fn create_router(state: Arc) -> Router { + let cors = CorsLayer::new() + .allow_origin(Any) + .allow_methods(Any) + .allow_headers(Any); + + Router::new() + // Health check + .route("/health", get(health_check)) + // Preview endpoints (mimics windmill-api) + .route("/api/w/:workspace/jobs/run/preview", post(run_preview)) + .route( + "/api/w/:workspace/jobs/run_wait_result/preview", + post(run_wait_result_preview), + ) + .route( + "/api/w/:workspace/jobs/run/preview_flow", + post(run_preview_flow), + ) + .route( + "/api/w/:workspace/jobs/run_wait_result/preview_flow", + post(run_wait_result_preview_flow), + ) + // Get job result + .route( + "/api/w/:workspace/jobs_u/completed/get_result/:job_id", + get(get_job_result), + ) + .layer(TraceLayer::new_for_http()) + .layer(cors) + .with_state(state) +} + +// === Request/Response Types === + +#[derive(Debug, Deserialize)] +struct PreviewPayload { + content: String, + language: String, + #[serde(default)] + args: serde_json::Value, + lock: Option, + tag: Option, +} + +#[derive(Debug, Deserialize)] +struct FlowPreviewPayload { + value: serde_json::Value, + #[serde(default)] + args: serde_json::Value, + tag: Option, +} + +#[derive(Debug, Serialize)] +struct JobCreatedResponse { + job_id: String, +} + +#[derive(Debug, Serialize)] +struct ErrorResponse { + error: String, +} + +// === Handlers === + +async fn health_check() -> &'static str { + "OK" +} + +/// Run a preview job (async - returns job ID) +async fn run_preview( + State(state): State>, + Path(_workspace): Path, + Json(payload): Json, +) -> impl IntoResponse { + let lang = match ScriptLang::from_str(&payload.language) { + Some(l) => l, + None => { + return ( + StatusCode::BAD_REQUEST, + Json(ErrorResponse { + error: format!("Unknown language: {}", payload.language), + }), + ) + .into_response(); + } + }; + + let req = PreviewRequest { + content: payload.content, + language: lang, + args: payload.args, + lock: payload.lock, + tag: payload.tag, + }; + + match push_preview(&state.db, req).await { + Ok(job_id) => ( + StatusCode::CREATED, + Json(JobCreatedResponse { + job_id: job_id.to_string(), + }), + ) + .into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: e.to_string(), + }), + ) + .into_response(), + } +} + +/// Run a preview job and wait for result +async fn run_wait_result_preview( + State(state): State>, + Path(_workspace): Path, + Json(payload): Json, +) -> impl IntoResponse { + let lang = match ScriptLang::from_str(&payload.language) { + Some(l) => l, + None => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Unknown language: {}", payload.language)})), + ) + .into_response(); + } + }; + + let req = PreviewRequest { + content: payload.content, + language: lang, + args: payload.args, + lock: payload.lock, + tag: payload.tag, + }; + + let job_id = match push_preview(&state.db, req).await { + Ok(id) => id, + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": e.to_string()})), + ) + .into_response(); + } + }; + + // Poll for result with timeout + wait_for_result(&state.db, job_id, Duration::from_secs(60)).await +} + +/// Run a flow preview job (async - returns job ID) +async fn run_preview_flow( + State(state): State>, + Path(_workspace): Path, + Json(payload): Json, +) -> impl IntoResponse { + let req = FlowPreviewRequest { + value: payload.value, + args: payload.args, + tag: payload.tag, + }; + + match push_flow_preview(&state.db, req).await { + Ok(job_id) => ( + StatusCode::CREATED, + Json(JobCreatedResponse { + job_id: job_id.to_string(), + }), + ) + .into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: e.to_string(), + }), + ) + .into_response(), + } +} + +/// Run a flow preview job and wait for result +async fn run_wait_result_preview_flow( + State(state): State>, + Path(_workspace): Path, + Json(payload): Json, +) -> impl IntoResponse { + let req = FlowPreviewRequest { + value: payload.value, + args: payload.args, + tag: payload.tag, + }; + + let job_id = match push_flow_preview(&state.db, req).await { + Ok(id) => id, + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": e.to_string()})), + ) + .into_response(); + } + }; + + // Poll for result with timeout + wait_for_result(&state.db, job_id, Duration::from_secs(120)).await +} + +/// Get the result of a completed job +async fn get_job_result( + State(state): State>, + Path((_workspace, job_id)): Path<(String, String)>, +) -> impl IntoResponse { + let job_id = match Uuid::parse_str(&job_id) { + Ok(id) => id, + Err(_) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "Invalid job ID"})), + ) + .into_response(); + } + }; + + match get_completed_job(&state.db, job_id).await { + Ok(Some(job)) => { + if job.status == JobStatus::Success { + (StatusCode::OK, Json(job.result)).into_response() + } else { + (StatusCode::INTERNAL_SERVER_ERROR, Json(job.result)).into_response() + } + } + Ok(None) => ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Job not found or not completed"})), + ) + .into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": e.to_string()})), + ) + .into_response(), + } +} + +/// Poll for job completion with timeout +async fn wait_for_result( + db: &LocalDb, + job_id: Uuid, + timeout: Duration, +) -> axum::response::Response { + let start = std::time::Instant::now(); + let fast_poll_duration = Duration::from_secs(2); + let fast_poll_interval = Duration::from_millis(50); + let slow_poll_interval = Duration::from_millis(200); + + loop { + if start.elapsed() > timeout { + return ( + StatusCode::REQUEST_TIMEOUT, + Json(serde_json::json!({"error": "Timeout waiting for job result"})), + ) + .into_response(); + } + + match get_completed_job(db, job_id).await { + Ok(Some(job)) => { + if job.status == JobStatus::Success { + return (StatusCode::OK, Json(job.result)).into_response(); + } else { + return (StatusCode::INTERNAL_SERVER_ERROR, Json(job.result)).into_response(); + } + } + Ok(None) => { + // Job not completed yet, keep polling + let interval = if start.elapsed() < fast_poll_duration { + fast_poll_interval + } else { + slow_poll_interval + }; + tokio::time::sleep(interval).await; + } + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": e.to_string()})), + ) + .into_response(); + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use axum::body::Body; + use axum::http::Request; + use tower::ServiceExt; + + #[tokio::test] + async fn test_health_check() { + let db = Arc::new(LocalDb::in_memory().await.unwrap()); + let state = Arc::new(AppState { db }); + let app = create_router(state); + + let response = app + .oneshot(Request::builder().uri("/health").body(Body::empty()).unwrap()) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + } +} diff --git a/backend/windmill-local/src/worker.rs b/backend/windmill-local/src/worker.rs new file mode 100644 index 0000000000..5041507728 --- /dev/null +++ b/backend/windmill-local/src/worker.rs @@ -0,0 +1,287 @@ +//! Worker for local mode +//! +//! A single embedded worker that pulls jobs from the queue and executes them. + +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::watch; +use chrono::Utc; + +use crate::db::LocalDb; +use crate::error::Result; +use crate::executor::{execute_script, ExecutionResult}; +use crate::jobs::{complete_job, JobKind, JobStatus, QueuedJob}; +use crate::queue::pull_job; + +/// Worker that processes jobs from the queue +pub struct Worker { + db: Arc, + /// Channel to signal shutdown + shutdown_rx: watch::Receiver, +} + +impl Worker { + /// Create a new worker + pub fn new(db: Arc, shutdown_rx: watch::Receiver) -> Self { + Self { db, shutdown_rx } + } + + /// Run the worker loop + pub async fn run(&mut self) -> Result<()> { + tracing::info!("Worker started"); + + loop { + // Check for shutdown signal + if *self.shutdown_rx.borrow() { + tracing::info!("Worker received shutdown signal"); + break; + } + + // Try to pull a job + match pull_job(&self.db).await { + Ok(Some(job)) => { + tracing::info!("Processing job: {} (kind: {:?})", job.id, job.kind); + if let Err(e) = self.process_job(job).await { + tracing::error!("Error processing job: {}", e); + } + } + Ok(None) => { + // No jobs available, wait a bit before polling again + tokio::select! { + _ = tokio::time::sleep(Duration::from_millis(100)) => {} + _ = self.shutdown_rx.changed() => {} + } + } + Err(e) => { + tracing::error!("Error pulling job: {}", e); + tokio::time::sleep(Duration::from_millis(500)).await; + } + } + } + + tracing::info!("Worker stopped"); + Ok(()) + } + + /// Process a single job + async fn process_job(&self, job: QueuedJob) -> Result<()> { + let started_at = Utc::now(); + + match job.kind { + JobKind::Preview => { + self.process_preview_job(job, started_at).await + } + JobKind::FlowPreview => { + self.process_flow_preview_job(job, started_at).await + } + _ => { + // Unsupported job kind + let error_result = serde_json::json!({ + "error": format!("Unsupported job kind in local mode: {:?}", job.kind) + }); + complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await + } + } + } + + /// Process a script preview job + async fn process_preview_job(&self, job: QueuedJob, started_at: chrono::DateTime) -> Result<()> { + let Some(code) = &job.raw_code else { + let error_result = serde_json::json!({"error": "No code provided for preview"}); + return complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await; + }; + + let Some(lang) = job.script_lang else { + let error_result = serde_json::json!({"error": "No language specified for preview"}); + return complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await; + }; + + // Execute the script + let exec_result = execute_script(lang, code, &job.args).await; + + match exec_result { + Ok(ExecutionResult { success, result, logs }) => { + tracing::debug!("Job {} logs:\n{}", job.id, logs); + let status = if success { JobStatus::Success } else { JobStatus::Failure }; + complete_job(&self.db, job.id, status, result, started_at).await + } + Err(e) => { + let error_result = serde_json::json!({"error": e.to_string()}); + complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await + } + } + } + + /// Process a flow preview job + /// + /// This is a simplified flow executor that handles basic linear flows. + /// A full implementation would need to handle branching, loops, etc. + async fn process_flow_preview_job(&self, job: QueuedJob, started_at: chrono::DateTime) -> Result<()> { + let Some(flow_value) = &job.raw_flow else { + let error_result = serde_json::json!({"error": "No flow definition provided"}); + return complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await; + }; + + // Extract modules from flow value + let modules = flow_value + .get("modules") + .and_then(|m| m.as_array()) + .cloned() + .unwrap_or_default(); + + if modules.is_empty() { + let error_result = serde_json::json!({"error": "Flow has no modules"}); + return complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await; + } + + // Execute modules sequentially (simplified - no branching support) + let mut current_result = job.args.clone(); + let mut flow_status = serde_json::json!({ + "modules": [], + "failure_module": serde_json::Value::Null + }); + + for (idx, module) in modules.iter().enumerate() { + let default_id = format!("module_{}", idx); + let module_id = module + .get("id") + .and_then(|id| id.as_str()) + .unwrap_or(&default_id); + + tracing::info!("Executing flow module: {}", module_id); + + // Update flow status + if let Some(modules_arr) = flow_status.get_mut("modules").and_then(|m| m.as_array_mut()) { + modules_arr.push(serde_json::json!({ + "id": module_id, + "type": "InProgress" + })); + } + + // Extract module value (the actual script/action) + let module_value = module.get("value"); + + match self.execute_flow_module(module_value, ¤t_result).await { + Ok(result) => { + current_result = result; + // Update status to success + if let Some(modules_arr) = flow_status.get_mut("modules").and_then(|m| m.as_array_mut()) { + if let Some(last) = modules_arr.last_mut() { + last["type"] = serde_json::json!("Success"); + last["result"] = current_result.clone(); + } + } + } + Err(e) => { + // Module failed + flow_status["failure_module"] = serde_json::json!({ + "id": module_id, + "error": e.to_string() + }); + let error_result = serde_json::json!({ + "error": e.to_string(), + "flow_status": flow_status + }); + return complete_job(&self.db, job.id, JobStatus::Failure, error_result, started_at).await; + } + } + } + + // Flow completed successfully + let final_result = serde_json::json!({ + "result": current_result, + "flow_status": flow_status + }); + complete_job(&self.db, job.id, JobStatus::Success, final_result, started_at).await + } + + /// Execute a single flow module + async fn execute_flow_module( + &self, + module_value: Option<&serde_json::Value>, + input: &serde_json::Value, + ) -> std::result::Result { + let Some(value) = module_value else { + return Err("Module has no value".to_string()); + }; + + // Check module type + let module_type = value.get("type").and_then(|t| t.as_str()).unwrap_or(""); + + match module_type { + "rawscript" => { + // Inline script + let code = value + .get("content") + .and_then(|c| c.as_str()) + .ok_or("rawscript module missing content")?; + + let lang_str = value + .get("language") + .and_then(|l| l.as_str()) + .unwrap_or("deno"); + + let lang = crate::jobs::ScriptLang::from_str(lang_str) + .ok_or_else(|| format!("Unknown language: {}", lang_str))?; + + let result = execute_script(lang, code, input) + .await + .map_err(|e| e.to_string())?; + + if result.success { + Ok(result.result) + } else { + Err(result.result.to_string()) + } + } + "identity" => { + // Pass through input + Ok(input.clone()) + } + _ => { + Err(format!("Unsupported module type in local mode: {}", module_type)) + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::jobs::{get_completed_job, push_preview, PreviewRequest, ScriptLang}; + + #[tokio::test] + async fn test_worker_processes_bash_preview() { + let db = Arc::new(LocalDb::in_memory().await.unwrap()); + let (shutdown_tx, shutdown_rx) = watch::channel(false); + + // Push a bash preview job + let req = PreviewRequest { + content: "echo 42".to_string(), + language: ScriptLang::Bash, + args: serde_json::json!({}), + lock: None, + tag: None, + }; + let job_id = push_preview(&db, req).await.unwrap(); + + // Create and run worker for one iteration + let mut worker = Worker::new(db.clone(), shutdown_rx); + + // Process one job then shutdown + tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(500)).await; + shutdown_tx.send(true).unwrap(); + }); + + worker.run().await.unwrap(); + + // Check the job completed + let completed = get_completed_job(&db, job_id).await.unwrap(); + assert!(completed.is_some()); + let completed = completed.unwrap(); + assert_eq!(completed.status, JobStatus::Success); + // Output "42" is parsed as JSON number + assert_eq!(completed.result, serde_json::json!(42)); + } +}