Compare commits

..

1 Commits

Author SHA1 Message Date
Ruben Fiszel
6ea6258910 all 2024-10-02 14:14:36 +02:00
318 changed files with 5337 additions and 17078 deletions

View File

@@ -1,6 +1,22 @@
ARG DEBIAN_IMAGE=debian:bookworm-slim
ARG RUST_IMAGE=rust:1.80-slim-bookworm
ARG PYTHON_IMAGE=python:3.11.4-slim-bookworm
FROM ${DEBIAN_IMAGE} as downloader
ARG TARGETPLATFORM
SHELL ["/bin/bash", "-c"]
RUN apt update -y
RUN apt install -y unzip curl
RUN [ "$TARGETPLATFORM" == "linux/amd64" ] && curl -Lsf https://github.com/denoland/deno/releases/download/v1.46.3/deno-x86_64-unknown-linux-gnu.zip -o deno.zip || true
RUN [ "$TARGETPLATFORM" == "linux/arm64" ] && curl -Lsf https://github.com/denoland/deno/releases/download/v1.46.3/deno-aarch64-unknown-linux-gnu.zip -o deno.zip || true
RUN unzip deno.zip && rm deno.zip
FROM ${RUST_IMAGE} as builder
@@ -15,7 +31,7 @@ ENV SQLX_OFFLINE=true
RUN mkdir -p /frontend/build
RUN apt-get update \
&& apt-get install -y ca-certificates tzdata libpq5 cmake unzip\
&& apt-get install -y ca-certificates tzdata libpq5 cmake\
make build-essential libssl-dev zlib1g-dev libbz2-dev libreadline-dev \
libsqlite3-dev wget curl llvm libncurses5-dev libncursesw5-dev xz-utils tk-dev libxml2-dev \
libxmlsec1-dev libffi-dev liblzma-dev mecab-ipadic-utf8 libgdbm-dev libc6-dev git libprotobuf-dev libnl-route-3-dev \
@@ -27,9 +43,6 @@ RUN wget https://golang.org/dl/go1.21.5.linux-amd64.tar.gz && tar -C /usr/local
ENV PATH="${PATH}:/usr/local/go/bin"
ENV GO_PATH=/usr/local/go/bin/go
# Install UV
RUN curl --proto '=https' --tlsv1.2 -LsSf https://github.com/astral-sh/uv/releases/download/0.4.18/uv-installer.sh | sh && mv /usr/local/cargo/bin/uv /usr/local/bin/uv
ENV TZ=Etc/UTC
ENV PYTHON_VERSION 3.11.4
@@ -40,14 +53,13 @@ RUN wget https://www.python.org/ftp/python/${PYTHON_VERSION}/Python-${PYTHON_VER
RUN /usr/local/bin/python3 -m pip install pip-tools
COPY --from=oven/bun:1.1.30 /usr/local/bin/bun /usr/bin/bun
COPY --from=oven/bun:1.1.27 /usr/local/bin/bun /usr/bin/bun
ARG TARGETPLATFORM
RUN curl -Lsf https://github.com/denoland/deno/releases/download/v2.0.0/deno-x86_64-unknown-linux-gnu.zip -o deno.zip
# RUN [ "$TARGETPLATFORM" == "linux/arm64" ] && curl -Lsf https://github.com/denoland/deno/releases/download/v2.0.0/deno-aarch64-unknown-linux-gnu.zip -o deno.zip || true
RUN [ "$TARGETPLATFORM" == "linux/amd64" ] && curl -Lsf https://github.com/denoland/deno/releases/download/v1.41.0/deno-x86_64-unknown-linux-gnu.zip -o deno.zip || true
RUN [ "$TARGETPLATFORM" == "linux/arm64" ] && curl -Lsf https://github.com/denoland/deno/releases/download/v1.41.0/deno-aarch64-unknown-linux-gnu.zip -o deno.zip || true
RUN unzip deno.zip && rm deno.zip && mv deno /usr/bin/deno
COPY --from=downloader --chmod=755 /deno /usr/bin/deno
RUN apt-get update \
&& apt-get install -y postgresql-client --allow-unauthenticated

View File

@@ -41,12 +41,8 @@ jobs:
- name: cargo test
timeout-minutes: 15
run:
/usr/bin/deno --version &&
/usr/bin/bun -v &&
go version &&
/usr/local/bin/python3 --version &&
mkdir frontend/build && cd backend && touch
windmill-api/openapi-deref.yaml &&
DATABASE_URL=postgres://postgres:changeme@postgres:5432/windmill
DISABLE_EMBEDDING=true RUST_LOG=info cargo test --features
enterprise,deno_core --all -- --nocapture
DISABLE_EMBEDDING=true RUST_LOG=info cargo test --features enterprise
--all -- --nocapture

View File

@@ -1,128 +0,0 @@
env:
REGISTRY: ghcr.io
IMAGE_NAME: ${{ github.repository }}
name: Build and publish windmill for RHEL9
on:
workflow_dispatch
permissions: write-all
jobs:
build_ee:
runs-on: ubicloud
steps:
- uses: actions/checkout@v4
with:
fetch-depth: 0
- name: Read EE repo commit hash
run: |
echo "ee_repo_ref=$(cat ./backend/ee-repo-ref.txt)" >> "$GITHUB_ENV"
- uses: actions/checkout@v4
with:
repository: windmill-labs/windmill-ee-private
path: ./windmill-ee-private
ref: ${{ env.ee_repo_ref }}
token: ${{ secrets.WINDMILL_EE_PRIVATE_ACCESS }}
fetch-depth: 0
# - name: Set up Docker Buildx
# uses: docker/setup-buildx-action@v2
- uses: depot/setup-action@v1
- name: Docker meta
id: meta-ee-public
uses: docker/metadata-action@v5
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee-rhel9
flavor: |
latest=false
tags: |
type=sha
- name: Login to registry
uses: docker/login-action@v3
with:
registry: ${{ env.REGISTRY }}
username: ${{ github.actor }}
password: ${{ secrets.GITHUB_TOKEN }}
- name: Substitute EE code
run: |
./backend/substitute_ee_code.sh --copy --dir ./windmill-ee-private
- name: Copy RHEL9 Dockerfile
run: |
cp ./docker/RHEL9/Dockerfile ./Dockerfile
- name: Build and push publicly ee amd64
uses: depot/build-push-action@v1
with:
context: .
platforms: linux/amd64
push: true
build-args: |
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core
secrets: |
rh_username=${{ secrets.RH_USERNAME }}
rh_password=${{ secrets.RH_PASSWORD }}
tags: |
${{ steps.meta-ee-public.outputs.tags }}-amd64
labels: |
${{ steps.meta-ee-public.outputs.labels }}-amd64
org.opencontainers.image.licenses=Windmill-Enterprise-License
- name: Build and push publicly ee arm64
uses: depot/build-push-action@v1
with:
context: .
platforms: linux/arm64
push: true
build-args: |
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core
secrets: |
rh_username=${{ secrets.RH_USERNAME }}
rh_password=${{ secrets.RH_PASSWORD }}
tags: |
${{ steps.meta-ee-public.outputs.tags }}-arm64
labels: |
${{ steps.meta-ee-public.outputs.labels }}-arm64
org.opencontainers.image.licenses=Windmill-Enterprise-License
- uses: shrink/actions-docker-extract@v3
id: extract-ee-amd64
with:
image: ${{ steps.meta-ee-public.outputs.tags}}-amd64
path: "/windmill/target/release/windmill"
- uses: shrink/actions-docker-extract@v3
id: extract-ee-arm64
with:
image: ${{ steps.meta-ee-public.outputs.tags}}-arm64
path: "/windmill/target/release/windmill"
- name: Rename binary with corresponding architecture
run: |
mv "${{ steps.extract-ee-amd64.outputs.destination }}/windmill" "${{ steps.extract-ee-amd64.outputs.destination }}/windmill-ee-amd64-rhel9"
mv "${{ steps.extract-ee-arm64.outputs.destination }}/windmill" "${{ steps.extract-ee-arm64.outputs.destination }}/windmill-ee-arm64-rhel9"
- uses: actions/upload-artifact@v4
with:
name: RHEL9-amd64 build
path: ${{ steps.extract-ee-amd64.outputs.destination }}/windmill-ee-amd64-rhel9
- uses: actions/upload-artifact@v4
with:
name: RHEL9-arm64 build
path: ${{ steps.extract-ee-arm64.outputs.destination }}/windmill-ee-arm64-rhel9
# - name: Attach binary to release
# uses: softprops/action-gh-release@v2
# if: startsWith(github.ref, 'refs/tags/')
# with:
# files: |
# ${{ steps.extract-ee-arm64.outputs.destination }}/windmill-ee-arm64-rhel9
# ${{ steps.extract-ee-amd64.outputs.destination }}/windmill-ee-amd64-rhel9

View File

@@ -62,7 +62,7 @@ jobs:
platforms: linux/amd64,linux/arm64
push: true
build-args: |
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc
tags: |
${{ steps.meta-ee-public.outputs.tags }}
labels: |

View File

@@ -1,60 +0,0 @@
name: Build and Publish Windows Worker
on:
push:
tags:
- "v*"
env:
CARGO_INCREMENTAL: 0
SQLX_OFFLINE: true
DISABLE_EMBEDDING: true
RUST_LOG: info
jobs:
cargo_build_windows:
runs-on: windows-latest
steps:
- uses: actions/checkout@v4
- name: Read EE repo commit hash
shell: pwsh
run: |
$ee_repo_ref = Get-Content .\backend\ee-repo-ref.txt
echo "ee_repo_ref=$ee_repo_ref" | Out-File -FilePath $env:GITHUB_ENV -Append
- name: Checkout windmill-ee-private repository
uses: actions/checkout@v4
with:
repository: windmill-labs/windmill-ee-private
path: ./windmill-ee-private
ref: ${{ env.ee_repo_ref }}
token: ${{ secrets.WINDMILL_EE_PRIVATE_ACCESS }}
fetch-depth: 0
- name: Substitute EE code
shell: bash
run: |
./backend/substitute_ee_code.sh --copy --dir ./windmill-ee-private
- name: Cargo build windows
timeout-minutes: 90
run: |
vcpkg.exe install openssl-windows:x64-windows
vcpkg.exe install openssl:x64-windows-static
vcpkg.exe integrate install
$env:VCPKGRS_DYNAMIC=1
$env:OPENSSL_DIR="${Env:VCPKG_INSTALLATION_ROOT}\installed\x64-windows-static"
mkdir frontend/build && cd backend
New-Item -Path . -Name "windmill-api/openapi-deref.yaml" -ItemType "File" -Force
cargo build --release --features=enterprise,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core
- name: Rename binary with corresponding architecture
run: |
Rename-Item -Path ".\backend\target\release\windmill.exe" -NewName "windmill-ee.exe"
- name: Attach binary to release
uses: softprops/action-gh-release@v2
with:
files: |
./backend/target/release/windmill-ee.exe

View File

@@ -67,7 +67,7 @@ jobs:
platforms: linux/amd64,linux/arm64
push: true
build-args: |
features=embedding,parquet,openidconnect,deno_core
features=embedding,parquet,openidconnect
tags: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}:dev
${{ steps.meta-public.outputs.tags }}

View File

@@ -1,10 +1,8 @@
env:
REGISTRY: ghcr.io
IMAGE_NAME:
${{ github.event_name != 'pull_request' && github.repository ||
IMAGE_NAME: ${{ github.event_name != 'pull_request' && github.repository ||
'windmill-labs/windmill-test' }}
DEV_SHA:
${{ github.event_name != 'pull_request' && 'dev' || format('pr-{0}',
DEV_SHA: ${{ github.event_name != 'pull_request' && 'dev' || format('pr-{0}',
github.event.number) }}
name: Build windmill:main
@@ -77,7 +75,7 @@ jobs:
platforms: linux/amd64,linux/arm64
push: true
build-args: |
features=embedding,parquet,openidconnect,jemalloc,deno_core
features=embedding,parquet,openidconnect,jemalloc
tags: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}:${{ env.DEV_SHA }}
${{ steps.meta-public.outputs.tags }}
@@ -138,7 +136,7 @@ jobs:
platforms: linux/amd64,linux/arm64
push: true
build-args: |
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy
tags: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee:${{ env.DEV_SHA }}
${{ steps.meta-ee-public.outputs.tags }}
@@ -200,7 +198,7 @@ jobs:
platforms: linux/amd64
push: true
build-args: |
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core
features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy
PYTHON_IMAGE=python:3.12.2-slim-bookworm
tags: |
${{ steps.meta-ee-public-py312.outputs.tags }}
@@ -394,7 +392,7 @@ jobs:
verify_ee_image_vulnerabilities:
runs-on: ubicloud
needs: [tag_latest_ee]
if: ${{ startsWith(github.ref, 'refs/tags/') }}
# if: ${{ startsWith(github.ref, 'refs/tags/') }}
steps:
- name: Checkout code
uses: actions/checkout@v4
@@ -589,6 +587,8 @@ jobs:
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee-cuda
flavor: |
latest=false
tags: |
type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}}
@@ -633,6 +633,8 @@ jobs:
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-slim
flavor: |
latest=false
tags: |
type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}}
@@ -676,6 +678,8 @@ jobs:
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee-slim
flavor: |
latest=false
tags: |
type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}}
@@ -720,6 +724,8 @@ jobs:
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-full
flavor: |
latest=false
tags: |
type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}}
@@ -763,6 +769,8 @@ jobs:
with:
images: |
${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee-full
flavor: |
latest=false
tags: |
type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}}

View File

@@ -1,173 +1,5 @@
# Changelog
## [1.409.2](https://github.com/windmill-labs/windmill/compare/v1.409.1...v1.409.2) (2024-10-16)
### Bug Fixes
* add extra args support for exception to bun scripts ([1466da3](https://github.com/windmill-labs/windmill/commit/1466da3999add0238b9c42ca13df52194f082fc0))
* fix script persistence in url + add support for extra error args in python ([3174024](https://github.com/windmill-labs/windmill/commit/3174024d8e6ecbe9f8c9e1ea055d611f652c0057))
## [1.409.1](https://github.com/windmill-labs/windmill/compare/v1.409.0...v1.409.1) (2024-10-16)
### Bug Fixes
* **apidocs:** fix generated openapi files ([d24e153](https://github.com/windmill-labs/windmill/commit/d24e1530655d27ffda9bb4c19471dc1431124cc2))
* **git-sync:** propagate update of folders with git sync ([6abb346](https://github.com/windmill-labs/windmill/commit/6abb346013da4a907a860713a8a67642985b8025))
## [1.409.0](https://github.com/windmill-labs/windmill/compare/v1.408.1...v1.409.0) (2024-10-16)
### Features
* **frontend:** unify all triggers UX and simplify flow settings ([#4259](https://github.com/windmill-labs/windmill/issues/4259)) ([91a3d06](https://github.com/windmill-labs/windmill/commit/91a3d065298cce7a882464fa0cd31d8f1ae9dda2))
* Scroll to element in virtual list when clicking on graph point ([#4532](https://github.com/windmill-labs/windmill/issues/4532)) ([7126ba1](https://github.com/windmill-labs/windmill/commit/7126ba12c7eb52d2cfbe8d83311b5592a5707bce))
* **sso:** adding the ability to define a custom display name for sso ([#4529](https://github.com/windmill-labs/windmill/issues/4529)) ([99c5b3e](https://github.com/windmill-labs/windmill/commit/99c5b3ecdacb1158c2cba5b891c4c3b8b70c3b6a))
### Bug Fixes
* Add indexer backup lock to fit the deployment model ([#4531](https://github.com/windmill-labs/windmill/issues/4531)) ([411bce7](https://github.com/windmill-labs/windmill/commit/411bce7e13aabfa53db80d55261ea4511f6d1ae9))
* **app:** accept connecting to non yet existing state output for convenience ([9eb1ecc](https://github.com/windmill-labs/windmill/commit/9eb1ecc9f3017e2f4284a827a2bcd21f36b1b8ac))
* **app:** improve absolute url handling in download button and downloadFile ([dcdbf1a](https://github.com/windmill-labs/windmill/commit/dcdbf1afb4d5a18e00b9bbb1eb0bef129ea5667f))
* **app:** make s3 uploads persistent across tabs change ([c3b536b](https://github.com/windmill-labs/windmill/commit/c3b536b1b8069898131768867a187b186b21e537))
* canceled jobs button reporting 0 jobs cancelled ([#4534](https://github.com/windmill-labs/windmill/issues/4534)) ([e736572](https://github.com/windmill-labs/windmill/commit/e736572db10929ae5e123c4f6cef73b8e90fc29b))
* **python-client:** improve get_job_status for running jobs ([a8c4ea2](https://github.com/windmill-labs/windmill/commit/a8c4ea2334d2535fa7d5d43f58d65565afe8f3e5))
* **ui:** dark mode support for queue metrics based critical alert ([#4535](https://github.com/windmill-labs/windmill/issues/4535)) ([f38b3d1](https://github.com/windmill-labs/windmill/commit/f38b3d14e8092ae58817511aea91a4e77725ead6))
## [1.408.1](https://github.com/windmill-labs/windmill/compare/v1.408.0...v1.408.1) (2024-10-12)
### Bug Fixes
* fix deno cache --allow-import on deno 2 ([42fe31f](https://github.com/windmill-labs/windmill/commit/42fe31f804c9e6643cd90167494e45270831e013))
## [1.408.0](https://github.com/windmill-labs/windmill/compare/v1.407.2...v1.408.0) (2024-10-12)
### Features
* **app builder:** file download helper ([#4511](https://github.com/windmill-labs/windmill/issues/4511)) ([f82f091](https://github.com/windmill-labs/windmill/commit/f82f09129096cfff975370d8bb7b6d832a2b8f9f))
### Bug Fixes
* **cli:** handle case where 'toString' is a schema field ([568cc66](https://github.com/windmill-labs/windmill/commit/568cc66932fb0470f5e89de7b02d94dba4050638))
* **frontend:** s3 file uploader works on public apps too ([982dde2](https://github.com/windmill-labs/windmill/commit/982dde2b9dfe6d9eda300c683af729d97a03cb4d))
* **frontend:** set unused schema property fields to null ([be11240](https://github.com/windmill-labs/windmill/commit/be112408e7c4601314726e9517c37daaeaa1bf09))
* improve workflow as code row-lock on db to handle more concurrency ([d2c4d3f](https://github.com/windmill-labs/windmill/commit/d2c4d3fa207379cb0b8ac180f6ccc759045580e8))
## [1.407.2](https://github.com/windmill-labs/windmill/compare/v1.407.1...v1.407.2) (2024-10-10)
### Bug Fixes
* improve default properties of new nodes of flows (suspend, branchone, branchall) ([d9bdc5a](https://github.com/windmill-labs/windmill/commit/d9bdc5a5b08dd4d0381304656af097315398c9d4))
## [1.407.1](https://github.com/windmill-labs/windmill/compare/v1.407.0...v1.407.1) (2024-10-10)
### Bug Fixes
* improve handling of empty lock files on deno 2.0 ([7ca5bf2](https://github.com/windmill-labs/windmill/commit/7ca5bf2faeff44a7543b1afa9369c140fcb71dfc))
## [1.407.0](https://github.com/windmill-labs/windmill/compare/v1.406.0...v1.407.0) (2024-10-10)
### Features
* upgrade to deno 2 ([26b11a0](https://github.com/windmill-labs/windmill/commit/26b11a00150acbe101abe4bb542f24379da0cc56))
### Bug Fixes
* update internal deno runtime to latest (deno 2.0) ([c3a5736](https://github.com/windmill-labs/windmill/commit/c3a57366419882ea2de1938bea592c795b1a1d03))
## [1.406.0](https://github.com/windmill-labs/windmill/compare/v1.405.5...v1.406.0) (2024-10-09)
### Features
* **frontend:** components can be moved inside containers by holding ctrl/cmd ([111bfc6](https://github.com/windmill-labs/windmill/commit/111bfc6a659037ae7029e8f557256e2fffcf979b))
* **monitoring:** Critical Alerts for Jobs Waiting in Queue [enterprise] ([#4491](https://github.com/windmill-labs/windmill/issues/4491)) ([d90d6c2](https://github.com/windmill-labs/windmill/commit/d90d6c2b896c5f99e00681656f376b180901f272))
### Bug Fixes
* **cli:** instance sync push does not require sync pull ([257f097](https://github.com/windmill-labs/windmill/commit/257f0971f86938da71b879f32d93473976eaa920))
* remove monaco-editor for app preview code path for faster app loads ([7b05033](https://github.com/windmill-labs/windmill/commit/7b0503332d1bdd7f5999a5ef99150f8c9f6f18be))
## [1.405.5](https://github.com/windmill-labs/windmill/compare/v1.405.4...v1.405.5) (2024-10-04)
### Bug Fixes
* windows.exe build with github workflow doesn't have openssl.dll bundled in ([#4489](https://github.com/windmill-labs/windmill/issues/4489)) ([284cb40](https://github.com/windmill-labs/windmill/commit/284cb4069c97efe59b5caf3effb68c8b30e02b73))
## [1.405.4](https://github.com/windmill-labs/windmill/compare/v1.405.3...v1.405.4) (2024-10-04)
### Bug Fixes
* **frontend:** correctly initialize step inputs on new inline script ([289ad51](https://github.com/windmill-labs/windmill/commit/289ad51374f0344582372572fc521f3b2bf12b33))
## [1.405.3](https://github.com/windmill-labs/windmill/compare/v1.405.2...v1.405.3) (2024-10-04)
### Bug Fixes
* fix id save on apps ([b034b07](https://github.com/windmill-labs/windmill/commit/b034b070c075a0fab74678bec1fd8b829d55b204))
## [1.405.2](https://github.com/windmill-labs/windmill/compare/v1.405.1...v1.405.2) (2024-10-03)
### Bug Fixes
* **cli:** fix opts.yes for instance sync ([26659ce](https://github.com/windmill-labs/windmill/commit/26659ce37d2887d5b98dbdbdbba27bab85d4fe3f))
* fix uv path ([19c62ba](https://github.com/windmill-labs/windmill/commit/19c62ba195b1df85c38c748dab7d9f137696a5c3))
## [1.405.1](https://github.com/windmill-labs/windmill/compare/v1.405.0...v1.405.1) (2024-10-03)
### Bug Fixes
* flow picker of flows + precache hub scripts as bundles ([c84e6fd](https://github.com/windmill-labs/windmill/commit/c84e6fd05de2bea426cae61fa25db0323b8770f5))
## [1.405.0](https://github.com/windmill-labs/windmill/compare/v1.404.1...v1.405.0) (2024-10-03)
### Features
* Replace `pip-compile` with `uv` ([#4460](https://github.com/windmill-labs/windmill/issues/4460)) ([b54c9ee](https://github.com/windmill-labs/windmill/commit/b54c9ee657cc88fabe694cae39dc0d3c1918fcbb))
* **worker:** support workers to run natively on windows ([#4446](https://github.com/windmill-labs/windmill/issues/4446)) ([f5c4727](https://github.com/windmill-labs/windmill/commit/f5c472727465dd95f5378bc08ee9bbb983f4d259))
### Bug Fixes
* **cli:** fix set client of instance when passing token and base url ([794c4cd](https://github.com/windmill-labs/windmill/commit/794c4cde3cd47042472dccdf4b60a012014dd26d))
## [1.404.1](https://github.com/windmill-labs/windmill/compare/v1.404.0...v1.404.1) (2024-10-03)
### Bug Fixes
* flow picker of flows ([92f61f0](https://github.com/windmill-labs/windmill/commit/92f61f07ed6d354407d26843e3a270b95bae90bc))
## [1.404.0](https://github.com/windmill-labs/windmill/compare/v1.403.1...v1.404.0) (2024-10-03)
### Features
* **frontend:** add quick access menu in flow editor ([#4415](https://github.com/windmill-labs/windmill/issues/4415)) ([45ccd45](https://github.com/windmill-labs/windmill/commit/45ccd45e306c66931880a9b8fd48bfe684c774ac))
### Bug Fixes
* **cli:** improve schedule path handling on windows ([9ac3b6b](https://github.com/windmill-labs/windmill/commit/9ac3b6b1d5d64d7467dd80506f8a8d772c4630bd))
* fix id editor for app ([8e58e43](https://github.com/windmill-labs/windmill/commit/8e58e4320a31d71c40a5ed352416a4c2dd3adb26))
* **frontend:** disable runnable field on route editor from detail panel ([#4469](https://github.com/windmill-labs/windmill/issues/4469)) ([3134f79](https://github.com/windmill-labs/windmill/commit/3134f79ced80aab86912643ab7a60dcf909ab104))
## [1.403.1](https://github.com/windmill-labs/windmill/compare/v1.403.0...v1.403.1) (2024-10-01)

View File

@@ -158,9 +158,6 @@ RUN set -eux; \
ENV PATH="${PATH}:/usr/local/go/bin"
ENV GO_PATH=/usr/local/go/bin/go
# Install UV
RUN curl --proto '=https' --tlsv1.2 -LsSf https://github.com/astral-sh/uv/releases/download/0.4.18/uv-installer.sh | sh && mv /root/.cargo/bin/uv /usr/local/bin/uv
RUN curl -sL https://deb.nodesource.com/setup_20.x | bash -
RUN apt-get -y update && apt-get install -y curl nodejs awscli && apt-get clean \
&& rm -rf /var/lib/apt/lists/*
@@ -175,9 +172,9 @@ RUN /usr/local/bin/python3 -m pip install pip-tools
COPY --from=builder /frontend/build /static_frontend
COPY --from=builder /windmill/target/release/windmill ${APP}/windmill
COPY --from=denoland/deno:2.0.0 --chmod=755 /usr/bin/deno /usr/bin/deno
COPY --from=denoland/deno:1.46.3 --chmod=755 /usr/bin/deno /usr/bin/deno
COPY --from=oven/bun:1.1.30 /usr/local/bin/bun /usr/bin/bun
COPY --from=oven/bun:1.1.27 /usr/local/bin/bun /usr/bin/bun
COPY --from=php:8.3.7-cli /usr/local/bin/php /usr/bin/php
COPY --from=composer:2.7.6 /usr/bin/composer /usr/bin/composer

View File

@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM healthchecks WHERE check_type = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text"
]
},
"nullable": []
},
"hash": "0ee63ef2dd5c88edba2a1f56d31f29876724922f148dad7af35b36efbf70207a"
}

View File

@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO concurrency_locks (id, last_locked_at, owner)\n VALUES ($1, now(), $2)\n ON CONFLICT (id)\n DO UPDATE SET\n last_locked_at = now(),\n owner = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar"
]
},
"nullable": []
},
"hash": "14abf759dae7ba5c38017ba6001927c6df0653a02b87bcea939066e39ebcf24d"
}

View File

@@ -1,24 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM schedule WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
null
]
},
"hash": "1ca5bc2d35c0498b587fd0618434def64233dc4f8fc3344d8d74be8e96ded659"
}

View File

@@ -1,12 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM healthchecks WHERE healthy = true AND created_at < NOW() - INTERVAL '14 days'",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "2041526bc58872d71f91f7698144039bd67f8e37895befa94a15b7e4019e114b"
}

View File

@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO healthchecks (check_type, healthy) VALUES ($1, false)",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar"
]
},
"nullable": []
},
"hash": "27920aaa55666ffc14a36a247f89ff7994ee40d3953b9f772d0e0ab999bccb7b"
}

View File

@@ -1,24 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM http_trigger WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
null
]
},
"hash": "31bc3dcea29be9cc0242771d25a232f173446d29c08fc29ddb8d55294f2c070e"
}

View File

@@ -1,22 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT created_at FROM healthchecks WHERE check_type = $1 ORDER BY created_at DESC LIMIT 1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "created_at",
"type_info": "Timestamptz"
}
],
"parameters": {
"Left": [
"Text"
]
},
"nullable": [
false
]
},
"hash": "34a45763bb4d14162f4cd3fa07cd8020f1f6085f4ee85f5eab3458637edf26cd"
}

View File

@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO token\n (token, email, label, expiration, super_admin, scopes, workspace_id)\n VALUES ($1, $2, $3, $4, $5, $6, $7)",
"query": "INSERT INTO token\n (token, email, label, expiration, super_admin, scopes)\n VALUES ($1, $2, $3, $4, $5, $6)",
"describe": {
"columns": [],
"parameters": {
@@ -10,11 +10,10 @@
"Varchar",
"Timestamptz",
"Bool",
"TextArray",
"Varchar"
"TextArray"
]
},
"nullable": []
},
"hash": "c624f15f3e321b1eecf123da9bf0b18e8c1d16ef25ffb9d04e5447d0d583d55c"
"hash": "34ad8a2a5bd89b9b8e25847a7e5e94ef99e35a178ad6328c1bcde2a6d6f88cb5"
}

View File

@@ -1,29 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT COUNT(*) as count, \n MIN(scheduled_for) as oldest_job\n FROM queue \n WHERE tag = $1 \n AND scheduled_for <= NOW() - $2::interval \n AND running = false\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
},
{
"ordinal": 1,
"name": "oldest_job",
"type_info": "Timestamptz"
}
],
"parameters": {
"Left": [
"Text",
"Interval"
]
},
"nullable": [
null,
null
]
},
"hash": "3ecb25b05d6c14b499f9b00af42ae74134728899f6b59c68b246979bc5143e30"
}

View File

@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE concurrency_locks SET\n last_locked_at = now()\n WHERE id = $1 AND owner = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": []
},
"hash": "57e270e032e8c04dda7b5c1ca949861756b3ad367a4a500728332a7cb91560a4"
}

View File

@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "5fc6b4a4dbb7875bdec76f876c18543435a95b019b20081f52f6ed6f4457e3c7"
}

View File

@@ -1,22 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT owner FROM concurrency_locks WHERE id = $1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "owner",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text"
]
},
"nullable": [
true
]
},
"hash": "5fd70c70ce52cbc51fa9124cb05f82b5951f17d1b7eade53c6d89253d55f8b9f"
}

View File

@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO windmill_migrations (name) VALUES ('bypassrls_1-2')",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "722a3096f03d25ef94292d53801d41037de4bc69dd434232029c731cbbcbc22f"
}

View File

@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO concurrency_locks (id, last_locked_at) VALUES ($1, NOW()) ON CONFLICT (id) DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar"
]
},
"nullable": []
},
"hash": "900ac59515e4283f4b57516210575dfe92f74a7220ed69e61899a6e0f053d9cd"
}

View File

@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM token WHERE label LIKE 'email-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "a7f5431e3b8960e9dc46fae69dd4391516d8b169186548ec44528c84078b80d8"
}

View File

@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE healthchecks SET healthy = true WHERE check_type = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text"
]
},
"nullable": []
},
"hash": "ad42118ccf6a9d2d1e072c4df064ddf964a5b3cd088fc162d0d8222325d4a5ea"
}

View File

@@ -18,8 +18,8 @@
"Left": []
},
"nullable": [
false,
true
true,
false
]
},
"hash": "b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76"

View File

@@ -52,8 +52,7 @@
"trigger",
"failure",
"command",
"approval",
"preprocessor"
"approval"
]
}
}

View File

@@ -1,24 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT schedule FROM schedule WHERE path = $1 AND script_path = $1 AND is_flow = $2 AND workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "schedule",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
false
]
},
"hash": "c060b8bbc5af7d2e7d0aaff64f0f62ec9db58611a99b0ba7f0375638b128ab89"
}

View File

@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE concurrency_locks SET last_locked_at = NOW() WHERE id = $1 AND last_locked_at < NOW() - INTERVAL '1 second' * $2 RETURNING 1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "?column?",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Text",
"Float8"
]
},
"nullable": [
null
]
},
"hash": "c4e1873bfc7b905e7299a021f4baa2a97e95f4797c5e11f37822e19828422b7e"
}

View File

@@ -1,59 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT label, concat(substring(token for 10)) as token_prefix, expiration, created_at, last_used_at, scopes, email FROM token WHERE workspace_id = $1 AND scopes @> ARRAY['run:script/' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "label",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "token_prefix",
"type_info": "Text"
},
{
"ordinal": 2,
"name": "expiration",
"type_info": "Timestamptz"
},
{
"ordinal": 3,
"name": "created_at",
"type_info": "Timestamptz"
},
{
"ordinal": 4,
"name": "last_used_at",
"type_info": "Timestamptz"
},
{
"ordinal": 5,
"name": "scopes",
"type_info": "TextArray"
},
{
"ordinal": 6,
"name": "email",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
true,
null,
true,
false,
false,
true,
true
]
},
"hash": "c7ee7ce64686cef41cebd99ad7ef31572fc1bf12e6ae473fd58fafb025989965"
}

View File

@@ -1,22 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT last_locked_at\n FROM concurrency_locks\n WHERE id = $1\n FOR UPDATE",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "last_locked_at",
"type_info": "Timestamp"
}
],
"parameters": {
"Left": [
"Text"
]
},
"nullable": [
false
]
},
"hash": "cecf1addc4aecb087a14786b2a9165895ca61ef042947c7314f66514d7f29edc"
}

View File

@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "d7a0f19f9e18d2ea49316012375ad78b69292ba091d69880945e42bebe890d66"
}

View File

@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT EXISTS(SELECT 1 FROM healthchecks WHERE check_type = $1 AND healthy = false)",
"query": "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'bypassrls_1-2')",
"describe": {
"columns": [
{
@@ -10,13 +10,11 @@
}
],
"parameters": {
"Left": [
"Text"
]
"Left": []
},
"nullable": [
null
]
},
"hash": "eb932b613a6dbb2cdff97e5512d42b538ba83115c0ea798be00b01659600f45a"
"hash": "eb1f916f9beea3eea83ce359f5305d0cfb0d6cdba9cc56c6139f57e46345f843"
}

View File

@@ -1,59 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT label, concat(substring(token for 10)) as token_prefix, expiration, created_at, last_used_at, scopes, email FROM token WHERE workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "label",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "token_prefix",
"type_info": "Text"
},
{
"ordinal": 2,
"name": "expiration",
"type_info": "Timestamptz"
},
{
"ordinal": 3,
"name": "created_at",
"type_info": "Timestamptz"
},
{
"ordinal": 4,
"name": "last_used_at",
"type_info": "Timestamptz"
},
{
"ordinal": 5,
"name": "scopes",
"type_info": "TextArray"
},
{
"ordinal": 6,
"name": "email",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
true,
null,
true,
false,
false,
true,
true
]
},
"hash": "eff32aeac25a75d06f73e08c26dd3fd25f6b85cbea870505751c6a82457ae1da"
}

View File

@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM token WHERE label LIKE 'email-%' AND workspace_id = $1 AND scopes @> ARRAY['run:script/' || $2]::text[]",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "f06e0e4fa358b26792df22fff48b71a6fcfa1e5603ea472892917c1accd1aafb"
}

948
backend/Cargo.lock generated

File diff suppressed because it is too large Load Diff

View File

@@ -1,6 +1,6 @@
[package]
name = "windmill"
version = "1.409.2"
version = "1.403.1"
authors.workspace = true
edition.workspace = true
@@ -27,7 +27,7 @@ members = [
]
[workspace.package]
version = "1.409.2"
version = "1.403.1"
authors = ["Ruben Fiszel <ruben@windmill.dev>"]
edition = "2021"
@@ -39,9 +39,6 @@ path = "./src/main.rs"
opt-level = 0
incremental = true
[profile.release]
lto = "thin"
[features]
default = []
enterprise = ["windmill-worker/enterprise", "windmill-queue/enterprise", "windmill-api/enterprise", "windmill-git-sync/enterprise", "windmill-common/prometheus", "windmill-common/enterprise", "windmill-indexer/enterprise"]
@@ -50,16 +47,16 @@ stripe = ["windmill-api/stripe"]
benchmark = ["windmill-api/benchmark", "windmill-worker/benchmark", "windmill-queue/benchmark", "windmill-common/benchmark"]
flamegraph = ["windmill-common/flamegraph", "windmill-worker/flamegraph"]
loki = ["windmill-common/loki"]
pg_embed = ["dep:pg-embed"]
embedding = ["windmill-api/embedding"]
parquet = ["windmill-api/parquet", "windmill-common/parquet", "windmill-worker/parquet", "windmill-indexer/parquet", "dep:object_store"]
prometheus = ["windmill-common/prometheus", "windmill-api/prometheus", "windmill-worker/prometheus", "windmill-queue/prometheus"]
flow_testing = ["windmill-worker/flow_testing"]
openidconnect = ["windmill-api/openidconnect", "windmill-common/openidconnect"]
openidconnect = ["windmill-api/openidconnect"]
cloud = ["windmill-queue/cloud", "windmill-worker/cloud"]
jemalloc = ["windmill-common/jemalloc", "dep:tikv-jemallocator", "dep:tikv-jemalloc-sys", "dep:tikv-jemalloc-ctl"]
tantivy = ["dep:windmill-indexer", "windmill-api/tantivy"]
sqlx = ["windmill-worker/sqlx"]
deno_core = ["windmill-worker/deno_core", "dep:deno_core"]
[dependencies]
anyhow.workspace = true
@@ -88,8 +85,9 @@ uuid.workspace = true
gethostname.workspace = true
serde_json.workspace = true
serde.workspace = true
deno_core = { workspace = true, optional = true }
deno_core.workspace = true
object_store = { workspace = true, optional = true }
pg-embed = {git = "https://github.com/faokunega/pg-embed", optional = true, default-features = false, features = ['rt_tokio']}
quote.workspace = true
@@ -107,7 +105,6 @@ serde.workspace = true
windmill-api-client.workspace = true
deno_core = { workspace = true, features = ["include_js_files_for_snapshotting", "unsafe_use_unprotected_platform"] }
[workspace.dependencies]
windmill-api = { path = "./windmill-api", default-features = false }
windmill-queue = { path = "./windmill-queue" }
@@ -175,24 +172,20 @@ tokio-util = { version = "^0", features = ["io"] }
json-pointer = "^0"
itertools = "^0"
regex = "^1"
deno_fetch = "0.195.0"
deno_tls = "0.158.0"
deno_console = "0.171.0"
deno_url = "0.171.0"
deno_webidl = "0.171.0"
deno_web = "0.202.0"
deno_net = "0.163.0"
deno_core = "0.311.0"
deno_ast = { version = "=0.42.2", features = ["transpiling"] }
swc_common = "=0.37.5"
swc_ecma_parser = "=0.149.1"
swc_ecma_ast = "=0.118.2"
swc_ecma_visit = "=0.104.8"
deno_fetch = "0.187.0"
deno_tls = "0.150.0"
deno_console = "0.163.0"
deno_url = "0.163.0"
deno_webidl = "0.163.0"
deno_web = "0.194.0"
deno_net = "0.155.0"
deno_core = "0.299.0"
deno_ast = { version = "=0.40.0", features = ["transpiling"] }
async-recursion = "^1"
swc_common = "=0.33.26"
swc_ecma_parser = "=0.144.3"
swc_ecma_ast = "=0.113.7"
swc_ecma_visit = "=0.99.1"
base64 = "0.21.0"
base32 = "^0"
hmac = "0.12.1"
@@ -277,7 +270,8 @@ tikv-jemallocator = { version = "0.5" }
tikv-jemalloc-sys = { version = "^0.5" }
tikv-jemalloc-ctl = { version = "^0.5" }
triomphe = "^0"
# 0.1.12 broken (nested dependency of swc_common)
triomphe = "<0.1.12"
tantivy = "0.22.0"

View File

@@ -0,0 +1,14 @@
CREATE POLICY admin_policy ON account TO windmill_admin USING (true);
CREATE POLICY admin_policy ON app TO windmill_admin USING (true);
CREATE POLICY admin_policy ON audit TO windmill_admin USING (true);
CREATE POLICY admin_policy ON capture TO windmill_admin USING (true);
CREATE POLICY admin_policy ON completed_job TO windmill_admin USING (true);
CREATE POLICY admin_policy ON flow TO windmill_admin USING (true);
CREATE POLICY admin_policy ON folder TO windmill_admin USING (true);
CREATE POLICY admin_policy ON queue TO windmill_admin USING (true);
CREATE POLICY admin_policy ON raw_app TO windmill_admin USING (true);
CREATE POLICY admin_policy ON resource TO windmill_admin USING (true);
CREATE POLICY admin_policy ON schedule TO windmill_admin USING (true);
CREATE POLICY admin_policy ON script TO windmill_admin USING (true);
CREATE POLICY admin_policy ON usr_to_group TO windmill_admin USING (true);
CREATE POLICY admin_policy ON variable TO windmill_admin USING (true);

View File

@@ -1 +1 @@
0428068e4fbbd1380a4d8bbaab5c8e7955decdb8
0f5f42d2f8f5f1af05c8086f3dc7ad38d83750df

View File

@@ -1 +0,0 @@
-- Add down migration script here

View File

@@ -1,2 +0,0 @@
-- Add up migration script here
ALTER TYPE SCRIPT_KIND ADD VALUE IF NOT EXISTS 'preprocessor';

View File

@@ -1 +0,0 @@
-- Add down migration script here

View File

@@ -1,24 +0,0 @@
-- Add up migration script here
DO
$$
DECLARE
tbl_name text;
policy_exists boolean;
tbl_names text[] := ARRAY['account', 'app', 'audit', 'capture', 'completed_job', 'flow', 'folder', 'http_trigger', 'queue', 'raw_app', 'resource', 'schedule', 'script', 'usr_to_group', 'variable'];
BEGIN
FOR tbl_name IN SELECT unnest(tbl_names)
LOOP
SELECT EXISTS (
SELECT 1
FROM pg_policies
WHERE schemaname = 'public'
AND tablename = tbl_name
AND policyname = 'admin_policy'
) INTO policy_exists;
IF NOT policy_exists THEN
EXECUTE format('CREATE POLICY admin_policy ON %I TO windmill_admin USING (true);', tbl_name);
END IF;
END LOOP;
END;
$$;

View File

@@ -1,2 +0,0 @@
-- Drop the alert_locks table
DROP TABLE IF EXISTS concurrency_locks;

View File

@@ -1,6 +0,0 @@
-- Create the alert_locks table
CREATE TABLE concurrency_locks (
id VARCHAR PRIMARY KEY,
last_locked_at TIMESTAMP NOT NULL,
owner VARCHAR NULL
);

View File

@@ -15,13 +15,13 @@ use windmill_parser::{
use swc_common::{sync::Lrc, FileName, SourceMap, SourceMapper, Span, Spanned};
use swc_ecma_ast::{
ArrayLit, AssignPat, BigInt, BindingIdent, Bool, Decl, ExportDecl, Expr, FnDecl, Ident,
IdentName, Lit, MemberExpr, MemberProp, ModuleDecl, ModuleItem, Number, ObjectLit, ObjectPat,
Param, Pat, Str, TsArrayType, TsEntityName, TsKeywordType, TsKeywordTypeKind, TsLit, TsLitType,
TsOptionalType, TsParenthesizedType, TsPropertySignature, TsType, TsTypeAnn, TsTypeElement,
TsTypeLit, TsTypeRef, TsUnionOrIntersectionType, TsUnionType,
ArrayLit, AssignPat, BigInt, BindingIdent, Bool, Decl, ExportDecl, Expr, FnDecl, Ident, Lit,
MemberExpr, MemberProp, ModuleDecl, ModuleItem, Number, ObjectLit, ObjectPat, Param, Pat, Str,
TsArrayType, TsEntityName, TsKeywordType, TsKeywordTypeKind, TsLit, TsLitType, TsOptionalType,
TsParenthesizedType, TsPropertySignature, TsType, TsTypeAnn, TsTypeElement, TsTypeLit,
TsTypeRef, TsUnionOrIntersectionType, TsUnionType,
};
use swc_ecma_parser::{lexer::Lexer, EsSyntax, Parser, StringInput, Syntax, TsSyntax};
use swc_ecma_parser::{lexer::Lexer, EsConfig, Parser, StringInput, Syntax, TsConfig};
use regex::Regex;
#[cfg(target_arch = "wasm32")]
@@ -48,9 +48,9 @@ impl Visit for ImportsFinder {
pub fn parse_expr_for_imports(code: &str) -> anyhow::Result<Vec<String>> {
let cm: Lrc<SourceMap> = Default::default();
let fm = cm.new_source_file(FileName::Custom("main.d.ts".into()).into(), code.into());
let fm = cm.new_source_file(FileName::Custom("main.d.ts".into()), code.into());
let lexer = Lexer::new(
Syntax::Typescript(TsSyntax::default()),
Syntax::Typescript(TsConfig::default()),
// EsVersion defaults to es5
Default::default(),
StringInput::from(&*fm),
@@ -69,7 +69,7 @@ pub fn parse_expr_for_imports(code: &str) -> anyhow::Result<Vec<String>> {
})?;
let mut visitor = ImportsFinder { imports: HashSet::new() };
visitor.visit_module(&expr);
swc_ecma_visit::visit_module(&mut visitor, &expr);
Ok(visitor.imports.into_iter().collect())
}
@@ -87,7 +87,7 @@ impl Visit for OutputFinder {
c.visit_with(self);
}
match m {
MemberExpr { obj, prop: MemberProp::Ident(IdentName { sym, .. }), .. } => {
MemberExpr { obj, prop: MemberProp::Ident(Ident { sym, .. }), .. } => {
match *obj.to_owned() {
Expr::Ident(Ident { sym: sym_i, .. }) => {
self.idents.insert((sym_i.to_string(), sym.to_string()));
@@ -102,10 +102,10 @@ impl Visit for OutputFinder {
pub fn parse_expr_for_ids(code: &str) -> anyhow::Result<Vec<(String, String)>> {
let cm: Lrc<SourceMap> = Default::default();
let fm = cm.new_source_file(FileName::Custom("main.ts".into()).into(), code.into());
let fm = cm.new_source_file(FileName::Custom("main.ts".into()), code.into());
let lexer = Lexer::new(
// We want to parse ecmascript
Syntax::Es(EsSyntax { jsx: false, ..Default::default() }),
Syntax::Es(EsConfig { jsx: false, ..Default::default() }),
// EsVersion defaults to es5
Default::default(),
StringInput::from(&*fm),
@@ -124,7 +124,7 @@ pub fn parse_expr_for_ids(code: &str) -> anyhow::Result<Vec<(String, String)>> {
})?;
let mut visitor = OutputFinder { idents: HashSet::new() };
visitor.visit_module(&expr);
swc_ecma_visit::visit_module(&mut visitor, &expr);
Ok(visitor.idents.into_iter().collect())
}
@@ -135,10 +135,10 @@ pub fn parse_deno_signature(
main_override: Option<String>,
) -> anyhow::Result<MainArgSignature> {
let cm: Lrc<SourceMap> = Default::default();
let fm = cm.new_source_file(FileName::Custom("main.ts".into()).into(), code.into());
let fm = cm.new_source_file(FileName::Custom("main.ts".into()), code.into());
let lexer = Lexer::new(
// We want to parse ecmascript
Syntax::Typescript(TsSyntax::default()),
Syntax::Typescript(TsConfig::default()),
// EsVersion defaults to es5
Default::default(),
StringInput::from(&*fm),

View File

@@ -399,11 +399,15 @@ fn parse_ansible_options(opts: &Vec<Yaml>) -> AnsiblePlaybookOptions {
if c > 0 && c <= 6 {
ret.verbosity = Some("v".repeat(c.min(6)));
}
}
}
_ => (),
_ => ()
}
}
}
}
@@ -418,10 +422,10 @@ fn count_consecutive_vs(s: &str) -> usize {
if c == 'v' {
current_count += 1;
if current_count == 6 {
return 6; // Stop early if we reach 6
return 6; // Stop early if we reach 6
}
} else {
current_count = 0; // Reset count if the character is not 'v'
current_count = 0; // Reset count if the character is not 'v'
}
max_count = max_count.max(current_count);
}

View File

@@ -67,7 +67,7 @@ use windmill_worker::{
get_hub_script_content_and_requirements, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR,
BUN_DEPSTAR_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM,
GO_BIN_CACHE_DIR, GO_CACHE_DIR, LOCK_CACHE_DIR, PIP_CACHE_DIR, POWERSHELL_CACHE_DIR,
RUST_CACHE_DIR, TAR_PIP_CACHE_DIR, TMP_LOGS_DIR, UV_CACHE_DIR,
RUST_CACHE_DIR, TAR_PIP_CACHE_DIR, TMP_LOGS_DIR,
};
use crate::monitor::{
@@ -92,6 +92,9 @@ const DEFAULT_SERVER_BIND_ADDR: Ipv4Addr = Ipv4Addr::new(0, 0, 0, 0);
mod ee;
mod monitor;
#[cfg(feature = "pg_embed")]
mod pg_embed;
#[inline(always)]
fn create_and_run_current_thread_inner<F, R>(future: F) -> R
where
@@ -115,8 +118,7 @@ where
}
pub fn main() -> anyhow::Result<()> {
#[cfg(feature = "deno_core")]
deno_core::JsRuntime::init_platform(None, false);
deno_core::JsRuntime::init_platform(None);
create_and_run_current_thread_inner(windmill_main())
}
@@ -135,7 +137,6 @@ async fn cache_hub_scripts(file_path: Option<String>) -> anyhow::Result<()> {
})?;
create_dir_all(HUB_CACHE_DIR).await?;
create_dir_all(BUN_BUNDLE_CACHE_DIR).await?;
for path in paths.values() {
tracing::info!("Caching hub script at {path}");
@@ -167,7 +168,7 @@ async fn cache_hub_scripts(file_path: Option<String>) -> anyhow::Result<()> {
create_dir_all(&job_dir).await?;
if let Some(lockfile) = res.lockfile {
let _ = windmill_worker::prepare_job_dir(&lockfile, &job_dir).await?;
let envs = windmill_worker::get_common_bun_proc_envs(None).await;
let _ = windmill_worker::install_bun_lockfile(
&mut 0,
&mut None,
@@ -176,31 +177,11 @@ async fn cache_hub_scripts(file_path: Option<String>) -> anyhow::Result<()> {
None,
&job_dir,
"cache_init",
envs.clone(),
windmill_worker::get_common_bun_proc_envs(None).await,
false,
&mut None,
)
.await?;
let _ = windmill_common::worker::write_file(&job_dir, "main.js", &res.content)?;
if let Err(e) = windmill_worker::prebundle_bun_script(
&res.content,
Some(lockfile),
&path,
&job_id,
"admins",
None,
&job_dir,
"",
"cache_init",
"",
&mut None,
)
.await
{
panic!("Error prebundling bun script: {e:#}");
}
} else {
tracing::warn!("No lockfile found for bun script {path}, skipping...");
}
@@ -360,6 +341,14 @@ async fn windmill_main() -> anyhow::Result<()> {
config
});
#[cfg(feature = "pg_embed")]
let _pg = {
let (db_url, pg) = pg_embed::start().await.expect("pg embed");
tracing::info!("Use embedded pg: {db_url}");
std::env::set_var("DATABASE_URL", db_url);
pg
};
tracing::info!("Connecting to database...");
let db = windmill_common::connect_db(server_mode, indexer_mode).await?;
tracing::info!("Database connected");
@@ -384,16 +373,8 @@ async fn windmill_main() -> anyhow::Result<()> {
let is_agent = mode == Mode::Agent;
if !is_agent {
let skip_migration = std::env::var("SKIP_MIGRATION")
.map(|val| val == "true")
.unwrap_or(false);
if !skip_migration {
// migration code to avoid break
windmill_api::migrate_db(&db).await?;
} else {
tracing::info!("SKIP_MIGRATION set, skipping db migration...")
}
// migration code to avoid break
windmill_api::migrate_db(&db).await?;
}
let (killpill_tx, mut killpill_rx) = tokio::sync::broadcast::channel::<()>(2);
@@ -476,7 +457,7 @@ Windmill Community Edition {GIT_VERSION}
#[cfg(feature = "tantivy")]
let (index_reader, index_writer) = if should_index_jobs {
let (r, w) = windmill_indexer::indexer_ee::init_index(&db).await?;
let (r, w) = windmill_indexer::indexer_ee::init_index().await?;
(Some(r), Some(w))
} else {
(None, None)
@@ -893,7 +874,6 @@ pub async fn run_workers<R: rsmq_async::RsmqConnection + Send + Sync + Clone + '
LOCK_CACHE_DIR,
TMP_LOGS_DIR,
PIP_CACHE_DIR,
UV_CACHE_DIR,
TAR_PIP_CACHE_DIR,
DENO_CACHE_DIR,
DENO_CACHE_DIR_DEPS,

View File

@@ -27,7 +27,7 @@ use windmill_api::{
DEFAULT_BODY_LIMIT, IS_SECURE, OAUTH_CLIENTS, REQUEST_SIZE_LIMIT, SAML_METADATA, SCIM_TOKEN,
};
#[cfg(feature = "enterprise")]
use windmill_common::ee::{worker_groups_alerts, jobs_waiting_alerts};
use windmill_common::ee::worker_groups_alerts;
use windmill_common::{
auth::JWT_SECRET,
ee::CriticalErrorChannel,
@@ -1061,20 +1061,12 @@ pub async fn monitor_db(
}
};
let jobs_waiting_alerts_f = async {
#[cfg(feature = "enterprise")]
if server_mode {
jobs_waiting_alerts(&db).await;
}
};
join!(
expired_items_f,
zombie_jobs_f,
expose_queue_metrics_f,
verify_license_key_f,
worker_groups_alerts_f,
jobs_waiting_alerts_f,
worker_groups_alerts_f
);
}

46
backend/src/pg_embed.rs Normal file
View File

@@ -0,0 +1,46 @@
use pg_embed::pg_enums::PgAuthMethod;
use pg_embed::pg_fetch::PgFetchSettings;
use pg_embed::postgres::{PgEmbed, PgSettings};
use std::path::PathBuf;
use std::time::Duration;
pub async fn start() -> anyhow::Result<(String, PgEmbed)> {
let pg_settings = PgSettings {
database_dir: PathBuf::from("/tmp/db"),
port: 6543,
user: "postgres".to_string(),
password: "password".to_string(),
auth_method: PgAuthMethod::Plain,
persistent: false,
timeout: Some(Duration::from_secs(15)),
migration_dir: None,
};
let fetch_settings = PgFetchSettings {
version: pg_embed::pg_fetch::PostgresVersion("15.3.0"),
..Default::default()
};
tracing::info!(
"Fetch settings: {:?} {:?}",
fetch_settings.operating_system,
fetch_settings.architecture
);
let mut pg = PgEmbed::new(pg_settings, fetch_settings).await?;
pg.setup().await.expect("pg setup");
pg.start_db().await.expect("pg start db");
//TODO: re-enable this to make it work
// if !pg.database_exists("windmill").await.expect("db exists") {
// pg.create_database("windmill")
// .await
// .expect("pg create database");
// }
let uri = pg.full_db_uri("windmill");
Ok((uri, pg))
}

View File

@@ -2745,7 +2745,7 @@ async fn test_flow_lock_all(db: Pool<Postgres>) {
"lock": null,
"path": null,
"type": "rawscript",
"content": "import * as wmill from \"https://deno.land/x/windmill@v1.50.0/mod.ts\"\n\nexport async function main() {\n return wmill\n}\n",
"content": "import * as wmill from \"https://deno.land/x/windmill@v1.50.0/mod.ts\"\n\nexport async function main() {\n return \"Hello\"\n}\n",
"language": "deno",
"input_transforms": {}
},

View File

@@ -17,7 +17,7 @@ benchmark = []
embedding = ["dep:tinyvector", "dep:hf-hub", "dep:tokenizers", "dep:candle-core", "dep:candle-transformers", "dep:candle-nn"]
parquet = ["dep:datafusion", "dep:object_store", "dep:url", "windmill-common/parquet"]
prometheus = ["windmill-common/prometheus", "windmill-queue/prometheus", "dep:prometheus"]
openidconnect = ["dep:openidconnect", "windmill-common/openidconnect"]
openidconnect = ["dep:openidconnect"]
tantivy = ["dep:windmill-indexer"]
[dependencies]

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -1,7 +1,7 @@
openapi: "3.0.3"
info:
version: 1.409.2
version: 1.403.1
title: Windmill API
contact:
@@ -2789,14 +2789,7 @@ paths:
oauth:
type: array
items:
type: object
properties:
type:
type: string
display_name:
type: string
required:
- type
type: string
saml:
type: string
required:
@@ -4020,42 +4013,6 @@ paths:
schema:
$ref: "#/components/schemas/Script"
/w/{workspace}/scripts/get_triggers_count/{path}:
get:
summary: get triggers count of script
operationId: getTriggersCountOfScript
tags:
- script
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/ScriptPath"
responses:
"200":
description: triggers count
content:
application/json:
schema:
$ref: "#/components/schemas/TriggersCount"
/w/{workspace}/scripts/list_tokens/{path}:
get:
summary: get tokens with script scope
operationId: listTokensOfScript
tags:
- script
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/ScriptPath"
responses:
"200":
description: tokens list
content:
application/json:
schema:
type: array
items:
$ref: "#/components/schemas/TruncatedToken"
/w/{workspace}/scripts/get/draft/{path}:
get:
summary: get script by path with draft
@@ -4660,43 +4617,6 @@ paths:
schema:
$ref: "#/components/schemas/Flow"
/w/{workspace}/flows/get_triggers_count/{path}:
get:
summary: get triggers count of flow
operationId: getTriggersCountOfFlow
tags:
- flow
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/ScriptPath"
responses:
"200":
description: triggers count
content:
application/json:
schema:
$ref: "#/components/schemas/TriggersCount"
/w/{workspace}/flows/list_tokens/{path}:
get:
summary: get tokens with flow scope
operationId: listTokensOfFlow
tags:
- flow
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/ScriptPath"
responses:
"200":
description: tokens list
content:
application/json:
schema:
type: array
items:
$ref: "#/components/schemas/TruncatedToken"
/w/{workspace}/flows/toggle_workspace_error_handler/{path}:
post:
summary: Toggle ON and OFF the workspace error handler for a given flow
@@ -10327,8 +10247,6 @@ components:
type: array
items:
type: string
email:
type: string
required:
- token_prefix
- created_at
@@ -10346,8 +10264,6 @@ components:
type: array
items:
type: string
workspace_id:
type: string
NewTokenImpersonate:
type: object
@@ -10359,8 +10275,6 @@ components:
format: date-time
impersonate_email:
type: string
workspace_id:
type: string
required:
- impersonate_email
@@ -11197,23 +11111,6 @@ components:
- requires_auth
- http_method
TriggersCount:
type: object
properties:
primary_schedule:
type: object
properties:
schedule:
type: string
schedule_count:
type: number
http_routes_count:
type: number
webhook_count:
type: number
email_count:
type: number
Group:
type: object
properties:

View File

@@ -199,6 +199,11 @@ pub async fn migrate(db: &DB) -> Result<(), Error> {
Err(err) => Err(err),
}?;
#[cfg(feature = "enterprise")]
if let Err(e) = windmill_migrations(&mut custom_migrator, db).await {
tracing::error!("Could not apply windmill custom migrations: {e:#}")
}
Ok(())
}
@@ -492,6 +497,33 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> {
Ok(())
}
#[cfg(feature = "enterprise")]
async fn windmill_migrations(migrator: &mut CustomMigrator, db: &DB) -> Result<(), Error> {
if std::env::var("MIGRATION_NO_BYPASSRLS").is_ok() {
migrator.lock().await?;
let has_done_migration = sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'bypassrls_1-2')",
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !has_done_migration {
let query = include_str!("../../custom_migrations/bypassrls_1.sql");
tracing::info!("Applying bypassrls_1.sql");
let mut tx: sqlx::Transaction<'_, Postgres> = db.begin().await?;
tx.execute(query).await?;
tracing::info!("Applied bypassrls_1.sql");
sqlx::query!("INSERT INTO windmill_migrations (name) VALUES ('bypassrls_1-2')")
.execute(&mut *tx)
.await?;
tx.commit().await?;
}
migrator.unlock().await?;
}
Ok(())
}
#[derive(Clone, Debug)]
pub struct ApiAuthed {
pub email: String,

View File

@@ -9,9 +9,6 @@
use std::collections::HashMap;
use crate::db::ApiAuthed;
use crate::triggers::{
get_triggers_count_internal, list_tokens_internal, TriggersCount, TruncatedTokenWithEmail,
};
use crate::utils::WithStarredInfoQuery;
use crate::{
db::DB,
@@ -56,8 +53,6 @@ pub fn workspaced_service() -> Router {
.route("/update/*path", post(update_flow))
.route("/archive/*path", post(archive_flow_by_path))
.route("/delete/*path", delete(delete_flow_by_path))
.route("/get_triggers_count/*path", get(get_triggers_count))
.route("/list_tokens/*path", get(list_tokens))
.route("/get/*path", get(get_flow_by_path))
.route("/get/draft/*path", get(get_flow_by_path_w_draft))
.route("/exists/*path", get(exists_flow_by_path))
@@ -879,22 +874,6 @@ async fn update_flow(
Ok(nf.path.to_string())
}
async fn get_triggers_count(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<TriggersCount> {
let path = path.to_path();
get_triggers_count_internal(&db, &w_id, &path, true).await
}
async fn list_tokens(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Vec<TruncatedTokenWithEmail>> {
let path = path.to_path();
list_tokens_internal(&db, &w_id, &path, true).await
}
async fn get_flow_by_path(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,

View File

@@ -275,10 +275,8 @@ pub fn require_is_owner(authed: &ApiAuthed, name: &str) -> Result<()> {
async fn update_folder(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Extension(rsmq): Extension<Option<rsmq_async::MultiplexedRsmq>>,
Path((w_id, name)): Path<(String, String)>,
Json(mut ng): Json<UpdateFolder>,
) -> Result<String> {
@@ -369,18 +367,6 @@ async fn update_folder(
}
}
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Folder { path: format!("f/{}", name) },
Some(format!("Folder '{}' updated", name)),
rsmq,
true,
)
.await?;
audit_log(
&mut *tx,
&authed,

View File

@@ -11,6 +11,7 @@ use axum::http::HeaderValue;
use quick_cache::sync::Cache;
use serde_json::value::RawValue;
use sqlx::Pool;
use windmill_common::error::JsonResult;
use std::collections::HashMap;
#[cfg(feature = "prometheus")]
use std::sync::atomic::Ordering;
@@ -18,13 +19,12 @@ use tokio::io::AsyncReadExt;
#[cfg(feature = "prometheus")]
use tokio::time::Instant;
use tower::ServiceBuilder;
use windmill_common::error::JsonResult;
use windmill_common::flow_status::{JobResult, RestartedFrom};
use windmill_common::jobs::{
format_completed_job_result, format_result, CompletedJobWithFormattedResult, FormattedResult,
ENTRYPOINT_OVERRIDE,
};
use windmill_common::worker::{CLOUD_HOSTED, TMP_DIR};
use windmill_common::worker::TMP_DIR;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use windmill_common::scripts::PREVIEW_IS_CODEBASE_HASH;
@@ -81,7 +81,7 @@ use windmill_common::{METRICS_DEBUG_ENABLED, METRICS_ENABLED};
use windmill_common::{get_latest_deployed_hash_for_path, BASE_URL};
use windmill_queue::{
cancel_job, get_queued_job, get_result_by_id_from_running_flow, job_is_complete, push,
DecodeQueries, PushArgs, PushArgsOwned, PushIsolationLevel,
DecodeQueries, PushArgs, PushArgsOwned, PushIsolationLevel, QueueTransaction,
};
#[cfg(feature = "prometheus")]
@@ -293,8 +293,11 @@ pub fn workspace_unauthed_service() -> Router {
pub fn global_root_service() -> Router {
Router::new()
.route("/db_clock", get(get_db_clock))
.route("/completed/count_by_tag", get(count_by_tag))
.route("/db_clock", get(get_db_clock))
.route(
"/completed/count_by_tag",
get(count_by_tag),
)
}
#[derive(Deserialize)]
@@ -544,8 +547,8 @@ pub async fn get_path_for_hash<'c>(
Ok(path)
}
pub async fn get_path_tag_limits_cache_for_hash(
tx: &DB,
pub async fn get_path_tag_limits_cache_for_hash<'c, R: rsmq_async::RsmqConnection + Send>(
tx: &mut QueueTransaction<'c, R>,
w_id: &str,
hash: i64,
) -> error::Result<(
@@ -1469,8 +1472,6 @@ async fn cancel_jobs(
}
}
uuids.extend(trivial_jobs);
Ok(Json(uuids))
}
@@ -2812,6 +2813,7 @@ pub async fn run_flow_by_path_inner(
let flow_path = flow_path.to_path();
check_scopes(&authed, || format!("run:flow/{flow_path}"))?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let (tag, dedicated_worker, has_preprocessor) = sqlx::query!(
"SELECT tag, dedicated_worker, flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor
@@ -2822,7 +2824,7 @@ pub async fn run_flow_by_path_inner(
flow_path,
w_id
)
.fetch_optional(&db)
.fetch_optional(&mut tx)
.await?
.map(|x| (x.tag, x.dedicated_worker, x.has_preprocessor))
.ok_or_else(|| {
@@ -2835,7 +2837,7 @@ pub async fn run_flow_by_path_inner(
check_tag_available_for_workspace(&w_id, &tag).await?;
let scheduled_for = run_query.get_scheduled_for(&db).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
tx,
@@ -2907,13 +2909,14 @@ pub async fn restart_flow(
) -> error::Result<(StatusCode, String)> {
check_license_key_valid().await?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let completed_job = sqlx::query_as::<_, CompletedJob>(
"SELECT *, result->'wm_labels' as labels from completed_job WHERE id = $1 and workspace_id = $2",
)
.bind(job_id)
.bind(&w_id)
.fetch_optional(&db)
.fetch_optional(&mut tx)
.await?
.with_context(|| "Unable to find completed job with the given job UUID")?;
@@ -2931,7 +2934,7 @@ pub async fn restart_flow(
let scheduled_for = run_query.get_scheduled_for(&db).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -3007,16 +3010,16 @@ pub async fn run_script_by_path_inner(
check_scopes(&authed, || format!("run:script/{script_path}"))?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let (job_payload, tag, _delete_after_use, timeout) =
script_path_to_payload(script_path, &db, &w_id, run_query.skip_preprocessor).await?;
script_path_to_payload(script_path, &mut tx, &w_id, run_query.skip_preprocessor).await?;
let scheduled_for = run_query.get_scheduled_for(&db).await?;
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -3049,11 +3052,6 @@ pub async fn run_script_by_path_inner(
Ok((StatusCode::CREATED, uuid.to_string()))
}
#[derive(Deserialize)]
pub struct WorkflowAsCodeQuery {
pub skip_update: Option<bool>,
}
pub async fn run_workflow_as_code(
authed: ApiAuthed,
Extension(db): Extension<DB>,
@@ -3061,38 +3059,15 @@ pub async fn run_workflow_as_code(
Extension(rsmq): Extension<Option<rsmq_async::MultiplexedRsmq>>,
Path((w_id, job_id, entrypoint)): Path<(String, Uuid, String)>,
Query(run_query): Query<RunJobQuery>,
Query(wkflow_query): Query<WorkflowAsCodeQuery>,
Json(task): Json<WorkflowTask>,
) -> error::Result<(StatusCode, String)> {
let mut i = 1;
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
#[cfg(feature = "enterprise")]
check_license_key_valid().await?;
check_tag_available_for_workspace(&w_id, &run_query.tag).await?;
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let job = get_queued_job(&job_id, &w_id, &db).await?;
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
let job = not_found_if_none(job, "Queued Job", &job_id.to_string())?;
let (job_payload, tag, _delete_after_use, timeout) = match job.job_kind {
JobKind::Preview => (
@@ -3117,7 +3092,7 @@ pub async fn run_workflow_as_code(
JobKind::Script => {
script_path_to_payload(
job.script_path(),
&db,
&mut tx,
&w_id,
run_query.skip_preprocessor,
)
@@ -3126,12 +3101,6 @@ pub async fn run_workflow_as_code(
_ => return Err(anyhow::anyhow!("Not supported").into()),
};
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
let mut extra = HashMap::new();
extra.insert(ENTRYPOINT_OVERRIDE.to_string(), to_raw_value(&entrypoint));
@@ -3140,21 +3109,7 @@ pub async fn run_workflow_as_code(
let tag = run_query.tag.clone().or(tag).or(Some(job.tag));
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, mut tx) = push(
&db,
@@ -3181,39 +3136,14 @@ pub async fn run_workflow_as_code(
Some(&authed.clone().into()),
)
.await?;
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
if !wkflow_query.skip_update.unwrap_or(false) {
sqlx::query!(
"UPDATE queue SET flow_status = jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], jsonb_set(jsonb_set('{}'::jsonb, '{scheduled_for}', to_jsonb(now()::text)), '{name}', to_jsonb($4::text))) WHERE id = $2 AND workspace_id = $3",
uuid.to_string(),
job_id,
w_id,
entrypoint
).execute(&mut tx).await?;
} else {
tracing::info!("Skipping update of flow status for job {job_id} in workspace {w_id}");
}
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
i += 1;
}
sqlx::query!(
"UPDATE queue SET flow_status = jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], jsonb_set(jsonb_set('{}'::jsonb, '{scheduled_for}', to_jsonb(now()::text)), '{name}', to_jsonb($4::text))) WHERE id = $2 AND workspace_id = $3",
uuid.to_string(),
job_id,
w_id,
entrypoint
).execute(&mut tx).await?;
tx.commit().await?;
if *CLOUD_HOSTED {
tracing::info!("workflow_as_code_tracing id {i} ");
}
Ok((StatusCode::CREATED, uuid.to_string()))
}
@@ -3577,13 +3507,15 @@ pub async fn run_wait_result_job_by_path_get(
let script_path = script_path.to_path();
check_scopes(&authed, || format!("run:script/{script_path}"))?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let (job_payload, tag, delete_after_use, timeout) =
script_path_to_payload(script_path, &db, &w_id, run_query.skip_preprocessor).await?;
script_path_to_payload(script_path, &mut tx, &w_id, run_query.skip_preprocessor).await?;
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -3700,13 +3632,15 @@ pub async fn run_wait_result_script_by_path_internal(
let script_path = script_path.to_path();
check_scopes(&authed, || format!("run:script/{script_path}"))?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let (job_payload, tag, delete_after_use, timeout) =
script_path_to_payload(script_path, &db, &w_id, run_query.skip_preprocessor).await?;
script_path_to_payload(script_path, &mut tx, &w_id, run_query.skip_preprocessor).await?;
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -3758,6 +3692,8 @@ pub async fn run_wait_result_script_by_hash(
check_queue_too_long(&db, run_query.queue_limit).await?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let hash = script_hash.0;
let (
path,
@@ -3772,7 +3708,7 @@ pub async fn run_wait_result_script_by_hash(
delete_after_use,
timeout,
has_preprocessor,
) = get_path_tag_limits_cache_for_hash(&db, &w_id, hash).await?;
) = get_path_tag_limits_cache_for_hash(&mut tx, &w_id, hash).await?;
if let Some(run_query_cache_ttl) = run_query.cache_ttl {
cache_ttl = Some(run_query_cache_ttl);
}
@@ -3781,7 +3717,7 @@ pub async fn run_wait_result_script_by_hash(
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -3863,6 +3799,7 @@ pub async fn run_wait_result_flow_by_path_internal(
let flow_path = flow_path.to_path();
check_scopes(&authed, || format!("run:flow/{flow_path}"))?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let scheduled_for = run_query.get_scheduled_for(&db).await?;
@@ -3875,7 +3812,7 @@ pub async fn run_wait_result_flow_by_path_internal(
flow_path,
w_id
)
.fetch_optional(&db)
.fetch_optional(&mut tx)
.await?
.map(|x| (x.tag, x.dedicated_worker, x.early_return, x.has_preprocessor))
.ok_or_else(|| {
@@ -3887,7 +3824,7 @@ pub async fn run_wait_result_flow_by_path_internal(
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -4193,7 +4130,7 @@ async fn run_dependencies_job(
JsonRawValue::from_string("true".to_string()).unwrap(),
);
if language == ScriptLang::Bun {
let annotation = windmill_common::worker::get_annotation_ts(&raw_code);
let annotation = windmill_common::worker::get_annotation(&raw_code);
hm.insert(
"npm_mode".to_string(),
JsonRawValue::from_string(annotation.npm_mode.to_string()).unwrap(),
@@ -4359,6 +4296,7 @@ async fn add_batch_jobs(
}
}
"flow" => {
let mut tx = PushIsolationLevel::IsolatedRoot(db.clone(), rsmq);
let mut uuids: Vec<Uuid> = Vec::new();
let payload = if let Some(ref fv) = batch_info.flow_value {
@@ -4376,7 +4314,6 @@ async fn add_batch_jobs(
))?
}
};
let mut tx = PushIsolationLevel::IsolatedRoot(db.clone(), rsmq);
for _ in 0..n {
let ehm = HashMap::new();
let (uuid, ntx) = push(
@@ -4577,6 +4514,7 @@ pub async fn run_job_by_hash_inner(
#[cfg(feature = "enterprise")]
check_license_key_valid().await?;
let mut tx: QueueTransaction<'_, _> = (rsmq, user_db.begin(&authed).await?).into();
let hash = script_hash.0;
let (
@@ -4592,7 +4530,7 @@ pub async fn run_job_by_hash_inner(
_delete_after_use, // not taken into account in async endpoints
timeout,
has_preprocessor,
) = get_path_tag_limits_cache_for_hash(&db, &w_id, hash).await?;
) = get_path_tag_limits_cache_for_hash(&mut tx, &w_id, hash).await?;
check_scopes(&authed, || format!("run:script/{path}"))?;
if let Some(run_query_cache_ttl) = run_query.cache_ttl {
cache_ttl = Some(run_query_cache_ttl);
@@ -4601,7 +4539,7 @@ pub async fn run_job_by_hash_inner(
let tag = run_query.tag.clone().or(tag);
check_tag_available_for_workspace(&w_id, &tag).await?;
let tx = PushIsolationLevel::Isolated(user_db, authed.clone().into(), rsmq);
let tx = PushIsolationLevel::Transaction(tx);
let (uuid, tx) = push(
&db,
@@ -4745,8 +4683,8 @@ async fn get_job_update(
.fetch_optional(&db)
.await?;
let progress: Option<i32> = if get_progress == Some(true) {
sqlx::query_scalar!(
let progress: Option<i32> = if get_progress == Some(true){
sqlx::query_scalar!(
"SELECT scalar_int FROM job_stats WHERE workspace_id = $1 AND job_id = $2 AND metric_id = $3",
&w_id,
job_id,
@@ -5177,6 +5115,8 @@ async fn get_completed_job_result(
Ok(Json(result).into_response())
}
#[derive(Deserialize)]
struct CountByTagQuery {
horizon_secs: Option<i64>,
@@ -5190,7 +5130,7 @@ struct TagCount {
}
async fn count_by_tag(
ApiAuthed { email, .. }: ApiAuthed,
ApiAuthed { email, ..}: ApiAuthed,
Extension(db): Extension<DB>,
Query(query): Query<CountByTagQuery>,
) -> JsonResult<Vec<TagCount>> {

View File

@@ -82,7 +82,6 @@ pub mod smtp_server_ee;
mod static_assets;
mod stripe_ee;
mod tracing_init;
mod triggers;
mod users;
mod utils;
mod variables;

View File

@@ -58,10 +58,7 @@ pub fn workspaced_service() -> Router {
.route("/type/exists/:name", get(exists_resource_type))
.route("/type/update/:name", post(update_resource_type))
.route("/type/delete/:name", delete(delete_resource_type))
.route(
"/file_resource_type_to_file_ext_map",
get(file_resource_ext_to_resource_type),
)
.route("/file_resource_type_to_file_ext_map", get(file_resource_ext_to_resource_type))
.route("/type/create", post(create_resource_type))
}

View File

@@ -9,9 +9,6 @@
use crate::{
db::{ApiAuthed, DB},
schedule::clear_schedule,
triggers::{
get_triggers_count_internal, list_tokens_internal, TriggersCount, TruncatedTokenWithEmail,
},
users::{maybe_refresh_folders, require_owner_of_path, AuthCache},
utils::WithStarredInfoQuery,
webhook_util::{WebhookMessage, WebhookShared},
@@ -56,7 +53,7 @@ use windmill_common::{
utils::{
not_found_if_none, paginate, query_elems_from_hub, require_admin, Pagination, StripPath,
},
worker::{get_annotation_ts, to_raw_value},
worker::{get_annotation, to_raw_value},
HUB_BASE_URL,
};
use windmill_git_sync::{handle_deployment_metadata, DeployedObject};
@@ -135,8 +132,6 @@ pub fn workspaced_service() -> Router {
.route("/archive/p/*path", post(archive_script_by_path))
.route("/get/draft/*path", get(get_script_by_path_w_draft))
.route("/get/p/*path", get(get_script_by_path))
.route("/get_triggers_count/*path", get(get_triggers_count))
.route("/list_tokens/*path", get(list_tokens))
.route("/raw/p/*path", get(raw_script_by_path))
.route("/raw_unpinned/p/*path", get(raw_script_by_path_unpinned))
.route("/exists/p/*path", get(exists_script_by_path))
@@ -606,7 +601,7 @@ async fn create_script_internal<'c>(
};
let lang = if &ns.language == &ScriptLang::Bun || &ns.language == &ScriptLang::Bunnative {
let anns = get_annotation_ts(&ns.content);
let anns = get_annotation(&ns.content);
if anns.native_mode {
ScriptLang::Bunnative
} else {
@@ -879,22 +874,6 @@ async fn get_script_by_path(
Ok(Json(script))
}
async fn list_tokens(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Vec<TruncatedTokenWithEmail>> {
let path = path.to_path();
list_tokens_internal(&db, &w_id, &path, false).await
}
async fn get_triggers_count(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<TriggersCount> {
let path = path.to_path();
get_triggers_count_internal(&db, &w_id, &path, false).await
}
async fn get_script_by_path_w_draft(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,

View File

@@ -9,7 +9,6 @@
use ::tracing::{field, Span};
use hyper::Response;
use tower_http::trace::{MakeSpan, OnFailure, OnResponse};
use uuid::Uuid;
lazy_static::lazy_static! {
static ref LOG_REQUESTS: bool = std::env::var("LOG_REQUESTS")
@@ -46,28 +45,17 @@ impl<B> OnFailure<B> for MyOnFailure {
// tracing::error!(latency = latency.as_millis(), "response")
}
}
lazy_static::lazy_static! {
static ref TRACING_HEADER: String = std::env::var("TRACING_HEADER")
.ok().unwrap_or_else(|| "x-tracing-id".to_string());
}
#[derive(Clone)]
pub struct MyMakeSpan {}
impl<B> MakeSpan<B> for MyMakeSpan {
fn make_span(&mut self, request: &hyper::Request<B>) -> Span {
let tracing_id = request
.headers()
.get(TRACING_HEADER.as_str())
.and_then(|x| x.to_str().map(|x| x.to_string()).ok())
.unwrap_or(Uuid::new_v4().to_string());
tracing::info_span!(
"request",
method = %request.method(),
uri = %request.uri(),
username = field::Empty,
workspace_id = field::Empty,
trace_id = tracing_id,
email = field::Empty,
)
}

View File

@@ -1,131 +0,0 @@
use axum::Json;
use serde::{Deserialize, Serialize};
use sqlx::FromRow;
use windmill_common::error::JsonResult;
use crate::db::DB;
#[derive(Serialize, Deserialize, Debug)]
pub struct TriggerPrimarySchedule {
schedule: String,
}
#[derive(Serialize, Deserialize, Debug)]
pub struct TriggersCount {
primary_schedule: Option<TriggerPrimarySchedule>,
schedule_count: i64,
http_routes_count: i64,
webhook_count: i64,
email_count: i64,
}
pub(crate) async fn get_triggers_count_internal(
db: &DB,
w_id: &str,
path: &str,
is_flow: bool,
) -> JsonResult<TriggersCount> {
let primary_schedule = sqlx::query_scalar!(
"SELECT schedule FROM schedule WHERE path = $1 AND script_path = $1 AND is_flow = $2 AND workspace_id = $3",
path,
is_flow,
w_id
)
.fetch_optional(db)
.await?;
let schedule_count = sqlx::query_scalar!(
"SELECT COUNT(*) FROM schedule WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
path,
is_flow,
w_id
)
.fetch_one(db)
.await?
.unwrap_or(0);
let http_routes_count = sqlx::query_scalar!(
"SELECT COUNT(*) FROM http_trigger WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
path,
is_flow,
w_id
)
.fetch_one(db)
.await?
.unwrap_or(0);
let webhook_count = (if is_flow {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
w_id,
path,
)
} else {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:' || $2]::text[]",
w_id,
path,
)
}).fetch_one(db)
.await?
.unwrap_or(0);
let email_count = (if is_flow {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'email-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
w_id,
path,
)
} else {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'email-%' AND workspace_id = $1 AND scopes @> ARRAY['run:script/' || $2]::text[]",
w_id,
path,
)
}).fetch_one(db)
.await?
.unwrap_or(0);
Ok(Json(TriggersCount {
primary_schedule: primary_schedule.map(|s| TriggerPrimarySchedule { schedule: s }),
schedule_count,
http_routes_count,
webhook_count,
email_count,
}))
}
#[derive(FromRow, Serialize)]
pub struct TruncatedTokenWithEmail {
pub label: Option<String>,
pub token_prefix: Option<String>,
pub expiration: Option<chrono::DateTime<chrono::Utc>>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub last_used_at: chrono::DateTime<chrono::Utc>,
pub scopes: Option<Vec<String>>,
pub email: Option<String>,
}
pub async fn list_tokens_internal(
db: &DB,
w_id: &str,
path: &str,
is_flow: bool,
) -> JsonResult<Vec<TruncatedTokenWithEmail>> {
let tokens = if is_flow {
sqlx::query_as!(
TruncatedTokenWithEmail,
"SELECT label, concat(substring(token for 10)) as token_prefix, expiration, created_at, last_used_at, scopes, email FROM token WHERE workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
w_id, path).fetch_all(db)
.await?
} else {
sqlx::query_as!(
TruncatedTokenWithEmail,
"SELECT label, concat(substring(token for 10)) as token_prefix, expiration, created_at, last_used_at, scopes, email FROM token WHERE workspace_id = $1 AND scopes @> ARRAY['run:script/' || $2]::text[]",
w_id, path)
.fetch_all(db)
.await?
};
Ok(Json(tokens))
}

View File

@@ -128,13 +128,7 @@ pub fn make_unauthed_service() -> Router {
fn username_override_from_label(label: Option<String>) -> Option<String> {
match label {
Some(label)
if label.starts_with("webhook-")
|| label.starts_with("http-")
|| label.starts_with("email-") =>
{
Some(label)
}
Some(label) if label.starts_with("webhook-") => Some(label),
Some(label) if label.starts_with("ephemeral-script-end-user-") => Some(
label
.trim_start_matches("ephemeral-script-end-user-")
@@ -276,10 +270,9 @@ impl AuthCache {
_ => {
let user_o = sqlx::query_as::<_, (Option<String>, Option<String>, bool, Option<Vec<String>>, Option<String>)>(
"UPDATE token SET last_used_at = now() WHERE token = $1 AND (expiration > NOW() \
OR expiration IS NULL) AND (workspace_id IS NULL OR workspace_id = $2) RETURNING owner, email, super_admin, scopes, label",
OR expiration IS NULL) RETURNING owner, email, super_admin, scopes, label",
)
.bind(token)
.bind(w_id.as_ref())
.fetch_optional(&self.db)
.await
.ok()
@@ -844,7 +837,6 @@ pub struct NewToken {
pub expiration: Option<chrono::DateTime<chrono::Utc>>,
pub impersonate_email: Option<String>,
pub scopes: Option<Vec<String>>,
pub workspace_id: Option<String>,
}
#[derive(Deserialize)]
@@ -2397,15 +2389,14 @@ async fn create_token(
.unwrap_or(false);
sqlx::query!(
"INSERT INTO token
(token, email, label, expiration, super_admin, scopes, workspace_id)
VALUES ($1, $2, $3, $4, $5, $6, $7)",
(token, email, label, expiration, super_admin, scopes)
VALUES ($1, $2, $3, $4, $5, $6)",
token,
authed.email,
new_token.label,
new_token.expiration,
is_super_admin,
new_token.scopes.as_ref().map(|x| x.as_slice()),
new_token.workspace_id,
new_token.scopes.as_ref().map(|x| x.as_slice())
)
.execute(&mut *tx)
.await?;

View File

@@ -13,7 +13,6 @@ flamegraph = ["dep:tracing-flame"]
loki = ["dep:tracing-loki"]
benchmark = []
parquet = ["dep:object_store", "dep:aws-config", "dep:aws-sdk-sts"]
openidconnect = ["dep:openidconnect"]
[lib]
name = "windmill_common"
@@ -58,7 +57,6 @@ futures-core.workspace = true
async-stream.workspace = true
const_format.workspace = true
crc.workspace = true
openidconnect = { workspace = true, optional = true}
[target.'cfg(not(target_env = "msvc"))'.dependencies]
tikv-jemalloc-ctl = { optional = true, workspace = true }

View File

@@ -0,0 +1,73 @@
#[cfg(feature = "enterprise")]
use crate::db::DB;
use crate::ee::LicensePlan::Community;
#[cfg(feature = "enterprise")]
use crate::error;
use serde::Deserialize;
use std::sync::Arc;
use tokio::sync::RwLock;
lazy_static::lazy_static! {
pub static ref LICENSE_KEY_VALID: Arc<RwLock<bool>> = Arc::new(RwLock::new(true));
pub static ref LICENSE_KEY_ID: Arc<RwLock<String>> = Arc::new(RwLock::new("".to_string()));
pub static ref LICENSE_KEY: Arc<RwLock<String>> = Arc::new(RwLock::new("".to_string()));
}
pub enum LicensePlan {
Community,
Pro,
Enterprise,
}
pub async fn get_license_plan() -> LicensePlan {
// Implementation is not open source
return Community;
}
#[derive(Deserialize)]
#[serde(untagged)]
pub enum CriticalErrorChannel {}
pub enum CriticalAlertKind {
#[cfg(feature = "enterprise")]
CriticalError,
#[cfg(feature = "enterprise")]
RecoveredCriticalError,
}
#[cfg(feature = "enterprise")]
pub async fn send_critical_alert(
_error_message: String,
_db: &DB,
_kind: CriticalAlertKind,
_channels: Option<Vec<CriticalErrorChannel>>,
) {
}
#[cfg(feature = "enterprise")]
pub async fn schedule_key_renewal(_http_client: &reqwest::Client, _db: &crate::db::DB) -> () {
// Implementation is not open source
}
#[cfg(feature = "enterprise")]
pub async fn renew_license_key(
_http_client: &reqwest::Client,
_db: &crate::db::DB,
_key: Option<String>,
_manual: bool,
) -> String {
// Implementation is not open source
"".to_string()
}
#[cfg(feature = "enterprise")]
pub async fn create_customer_portal_session(
_http_client: &reqwest::Client,
_key: Option<String>,
) -> error::Result<String> {
// Implementation is not open source
Ok("".to_string())
}
#[cfg(feature = "enterprise")]
pub async fn worker_groups_alerts(_db: &DB) {}

View File

@@ -0,0 +1,76 @@
#[cfg(feature = "enterprise")]
use crate::db::DB;
use crate::ee::LicensePlan::Community;
#[cfg(feature = "enterprise")]
use crate::error;
use serde::Deserialize;
use std::sync::Arc;
use tokio::sync::RwLock;
lazy_static::lazy_static! {
pub static ref LICENSE_KEY_VALID: Arc<RwLock<bool>> = Arc::new(RwLock::new(true));
pub static ref LICENSE_KEY_ID: Arc<RwLock<String>> = Arc::new(RwLock::new("".to_string()));
pub static ref LICENSE_KEY: Arc<RwLock<String>> = Arc::new(RwLock::new("".to_string()));
}
pub enum LicensePlan {
Community,
Pro,
Enterprise,
}
pub async fn get_license_plan() -> LicensePlan {
// Implementation is not open source
return Community;
}
#[derive(Deserialize)]
#[serde(untagged)]
pub enum CriticalErrorChannel {
Email { email: String },
Slack { slack_channel: String },
}
pub enum CriticalAlertKind {
#[cfg(feature = "enterprise")]
CriticalError,
#[cfg(feature = "enterprise")]
RecoveredCriticalError,
}
#[cfg(feature = "enterprise")]
pub async fn send_critical_alert(
_error_message: String,
_db: &DB,
_kind: CriticalAlertKind,
_channels: Option<Vec<CriticalErrorChannel>>,
) {
}
#[cfg(feature = "enterprise")]
pub async fn schedule_key_renewal(_http_client: &reqwest::Client, _db: &crate::db::DB) -> () {
// Implementation is not open source
}
#[cfg(feature = "enterprise")]
pub async fn renew_license_key(
_http_client: &reqwest::Client,
_db: &crate::db::DB,
_key: Option<String>,
_manual: bool,
) -> String {
// Implementation is not open source
"".to_string()
}
#[cfg(feature = "enterprise")]
pub async fn create_customer_portal_session(
_http_client: &reqwest::Client,
_key: Option<String>,
) -> error::Result<String> {
// Implementation is not open source
Ok("".to_string())
}
#[cfg(feature = "enterprise")]
pub async fn worker_groups_alerts(_db: &DB) {}

View File

@@ -32,8 +32,6 @@ pub mod job_s3_helpers_ee;
pub mod jobs;
pub mod more_serde;
pub mod oauth2;
#[cfg(feature = "openidconnect")]
pub mod oidc_ee;
pub mod s3_helpers;
pub mod auth;

View File

@@ -1,170 +0,0 @@
/*
* Author: Ruben Fiszel
* Copyright: Windmill Labs, Inc 2023
* This file and its contents are licensed under the AGPLv3 License.
* Please see the included NOTICE for copyright information and
* LICENSE-AGPL for a copy of the license.
*/
#[cfg(feature = "openidconnect")]
use anyhow;
#[cfg(feature = "openidconnect")]
use std::process::Command;
#[cfg(feature = "openidconnect")]
use openidconnect::{
core::{CoreJwsSigningAlgorithm, CoreRsaPrivateSigningKey},
IssuerUrl, JsonWebKeyId,
};
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
use openidconnect::{
core::{
CoreClaimName, CoreJsonWebKeySet, CoreProviderMetadata, CoreResponseType,
CoreSubjectIdentifierType,
},
AuthUrl, EmptyAdditionalProviderMetadata, JsonWebKeySetUrl, ResponseTypes,
};
#[cfg(feature = "openidconnect")]
use openidconnect::AdditionalClaims;
#[cfg(feature = "openidconnect")]
use crate::db::DB;
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
use axum::extract::Path;
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
use axum::routing::{get, post};
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
use axum::Extension;
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
use axum::Json;
use axum::Router;
use serde::{Deserialize, Serialize};
#[cfg(all(feature = "enterprise", feature = "openidconnect"))]
pub async fn generate_id_token<T: AdditionalClaims>(
db: &DB,
claim: T,
audience: String,
identifier: String,
email: Option<String>,
) -> crate::error::Result<String> {
use chrono::{Duration, Utc};
use openidconnect::{
core::{CoreGenderClaim, CoreJsonWebKeyType, CoreJweContentEncryptionAlgorithm},
Audience, EndUserEmail, IdToken, IdTokenClaims, StandardClaims, SubjectIdentifier,
};
let private_key = get_private_key(&db).await?;
let issue_url = format!("{}/api/oidc/", crate::BASE_URL.read().await.clone());
let id_token = IdToken::<
T,
CoreGenderClaim,
CoreJweContentEncryptionAlgorithm,
CoreJwsSigningAlgorithm,
CoreJsonWebKeyType,
>::new(
IdTokenClaims::<T, CoreGenderClaim>::new(
// Specify the issuer URL for the OpenID Connect Provider.
IssuerUrl::new(issue_url)
.map_err(|e| anyhow::anyhow!("Failed to generate IssueUrl: {}", e))?,
// The audience is usually a single entry with the client ID of the client for whom
// the ID token is intended. This is a required claim.
vec![Audience::new(audience)],
// The ID token expiration is usually much shorter than that of the access or refresh
// tokens issued to clients.
Utc::now() + Duration::try_hours(48).unwrap(),
// The issue time is usually the current time.
Utc::now(),
// Set the standard claims defined by the OpenID Connect Core spec.
StandardClaims::new(
// Stable subject identifiers are recommended in place of e-mail addresses or other
// potentially unstable identifiers. This is the only required claim.
SubjectIdentifier::new(identifier),
)
// Optional: specify the user's e-mail address. This should only be provided if the
// client has been granted the 'profile' or 'email' scopes.
.set_email(email.map(|x| EndUserEmail::new(x)))
// Optional: specify whether the provider has verified the user's e-mail address.
.set_email_verified(Some(true)),
// OpenID Connect Providers may supply custom claims by providing a struct that
// implements the AdditionalClaims trait. This requires manually using the
// generic IdTokenClaims struct rather than the CoreIdTokenClaims type alias,
// however.
claim,
),
// The private key used for signing the ID token. For confidential clients (those able
// to maintain a client secret), a CoreHmacKey can also be used, in conjunction
// with one of the CoreJwsSigningAlgorithm::HmacSha* signing algorithms. When using an
// HMAC-based signing algorithm, the UTF-8 representation of the client secret should
// be used as the HMAC key.
&CoreRsaPrivateSigningKey::from_pem(
&private_key,
Some(JsonWebKeyId::new("windmill".to_string())),
)
.map_err(|e| anyhow::anyhow!("Invalid private key: {}", e))?,
// Uses the RS256 signature algorithm. This crate supports any RS*, PS*, or HS*
// signature algorithm.
CoreJwsSigningAlgorithm::RsaSsaPkcs1V15Sha256,
// When returning the ID token alongside an access token (e.g., in the Authorization Code
// flow), it is recommended to pass the access token here to set the `at_hash` claim
// automatically.
None,
// When returning the ID token alongside an authorization code (e.g., in the implicit
// flow), it is recommended to pass the authorization code here to set the `c_hash` claim
// automatically.
None,
)
.map_err(|e| anyhow::anyhow!("Failed to generate token: {}", e))?;
Ok(id_token.to_string())
}
#[cfg(feature = "openidconnect")]
pub async fn get_private_key(db: &DB) -> anyhow::Result<String> {
let key = sqlx::query_scalar!(
"SELECT value->>'private_key' FROM global_settings WHERE name = 'rsa_keys'",
)
.fetch_optional(db)
.await?
.flatten();
if let Some(key) = key {
return Ok(key);
} else {
let keys = gen_pems(db).await?;
return Ok(keys.private_key);
}
}
#[cfg(feature = "openidconnect")]
async fn gen_pems(db: &DB) -> anyhow::Result<Keys> {
let private_key_cmd = Command::new("openssl")
.arg("genrsa")
.arg("--traditional")
.arg("2048")
.output()
.expect("failed to execute process");
let private_key = String::from_utf8(private_key_cmd.stdout).unwrap();
tracing::debug!("Generated private key: {}", private_key);
let keys = Keys { private_key };
sqlx::query!(
"INSERT INTO global_settings (name, value) VALUES ('rsa_keys', $1)",
serde_json::to_value(&keys).unwrap()
)
.execute(db)
.await?;
Ok(keys)
}
#[derive(Debug, Clone, serde::Serialize)]
struct Keys {
private_key: String,
}

View File

@@ -22,7 +22,6 @@ use tokio::sync::RwLock;
lazy_static::lazy_static! {
pub static ref OBJECT_STORE_CACHE_SETTINGS: Arc<RwLock<Option<Arc<dyn ObjectStore>>>> = Arc::new(RwLock::new(None));
pub static ref OBJECT_STORE_OIDC_SETTINGS: Arc<RwLock<Option<Arc<S3AwsOidcResource>>>> = Arc::new(RwLock::new(None));
}
#[derive(Serialize, Deserialize, Debug)]
@@ -357,42 +356,18 @@ pub enum ObjectStoreSettings {
pub enum ObjectSettings {
S3(S3Settings),
Azure(AzureBlobResource),
AwsOidc(S3AwsOidcResource),
}
#[cfg(feature = "parquet")]
pub async fn build_object_store_from_settings(
settings: ObjectSettings,
) -> error::Result<Arc<dyn ObjectStore>> {
use crate::oidc_ee::generate_id_token;
match settings {
ObjectSettings::S3(s3_settings) => build_s3_client_from_settings(s3_settings).await,
ObjectSettings::Azure(azure_settings) => {
let azure_blob_resource = azure_settings;
build_azure_blob_client(&azure_blob_resource)
}
ObjectSettings::AwsOidc(aws_oidc_settings) => {
#[cfg(feature = "openidconnect")]
{
let token_fn = |audience: String| async move {
generate_id_token(
db,
claim,
aws_oidc_settings.audience,
"windmill_instance",
"instance_storage@windmill.dev",
)
};
todo!()
}
#[cfg(not(feature = "openidconnect"))]
{
return Err(error::Error::InternalErr(
"OpenID Connect is not enabled".to_string(),
));
}
}
}
}

View File

@@ -303,14 +303,14 @@ fn parse_file<T: FromStr>(path: &str) -> Option<T> {
.flatten()
}
pub struct TypeScriptAnnotations {
pub struct Annotations {
pub npm_mode: bool,
pub nodejs_mode: bool,
pub native_mode: bool,
pub nobundling: bool,
}
pub fn get_annotation_ts(inner_content: &str) -> TypeScriptAnnotations {
pub fn get_annotation(inner_content: &str) -> Annotations {
let annotations = inner_content
.lines()
.take_while(|x| x.starts_with("//"))
@@ -324,25 +324,7 @@ pub fn get_annotation_ts(inner_content: &str) -> TypeScriptAnnotations {
let nobundling: bool =
annotations.contains(&"nobundling".to_string()) || nodejs_mode || *DISABLE_BUNDLING;
TypeScriptAnnotations { npm_mode, nodejs_mode, native_mode, nobundling }
}
pub struct PythonAnnotations {
pub no_uv: bool,
pub no_cache: bool,
}
pub fn get_annotation_python(inner_content: &str) -> PythonAnnotations {
let annotations = inner_content
.lines()
.take_while(|x| x.starts_with("#"))
.map(|x| x.to_string().replace("#", "").trim().to_string())
.collect_vec();
let no_uv: bool = annotations.contains(&"no_uv".to_string());
let no_cache: bool = annotations.contains(&"no_cache".to_string());
PythonAnnotations { no_uv, no_cache }
Annotations { npm_mode, nodejs_mode, native_mode, nobundling }
}
pub struct SqlAnnotations {
@@ -467,13 +449,10 @@ pub async fn save_cache(
fn write_binary_file(main_path: &str, byts: &mut bytes::Bytes) -> error::Result<()> {
use std::fs::{File, Permissions};
use std::io::Write;
#[cfg(unix)]
use std::os::unix::fs::PermissionsExt;
let mut file = File::create(main_path)?;
file.write_all(byts)?;
#[cfg(unix)]
file.set_permissions(Permissions::from_mode(0o755))?;
file.flush()?;
Ok(())

View File

@@ -29,4 +29,3 @@ tempfile.workspace = true
bytes.workspace = true
object_store = { workspace = true, optional = true}
tokio-tar.workspace = true
lazy_static.workspace = true

View File

@@ -11,14 +11,13 @@ path = "src/lib.rs"
[features]
default = []
prometheus = ["dep:prometheus", "windmill-common/prometheus"]
enterprise = ["windmill-queue/enterprise", "windmill-git-sync/enterprise", "windmill-common/enterprise", "dep:gcp_auth", "dep:pem", "dep:tiberius", "dep:tokio-util"]
enterprise = ["windmill-queue/enterprise", "windmill-git-sync/enterprise", "windmill-common/enterprise", "dep:gcp_auth", "dep:pem", "dep:tiberius", "dep:tokio-util", "dep:openidconnect"]
benchmark = ["windmill-queue/benchmark", "windmill-common/benchmark"]
flamegraph = []
parquet = ["windmill-common/parquet", "dep:object_store"]
flow_testing = []
cloud = []
sqlx = []
deno_core = ["dep:deno_fetch", "dep:deno_webidl", "dep:deno_web", "dep:deno_net", "dep:deno_console", "dep:deno_url", "dep:deno_core", "dep:deno_ast", "dep:deno_tls"]
[dependencies]
windmill-queue.workspace = true
@@ -60,15 +59,15 @@ once_cell.workspace = true
rsmq_async.workspace = true
tokio-postgres.workspace = true
bit-vec.workspace = true
deno_fetch = { workspace = true, optional = true }
deno_webidl = { workspace = true, optional = true }
deno_web = { workspace = true, optional = true }
deno_net = { workspace = true, optional = true }
deno_console = { workspace = true, optional = true }
deno_url = { workspace = true, optional = true }
deno_core = { workspace = true, optional = true }
deno_ast = { workspace = true, optional = true }
deno_tls = { workspace = true, optional = true }
deno_fetch.workspace = true
deno_webidl.workspace = true
deno_web.workspace = true
deno_net.workspace = true
deno_console.workspace = true
deno_url.workspace = true
deno_core.workspace = true
deno_ast.workspace = true
deno_tls.workspace = true
postgres-native-tls.workspace = true
native-tls.workspace = true
mysql_async.workspace = true
@@ -85,21 +84,18 @@ reqwest.workspace = true
hex.workspace = true
tiberius = { workspace = true, optional = true }
tokio-util = { workspace = true, optional = true }
openidconnect = { workspace = true, optional = true}
tar.workspace = true
object_store = { workspace = true, optional = true}
convert_case.workspace = true
yaml-rust.workspace = true
swc_ecma_parser.workspace = true
[build-dependencies]
deno_fetch = { workspace = true, optional = true }
deno_webidl = { workspace = true, optional = true }
deno_web = { workspace = true, optional = true }
deno_net = { workspace = true, optional = true }
deno_console = { workspace = true, optional = true }
deno_url = { workspace = true, optional = true }
deno_core = { workspace = true, optional = true }
deno_ast = { workspace = true, optional = true }
deno_tls = { workspace = true, optional = true }
deno_fetch.workspace = true
deno_webidl.workspace = true
deno_web.workspace = true
deno_console.workspace = true
deno_url.workspace = true
deno_core.workspace = true
deno_net.workspace = true
zstd.workspace = true

View File

@@ -1,24 +1,13 @@
#[cfg(feature = "deno_core")]
use deno_fetch::FetchPermissions;
#[cfg(feature = "deno_core")]
use deno_net::NetPermissions;
#[cfg(feature = "deno_core")]
use deno_web::{BlobStore, TimersPermission};
#[cfg(feature = "deno_core")]
use std::borrow::Cow;
#[cfg(feature = "deno_core")]
use std::env;
#[cfg(feature = "deno_core")]
use std::io::Write;
#[cfg(feature = "deno_core")]
use std::path::{Path, PathBuf};
#[cfg(feature = "deno_core")]
use std::path::PathBuf;
use std::sync::Arc;
// #[cfg(feature = "deno_core")]
pub struct PermissionsContainer;
#[cfg(feature = "deno_core")]
impl FetchPermissions for PermissionsContainer {
#[inline(always)]
fn check_net_url(
@@ -26,20 +15,19 @@ impl FetchPermissions for PermissionsContainer {
_url: &deno_core::url::Url,
_api_name: &str,
) -> Result<(), deno_core::error::AnyError> {
unreachable!("snapshotting")
Ok(())
}
#[inline(always)]
fn check_read<'a>(
fn check_read(
&mut self,
_p: &'a std::path::Path,
_p: &std::path::Path,
_api_name: &str,
) -> Result<Cow<'a, Path>, deno_core::error::AnyError> {
unreachable!("snapshotting")
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
}
#[cfg(feature = "deno_core")]
impl TimersPermission for PermissionsContainer {
#[inline(always)]
fn allow_hrtime(&mut self) -> bool {
@@ -47,22 +35,21 @@ impl TimersPermission for PermissionsContainer {
}
}
#[cfg(feature = "deno_core")]
impl NetPermissions for PermissionsContainer {
fn check_read<'a>(
fn check_read(
&mut self,
_p: &'a str,
_p: &std::path::Path,
_api_name: &str,
) -> Result<PathBuf, deno_core::error::AnyError> {
unreachable!("snapshotting")
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
fn check_write<'a>(
fn check_write(
&mut self,
_p: &'a str,
_p: &std::path::Path,
_api_name: &str,
) -> Result<PathBuf, deno_core::error::AnyError> {
unreachable!("snapshotting")
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
fn check_net<T: AsRef<str>>(
@@ -70,26 +57,16 @@ impl NetPermissions for PermissionsContainer {
_host: &(T, Option<u16>),
_api_name: &str,
) -> Result<(), deno_core::error::AnyError> {
unreachable!("snapshotting")
}
fn check_write_path<'a>(
&mut self,
_: &'a Path,
_: &str,
) -> Result<Cow<'a, Path>, deno_core::anyhow::Error> {
todo!()
Ok(())
}
}
#[cfg(feature = "deno_core")]
deno_core::extension!(
fetch,
esm_entry_point = "ext:fetch/src/runtime.js",
esm = ["src/runtime.js"],
);
#[cfg(feature = "deno_core")]
fn main() {
println!("cargo:rustc-env=TARGET={}", env::var("TARGET").unwrap());
println!("cargo:rustc-env=PROFILE={}", env::var("PROFILE").unwrap());
@@ -145,6 +122,3 @@ fn main() {
println!("cargo:rerun-if-changed={}", path.display());
}
}
#[cfg(not(feature = "deno_core"))]
fn main() {}

View File

@@ -97,13 +97,6 @@ mount {
is_bind: true
}
mount {
dst: "/dev/shm"
fstype: "tmpfs"
rw: true
is_bind: false
}
mount {
src: "/dev/random"
dst: "/dev/random"

View File

@@ -1,4 +1,3 @@
#[cfg(unix)]
use std::{
collections::HashMap,
os::unix::fs::PermissionsExt,
@@ -6,13 +5,6 @@ use std::{
process::Stdio,
};
#[cfg(windows)]
use std::{
collections::HashMap,
path::{Path, PathBuf},
process::Stdio,
};
use anyhow::anyhow;
use itertools::Itertools;
use serde_json::value::RawValue;
@@ -33,7 +25,7 @@ use crate::{
OccupancyMetrics,
},
handle_child::handle_child,
python_executor::{create_dependencies_dir, handle_python_reqs, uv_pip_compile},
python_executor::{create_dependencies_dir, handle_python_reqs, pip_compile},
AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV,
TZ_ENV,
};
@@ -80,7 +72,7 @@ async fn handle_ansible_python_deps(
if requirements.is_empty() {
"".to_string()
} else {
uv_pip_compile(
pip_compile(
job_id,
&requirements,
mem_peak,
@@ -90,8 +82,6 @@ async fn handle_ansible_python_deps(
worker_name,
w_id,
&mut Some(occupancy_metrics),
false,
false,
)
.await
.map_err(|e| {
@@ -388,7 +378,6 @@ fi
let file = write_file(job_dir, "wrapper.sh", &wrapper)?;
#[cfg(unix)]
file.metadata()?.permissions().set_mode(0o777);
// let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());

View File

@@ -32,9 +32,6 @@ use crate::{
POWERSHELL_CACHE_DIR, POWERSHELL_PATH, TZ_ENV,
};
#[cfg(windows)]
use crate::SYSTEM_ROOT;
lazy_static::lazy_static! {
pub static ref ANSI_ESCAPE_RE: Regex = Regex::new(r"\x1b\[[0-9;]*m").unwrap();
@@ -229,19 +226,13 @@ pub async fn handle_powershell_job(
.collect::<Vec<_>>()
};
#[cfg(windows)]
let split_char = '\\';
#[cfg(unix)]
let split_char = '/';
let installed_modules = fs::read_dir(POWERSHELL_CACHE_DIR)?
.filter_map(|x| {
x.ok().map(|x| {
x.path()
.display()
.to_string()
.split(split_char)
.split('/')
.last()
.unwrap_or_default()
.to_lowercase()
@@ -298,83 +289,34 @@ pub async fn handle_powershell_job(
append_logs(&job.id, &job.workspace_id, logs2, db).await;
// make sure default (only allhostsallusers) modules are loaded, disable autoload (cache can be large to explore especially on cloud) and add /tmp/windmill/cache to PSModulePath
#[cfg(unix)]
let profile = format!(
"$PSModuleAutoloadingPreference = 'None'
$PSModulePathBackup = $env:PSModulePath
$env:PSModulePath = \"$PSHome/Modules\"
$env:PSModulePath = ($Env:PSModulePath -split ':')[-1]
Get-Module -ListAvailable | Import-Module
$env:PSModulePath = \"{}:$PSModulePathBackup\"",
POWERSHELL_CACHE_DIR
);
#[cfg(windows)]
let profile = format!(
"$PSModuleAutoloadingPreference = 'None'
$PSModulePathBackup = $env:PSModulePath
$env:PSModulePath = \"C:\\Program Files\\PowerShell\\7\\Modules\"
Get-Module -ListAvailable | Import-Module
$env:PSModulePath = \"{};$PSModulePathBackup\"",
POWERSHELL_CACHE_DIR
);
// NOTE: powershell error handling / termination is quite tricky compared to bash
// here we're trying to catch terminating errors and propagate the exit code
// to the caller such that the job will be marked as failed. It's up to the user
// to catch specific errors in their script not caught by the below as there is no
// generic set -eu as in bash
let strict_termination_start = "$ErrorActionPreference = 'Stop'\n\
Set-StrictMode -Version Latest\n\
try {\n";
let strict_termination_end = "\n\
} catch {\n\
Write-Output \"An error occurred:\n\"\
Write-Output $_
exit 1\n\
}\n";
// make sure param() is first
let param_match = windmill_parser_bash::RE_POWERSHELL_PARAM.find(&content);
let content: String = if let Some(param_match) = param_match {
let param_match = param_match.as_str();
format!(
"{}\n{}\n{}\n{}\n{}",
"{}\n{}\n{}",
param_match,
profile,
strict_termination_start,
content.replace(param_match, ""),
strict_termination_end
content.replace(param_match, "")
)
} else {
format!("{}\n{}", profile, content)
};
write_file(job_dir, "main.ps1", content.as_str())?;
#[cfg(unix)]
write_file(
job_dir,
"wrapper.sh",
&format!("set -o pipefail\nset -e\nmkfifo bp\ncat bp | tail -1 > ./result2.out &\n{} -F ./main.ps1 \"$@\" 2>&1 | tee bp\nwait $!", POWERSHELL_PATH.as_str()),
)?;
#[cfg(windows)]
write_file(
job_dir,
"wrapper.ps1",
&format!(
"param([string[]]$args)\n\
$ErrorActionPreference = 'Stop'\n\
$pipe = New-TemporaryFile\n\
& \"{}\" -File ./main.ps1 @args 2>&1 | Tee-Object -FilePath $pipe\n\
Get-Content -Path $pipe | Select-Object -Last 1 | Set-Content -Path './result2.out'\n\
Remove-Item $pipe\n\
exit $LASTEXITCODE\n",
POWERSHELL_PATH.as_str()
),
)?;
let token = client.get_token().await;
let mut reserved_variables = get_reserved_variables(job, &token, db).await?;
reserved_variables.insert("RUST_LOG".to_string(), "info".to_string());
@@ -413,24 +355,10 @@ $env:PSModulePath = \"{};$PSModulePathBackup\"",
.stderr(Stdio::piped())
.spawn()?
} else {
let mut cmd;
let mut cmd_args;
#[cfg(unix)]
{
cmd_args = vec!["wrapper.sh"];
cmd_args.extend(pwsh_args.iter().map(|x| x.as_str()));
cmd = Command::new(BIN_BASH.as_str());
}
#[cfg(windows)]
{
cmd_args = vec![r".\wrapper.ps1".to_string()];
cmd_args.extend(pwsh_args.iter().map(|x| x.replace("--", "-")));
cmd = Command::new(POWERSHELL_PATH.as_str());
}
cmd.current_dir(job_dir)
let mut cmd_args = vec!["wrapper.sh"];
cmd_args.extend(pwsh_args.iter().map(|x| x.as_str()));
Command::new(BIN_BASH.as_str())
.current_dir(job_dir)
.env_clear()
.envs(envs)
.envs(reserved_variables)
@@ -438,53 +366,11 @@ $env:PSModulePath = \"{};$PSModulePathBackup\"",
.env("PATH", PATH_ENV.as_str())
.env("BASE_INTERNAL_URL", base_internal_url)
.env("HOME", HOME_ENV.as_str())
.args(&cmd_args)
.args(cmd_args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
{
cmd.env("SystemRoot", SYSTEM_ROOT.as_str())
.env(
"LOCALAPPDATA",
std::env::var("LOCALAPPDATA")
.unwrap_or_else(|_| format!("{}\\AppData\\Local", HOME_ENV.as_str())),
)
.env(
"ProgramData",
std::env::var("ProgramData")
.unwrap_or_else(|_| String::from("C:\\ProgramData")),
)
.env(
"ProgramFiles",
std::env::var("ProgramFiles")
.unwrap_or_else(|_| String::from("C:\\Program Files")),
)
.env(
"ProgramFiles(x86)",
std::env::var("ProgramFiles(x86)")
.unwrap_or_else(|_| String::from("C:\\Program Files (x86)")),
)
.env(
"ProgramW6432",
std::env::var("ProgramW6432")
.unwrap_or_else(|_| String::from("C:\\Program Files")),
)
.env(
"TMP",
std::env::var("TMP").unwrap_or_else(|_| String::from("/tmp")),
)
.env(
"PATHEXT",
std::env::var("PATHEXT").unwrap_or_else(|_| {
String::from(".COM;.EXE;.BAT;.CMD;.VBS;.VBE;.JS;.JSE;.WSF;.WSH;.MSC;.CPL")
}),
);
}
cmd.spawn()?
.stderr(Stdio::piped())
.spawn()?
};
handle_child(
&job.id,
db,

View File

@@ -1,14 +1,8 @@
#[cfg(feature = "deno_core")]
use std::time::Instant;
use std::{collections::HashMap, fs, io, path::Path, process::Stdio};
use std::{collections::HashMap, fs, io, path::Path, process::Stdio, time::Instant};
use base64::Engine;
use itertools::Itertools;
#[cfg(not(feature = "deno_core"))]
use serde_json::value::to_raw_value;
use serde_json::value::RawValue;
use sha2::Digest;
use uuid::Uuid;
use windmill_parser_ts::remove_pinned_imports;
@@ -29,9 +23,6 @@ use crate::{
NODE_PATH, NPM_CONFIG_REGISTRY, NPM_PATH, NSJAIL_PATH, PATH_ENV, TZ_ENV,
};
#[cfg(windows)]
use crate::SYSTEM_ROOT;
use tokio::{fs::File, process::Command};
use tokio::io::AsyncReadExt;
@@ -47,7 +38,7 @@ use windmill_common::{
get_latest_hash_for_path,
jobs::{QueuedJob, PREPROCESSOR_FAKE_ENTRYPOINT},
scripts::ScriptLang,
worker::{exists_in_cache, get_annotation_ts, save_cache, write_file},
worker::{exists_in_cache, get_annotation, save_cache, write_file},
DB,
};
@@ -121,9 +112,6 @@ pub async fn gen_bun_lockfile(
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
child_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
let mut child_process = start_child_process(child_cmd, &*BUN_PATH).await?;
if let Some(db) = db {
@@ -257,9 +245,6 @@ pub async fn install_bun_lockfile(
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
child_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
let mut npm_logs = if npm_mode {
"NPM mode\n".to_string()
} else {
@@ -468,10 +453,6 @@ pub async fn generate_wrapper_mjs(
.args(vec!["run", "node_builder.ts"])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
child.env("SystemRoot", SYSTEM_ROOT.as_str());
let child_process = start_child_process(child, &*BUN_PATH).await?;
handle_child(
job_id,
@@ -506,7 +487,7 @@ pub async fn generate_bun_bundle(
mem_peak: &mut i32,
canceled_by: &mut Option<CanceledBy>,
common_bun_proc_envs: &HashMap<String, String>,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
occupancy_metrics: &mut OccupancyMetrics,
) -> Result<()> {
let mut child = Command::new(&*BUN_PATH);
child
@@ -517,10 +498,6 @@ pub async fn generate_bun_bundle(
.args(vec!["run", "node_builder.ts"])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
child.env("SystemRoot", SYSTEM_ROOT.as_str());
let mut child_process = start_child_process(child, &*BUN_PATH).await?;
if let Some(db) = db {
handle_child(
@@ -535,7 +512,7 @@ pub async fn generate_bun_bundle(
"bun build",
timeout,
false,
occupancy_metrics,
&mut Some(occupancy_metrics),
)
.await?;
} else {
@@ -563,11 +540,7 @@ pub async fn pull_codebase(w_id: &str, id: &str, job_dir: &str) -> Result<()> {
if is_tar {
extract_tar(fs::read(bun_cache_path)?.into(), job_dir).await?;
} else {
#[cfg(unix)]
tokio::fs::symlink(&bun_cache_path, dst).await?;
#[cfg(windows)]
std::os::windows::fs::symlink_dir(&bun_cache_path, &dst)?;
}
} else if let Some(os) = windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS
.read()
@@ -580,11 +553,7 @@ pub async fn pull_codebase(w_id: &str, id: &str, job_dir: &str) -> Result<()> {
if is_tar {
extract_tar(bytes, job_dir).await?;
} else {
#[cfg(unix)]
tokio::fs::symlink(bun_cache_path, dst).await?;
#[cfg(windows)]
std::os::windows::fs::symlink_dir(&bun_cache_path, &dst)?;
}
// extract_tar(bytes, job_dir).await?;
@@ -650,7 +619,7 @@ pub async fn prebundle_bun_script(
base_internal_url: &str,
worker_name: &str,
token: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
occupancy_metrics: &mut OccupancyMetrics,
) -> Result<()> {
let (local_path, remote_path) = compute_bundle_local_and_remote_path(
inner_content,
@@ -663,7 +632,7 @@ pub async fn prebundle_bun_script(
if exists_in_cache(&local_path, &remote_path).await {
return Ok(());
}
let annotation = get_annotation_ts(inner_content);
let annotation = get_annotation(inner_content);
if annotation.nobundling {
return Ok(());
}
@@ -753,10 +722,6 @@ async fn compute_bundle_local_and_remote_path(
let hash = windmill_common::utils::calculate_hash(&input_src);
let local_path = format!("{BUN_BUNDLE_CACHE_DIR}/{hash}");
#[cfg(windows)]
let local_path = local_path.replace("/tmp", r"C:\tmp").replace("/", r"\");
let remote_path = format!("{BUN_BUNDLE_OBJECT_STORE_PREFIX}{hash}");
(local_path, remote_path)
}
@@ -800,7 +765,7 @@ pub async fn handle_bun_job(
new_args: &mut Option<HashMap<String, Box<RawValue>>>,
occupancy_metrics: &mut OccupancyMetrics,
) -> error::Result<Box<RawValue>> {
let mut annotation = windmill_common::worker::get_annotation_ts(inner_content);
let mut annotation = windmill_common::worker::get_annotation(inner_content);
let (mut has_bundle_cache, cache_logs, local_path, remote_path) =
if requirements_o.is_some() && !annotation.nobundling && codebase.is_none() {
@@ -850,23 +815,10 @@ pub async fn handle_bun_job(
));
}
let mut gbuntar_name: Option<String> = None;
let mut gbuntar_name = None;
if has_bundle_cache {
let target;
let symlink;
#[cfg(unix)]
{
target = format!("{job_dir}/main.js");
symlink = std::os::unix::fs::symlink(&local_path, &target);
}
#[cfg(windows)]
{
target = format!("{job_dir}\\main.js");
symlink = std::os::windows::fs::symlink_dir(&local_path, &target);
}
symlink.map_err(|e| {
let target = format!("{job_dir}/main.js");
std::os::unix::fs::symlink(&local_path, &target).map_err(|e| {
error::Error::ExecutionErr(format!(
"could not copy cached binary from {local_path} to {job_dir}/main: {e:?}"
))
@@ -1105,16 +1057,6 @@ try {{
if (step_id) {{
err["step_id"] = step_id;
}}
const extra = {{}};
Object.getOwnPropertyNames(e).forEach((key) => {{
if (['line', 'name', 'stack', 'column', 'message', 'sourceURL', 'originalLine', 'originalColumn'].includes(key)) {{
return;
}}
extra[key] = e[key];
}});
if (Object.keys(extra).length > 0) {{
err["extra"] = extra;
}}
await fs.writeFile("result.json", JSON.stringify(err));
process.exit(1);
}}
@@ -1200,7 +1142,7 @@ try {{
mem_peak,
canceled_by,
&common_bun_proc_envs,
&mut Some(occupancy_metrics),
occupancy_metrics,
)
.await?;
if !local_path.is_empty() {
@@ -1248,58 +1190,51 @@ try {{
}
}
if annotation.native_mode {
#[cfg(not(feature = "deno_core"))]
return Ok(to_raw_value("").unwrap());
#[cfg(feature = "deno_core")]
{
let env_code = format!(
let env_code = format!(
"const process = {{ env: {{}} }};\nconst BASE_URL = '{base_internal_url}';\nconst BASE_INTERNAL_URL = '{base_internal_url}';\nprocess.env['BASE_URL'] = BASE_URL;process.env['BASE_INTERNAL_URL'] = BASE_INTERNAL_URL;\n{}",
reserved_variables
.iter()
.map(|(k, v)| format!("process.env['{}'] = '{}';\n", k, v))
.collect::<Vec<String>>()
.join("\n"));
let js_code = read_file_content(&format!("{job_dir}/main.js")).await?;
let started_at = Instant::now();
let args = crate::common::build_args_map(job, client, db)
.await?
.map(sqlx::types::Json);
let job_args = if args.is_some() {
args.as_ref()
} else {
job.args.as_ref()
};
let result = crate::js_eval::eval_fetch_timeout(
env_code,
inner_content.clone(),
js_code,
job_args,
job.id,
job.timeout,
db,
mem_peak,
canceled_by,
worker_name,
&job.workspace_id,
false,
occupancy_metrics,
)
.await?;
tracing::info!(
"Executed native code in {}ms",
started_at.elapsed().as_millis()
);
append_logs(
&job.id,
&job.workspace_id,
format!("{}\n{}", init_logs, result.1),
db,
)
.await;
return Ok(result.0);
}
let js_code = read_file_content(&format!("{job_dir}/main.js")).await?;
let started_at = Instant::now();
let args = crate::common::build_args_map(job, client, db)
.await?
.map(sqlx::types::Json);
let job_args = if args.is_some() {
args.as_ref()
} else {
job.args.as_ref()
};
let result = crate::js_eval::eval_fetch_timeout(
env_code,
inner_content.clone(),
js_code,
job_args,
job.id,
job.timeout,
db,
mem_peak,
canceled_by,
worker_name,
&job.workspace_id,
false,
occupancy_metrics,
)
.await?;
tracing::info!(
"Executed native code in {}ms",
started_at.elapsed().as_millis()
);
append_logs(
&job.id,
&job.workspace_id,
format!("{}\n{}", init_logs, result.1),
db,
)
.await;
return Ok(result.0);
}
append_logs(&job.id, &job.workspace_id, init_logs, db).await;
@@ -1391,10 +1326,6 @@ try {{
.args(vec!["--preserve-symlinks", &script_path])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
bun_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
bun_cmd
} else {
let script_path = format!("{job_dir}/wrapper.mjs");
@@ -1421,13 +1352,8 @@ try {{
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
bun_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
bun_cmd
};
start_child_process(
cmd,
if annotation.nodejs_mode {
@@ -1445,7 +1371,7 @@ try {{
mem_peak,
canceled_by,
child,
!*DISABLE_NSJAIL,
false,
worker_name,
&job.workspace_id,
"bun run",
@@ -1535,7 +1461,7 @@ pub async fn start_worker(
let common_bun_proc_envs: HashMap<String, String> =
get_common_bun_proc_envs(Some(&base_internal_url)).await;
let mut annotation = windmill_common::worker::get_annotation_ts(inner_content);
let mut annotation = windmill_common::worker::get_annotation(inner_content);
//TODO: remove this when bun dedicated workers work without issues
annotation.nodejs_mode = true;

View File

@@ -125,8 +125,7 @@ pub async fn generate_deno_lock(
"--unstable-worker-options",
"--unstable-http",
"--lock=lock.json",
"--frozen=false",
"--allow-import",
"--lock-write",
"--import-map",
&import_map_path,
"main.ts",
@@ -157,13 +156,10 @@ pub async fn generate_deno_lock(
}
let path_lock = format!("{job_dir}/lock.json");
if let Ok(mut file) = File::open(path_lock).await {
let mut req_content = "".to_string();
file.read_to_string(&mut req_content).await?;
Ok(req_content)
} else {
Ok("".to_string())
}
let mut file = File::open(path_lock).await?;
let mut req_content = "".to_string();
file.read_to_string(&mut req_content).await?;
Ok(req_content)
}
#[tracing::instrument(level = "trace", skip_all)]
@@ -370,10 +366,10 @@ try {{
} else if !*DISABLE_NSJAIL {
args.push("--allow-net");
args.push("--allow-sys");
args.push("--allow-hrtime");
args.push(allow_read.as_str());
args.push("--allow-write=./");
args.push("--allow-env");
args.push("--allow-import");
args.push("--allow-run=git,/usr/bin/chromium");
} else {
args.push("-A");

View File

@@ -225,12 +225,7 @@ func Run(req Req) (interface{{}}, error){{
}
} else {
let target = format!("{job_dir}/main");
#[cfg(unix)]
let symlink = std::os::unix::fs::symlink(&bin_path, &target);
#[cfg(windows)]
let symlink = std::os::windows::fs::symlink_dir(&bin_path, &target);
symlink.map_err(|e| {
std::os::unix::fs::symlink(&bin_path, &target).map_err(|e| {
Error::ExecutionErr(format!(
"could not copy cached binary from {bin_path} to {job_dir}/main: {e:?}"
))

View File

@@ -6,11 +6,7 @@ use nix::sys::signal::{self, Signal};
use nix::unistd::Pid;
use sqlx::{Pool, Postgres};
#[cfg(windows)]
use std::process::Stdio;
use tokio::fs::File;
#[cfg(windows)]
use tokio::process::Command;
use windmill_common::error::to_anyhow;
use windmill_common::error::{self, Error};
@@ -53,33 +49,6 @@ use crate::{MAX_RESULT_SIZE, MAX_WAIT_FOR_SIGINT, MAX_WAIT_FOR_SIGTERM};
lazy_static::lazy_static! {
pub static ref SLOW_LOGS: bool = std::env::var("SLOW_LOGS").ok().is_some_and(|x| x == "1" || x == "true");
}
// - kill windows process along with all child processes
#[cfg(windows)]
async fn kill_process_tree(pid: Option<u32>) -> Result<(), String> {
let pid = match pid {
Some(pid) => pid,
None => return Err("No PID provided to kill.".to_string()),
};
let output = Command::new("cmd")
.args(&["/C", "taskkill", "/PID", &pid.to_string(), "/T", "/F"])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()
.await
.map_err(|e| format!("Failed to execute taskkill: {}", e))?;
if output.status.success() {
Ok(())
} else {
Err(format!(
"Failed to kill process tree. Error: {}",
String::from_utf8_lossy(&output.stderr)
))
}
}
/// - wait until child exits and return with exit status
/// - read lines from stdout and stderr and append them to the "queue"."logs"
/// quitting early if output exceedes MAX_LOG_SIZE characters (not bytes)
@@ -227,26 +196,9 @@ pub async fn handle_child(
}
}
}
#[cfg(windows)]
{
let pid_to_kill = child.id();
match kill_process_tree(pid_to_kill).await {
Ok(_) => tracing::debug!(
"successfully killed process tree with PID: {:?}",
pid_to_kill
),
Err(e) => tracing::error!("failed to kill process tree: {:?}", e),
};
set_reason.await;
return Ok(Err(kill_reason));
}
#[cfg(unix)]
{
/* send SIGKILL and reap child process */
let (_, kill) = future::join(set_reason, child.kill()).await;
kill.map(|()| Err(kill_reason))
}
/* send SIGKILL and reap child process */
let (_, kill) = future::join(set_reason, child.kill()).await;
kill.map(|()| Err(kill_reason))
};
/* a future that reads output from the child and appends to the database */
@@ -414,30 +366,9 @@ async fn get_mem_peak(pid: Option<u32>, nsjail: bool) -> i32 {
return -1;
}
let pid = if nsjail {
// Read /proc/<nsjail_pid>/task/<nsjail_pid>/children and extract pid
let nsjail_pid = pid.unwrap();
let children_path = format!("/proc/{}/task/{}/children", nsjail_pid, nsjail_pid);
if let Ok(mut file) = File::open(children_path).await {
let mut contents = String::new();
if tokio::io::AsyncReadExt::read_to_string(&mut file, &mut contents)
.await
.is_ok()
{
if let Some(child_pid) = contents.split_whitespace().next() {
if let Ok(child_pid) = child_pid.parse::<u32>() {
child_pid
} else {
return -1;
}
} else {
return -1;
}
} else {
return -1;
}
} else {
return -1;
}
// This is a bit hacky, but the process id of the nsjail process is the pid of nsjail + 1.
// Ideally, we would get the number from fork() itself. This works in MOST cases.
pid.unwrap() + 1
} else {
pid.unwrap()
};

View File

@@ -1,5 +1,5 @@
use deno_ast::swc::parser::lexer::util::CharExt;
use itertools::Itertools;
use swc_ecma_parser::lexer::util::CharExt;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use object_store::path::Path;

View File

@@ -6,74 +6,56 @@
* LICENSE-AGPL for a copy of the license.
*/
#[cfg(feature = "deno_core")]
use std::{
borrow::Cow,
cell::RefCell,
collections::HashMap,
env,
io::{self, BufReader},
path::PathBuf,
rc::Rc,
sync::Arc,
};
use std::{collections::HashMap, sync::Arc};
#[cfg(feature = "deno_core")]
use deno_ast::ParseParams;
#[cfg(feature = "deno_core")]
use deno_core::{
error::AnyError,
op2, serde_v8, url,
v8::{self, IsolateHandle},
Extension, JsRuntime, OpState, PollEventLoopOptions, RuntimeOptions,
};
#[cfg(feature = "deno_core")]
use deno_fetch::FetchPermissions;
#[cfg(feature = "deno_core")]
use deno_net::NetPermissions;
#[cfg(feature = "deno_core")]
use deno_tls::{rustls::RootCertStore, rustls_pemfile};
#[cfg(feature = "deno_core")]
use deno_web::{BlobStore, TimersPermission};
#[cfg(feature = "deno_core")]
use itertools::Itertools;
use lazy_static::lazy_static;
use regex::Regex;
use serde_json::value::RawValue;
use sqlx::types::Json;
#[cfg(feature = "deno_core")]
use tokio::{
sync::{mpsc, oneshot},
time::timeout,
};
use uuid::Uuid;
#[cfg(feature = "deno_core")]
use windmill_common::error::Error;
use windmill_common::{flow_status::JobResult, DB};
use windmill_common::{error::Error, flow_status::JobResult, DB};
use windmill_queue::CanceledBy;
use crate::{common::OccupancyMetrics, AuthedClient};
#[cfg(feature = "deno_core")]
use crate::{common::unsafe_raw, handle_child::run_future_with_polling_update_job_poller};
use crate::{
common::{unsafe_raw, OccupancyMetrics},
handle_child::run_future_with_polling_update_job_poller,
AuthedClient,
};
#[derive(Debug, Clone)]
pub struct IdContext {
pub flow_job: Uuid,
#[allow(dead_code)]
pub steps_results: HashMap<String, JobResult>,
pub previous_id: String,
}
#[cfg(feature = "deno_core")]
pub struct ContainerRootCertStoreProvider {
root_cert_store: RootCertStore,
}
#[cfg(feature = "deno_core")]
impl ContainerRootCertStoreProvider {
fn new() -> ContainerRootCertStoreProvider {
return ContainerRootCertStoreProvider {
@@ -91,17 +73,14 @@ impl ContainerRootCertStoreProvider {
}
}
#[cfg(feature = "deno_core")]
impl deno_tls::RootCertStoreProvider for ContainerRootCertStoreProvider {
fn get_or_try_init(&self) -> Result<&RootCertStore, AnyError> {
Ok(&self.root_cert_store)
}
}
#[cfg(feature = "deno_core")]
pub struct PermissionsContainer;
#[cfg(feature = "deno_core")]
impl FetchPermissions for PermissionsContainer {
#[inline(always)]
fn check_net_url(
@@ -113,16 +92,15 @@ impl FetchPermissions for PermissionsContainer {
}
#[inline(always)]
fn check_read<'a>(
fn check_read(
&mut self,
p: &'a std::path::Path,
_p: &std::path::Path,
_api_name: &str,
) -> Result<Cow<'a, std::path::Path>, anyhow::Error> {
Ok(Cow::Borrowed(p))
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
}
#[cfg(feature = "deno_core")]
impl TimersPermission for PermissionsContainer {
#[inline(always)]
fn allow_hrtime(&mut self) -> bool {
@@ -130,22 +108,21 @@ impl TimersPermission for PermissionsContainer {
}
}
#[cfg(feature = "deno_core")]
impl NetPermissions for PermissionsContainer {
fn check_read<'a>(
fn check_read(
&mut self,
p: &'a str,
_p: &std::path::Path,
_api_name: &str,
) -> Result<PathBuf, deno_core::error::AnyError> {
Ok(PathBuf::from(p))
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
fn check_write<'a>(
fn check_write(
&mut self,
p: &'a str,
_p: &std::path::Path,
_api_name: &str,
) -> Result<PathBuf, deno_core::error::AnyError> {
Ok(PathBuf::from(p))
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
fn check_net<T: AsRef<str>>(
@@ -155,17 +132,8 @@ impl NetPermissions for PermissionsContainer {
) -> Result<(), deno_core::error::AnyError> {
Ok(())
}
fn check_write_path<'a>(
&mut self,
p: &'a std::path::Path,
_api_name: &str,
) -> Result<std::borrow::Cow<'a, std::path::Path>, AnyError> {
Ok(Cow::Borrowed(p))
}
}
#[cfg(feature = "deno_core")]
pub struct OptAuthedClient(Option<AuthedClient>);
pub async fn eval_timeout(
@@ -174,7 +142,7 @@ pub async fn eval_timeout(
flow_input: Option<mappable_rc::Marc<HashMap<String, Box<RawValue>>>>,
authed_client: Option<&AuthedClient>,
by_id: Option<IdContext>,
#[allow(unused_variables)] ctx: Option<Vec<(String, String)>>,
ctx: Option<Vec<(String, String)>>,
) -> anyhow::Result<Box<RawValue>> {
let expr = expr.trim().to_string();
@@ -244,133 +212,121 @@ pub async fn eval_timeout(
}
}
#[cfg(not(feature = "deno_core"))]
{
#[allow(unreachable_code)]
return Err(anyhow::anyhow!("Deno core is not enabled".to_string()).into());
}
let expr2 = expr.clone();
let (sender, mut receiver) = oneshot::channel::<IsolateHandle>();
let has_client = authed_client.is_some();
let authed_client = authed_client.cloned();
timeout(
std::time::Duration::from_millis(10000),
tokio::task::spawn_blocking(move || {
let mut ops = vec![op_get_context()];
#[cfg(feature = "deno_core")]
{
let expr2 = expr.clone();
let (sender, mut receiver) = oneshot::channel::<IsolateHandle>();
let has_client = authed_client.is_some();
let authed_client = authed_client.cloned();
return timeout(
std::time::Duration::from_millis(10000),
tokio::task::spawn_blocking(move || {
let mut ops = vec![op_get_context()];
if authed_client.is_some() {
ops.extend([
// An op for summing an array of numbers
// The op-layer automatically deserializes inputs
// and serializes the returned Result & value
op_variable(),
op_resource(),
])
}
if authed_client.is_some() {
ops.extend([
// An op for summing an array of numbers
// The op-layer automatically deserializes inputs
// and serializes the returned Result & value
op_variable(),
op_resource(),
])
}
if by_id.is_some() && authed_client.is_some() {
ops.push(op_get_result());
ops.push(op_get_id());
}
if by_id.is_some() && authed_client.is_some() {
ops.push(op_get_result());
ops.push(op_get_id());
}
let ext = Extension { name: "js_eval", ops: ops.into(), ..Default::default() };
let exts = vec![ext];
// Use our snapshot to provision our new runtime
let options = RuntimeOptions {
extensions: exts,
// startup_snapshot: Some(Snapshot::Static(buffer)),
..Default::default()
};
let mut context_keys = transform_context
.keys()
.filter(|x| expr.contains(&x.to_string()))
.map(|x| x.clone())
.collect_vec();
if !context_keys.contains(&"previous_result".to_string())
&& (p_ids.is_some() && p_ids.as_ref().unwrap().iter().any(|x| expr.contains(x)))
|| expr.contains("error")
{
// tracing::error!("PREVIOUS_RESULT");
context_keys.push("previous_result".to_string());
}
let has_flow_input = expr.contains("flow_input");
if has_flow_input {
context_keys.push("flow_input".to_string())
}
let mut js_runtime = JsRuntime::new(options);
{
let op_state = js_runtime.op_state();
let mut op_state = op_state.borrow_mut();
let mut client = authed_client.clone();
if let Some(client) = client.as_mut() {
client.force_client = Some(
reqwest::ClientBuilder::new()
.user_agent("windmill/beta")
.danger_accept_invalid_certs(
std::env::var("ACCEPT_INVALID_CERTS").is_ok(),
)
.build()
.unwrap(),
);
}
op_state.put(OptAuthedClient(client));
op_state.put(TransformContext {
flow_input: if has_flow_input { flow_input } else { None },
envs: transform_context
.into_iter()
.filter(|(a, _)| context_keys.contains(a))
.collect(),
})
}
sender
.send(js_runtime.v8_isolate().thread_safe_handle())
.map_err(|_| {
Error::ExecutionErr("impossible to send v8 isolate".to_string())
})?;
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
// pretty frail but this it to make the expr more user friendly and not require the user to write await
let expr = ["variable", "resource"]
.into_iter()
.fold(expr, replace_with_await);
let expr = replace_with_await_result(expr);
let r = runtime.block_on(eval(
&mut js_runtime,
&expr,
context_keys,
by_id,
has_client,
ctx,
))?;
Ok(r) as anyhow::Result<Box<RawValue>>
}),
)
.await
.map_err(|_| {
if let Ok(isolate) = receiver.try_recv() {
isolate.terminate_execution();
let ext = Extension { name: "js_eval", ops: ops.into(), ..Default::default() };
let exts = vec![ext];
// Use our snapshot to provision our new runtime
let options = RuntimeOptions {
extensions: exts,
// startup_snapshot: Some(Snapshot::Static(buffer)),
..Default::default()
};
Error::ExecutionErr(format!(
"The expression of evaluation `{expr2}` took too long to execute (>10000ms)"
))
})??;
}
let mut context_keys = transform_context
.keys()
.filter(|x| expr.contains(&x.to_string()))
.map(|x| x.clone())
.collect_vec();
if !context_keys.contains(&"previous_result".to_string())
&& (p_ids.is_some() && p_ids.as_ref().unwrap().iter().any(|x| expr.contains(x)))
|| expr.contains("error")
{
// tracing::error!("PREVIOUS_RESULT");
context_keys.push("previous_result".to_string());
}
let has_flow_input = expr.contains("flow_input");
if has_flow_input {
context_keys.push("flow_input".to_string())
}
let mut js_runtime = JsRuntime::new(options);
{
let op_state = js_runtime.op_state();
let mut op_state = op_state.borrow_mut();
let mut client = authed_client.clone();
if let Some(client) = client.as_mut() {
client.force_client = Some(
reqwest::ClientBuilder::new()
.user_agent("windmill/beta")
.danger_accept_invalid_certs(
std::env::var("ACCEPT_INVALID_CERTS").is_ok(),
)
.build()
.unwrap(),
);
}
op_state.put(OptAuthedClient(client));
op_state.put(TransformContext {
flow_input: if has_flow_input { flow_input } else { None },
envs: transform_context
.into_iter()
.filter(|(a, _)| context_keys.contains(a))
.collect(),
})
}
sender
.send(js_runtime.v8_isolate().thread_safe_handle())
.map_err(|_| Error::ExecutionErr("impossible to send v8 isolate".to_string()))?;
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
// pretty frail but this it to make the expr more user friendly and not require the user to write await
let expr = ["variable", "resource"]
.into_iter()
.fold(expr, replace_with_await);
let expr = replace_with_await_result(expr);
let r = runtime.block_on(eval(
&mut js_runtime,
&expr,
context_keys,
by_id,
has_client,
ctx,
))?;
Ok(r) as anyhow::Result<Box<RawValue>>
}),
)
.await
.map_err(|_| {
if let Ok(isolate) = receiver.try_recv() {
isolate.terminate_execution();
};
Error::ExecutionErr(format!(
"The expression of evaluation `{expr2}` took too long to execute (>10000ms)"
))
})??
}
#[cfg(feature = "deno_core")]
fn replace_with_await(expr: String, fn_name: &str) -> String {
let sep = format!("{}(", fn_name);
let mut split = expr.split(&sep);
@@ -389,12 +345,10 @@ lazy_static! {
Regex::new(r"^(https?)://(([^:@\s]+):([^:@\s]+)@)?([^:@\s]+)(:(\d+))?$").unwrap();
}
#[cfg(feature = "deno_core")]
fn replace_with_await_result(expr: String) -> String {
RE.replace_all(&expr, "(await $r)").to_string()
}
#[cfg(feature = "deno_core")]
fn add_closing_bracket(s: &str) -> String {
let mut s = s.to_string();
let mut level = 1;
@@ -414,7 +368,6 @@ fn add_closing_bracket(s: &str) -> String {
s
}
#[cfg(feature = "deno_core")]
async fn eval(
context: &mut JsRuntime,
expr: &str,
@@ -555,7 +508,6 @@ function get_from_env(name) {{
// }
// TODO: Can we a) share the api configuration here somehow or b) just implement this natively in deno, via the deno client?
#[cfg(feature = "deno_core")]
#[op2(async)]
#[string]
async fn op_variable(
@@ -570,7 +522,6 @@ async fn op_variable(
}
}
#[cfg(feature = "deno_core")]
#[op2(async)]
#[string]
async fn op_get_result(
@@ -589,7 +540,6 @@ async fn op_get_result(
}
}
#[cfg(feature = "deno_core")]
#[op2(async)]
#[string]
async fn op_get_id(
@@ -613,7 +563,6 @@ async fn op_get_id(
}
}
#[cfg(feature = "deno_core")]
#[op2(async)]
#[string]
async fn op_resource(
@@ -631,13 +580,11 @@ async fn op_resource(
}
}
#[cfg(feature = "deno_core")]
pub struct TransformContext {
pub envs: HashMap<String, Arc<Box<RawValue>>>,
pub flow_input: Option<mappable_rc::Marc<HashMap<String, Box<RawValue>>>>,
}
#[cfg(feature = "deno_core")]
#[op2]
#[string]
fn op_get_context(op_state: Rc<RefCell<OpState>>, #[string] id: &str) -> String {
@@ -658,7 +605,6 @@ fn op_get_context(op_state: Rc<RefCell<OpState>>, #[string] id: &str) -> String
}
}
#[cfg(feature = "deno_core")]
pub fn transpile_ts(expr: String) -> anyhow::Result<String> {
let parsed = deno_ast::parse_module(ParseParams {
specifier: url::Url::parse("file:///eval.ts")?,
@@ -675,30 +621,21 @@ pub fn transpile_ts(expr: String) -> anyhow::Result<String> {
.text)
}
#[cfg(not(feature = "deno_core"))]
pub fn transpile_ts(_expr: String) -> anyhow::Result<String> {
Ok("require deno".to_string())
}
#[cfg(feature = "deno_core")]
static RUNTIME_SNAPSHOT: &[u8] = include_bytes!(concat!(env!("OUT_DIR"), "/FETCH_SNAPSHOT.bin"));
#[cfg(feature = "deno_core")]
pub struct MainArgs {
args: Vec<Option<Box<RawValue>>>,
}
#[cfg(feature = "deno_core")]
pub struct LogString {
pub s: String,
}
#[cfg(feature = "deno_core")]
pub struct NativeAnnotation {
pub useragent: Option<String>,
pub proxy: Option<(String, Option<(String, String)>)>,
}
#[cfg(feature = "deno_core")]
pub fn get_annotation(inner_content: &str) -> NativeAnnotation {
let mut res = NativeAnnotation { useragent: None, proxy: None };
@@ -718,7 +655,6 @@ pub fn get_annotation(inner_content: &str) -> NativeAnnotation {
res
}
#[cfg(feature = "deno_core")]
fn capture_proxy(s: &str) -> Option<(String, Option<(String, String)>)> {
RE_PROXY.captures(s).map(|x| {
(
@@ -739,27 +675,7 @@ fn capture_proxy(s: &str) -> Option<(String, Option<(String, String)>)> {
)
})
}
#[cfg(not(feature = "deno_core"))]
pub async fn eval_fetch_timeout(
_env_code: String,
_ts_expr: String,
_js_expr: String,
_args: Option<&Json<HashMap<String, Box<RawValue>>>>,
_job_id: Uuid,
_job_timeout: Option<i32>,
_db: &DB,
_mem_peak: &mut i32,
_canceled_by: &mut Option<CanceledBy>,
_worker_name: &str,
_w_id: &str,
_load_client: bool,
_occupation_metrics: &mut OccupancyMetrics,
) -> anyhow::Result<(Box<RawValue>, String)> {
use serde_json::value::to_raw_value;
Ok((to_raw_value("require deno_core").unwrap(), "".to_string()))
}
#[cfg(feature = "deno_core")]
pub async fn eval_fetch_timeout(
env_code: String,
ts_expr: String,
@@ -935,10 +851,8 @@ pub async fn eval_fetch_timeout(
Ok((res, format!("{extra_logs}{logs}")))
}
#[cfg(feature = "deno_core")]
const WINDMILL_CLIENT: &str = include_str!("./windmill-client.js");
#[cfg(feature = "deno_core")]
async fn eval_fetch(
js_runtime: &mut JsRuntime,
expr: &str,
@@ -984,7 +898,6 @@ import("file:///eval.ts").then((module) => module.main(...args)).then(JSON.strin
Ok(unsafe_raw(r.unwrap_or_else(|| "null".to_string())))
}
#[cfg(feature = "deno_core")]
#[op2]
#[serde]
fn op_get_static_args(op_state: Rc<RefCell<OpState>>) -> Vec<Option<String>> {
@@ -997,7 +910,6 @@ fn op_get_static_args(op_state: Rc<RefCell<OpState>>) -> Vec<Option<String>> {
.collect_vec()
}
#[cfg(feature = "deno_core")]
#[op2(fast)]
fn op_log(op_state: Rc<RefCell<OpState>>, #[string] log: &str) {
// tracing::error!("log: |{}|", log);
@@ -1008,7 +920,6 @@ fn op_log(op_state: Rc<RefCell<OpState>>, #[string] log: &str) {
.push_str(log);
}
#[cfg(feature = "deno_core")]
#[cfg(test)]
mod tests {

View File

@@ -34,7 +34,5 @@ pub use worker::*;
pub use result_processor::handle_job_error;
pub use bun_executor::{
get_common_bun_proc_envs, install_bun_lockfile, prebundle_bun_script, prepare_job_dir,
};
pub use bun_executor::{get_common_bun_proc_envs, install_bun_lockfile, prepare_job_dir};
pub use deno_executor::generate_deno_lock;

View File

@@ -345,7 +345,7 @@ try {{
mem_peak,
canceled_by,
child,
!*DISABLE_NSJAIL,
false,
worker_name,
&job.workspace_id,
"php run",

View File

@@ -29,9 +29,6 @@ lazy_static::lazy_static! {
static ref PYTHON_PATH: String =
std::env::var("PYTHON_PATH").unwrap_or_else(|_| "/usr/local/bin/python3".to_string());
static ref UV_PATH: String =
std::env::var("UV_PATH").unwrap_or_else(|_| "/usr/local/bin/uv".to_string());
static ref FLOCK_PATH: String =
std::env::var("FLOCK_PATH").unwrap_or_else(|_| "/usr/bin/flock".to_string());
static ref NON_ALPHANUM_CHAR: Regex = regex::Regex::new(r"[^0-9A-Za-z=.-]").unwrap();
@@ -39,9 +36,6 @@ lazy_static::lazy_static! {
static ref PIP_TRUSTED_HOST: Option<String> = std::env::var("PIP_TRUSTED_HOST").ok();
static ref PIP_INDEX_CERT: Option<String> = std::env::var("PIP_INDEX_CERT").ok();
static ref USE_PIP_COMPILE: bool = std::env::var("USE_PIP_COMPILE")
.ok().map(|flag| flag == "true").unwrap_or(false);
static ref RELATIVE_IMPORT_REGEX: Regex = Regex::new(r#"(import|from)\s(((u|f)\.)|\.)"#).unwrap();
@@ -67,12 +61,9 @@ use crate::{
handle_child::handle_child,
AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, HTTPS_PROXY, HTTP_PROXY,
LOCK_CACHE_DIR, NO_PROXY, NSJAIL_PATH, PATH_ENV, PIP_CACHE_DIR, PIP_EXTRA_INDEX_URL,
PIP_INDEX_URL, TZ_ENV, UV_CACHE_DIR,
PIP_INDEX_URL, TZ_ENV,
};
#[cfg(windows)]
use crate::SYSTEM_ROOT;
pub async fn create_dependencies_dir(job_dir: &str) {
DirBuilder::new()
.recursive(true)
@@ -102,7 +93,7 @@ pub fn handle_ephemeral_token(x: String) -> String {
x
}
pub async fn uv_pip_compile(
pub async fn pip_compile(
job_id: &Uuid,
requirements: &str,
mem_peak: &mut i32,
@@ -112,10 +103,6 @@ pub async fn uv_pip_compile(
worker_name: &str,
w_id: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
// Fallback to pip-compile. Will be removed in future
mut no_uv: bool,
// Debug-only flag
no_cache: bool,
) -> error::Result<String> {
let mut logs = String::new();
logs.push_str(&format!("\nresolving dependencies..."));
@@ -152,184 +139,83 @@ pub async fn uv_pip_compile(
#[cfg(feature = "enterprise")]
let requirements = replace_pip_secret(db, w_id, &requirements, worker_name, job_id).await?;
let mut req_hash = format!("py-{}", calculate_hash(&requirements));
if no_uv || *USE_PIP_COMPILE {
logs.push_str(&format!("\nFallback to pip-compile (Deprecated!)"));
// Set no_uv if not setted
no_uv = true;
// Make sure that if we put #no_uv (switch to pip-compile) to python code or used `USE_PIP_COMPILE=true` variable.
// Windmill will recalculate lockfile using pip-compile and dont take potentially broken lockfile (generated by uv) from cache (our db).
// It will recalculate lockfile even if inputs have not been changed.
req_hash.push_str("-no_uv");
// Will be in format:
// py-000..000-no_uv
}
if !no_cache {
if let Some(cached) = sqlx::query_scalar!(
"SELECT lockfile FROM pip_resolution_cache WHERE hash = $1",
req_hash
)
.fetch_optional(db)
.await?
{
logs.push_str(&format!("\nfound cached resolution: {req_hash}"));
return Ok(cached);
}
let req_hash = format!("py-{}", calculate_hash(&requirements));
if let Some(cached) = sqlx::query_scalar!(
"SELECT lockfile FROM pip_resolution_cache WHERE hash = $1",
req_hash
)
.fetch_optional(db)
.await?
{
logs.push_str(&format!("\nfound cached resolution: {req_hash}"));
return Ok(cached);
}
let file = "requirements.in";
write_file(job_dir, file, &requirements)?;
// Fallback pip-compile. Will be removed in future
if no_uv {
tracing::debug!("Fallback to pip-compile");
let mut args = vec![
"-q",
"--no-header",
file,
"--resolver=backtracking",
"--strip-extras",
];
let mut pip_args = vec![];
let pip_extra_index_url = PIP_EXTRA_INDEX_URL
.read()
.await
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_extra_index_url.as_ref() {
args.extend(["--extra-index-url", url, "--no-emit-index-url"]);
pip_args.push(format!("--extra-index-url {}", url));
}
let pip_index_url = PIP_INDEX_URL
.read()
.await
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_index_url.as_ref() {
args.extend(["--index-url", url, "--no-emit-index-url"]);
pip_args.push(format!("--index-url {}", url));
}
if let Some(host) = PIP_TRUSTED_HOST.as_ref() {
args.extend(["--trusted-host", host]);
}
if let Some(cert_path) = PIP_INDEX_CERT.as_ref() {
args.extend(["--cert", cert_path]);
}
let pip_args_str = pip_args.join(" ");
if pip_args.len() > 0 {
args.extend(["--pip-args", &pip_args_str]);
}
tracing::debug!("pip-compile args: {:?}", args);
let mut child_cmd = Command::new("pip-compile");
child_cmd
.current_dir(job_dir)
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let child_process = start_child_process(child_cmd, "pip-compile").await?;
append_logs(&job_id, &w_id, logs, db).await;
handle_child(
job_id,
db,
mem_peak,
canceled_by,
child_process,
false,
worker_name,
&w_id,
"pip-compile",
None,
false,
occupancy_metrics,
)
let mut args = vec![
"-q",
"--no-header",
file,
"--resolver=backtracking",
"--strip-extras",
];
let mut pip_args = vec![];
let pip_extra_index_url = PIP_EXTRA_INDEX_URL
.read()
.await
.map_err(|e| Error::ExecutionErr(format!("Lock file generation failed: {e:?}")))?;
} else {
let mut args = vec![
"pip",
"compile",
"-q",
"--no-header",
file,
"--strip-extras",
"-o",
"requirements.txt",
// Prefer main index over extra
// https://docs.astral.sh/uv/pip/compatibility/#packages-that-exist-on-multiple-indexes
// TODO: Use env variable that can be toggled from UI
"--index-strategy",
"unsafe-best-match",
// Target to /tmp/windmill/cache/uv
"--cache-dir",
UV_CACHE_DIR,
// We dont want UV to manage python installations
"--python-preference",
"only-system",
"--no-python-downloads",
];
if no_cache {
args.extend(["--no-cache"]);
}
let pip_extra_index_url = PIP_EXTRA_INDEX_URL
.read()
.await
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_extra_index_url.as_ref() {
args.extend(["--extra-index-url", url]);
}
let pip_index_url = PIP_INDEX_URL
.read()
.await
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_index_url.as_ref() {
args.extend(["--index-url", url]);
}
if let Some(host) = PIP_TRUSTED_HOST.as_ref() {
args.extend(["--trusted-host", host]);
}
if let Some(cert_path) = PIP_INDEX_CERT.as_ref() {
args.extend(["--cert", cert_path]);
}
tracing::debug!("uv args: {:?}", args);
#[cfg(windows)]
let uv_cmd = "uv";
#[cfg(unix)]
let uv_cmd = UV_PATH.as_str();
let mut child_cmd = Command::new(uv_cmd);
child_cmd
.current_dir(job_dir)
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let child_process = start_child_process(child_cmd, "/usr/local/bin/uv").await?;
append_logs(&job_id, &w_id, logs, db).await;
handle_child(
job_id,
db,
mem_peak,
canceled_by,
child_process,
false,
worker_name,
&w_id,
// TODO: Rename to uv-pip-compile?
"uv",
None,
false,
occupancy_metrics,
)
.await
.map_err(|e| Error::ExecutionErr(format!("Lock file generation failed: {e:?}")))?;
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_extra_index_url.as_ref() {
args.extend(["--extra-index-url", url, "--no-emit-index-url"]);
pip_args.push(format!("--extra-index-url {}", url));
}
let pip_index_url = PIP_INDEX_URL
.read()
.await
.clone()
.map(handle_ephemeral_token);
if let Some(url) = pip_index_url.as_ref() {
args.extend(["--index-url", url, "--no-emit-index-url"]);
pip_args.push(format!("--index-url {}", url));
}
if let Some(host) = PIP_TRUSTED_HOST.as_ref() {
args.extend(["--trusted-host", host]);
}
if let Some(cert_path) = PIP_INDEX_CERT.as_ref() {
args.extend(["--cert", cert_path]);
}
let pip_args_str = pip_args.join(" ");
if pip_args.len() > 0 {
args.extend(["--pip-args", &pip_args_str]);
}
tracing::debug!("pip-compile args: {:?}", args);
let mut child_cmd = Command::new("pip-compile");
child_cmd
.current_dir(job_dir)
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let child_process = start_child_process(child_cmd, "pip-compile").await?;
append_logs(&job_id, &w_id, logs, db).await;
handle_child(
job_id,
db,
mem_peak,
canceled_by,
child_process,
false,
worker_name,
&w_id,
"pip-compile",
None,
false,
occupancy_metrics,
)
.await
.map_err(|e| Error::ExecutionErr(format!("Lock file generation failed: {e:?}")))?;
let path_lock = format!("{job_dir}/requirements.txt");
let mut file = File::open(path_lock).await?;
let mut req_content = "".to_string();
@@ -498,10 +384,7 @@ except BaseException as e:
exc_type, exc_value, exc_traceback = sys.exc_info()
tb = traceback.format_tb(exc_traceback)
with open(result_json, 'w') as f:
err = {{ "message": str(e), "name": e.__class__.__name__, "stack": '\n'.join(tb[1:]) }}
extra = e.__dict__
if extra and len(extra) > 0:
err['extra'] = extra
err = {{ "message": str(e), "name": e.__class__.__name__, "stack": '\n'.join(tb[1:]) }}
flow_node_id = os.environ.get('WM_FLOW_STEP_ID')
if flow_node_id:
err['step_id'] = flow_node_id
@@ -516,9 +399,6 @@ except BaseException as e:
let mut reserved_variables = get_reserved_variables(job, &client.token, db).await?;
let additional_python_paths_folders = additional_python_paths.iter().join(":");
#[cfg(windows)]
let additional_python_paths_folders = additional_python_paths_folders.replace(":", ";");
if !*DISABLE_NSJAIL {
let shared_deps = additional_python_paths
.into_iter()
@@ -595,10 +475,6 @@ mount {{
.args(vec!["-u", "-m", "wrapper"])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
python_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
start_child_process(python_cmd, PYTHON_PATH.as_str()).await?
};
@@ -898,7 +774,6 @@ async fn handle_python_deps(
let requirements = match requirements_o {
Some(r) => r,
None => {
let annotation = windmill_common::worker::get_annotation_python(inner_content);
let mut already_visited = vec![];
let requirements = windmill_parser_py_imports::parse_python_imports(
@@ -913,7 +788,7 @@ async fn handle_python_deps(
if requirements.is_empty() {
"".to_string()
} else {
uv_pip_compile(
pip_compile(
job_id,
&requirements,
mem_peak,
@@ -923,8 +798,6 @@ async fn handle_python_deps(
worker_name,
w_id,
occupancy_metrics,
annotation.no_uv,
annotation.no_cache,
)
.await
.map_err(|e| {
@@ -1143,12 +1016,7 @@ pub async fn handle_python_reqs(
start_child_process(nsjail_cmd, NSJAIL_PATH.as_str()).await?
} else {
let fssafe_req = NON_ALPHANUM_CHAR.replace_all(&req, "_").to_string();
#[cfg(unix)]
let req = format!("'{}'", req);
#[cfg(windows)]
let req = format!("{}", req);
let mut command_args = vec![
PYTHON_PATH.as_str(),
"-m",
@@ -1204,35 +1072,19 @@ pub async fn handle_python_reqs(
tracing::debug!("pip install command: {:?}", command_args);
#[cfg(unix)]
{
let mut flock_cmd = Command::new(FLOCK_PATH.as_str());
flock_cmd
.env_clear()
.envs(envs)
.args([
"-x",
&format!("{}/pip-{}.lock", LOCK_CACHE_DIR, fssafe_req),
"--command",
&command_args.join(" "),
])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
start_child_process(flock_cmd, FLOCK_PATH.as_str()).await?
}
#[cfg(windows)]
{
let mut pip_cmd = Command::new(PYTHON_PATH.as_str());
pip_cmd
.env_clear()
.envs(envs)
.env("SystemRoot", SYSTEM_ROOT.as_str())
.args(&command_args[1..])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
start_child_process(pip_cmd, PYTHON_PATH.as_str()).await?
}
let mut flock_cmd = Command::new(FLOCK_PATH.as_str());
flock_cmd
.env_clear()
.envs(envs)
.args([
"-x",
&format!("{}/pip-{}.lock", LOCK_CACHE_DIR, fssafe_req),
"--command",
&command_args.join(" "),
])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
start_child_process(flock_cmd, FLOCK_PATH.as_str()).await?
};
let child = handle_child(

View File

@@ -23,28 +23,12 @@ use crate::{
RUST_CACHE_DIR, TZ_ENV,
};
#[cfg(windows)]
use crate::SYSTEM_ROOT;
const NSJAIL_CONFIG_RUN_RUST_CONTENT: &str = include_str!("../nsjail/run.rust.config.proto");
lazy_static::lazy_static! {
static ref HOME_DIR: String = std::env::var("HOME").expect("Could not find the HOME environment variable");
static ref CARGO_HOME: String = std::env::var("CARGO_HOME").unwrap_or_else(|_| { CARGO_HOME_DEFAULT.clone() });
static ref RUSTUP_HOME: String = std::env::var("RUSTUP_HOME").unwrap_or_else(|_| { RUSTUP_HOME_DEFAULT.clone() });
static ref CARGO_PATH: String = format!("{}/bin/cargo", CARGO_HOME.as_str());
}
#[cfg(windows)]
lazy_static::lazy_static! {
static ref CARGO_HOME_DEFAULT: String = format!("{}\\.cargo", *HOME_DIR);
static ref RUSTUP_HOME_DEFAULT: String = format!("{}\\.rustup", *HOME_DIR);
}
#[cfg(unix)]
lazy_static::lazy_static! {
static ref CARGO_HOME_DEFAULT: String = "/usr/local/cargo".to_string();
static ref RUSTUP_HOME_DEFAULT: String = "/usr/local/rustup".to_string();
static ref CARGO_HOME: String = std::env::var("CARGO_HOME").unwrap_or_else(|_| "/usr/local/cargo".to_string());
static ref RUSTUP_HOME: String = std::env::var("RUSTUP_HOME").unwrap_or_else(|_| "/usr/local/rustup".to_string());
static ref CARGO_PATH: String = format!("{}/bin/cargo", std::env::var("CARGO_HOME").unwrap_or("/usr/local/cargo/bin/cargo".to_string()));
}
const RUST_OBJECT_STORE_PREFIX: &str = "rustbin/";
@@ -142,14 +126,6 @@ pub async fn generate_cargo_lockfile(
.args(vec!["generate-lockfile"])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
{
gen_lockfile_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
gen_lockfile_cmd.env(
"TMP",
std::env::var("TMP").unwrap_or_else(|_| "C:\\tmp".to_string()),
);
}
let gen_lockfile_process = start_child_process(gen_lockfile_cmd, CARGO_PATH.as_str()).await?;
handle_child(
job_id,
@@ -200,16 +176,6 @@ pub async fn build_rust_crate(
.args(vec!["build", "--release"])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
{
build_rust_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
build_rust_cmd.env(
"TMP",
std::env::var("TMP").unwrap_or_else(|_| "C:\\tmp".to_string()),
);
}
let build_rust_process = start_child_process(build_rust_cmd, CARGO_PATH.as_str()).await?;
handle_child(
job_id,
@@ -313,13 +279,7 @@ pub async fn handle_rust_job(
let cache_logs = if cache {
let target = format!("{job_dir}/main");
#[cfg(unix)]
let symlink = std::os::unix::fs::symlink(&bin_path, &target);
#[cfg(windows)]
let symlink = std::os::windows::fs::symlink_dir(&bin_path, &target);
symlink.map_err(|e| {
std::os::unix::fs::symlink(&bin_path, &target).map_err(|e| {
Error::ExecutionErr(format!(
"could not copy cached binary from {bin_path} to {job_dir}/main: {e:?}"
))
@@ -400,9 +360,6 @@ pub async fn handle_rust_job(
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
run_rust.env("SystemRoot", SYSTEM_ROOT.as_str());
start_child_process(run_rust, compiled_executable_name).await?
};
handle_child(

View File

@@ -237,7 +237,6 @@ pub const ROOT_CACHE_NOMOUNT_DIR: &str = concatcp!(TMP_DIR, "/cache_nomount/");
pub const LOCK_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "lock");
pub const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip");
pub const UV_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "uv");
pub const TAR_PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "tar/pip");
pub const DENO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "deno");
pub const DENO_CACHE_DIR_DEPS: &str = concatcp!(ROOT_CACHE_DIR, "deno/deps");
@@ -349,6 +348,7 @@ lazy_static::lazy_static! {
pub static ref PIP_INDEX_URL: Arc<RwLock<Option<String>>> = Arc::new(RwLock::new(None));
pub static ref JOB_DEFAULT_TIMEOUT: Arc<RwLock<Option<i32>>> = Arc::new(RwLock::new(None));
static ref MAX_TIMEOUT: u64 = std::env::var("TIMEOUT")
.ok()
.and_then(|x| x.parse::<u64>().ok())
@@ -389,12 +389,6 @@ lazy_static::lazy_static! {
}
#[cfg(windows)]
lazy_static::lazy_static! {
pub static ref SYSTEM_ROOT: String = std::env::var("SystemRoot").unwrap_or_else(|_| "C:\\Windows".to_string());
}
//only matter if CLOUD_HOSTED
pub const MAX_RESULT_SIZE: usize = 1024 * 1024 * 2; // 2MB

View File

@@ -12,7 +12,7 @@ use windmill_common::flows::{FlowModule, FlowModuleValue};
use windmill_common::get_latest_deployed_hash_for_path;
use windmill_common::jobs::JobPayload;
use windmill_common::scripts::ScriptHash;
use windmill_common::worker::{get_annotation_ts, to_raw_value, to_raw_value_owned, write_file};
use windmill_common::worker::{get_annotation, to_raw_value, to_raw_value_owned, write_file};
use windmill_common::{
error::{self, to_anyhow},
flows::FlowValue,
@@ -26,7 +26,7 @@ use windmill_parser_ts::parse_expr_for_imports;
use windmill_queue::{append_logs, CanceledBy, PushIsolationLevel};
use crate::common::OccupancyMetrics;
use crate::python_executor::{create_dependencies_dir, handle_python_reqs, uv_pip_compile};
use crate::python_executor::{create_dependencies_dir, handle_python_reqs, pip_compile};
use crate::rust_executor::{build_rust_crate, compute_rust_hash, generate_cargo_lockfile};
use crate::{
bun_executor::gen_bun_lockfile,
@@ -953,7 +953,7 @@ async fn lock_modules<'c>(
}
if language == ScriptLang::Bun || language == ScriptLang::Bunnative {
let anns = get_annotation_ts(&content);
let anns = get_annotation(&content);
if anns.native_mode && language == ScriptLang::Bun {
language = ScriptLang::Bunnative;
} else if !anns.native_mode && language == ScriptLang::Bunnative {
@@ -1003,7 +1003,7 @@ async fn lock_modules<'c>(
fn skip_creating_new_lock(language: &ScriptLang, content: &str) -> bool {
if language == &ScriptLang::Bun || language == &ScriptLang::Bunnative {
let anns = get_annotation_ts(&content);
let anns = get_annotation(&content);
if anns.native_mode && language == &ScriptLang::Bun {
return false;
} else if !anns.native_mode && language == &ScriptLang::Bunnative {
@@ -1077,7 +1077,7 @@ async fn lock_modules_app(
match new_lock {
Ok(new_lock) => {
append_logs(&job.id, &job.workspace_id, logs, db).await;
let anns = get_annotation_ts(&content);
let anns = get_annotation(&content);
let nlang = if anns.native_mode && language == ScriptLang::Bun {
Some(ScriptLang::Bunnative)
} else if !anns.native_mode && language == ScriptLang::Bunnative
@@ -1280,7 +1280,7 @@ async fn python_dep(
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
) -> std::result::Result<String, Error> {
create_dependencies_dir(job_dir).await;
let req: std::result::Result<String, Error> = uv_pip_compile(
let req: std::result::Result<String, Error> = pip_compile(
job_id,
&reqs,
mem_peak,
@@ -1290,8 +1290,6 @@ async fn python_dep(
worker_name,
w_id,
occupancy_metrics,
false,
false,
)
.await;
// install the dependencies to pre-fill the cache
@@ -1436,9 +1434,8 @@ async fn capture_dependency_job(
.await
}
ScriptLang::Bun | ScriptLang::Bunnative => {
let npm_mode = npm_mode.unwrap_or_else(|| {
windmill_common::worker::get_annotation_ts(job_raw_code).npm_mode
});
let npm_mode = npm_mode
.unwrap_or_else(|| windmill_common::worker::get_annotation(job_raw_code).npm_mode);
if !raw_deps {
let _ = write_file(job_dir, "main.ts", job_raw_code)?;
}
@@ -1475,7 +1472,7 @@ async fn capture_dependency_job(
base_internal_url,
worker_name,
&token,
&mut Some(occupancy_metrics),
occupancy_metrics,
)
.await?;
}

View File

@@ -2,7 +2,7 @@ import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts";
import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts";
import * as api from "https://deno.land/x/windmill@v1.174.0/windmill-api/index.ts";
export const VERSION = "v1.409.2";
export const VERSION = "v1.403.1";
export async function login(email: string, password: string): Promise<string> {
return await windmill.UserService.login({

View File

@@ -1,6 +1,6 @@
// deno-lint-ignore-file no-explicit-any
import { requireLogin, resolveWorkspace, validatePath } from "./context.ts";
import { colors, Command, log, SEP, Table, yamlParseFile } from "./deps.ts";
import { colors, Command, log, SEP, Table, yamlParse } from "./deps.ts";
import * as wmill from "./gen/services.gen.ts";
import { ListableApp, Policy } from "./gen/types.gen.ts";
@@ -25,7 +25,7 @@ export async function pushApp(
return;
}
alreadySynced.push(localPath);
remotePath = remotePath.replaceAll(SEP, "/");
remotePath.replaceAll(SEP, "/");
let app: any = undefined;
// deleting old app if it exists in raw mode
try {
@@ -40,8 +40,8 @@ export async function pushApp(
if (!localPath.endsWith(SEP)) {
localPath += SEP;
}
const path = localPath + "app.yaml";
const localApp = (await yamlParseFile(path)) as AppFile;
const localAppRaw = await Deno.readTextFile(localPath + "app.yaml");
const localApp = yamlParse(localAppRaw) as AppFile;
function replaceInlineScripts(rec: any) {
if (!rec) {

View File

@@ -1,4 +1,4 @@
import { log, yamlParseFile } from "./deps.ts";
import { log, yamlParse } from "./deps.ts";
export interface SyncOptions {
stateful?: boolean;
@@ -40,7 +40,9 @@ export interface Codebase {
export async function readConfigFile(): Promise<SyncOptions> {
try {
const conf = (await yamlParseFile("wmill.yaml")) as SyncOptions;
const conf = yamlParse(
await Deno.readTextFile("wmill.yaml")
) as SyncOptions;
if (conf?.defaultTs == undefined) {
log.warn(
"No defaultTs defined in your wmill.yaml. Using 'bun' as default."

View File

@@ -21,29 +21,7 @@ export { copy } from "jsr:@std/io/copy";
export { readAll } from "jsr:@std/io/read-all";
export * as log from "jsr:@std/log";
export { stringify as yamlStringify } from "jsr:@std/yaml";
import { parse as yamlParse, ParseOptions } from "jsr:@std/yaml";
export async function yamlParseFile(path: string, options: ParseOptions = {}) {
try {
return yamlParse(await Deno.readTextFile(path), options);
} catch (e) {
throw new Error(`Error parsing yaml ${path}`, { cause: e });
}
}
export function yamlParseContent(
path: string,
content: string,
options: ParseOptions = {}
) {
try {
return yamlParse(content, options);
} catch (e) {
throw new Error(`Error parsing yaml ${path}`, { cause: e });
}
}
export { stringify as yamlStringify, parse as yamlParse } from "jsr:@std/yaml";
// other

View File

@@ -1,7 +1,7 @@
// deno-lint-ignore-file no-explicit-any
import { GlobalOptions, isSuperset } from "./types.ts";
import { Confirm, SEP, log, yamlStringify } from "./deps.ts";
import { colors, Command, Table, yamlParseFile } from "./deps.ts";
import { colors, Command, Table, yamlParse } from "./deps.ts";
import * as wmill from "./gen/services.gen.ts";
import { requireLogin, resolveWorkspace, validatePath } from "./context.ts";
@@ -29,23 +29,21 @@ export function replaceInlineScripts(
) {
modules.forEach((m) => {
if (m.value.type == "rawscript") {
if (m.value.content.startsWith("!inline")) {
const path = m.value.content.split(" ")[1];
m.value.content = Deno.readTextFileSync(localPath + path);
const lock = m.value.lock;
if (removeLocks && removeLocks.includes(path)) {
m.value.lock = undefined;
} else if (
lock &&
typeof lock == "string" &&
lock.trimStart().startsWith("!inline ")
) {
const path = lock.split(" ")[1];
try {
m.value.lock = readInlinePathSync(localPath + path);
} catch {
log.error(`Lock file ${path} not found`);
}
const path = m.value.content.split(" ")[1];
m.value.content = Deno.readTextFileSync(localPath + path);
const lock = m.value.lock;
if (removeLocks && removeLocks.includes(path)) {
m.value.lock = undefined;
} else if (
lock &&
typeof lock == "string" &&
lock.trimStart().startsWith("!inline ")
) {
const path = lock.split(" ")[1];
try {
m.value.lock = readInlinePathSync(localPath + path);
} catch {
log.error(`Lock file ${path} not found`);
}
}
} else if (m.value.type == "forloopflow") {
@@ -89,7 +87,8 @@ export async function pushFlow(
if (!localPath.endsWith(SEP)) {
localPath += SEP;
}
const localFlow = (await yamlParseFile(localPath + "flow.yaml")) as FlowFile;
const localFlowRaw = await Deno.readTextFile(localPath + "flow.yaml");
const localFlow = yamlParse(localFlowRaw) as FlowFile;
replaceInlineScripts(localFlow.value.modules, localPath, undefined);
@@ -121,7 +120,6 @@ export async function pushFlow(
});
} catch (e) {
throw new Error(
//@ts-ignore
`Failed to create flow ${remotePath}: ${e.body ?? e.message}`
);
}

View File

@@ -73,7 +73,6 @@ export async function pushFolder(
},
});
} catch (e) {
//@ts-ignore
console.error(e.body);
throw e;
}
@@ -88,7 +87,6 @@ export async function pushFolder(
},
});
} catch (e) {
//@ts-ignore
throw Error(`Failed to create folder ${name}: ${e.body ?? e.message}`);
}
}

View File

@@ -54,7 +54,7 @@ export const OpenAPI: OpenAPIConfig = {
PASSWORD: undefined,
TOKEN: getEnv("WM_TOKEN"),
USERNAME: undefined,
VERSION: '1.407.2',
VERSION: '1.401.0',
WITH_CREDENTIALS: true,
interceptors: {
request: new Interceptors(),

View File

@@ -1015,9 +1015,6 @@ export type FlowModule = {
skip_if_stopped?: boolean;
expr: string;
};
skip_if?: {
expr: string;
};
sleep?: InputTransform;
cache_ttl?: number;
timeout?: number;
@@ -1179,7 +1176,6 @@ export type FlowStatusModule = {
approver: string;
}>;
failed_retries?: Array<(string)>;
skipped?: boolean;
};
export type type4 = 'WaitingForPriorSteps' | 'WaitingForEvents' | 'WaitingForExecutor' | 'InProgress' | 'Success' | 'Failure';

View File

@@ -3,7 +3,7 @@ import {
path,
Confirm,
yamlStringify,
yamlParseFile,
yamlParse,
Command,
setClient,
Table,
@@ -25,8 +25,8 @@ import {
import {
add as workspaceSetup,
addWorkspace,
allWorkspaces,
removeWorkspace,
setActiveWorkspace,
} from "./workspace.ts";
import {
pushInstanceSettings,
@@ -37,7 +37,6 @@ import {
} from "./settings.ts";
import { deepEqual } from "./utils.ts";
import { GlobalOptions } from "./types.ts";
import { getActiveWorkspace } from "./workspace.ts";
export interface Instance {
remote: string;
@@ -191,12 +190,6 @@ export async function pickInstance(
const instances = await allInstances();
if (opts.baseUrl && opts.token) {
log.info("Using instance fully defined by --base-url and --token");
setClient(
opts.token,
opts.baseUrl.endsWith("/") ? opts.baseUrl.slice(0, -1) : opts.baseUrl
);
return {
name: "custom",
remote: opts.baseUrl,
@@ -296,14 +289,15 @@ async function instancePull(opts: GlobalOptions & InstanceSyncOptions) {
if (opts.includeWorkspaces) {
log.info("\nPulling all workspaces");
const rootDir = Deno.cwd();
const localWorkspaces = await getLocalWorkspaces(rootDir, instance.prefix);
const previousActiveWorkspace = await getActiveWorkspace(undefined);
const remoteWorkspaces = await wmill.listWorkspacesAsSuperAdmin({
page: 1,
perPage: 1000,
});
let localWorkspaces = await allWorkspaces();
localWorkspaces = localWorkspaces.filter((w) =>
w.name.startsWith(instance.prefix + "_")
);
const rootDir = Deno.cwd();
for (const remoteWorkspace of remoteWorkspaces) {
log.info("\nPulling workspace " + remoteWorkspace.id);
const workspaceName = instance.prefix + "_" + remoteWorkspace.id;
@@ -333,37 +327,31 @@ async function instancePull(opts: GlobalOptions & InstanceSyncOptions) {
includeSettings: true,
includeUsers: true,
includeKey: true,
yes: opts.yes,
});
}
const localWorkspacesToDelete = localWorkspaces.filter(
(w) => !remoteWorkspaces.find((r) => r.id === w.id)
(w) => !remoteWorkspaces.find((r) => r.id === w.workspaceId)
);
if (localWorkspacesToDelete.length > 0) {
const confirmDelete =
opts.yes ||
(await Confirm.prompt({
message:
"Do you want to delete the local copy of workspaces that don't exist anymore on the instance?\n" +
localWorkspacesToDelete.map((w) => w).join(", "),
default: true,
}));
const confirmDelete = await Confirm.prompt({
message:
"Do you want to delete the local copy of workspaces that don't exist anymore on the instance?\n" +
localWorkspacesToDelete.map((w) => w.workspaceId).join(", "),
default: true,
});
if (confirmDelete) {
for (const workspace of localWorkspacesToDelete) {
await removeWorkspace(workspace.id, false, {});
await Deno.remove(path.join(rootDir, workspace.dir), {
await removeWorkspace(workspace.name, false, {});
await Deno.remove(path.join(rootDir, workspace.name), {
recursive: true,
});
}
}
}
if (previousActiveWorkspace) {
await setActiveWorkspace(previousActiveWorkspace?.name);
}
log.info(colors.green.underline.bold("All workspaces pulled"));
}
}
@@ -423,8 +411,6 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
if (opts.includeWorkspaces) {
instances = await allInstances();
const rootDir = Deno.cwd();
const localPrefix = (await Select.prompt({
message: "What is the prefix of the local workspaces you want to sync?",
options: [
@@ -436,22 +422,22 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
default: instance.prefix as unknown,
})) as unknown as string;
const remoteWorkspaces = await wmill.listWorkspacesAsSuperAdmin({
page: 1,
perPage: 1000,
});
const rootDir = Deno.cwd();
const previousActiveWorkspace = await getActiveWorkspace(undefined);
const localWorkspaces = await getLocalWorkspaces(rootDir, localPrefix);
let localWorkspaces = Deno.readDir(".");
localWorkspaces = localWorkspaces.filter((w) =>
w.name.startsWith(localPrefix + "_")
);
log.info(
`\nPushing all workspaces: ${localWorkspaces.map((x) => x.id).join(", ")}`
`\nPushing all workspaces: ${localWorkspaces.join(
", "
)} with prefix ${localPrefix}`
);
for (const localWorkspace of localWorkspaces) {
log.info("\nPushing workspace " + localWorkspace.id);
log.info("\nPushing workspace " + localWorkspace.workspaceId);
try {
await Deno.chdir(path.join(rootDir, localWorkspace.dir));
await Deno.chdir(path.join(rootDir, localWorkspace.name));
} catch (_) {
throw new Error(
"Workspace folder not found, are you in the right directory?"
@@ -459,9 +445,9 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
}
try {
const workspaceSettings = (await yamlParseFile(
"settings.yaml"
)) as SimplifiedSettings;
const workspaceSettings = yamlParse(
await Deno.readTextFile("settings.yaml")
) as SimplifiedSettings;
await workspaceSetup(
{
token: instance.token,
@@ -471,8 +457,8 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
createWorkspaceName: workspaceSettings.name,
createUsername: undefined,
},
localWorkspace.dir,
localWorkspace.id,
localWorkspace.name,
localWorkspace.workspaceId,
instance.remote
);
} catch (_) {
@@ -482,7 +468,7 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
continue;
}
await push({
workspace: localWorkspace.dir,
workspace: localWorkspace.name,
token: undefined,
baseUrl: undefined,
includeGroups: true,
@@ -490,22 +476,23 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
includeSettings: true,
includeUsers: true,
includeKey: true,
yes: opts.yes,
});
}
const remoteWorkspaces = await wmill.listWorkspacesAsSuperAdmin({
page: 1,
perPage: 1000,
});
const workspacesToDelete = remoteWorkspaces.filter(
(w) => !localWorkspaces.find((l) => l.id === w.id)
(w) => !localWorkspaces.find((l) => l.workspaceId === w.id)
);
if (workspacesToDelete.length > 0) {
const confirmDelete =
opts.yes ||
(await Confirm.prompt({
message:
"Do you want to delete the following remote workspaces that don't exist locally?\n" +
workspacesToDelete.map((w) => w.id).join(", "),
default: true,
}));
const confirmDelete = await Confirm.prompt({
message:
"Do you want to delete the following remote workspaces that don't exist locally?\n" +
workspacesToDelete.map((w) => w.id).join(", "),
default: true,
});
if (confirmDelete) {
for (const workspace of workspacesToDelete) {
@@ -514,28 +501,10 @@ async function instancePush(opts: GlobalOptions & InstanceSyncOptions) {
}
}
}
if (previousActiveWorkspace) {
await setActiveWorkspace(previousActiveWorkspace?.name);
}
log.info(colors.green.underline.bold("All workspaces pushed"));
}
}
async function getLocalWorkspaces(rootDir: string, localPrefix: string) {
const localWorkspaces: { dir: string; id: string }[] = [];
for await (const dir of Deno.readDir(rootDir)) {
const dirName = dir.name;
if (dirName.startsWith(localPrefix + "_")) {
localWorkspaces.push({
dir: dirName,
id: dirName.substring(localPrefix.length + 1),
});
}
}
return localWorkspaces;
}
async function switchI(opts: {}, instanceName: string) {
const all = await allInstances();
if (all.findIndex((x) => x.name === instanceName) === -1) {
@@ -578,7 +547,6 @@ async function whoami(opts: {}) {
log.info(JSON.stringify(whoamiInfo, null, 2));
} catch (error) {
log.error(
//@ts-ignore
colors.red(`Failed to retrieve whoami information: ${error.message}`)
);
}
@@ -621,7 +589,7 @@ const command = new Command()
.description("Remove an instance")
.complete("instance", async () => (await allInstances()).map((x) => x.name))
.arguments("<instance:string:instance>")
.action(async (instance: any) => {
.action(async (instance) => {
const instances = await allInstances();
const choice = (await Select.prompt({

View File

@@ -60,7 +60,7 @@ export {
// }
// });
export const VERSION = "1.409.2";
export const VERSION = "1.403.1";
const command = new Command()
.name("wmill")

Some files were not shown because too many files have changed in this diff Show More