Compare commits

...
Author SHA1 Message Date
Ruben Fiszelandrubenfiszel 0ceb72f012 chore(main): release 1.535.0 (#6460)
* chore(main): release 1.535.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-25 14:43:29 +00:00
centdix eca3109ec6 feat(aichat): show diff mode on inline scripts changes (#6454)
* draft

* simplify logic

* convert diffeditor to svelte5

* add buttons to accept or reject

* set code on reject

* nit

* fix

* cleaning

* nit
2025-08-25 14:38:39 +00:00
Ruben Fiszel d3288947b2 fix: fix opening advanced popup for run resetting tag to default 2025-08-25 14:21:41 +00:00
Ruben Fiszelandrubenfiszel a691ae2883 chore(main): release 1.534.1 (#6458)
* chore(main): release 1.534.1

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-25 13:37:53 +00:00
Diego ImbertandGitHub Action 16d233bf46 fix: add alias to subquery for older postgres versions (#6455)
* add alias to subquery for older postgres versions

* Update SQLx metadata

---------

Co-authored-by: GitHub Action <action@github.com>
2025-08-25 13:19:10 +00:00
Guilhem fc20b7bd91 fix(frontend): fix test step behavior (#6427)
* fix flowStateStore val

* handle run preview multiple keyboard actions

* Synchronise input args and prview args

* Fix arg update one step load

* fix input ste manually not reseted after preview

* rename test steps to stepsInputArgs

* simplify job result update

* fix job preview logic

* fix import

* nit

* clean

* fix test job not displaying when data is pinned

* remove job history loader display delay

* nit

* nit

* add error handler to steps input args comparison function

* prevent result node to display connection
2025-08-25 12:49:46 +00:00
Ruben Fiszelandrubenfiszel 082312000f chore(main): release 1.534.0 (#6452)
* chore(main): release 1.534.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-25 12:46:19 +00:00
pyranota ef93e9ec8b tooling: update nix instructions on starting db (#6457) 2025-08-25 12:45:53 +00:00
hugocasa 8d31c2ab0d feat(backend): support unencrypted connection to mssql (#6453) 2025-08-25 12:39:37 +00:00
hugocasa 1074b22900 tooling: dev docker db script and readme nits (#6456)
* tooling: dev docker db script and readme nits

* nits

* nits
2025-08-25 12:38:35 +00:00
Ruben Fiszel 3845744492 nit 2025-08-25 12:35:04 +00:00
Ruben Fiszel 1073eb0e68 fix(flow): test this step preload step input evaluation 2025-08-25 12:12:12 +00:00
centdix e951c896b8 fix(aichat): fix wrong current model logic (#6451)
* fix model selection

* fix for context window

* cleaning

* fix lint

* fix
2025-08-25 11:47:02 +02:00
Ruben Fiszel 97ed4a539b nits cleanup + faster script index #6450 2025-08-23 22:42:39 +00:00
Ruben Fiszelandrubenfiszel 8964896c13 chore(main): release 1.533.1 (#6449)
* chore(main): release 1.533.1

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-23 08:02:32 +00:00
Ruben Fiszel 99666426ec nits 2025-08-23 07:54:59 +00:00
Ruben Fiszel 0ae8f44773 fix(app): fix oneOf selected undefined freeze 2025-08-23 07:52:56 +00:00
Ruben Fiszelandrubenfiszel 5b338bb749 chore(main): release 1.533.0 (#6448)
* chore(main): release 1.533.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-23 02:53:29 +01:00
centdix d948ff5d0d nits(aichat): add keyboard navigation in context list (#6443)
* better flow module code peice

* better keyboard nav for availablecontextlist

* cleaning

* escape to close + cleaning

* nit tab handling

* fix module extraction

* remove code category

* fixes

* comment

* fix

* fix tool params display
2025-08-23 02:50:17 +01:00
Alexander Petricandellipsis-dev[bot] a41b9e47e2 feat: CLI improvements (#6446)
* feat: branch specific items for cli

* error on wmill.yaml parsing errors

* also search for wmill.yaml in parent dirs when git

* git_branches -> gitBranches

* Update cli/src/core/specific_items.ts

Co-authored-by: ellipsis-dev[bot] <65095814+ellipsis-dev[bot]@users.noreply.github.com>

* sanitize branch name (regex + fs path)

* improve sanitatino

* robust relative paths

* hubpath

---------

Co-authored-by: ellipsis-dev[bot] <65095814+ellipsis-dev[bot]@users.noreply.github.com>
2025-08-23 02:45:00 +01:00
hugocasa c13747cda9 fix(frontend): ai agent flow status + UI nits (#6447)
* fix(frontend): ai agent flow status

* nit: prevent undefined node issue

* feat: UI nits + flow status select iter fix

* nit ai agent color in picker
2025-08-22 18:51:13 +01:00
Ruben Fiszel 964351e211 update monaco-vscode-api (#6445) 2025-08-22 11:54:26 +00:00
Ruben Fiszelandrubenfiszel 2046b64ec8 chore(main): release 1.532.0 (#6439)
* chore(main): release 1.532.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-22 11:05:42 +01:00
hugocasa 7da79a8bc5 feat: json schema resource (#6433)
* feat: json schema resource

* feat: cache

* fix cleanup

* fix: use format instead of custom property
2025-08-22 10:39:13 +01:00
Diego Imbert 4d8777b278 Fix read undefined when renaming flow step (#6440) 2025-08-22 09:24:43 +00:00
centdix 73272f16fd feat(aichat): allow adding contexts to flow mode (#6424)
* add new feature instructions

* add db as context for flow mode

* add diff

* cleaner diff

* add modules as available context

* convert to svelte 5

* auto add selected module to context

* change flowinline ai button + nits

* handle adding selected lines

* clean context handling

* apply code pieces

* new chat when changing mode

* clean

* show code for code steps

* add last saved flow

* fix size

* categorize context

* optionnaly categorize

* fix module finding

* logs

* nit prompt

* fix

* fix

* fix test tool for script

* clean
2025-08-22 08:46:03 +00:00
Ruben Fiszelandrubenfiszel ee5e39a3d5 chore(main): release 1.531.0 (#6429)
* chore(main): release 1.531.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
2025-08-22 08:01:58 +00:00
Ruben Fiszel 9df008b9f8 fix: s3 result presigned not working with list 2025-08-22 07:56:02 +00:00
124 changed files with 4910 additions and 3222 deletions
+76
View File
@@ -1,5 +1,81 @@
# Changelog
## [1.535.0](https://github.com/windmill-labs/windmill/compare/v1.534.1...v1.535.0) (2025-08-25)
### Features
* **aichat:** show diff mode on inline scripts changes ([#6454](https://github.com/windmill-labs/windmill/issues/6454)) ([eca3109](https://github.com/windmill-labs/windmill/commit/eca3109ec63967e3041521bf74d34a12c70f5ff8))
### Bug Fixes
* fix opening advanced popup for run resetting tag to default ([d328894](https://github.com/windmill-labs/windmill/commit/d3288947b2d2539b2f3302059a9aad2841275a28))
## [1.534.1](https://github.com/windmill-labs/windmill/compare/v1.534.0...v1.534.1) (2025-08-25)
### Bug Fixes
* add alias to subquery for older postgres versions ([#6455](https://github.com/windmill-labs/windmill/issues/6455)) ([16d233b](https://github.com/windmill-labs/windmill/commit/16d233bf466fd818ec1f9235377e4c1a8239d98c))
* **frontend:** fix test step behavior ([#6427](https://github.com/windmill-labs/windmill/issues/6427)) ([fc20b7b](https://github.com/windmill-labs/windmill/commit/fc20b7bd91d33115aacb38cc46394f9c6465aa0f))
## [1.534.0](https://github.com/windmill-labs/windmill/compare/v1.533.1...v1.534.0) (2025-08-25)
### Features
* **backend:** support unencrypted connection to mssql ([#6453](https://github.com/windmill-labs/windmill/issues/6453)) ([8d31c2a](https://github.com/windmill-labs/windmill/commit/8d31c2ab0d34036dc8057611857a5d72aad8598f))
### Bug Fixes
* **aichat:** fix wrong current model logic ([#6451](https://github.com/windmill-labs/windmill/issues/6451)) ([e951c89](https://github.com/windmill-labs/windmill/commit/e951c896b865df48d331968953c9e44848236516))
* **flow:** test this step preload step input evaluation ([1073eb0](https://github.com/windmill-labs/windmill/commit/1073eb0e682e7bd253c6d62225361b487d7f6d2f))
## [1.533.1](https://github.com/windmill-labs/windmill/compare/v1.533.0...v1.533.1) (2025-08-23)
### Bug Fixes
* **app:** fix oneOf selected undefined freeze ([0ae8f44](https://github.com/windmill-labs/windmill/commit/0ae8f44773adb0576e1b63e858b506ba9a9fe7b3))
## [1.533.0](https://github.com/windmill-labs/windmill/compare/v1.532.0...v1.533.0) (2025-08-23)
### Features
* CLI improvements ([#6446](https://github.com/windmill-labs/windmill/issues/6446)) ([a41b9e4](https://github.com/windmill-labs/windmill/commit/a41b9e47e233ebaa2baafb5cca1187bb85d6f8f4))
### Bug Fixes
* **frontend:** ai agent flow status + UI nits ([#6447](https://github.com/windmill-labs/windmill/issues/6447)) ([c13747c](https://github.com/windmill-labs/windmill/commit/c13747cda9449369288e8d078b60542ea79a49bf))
## [1.532.0](https://github.com/windmill-labs/windmill/compare/v1.531.0...v1.532.0) (2025-08-22)
### Features
* **aichat:** allow adding contexts to flow mode ([#6424](https://github.com/windmill-labs/windmill/issues/6424)) ([73272f1](https://github.com/windmill-labs/windmill/commit/73272f16fddc355703b04f2c3458520753d1e19c))
* json schema resource ([#6433](https://github.com/windmill-labs/windmill/issues/6433)) ([7da79a8](https://github.com/windmill-labs/windmill/commit/7da79a8bc525fc6b89748ad0af25c2bac4ca2ef3))
## [1.531.0](https://github.com/windmill-labs/windmill/compare/v1.530.0...v1.531.0) (2025-08-22)
### Features
* ai agent steps ([#6393](https://github.com/windmill-labs/windmill/issues/6393)) ([958e8af](https://github.com/windmill-labs/windmill/commit/958e8af78290cf859f98c45c012ed41e3bada39e))
* bump Go version from 1.22.0 to 1.25.0 [#6415](https://github.com/windmill-labs/windmill/issues/6415) ([c92bfe6](https://github.com/windmill-labs/windmill/commit/c92bfe6601fd96f6d74860f52f9307e02961ac21))
### Bug Fixes
* **app:** fix ctrl drag for insertion into subgrids ([51ea947](https://github.com/windmill-labs/windmill/commit/51ea9473ef23c6871699e69bbe79772a4d50d3b8))
* **frontend:** graph cache of ai agent step tools ([#6431](https://github.com/windmill-labs/windmill/issues/6431)) ([28f1d61](https://github.com/windmill-labs/windmill/commit/28f1d611643459d42531fa217c185408eb97d6d1))
* make relevant sidebar menu items a instead of button ([06d078e](https://github.com/windmill-labs/windmill/commit/06d078ebfa8f70b66bc764eae70d33c8c57b4012))
* s3 result presigned not working with list ([9df008b](https://github.com/windmill-labs/windmill/commit/9df008b9f8fe58692463e4b9da0538935e458b10))
## [1.530.0](https://github.com/windmill-labs/windmill/compare/v1.529.0...v1.530.0) (2025-08-20)
+11
View File
@@ -4,6 +4,17 @@
Windmill is an open-source developer platform for building internal tools, workflows, API integrations, background jobs, workflows, and user interfaces. See @windmill-overview.mdc for full platform details.
## New Feature Implementation Guidelines
When implementing new features in Windmill, follow these best practices:
- **Clean Code First**: Write clean, readable, and maintainable code. Prioritize clarity over cleverness.
- **Avoid Duplication at All Costs**: Before writing new code, thoroughly search for existing implementations that can be reused or extended.
- **Adapt Existing Code**: Refactor and generalize existing code when necessary to avoid logic duplication. Extract common patterns into reusable utilities.
- **Follow Established Patterns**: Study existing code patterns in the codebase and maintain consistency with established conventions.
- **Single Responsibility**: Each function, component, and module should have a single, well-defined responsibility.
- **Incremental Implementation**: Break large features into smaller, reviewable chunks that can be implemented and tested incrementally.
## Language-Specific Guides
- Backend (Rust): @backend/rust-best-practices.mdc + @backend/summarized_schema.txt
+48 -51
View File
@@ -332,40 +332,40 @@ you to have it being synced automatically everyday.
## Environment Variables
| Environment Variable name | Default | Description | Api Server/Worker/All |
| ----------------------------------- | ---------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------- |
| DATABASE_URL | | The Postgres database url. | All |
| WORKER_GROUP | default | The worker group the worker belongs to and get its configuration pulled from | Worker |
| MODE | standalone | The mode if the binary. Possible values: standalone, worker, server, agent | All |
| METRICS_ADDR | None | (ee only) The socket addr at which to expose Prometheus metrics at the /metrics path. Set to "true" to expose it on port 8001 | All |
| JSON_FMT | false | Output the logs in json format instead of logfmt | All |
| BASE_URL | http://localhost:8000 | The base url that is exposed publicly to access your instance. Is overriden by the instance settings if any. | Server |
| ZOMBIE_JOB_TIMEOUT | 30 | The timeout after which a job is considered to be zombie if the worker did not send pings about processing the job (every server check for zombie jobs every 30s) | Server |
| RESTART_ZOMBIE_JOBS | true | If true then a zombie job is restarted (in-place with the same uuid and some logs), if false the zombie job is failed | Server |
| SLEEP_QUEUE | 50 | The number of ms to sleep in between the last check for new jobs in the DB. It is multiplied by NUM_WORKERS such that in average, for one worker instance, there is one pull every SLEEP_QUEUE ms. | Worker |
| KEEP_JOB_DIR | false | Keep the job directory after the job is done. Useful for debugging. | Worker |
| LICENSE_KEY (EE only) | None | License key checked at startup for the Enterprise Edition of Windmill | Worker |
| SLACK_SIGNING_SECRET | None | The signing secret of your Slack app. See [Slack documentation](https://api.slack.com/authentication/verifying-requests-from-slack) | Server |
| COOKIE_DOMAIN | None | The domain of the cookie. If not set, the cookie will be set by the browser based on the full origin | Server |
| DENO_PATH | /usr/bin/deno | The path to the deno binary. | Worker |
| PYTHON_PATH | | The path to the python binary if wanting to not have it managed by uv. | Worker |
| GO_PATH | /usr/bin/go | The path to the go binary. | Worker |
| GOPRIVATE | | The GOPRIVATE env variable to use private go modules | Worker |
| GOPROXY | | The GOPROXY env variable to use | Worker |
| NETRC | | The netrc content to use a private go registry | Worker |
| PY_CONCURRENT_DOWNLOADS | 20 | Sets the maximum number of in-flight concurrent python downloads that windmill will perform at any given time. | Worker |
| PATH | None | The path environment variable, usually inherited | Worker |
| HOME | None | The home directory to use for Go and Bash , usually inherited | Worker |
| DATABASE_CONNECTIONS | 50 (Server)/3 (Worker) | The max number of connections in the database connection pool | All |
| SUPERADMIN_SECRET | None | A token that would let the caller act as a virtual superadmin superadmin@windmill.dev | Server |
| TIMEOUT_WAIT_RESULT | 20 | The number of seconds to wait before timeout on the 'run_wait_result' endpoint | Worker |
| QUEUE_LIMIT_WAIT_RESULT | None | The number of max jobs in the queue before rejecting immediately the request in 'run_wait_result' endpoint. Takes precedence on the query arg. If none is specified, there are no limit. | Worker |
| DENO_AUTH_TOKENS | None | Custom DENO_AUTH_TOKENS to pass to worker to allow the use of private modules | Worker |
| DISABLE_RESPONSE_LOGS | false | Disable response logs | Server |
| CREATE_WORKSPACE_REQUIRE_SUPERADMIN | true | If true, only superadmins can create new workspaces | Server |
| MIN_FREE_DISK_SPACE_MB | 15000 | Minimum amount of free space on worker. Sends critical alert if worker has less free space. | Worker |
| RUN_UPDATE_CA_CERTIFICATE_AT_START | false | If true, runs CA certificate update command at startup before other initialization | All |
| RUN_UPDATE_CA_CERTIFICATE_PATH | /usr/sbin/update-ca-certificates | Path to the CA certificate update command/script to run when RUN_UPDATE_CA_CERTIFICATE_AT_START is true | All |
| Environment Variable name | Default | Description | Api Server/Worker/All |
| ----------------------------------- | -------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------- |
| DATABASE_URL | | The Postgres database url. | All |
| WORKER_GROUP | default | The worker group the worker belongs to and get its configuration pulled from | Worker |
| MODE | standalone | The mode if the binary. Possible values: standalone, worker, server, agent | All |
| METRICS_ADDR | None | (ee only) The socket addr at which to expose Prometheus metrics at the /metrics path. Set to "true" to expose it on port 8001 | All |
| JSON_FMT | false | Output the logs in json format instead of logfmt | All |
| BASE_URL | http://localhost:8000 | The base url that is exposed publicly to access your instance. Is overriden by the instance settings if any. | Server |
| ZOMBIE_JOB_TIMEOUT | 30 | The timeout after which a job is considered to be zombie if the worker did not send pings about processing the job (every server check for zombie jobs every 30s) | Server |
| RESTART_ZOMBIE_JOBS | true | If true then a zombie job is restarted (in-place with the same uuid and some logs), if false the zombie job is failed | Server |
| SLEEP_QUEUE | 50 | The number of ms to sleep in between the last check for new jobs in the DB. It is multiplied by NUM_WORKERS such that in average, for one worker instance, there is one pull every SLEEP_QUEUE ms. | Worker |
| KEEP_JOB_DIR | false | Keep the job directory after the job is done. Useful for debugging. | Worker |
| LICENSE_KEY (EE only) | None | License key checked at startup for the Enterprise Edition of Windmill | Worker |
| SLACK_SIGNING_SECRET | None | The signing secret of your Slack app. See [Slack documentation](https://api.slack.com/authentication/verifying-requests-from-slack) | Server |
| COOKIE_DOMAIN | None | The domain of the cookie. If not set, the cookie will be set by the browser based on the full origin | Server |
| DENO_PATH | /usr/bin/deno | The path to the deno binary. | Worker |
| PYTHON_PATH | | The path to the python binary if wanting to not have it managed by uv. | Worker |
| GO_PATH | /usr/bin/go | The path to the go binary. | Worker |
| GOPRIVATE | | The GOPRIVATE env variable to use private go modules | Worker |
| GOPROXY | | The GOPROXY env variable to use | Worker |
| NETRC | | The netrc content to use a private go registry | Worker |
| PY_CONCURRENT_DOWNLOADS | 20 | Sets the maximum number of in-flight concurrent python downloads that windmill will perform at any given time. | Worker |
| PATH | None | The path environment variable, usually inherited | Worker |
| HOME | None | The home directory to use for Go and Bash , usually inherited | Worker |
| DATABASE_CONNECTIONS | 50 (Server)/3 (Worker) | The max number of connections in the database connection pool | All |
| SUPERADMIN_SECRET | None | A token that would let the caller act as a virtual superadmin superadmin@windmill.dev | Server |
| TIMEOUT_WAIT_RESULT | 20 | The number of seconds to wait before timeout on the 'run_wait_result' endpoint | Worker |
| QUEUE_LIMIT_WAIT_RESULT | None | The number of max jobs in the queue before rejecting immediately the request in 'run_wait_result' endpoint. Takes precedence on the query arg. If none is specified, there are no limit. | Worker |
| DENO_AUTH_TOKENS | None | Custom DENO_AUTH_TOKENS to pass to worker to allow the use of private modules | Worker |
| DISABLE_RESPONSE_LOGS | false | Disable response logs | Server |
| CREATE_WORKSPACE_REQUIRE_SUPERADMIN | true | If true, only superadmins can create new workspaces | Server |
| MIN_FREE_DISK_SPACE_MB | 15000 | Minimum amount of free space on worker. Sends critical alert if worker has less free space. | Worker |
| RUN_UPDATE_CA_CERTIFICATE_AT_START | false | If true, runs CA certificate update command at startup before other initialization | All |
| RUN_UPDATE_CA_CERTIFICATE_PATH | /usr/sbin/update-ca-certificates | Path to the CA certificate update command/script to run when RUN_UPDATE_CA_CERTIFICATE_AT_START is true | All |
## Run a local dev setup
@@ -374,7 +374,6 @@ Using [Nix](./frontend/README_DEV.md#nix) (Recommended).
See the [./frontend/README_DEV.md](./frontend/README_DEV.md) file for all
running options.
### only Frontend
This will use the backend of <https://app.windmill.dev> but your own frontend
@@ -400,29 +399,27 @@ npm run generate-backend-client-mac
See the [./frontend/README_DEV.md](./frontend/README_DEV.md) file for all
running options.
1. Create a Postgres Database for Windmill and create an admin role inside your
Postgres setup. The easiest way to get a working db is to run
1. Start a local Postgres database using for instance the `start-dev-db.sh` script which will make a database available at `postgres://postgres:changeme@localhost:5432/windmill`
Then run the migrations using the following command:
```
cargo install sqlx-cli
env DATABASE_URL=<YOUR_DATABASE_URL> sqlx migrate run
```
This will also avoid compile time issue with sqlx's `query!` macro
2. Install [nsjail](https://github.com/google/nsjail) and have it accessible in
This will also avoid compile time issue with sqlx's `query!` macro.
2. (optional, linux only) Install [nsjail](https://github.com/google/nsjail) and have it accessible in
your PATH
3. Install deno and python3, have the bins at `/usr/bin/deno` and
`/usr/local/bin/python3`
4. Install [caddy](https://caddyserver.com)
5. Install the [lld linker](https://lld.llvm.org/)
6. Go to `frontend/`:
1. `npm install`, `npm run generate-backend-client` then `npm run dev`
3. Install bun, deno and python3 (+ any languages you want to use), have the bins at `/usr/bin/bun`,`/usr/bin/deno`, and
`/usr/local/bin/python3` or set the corresponding environment variables.
4. (optional) Install the [lld linker](https://lld.llvm.org/)
5. Go to `frontend/`:
1. `npm install`, `npm run generate-backend-client` then `REMOTE=http://localhost:8000 npm run dev`
2. You might need to set some extra heap space for the node runtime
`export NODE_OPTIONS="--max-old-space-size=4096"`
3. In another shell `npm run build` otherwise the backend will not find the
`frontend/build` folder and will not compile.
4. In another shell `sudo caddy run --config Caddyfile`
7. Go to `backend/`:
`env DATABASE_URL=<DATABASE_URL_TO_YOUR_WINDMILL_DB> RUST_LOG=info cargo run`
8. Et voilà, windmill should be available at `http://localhost/`
3. Create an empty `frontend/build` folder using `mkdir frontend/build`
6. Go to `backend/`:
1. `env DATABASE_URL=<YOUR_DATABASE_URL> RUST_LOG=info cargo run`
2. You can specify any feature flag you want to enable, for example `cargo run --features python` to enable the python executor.
7. Et voilà, windmill should be available at `http://localhost:3000`
## Contributors
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT elem FROM (SELECT unnest($1::TEXT[]) AS elem)\n WHERE elem NOT IN (SELECT datname FROM pg_catalog.pg_database);",
"query": "SELECT elem FROM (SELECT unnest($1::TEXT[]) AS elem) AS e\n WHERE elem NOT IN (SELECT datname FROM pg_catalog.pg_database);",
"describe": {
"columns": [
{
@@ -18,5 +18,5 @@
null
]
},
"hash": "543859cf1c8d9e3bf2c2b23d21d096d01fe7a72d749229f21634b549c6b1241a"
"hash": "298f8609319a2928257fd5be60bb37f292c786d2348efe11d19868e5dc8fba11"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT workspace_id as workspace, path, summary, description, schema FROM script as o WHERE created_at = (select max(created_at) from script where o.path = path and workspace_id = $1) and workspace_id = $1",
"query": "SELECT workspace_id as workspace, path, summary, description, schema FROM script as o \n WHERE created_at = (select max(created_at) from script where o.path = path and workspace_id = $1 AND archived = false) \n AND workspace_id = $1 and archived = false",
"describe": {
"columns": [
{
@@ -42,5 +42,5 @@
true
]
},
"hash": "53e7243abd724816fb8d09c63b7ffa65f1cd622a989f5cefedbbf3c143b387c4"
"hash": "2d5f58dd2aff3bd49f3891ae76df23e2aa39891931516426f65b229314a0cee1"
}
@@ -18,8 +18,8 @@
"Left": []
},
"nullable": [
true,
false
false,
true
]
},
"hash": "b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76"
@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "CREATE INDEX CONCURRENTLY IF NOT EXISTS script_not_archived ON script (workspace_id, path, created_at DESC) where archived = false;",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "c1058d8816d139c63dd9c4a075ab63efc585942b30c9e853f2a5cff4cc9916cd"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT content FROM script WHERE path = $1 AND workspace_id = $2\n AND created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND\n workspace_id = $2)\n ",
"query": "\n SELECT content FROM script WHERE path = $1 AND workspace_id = $2\n AND archived = false ORDER BY created_at DESC LIMIT 1\n ",
"describe": {
"columns": [
{
@@ -19,5 +19,5 @@
false
]
},
"hash": "443bd83bcea1d37c79cb080095343c98104529879f991c49585cd181e34aa827"
"hash": "c7cae4cf872fce0a989cf89aa35929218a9d459ee1c2b36a28b110e9741ab623"
}
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "SELECT value from resource WHERE path = $1 AND workspace_id = $2 AND resource_type = 'json_schema'",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "value",
"type_info": "Jsonb"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
true
]
},
"hash": "d4c963fa653652b7a3e8529cbf0d0fca091d7c1cb0924f6f9343544abb2666a5"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "select hash, tag, concurrency_key, concurrent_limit, concurrency_time_window_s, cache_ttl, language as \"language: ScriptLang\", dedicated_worker, priority, timeout, on_behalf_of_email, created_by FROM script where path = $1 AND workspace_id = $2 AND\n created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND\n deleted = false AND archived = false)",
"query": "select hash, tag, concurrency_key, concurrent_limit, concurrency_time_window_s, cache_ttl, language as \"language: ScriptLang\", dedicated_worker, priority, timeout, on_behalf_of_email, created_by FROM script\n WHERE path = $1 AND workspace_id = $2 AND archived = false AND (lock IS NOT NULL OR $3 = false)\n ORDER BY created_at DESC LIMIT 1",
"describe": {
"columns": [
{
@@ -98,7 +98,8 @@
"parameters": {
"Left": [
"Text",
"Text"
"Text",
"Bool"
]
},
"nullable": [
@@ -116,5 +117,5 @@
false
]
},
"hash": "61c29d684e8e683e839a6d7210b3b9b96854e5bfd752e45922c358c42ebea0c4"
"hash": "e64f7044c74e96c2338580562f6b087805dad2b6fb1aa194ac2a2026fa24ecd0"
}
+173 -173
View File
File diff suppressed because it is too large Load Diff
+2 -2
View File
@@ -1,6 +1,6 @@
[package]
name = "windmill"
version = "1.530.0"
version = "1.535.0"
authors.workspace = true
edition.workspace = true
@@ -33,7 +33,7 @@ members = [
]
[workspace.package]
version = "1.530.0"
version = "1.535.0"
authors = ["Ruben Fiszel <ruben@windmill.dev>"]
edition = "2021"
@@ -489,8 +489,7 @@ async fn parse_python_imports_inner(
let code = sqlx::query_scalar!(
r#"
SELECT content FROM script WHERE path = $1 AND workspace_id = $2
AND created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND
workspace_id = $2)
AND archived = false ORDER BY created_at DESC LIMIT 1
"#,
&rpath,
w_id
-2
View File
@@ -2571,8 +2571,6 @@ pub async fn reload_app_workspaced_route_setting(conn: &DB) -> error::Result<()>
let app_workspaced_route =
load_value_from_global_settings(conn, APP_WORKSPACED_ROUTE_SETTING).await?;
println!("Updating...");
let ws_route = match app_workspaced_route {
Some(serde_json::Value::Bool(ws_route)) => ws_route,
None => false,
+1 -1
View File
@@ -1,7 +1,7 @@
openapi: "3.0.3"
info:
version: 1.530.0
version: 1.535.0
title: Windmill API
contact:
+20 -14
View File
@@ -762,20 +762,26 @@ async fn get_public_resource(
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Option<serde_json::Value>> {
let path = path.to_path();
if !path.starts_with("f/app_themes/") {
return Err(Error::BadRequest(
"Only app themes are public resources".to_string(),
));
}
let res = sqlx::query_scalar!(
"SELECT value from resource WHERE path = $1 AND workspace_id = $2",
path.to_owned(),
&w_id
)
.fetch_optional(&db)
.await?
.flatten();
Ok(Json(res))
let res = if path.starts_with("f/app_themes/") {
sqlx::query_scalar!(
"SELECT value from resource WHERE path = $1 AND workspace_id = $2",
path.to_owned(),
&w_id
)
.fetch_optional(&db)
.await?
} else {
sqlx::query_scalar!(
"SELECT value from resource WHERE path = $1 AND workspace_id = $2 AND resource_type = 'json_schema'",
path.to_owned(),
&w_id
)
.fetch_optional(&db)
.await?
};
Ok(Json(res.flatten()))
}
async fn get_secret_id(
+3 -584
View File
@@ -6,8 +6,6 @@
* LICENSE-AGPL for a copy of the license.
*/
use std::time::Duration;
use futures::FutureExt;
use sqlx::{
migrate::{Migrate, MigrateError},
@@ -17,11 +15,11 @@ use sqlx::{
use tokio::task::JoinHandle;
use windmill_audit::audit_oss::{AuditAuthor, AuditAuthorable};
use windmill_common::utils::generate_lock_id;
use windmill_common::{
db::{Authable, Authed},
error::Error,
};
use windmill_common::{utils::generate_lock_id, worker::MIN_VERSION_IS_AT_LEAST_1_461};
pub type DB = Pool<Postgres>;
@@ -59,7 +57,7 @@ lazy_static::lazy_static! {
].into_iter().collect();
}
struct CustomMigrator {
pub struct CustomMigrator {
inner: PoolConnection<Postgres>,
}
impl Migrate for CustomMigrator {
@@ -243,586 +241,7 @@ pub async fn migrate(db: &DB) -> Result<Option<JoinHandle<()>>, Error> {
Err(err) => Err(err),
}?;
if let Err(err) = fix_flow_versioning_migration(&mut custom_migrator, db).await {
tracing::error!("Could not apply flow versioning fix migration: {err:#}");
}
let db2 = db.clone();
let _ = tokio::task::spawn(async move {
if let Err(err) = fix_job_completed_index(&db2).await {
tracing::error!("Could not apply job completed index fix migration: {err:#}");
}
});
let mut jh = None;
if !has_done_migration(db, "v2_finalize_job_completed").await {
let db2 = db.clone();
let v2jh = tokio::task::spawn(async move {
loop {
if !*MIN_VERSION_IS_AT_LEAST_1_461.read().await {
tracing::info!("Waiting for all workers to be at least version 1.461 before applying v2 finalize migration, sleeping for 5s...");
tokio::time::sleep(Duration::from_secs(5)).await;
continue;
}
if let Err(err) = v2_finalize(&db2).await {
tracing::error!(
"{err:#}: Could not apply v2 finalize migration, retry in 30s.."
);
tokio::time::sleep(Duration::from_secs(30)).await;
continue;
}
tracing::info!("v2 finalization step successfully applied.");
break;
}
});
jh = Some(v2jh)
}
Ok(jh)
}
async fn fix_flow_versioning_migration(
migrator: &mut CustomMigrator,
db: &DB,
) -> Result<(), Error> {
let has_done_migration = sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_flow_versioning_2')",
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !has_done_migration {
migrator.lock().await?;
if migrator
.list_applied_migrations()
.await?
.iter()
.any(|x| x.version == 20240630102146)
{
let has_done_migration = sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_flow_versioning_2')",
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !has_done_migration {
let query = include_str!("../../custom_migrations/fix_flow_versioning_2.sql");
tracing::info!("Applying fix_flow_versioning_2.sql");
let mut tx: sqlx::Transaction<'_, Postgres> = db.begin().await?;
tx.execute(query).await?;
tracing::info!("Applied fix_flow_versioning_2.sql");
sqlx::query!(
"INSERT INTO windmill_migrations (name) VALUES ('fix_flow_versioning_2')"
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
}
}
migrator.unlock().await?;
}
Ok(())
}
async fn has_done_migration(db: &DB, migration_job_name: &str) -> bool {
sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)",
migration_job_name
)
.fetch_one(db)
.await
.ok()
.flatten()
.unwrap_or(false)
}
macro_rules! run_windmill_migration {
($migration_job_name:expr, $db:expr, |$tx:ident| $code:block) => {
{
let migration_job_name = $migration_job_name;
let db: &Pool<Postgres> = $db;
let has_done = has_done_migration(db, migration_job_name).await;
if !has_done {
tracing::info!("Applying {migration_job_name} migration");
let mut $tx = db.begin().await?;
let mut r = false;
while !r {
r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)")
.fetch_one(&mut *$tx)
.await
.map_err(|e| {
tracing::error!("Error acquiring {migration_job_name} lock: {e:#}");
sqlx::migrate::MigrateError::Execute(e)
})?
.unwrap_or(false);
if !r {
tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)");
drop($tx);
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
$tx = db.begin().await?;
}
}
tracing::info!("acquired lock for {migration_job_name}");
let has_done = has_done_migration(db, migration_job_name).await;
if !has_done {
$code
sqlx::query!(
"INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING",
migration_job_name
)
.execute(&mut *$tx)
.await?;
tracing::info!("Finished applying {migration_job_name} migration");
} else {
tracing::debug!("migration {migration_job_name} already done");
}
let _ = sqlx::query("SELECT pg_advisory_unlock(4242)")
.execute(&mut *$tx)
.await?;
$tx.commit().await?;
tracing::info!("released lock for {migration_job_name}");
} else {
tracing::debug!("migration {migration_job_name} already done");
}
}
};
}
async fn v2_finalize(db: &DB) -> Result<(), Error> {
run_windmill_migration!("v2_finalize_disable_sync_III", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue DISABLE ROW LEVEL SECURITY;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_2", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed DISABLE ROW LEVEL SECURITY;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_3", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_after_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_4", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_completed_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_completed_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_5", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_queue_after_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_6", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_runtime IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_7", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_status IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_status_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_status_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_8", db, |tx| {
tx.execute(
r#"
DROP VIEW IF EXISTS completed_job, completed_job_view, job, queue, queue_view CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_job_queue", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
DROP COLUMN IF EXISTS __script_hash CASCADE,
DROP COLUMN IF EXISTS __script_path CASCADE,
DROP COLUMN IF EXISTS __args CASCADE,
DROP COLUMN IF EXISTS __logs CASCADE,
DROP COLUMN IF EXISTS __raw_code CASCADE,
DROP COLUMN IF EXISTS __canceled CASCADE,
DROP COLUMN IF EXISTS __last_ping CASCADE,
DROP COLUMN IF EXISTS __job_kind CASCADE,
DROP COLUMN IF EXISTS __env_id CASCADE,
DROP COLUMN IF EXISTS __schedule_path CASCADE,
DROP COLUMN IF EXISTS __permissioned_as CASCADE,
DROP COLUMN IF EXISTS __flow_status CASCADE,
DROP COLUMN IF EXISTS __raw_flow CASCADE,
DROP COLUMN IF EXISTS __is_flow_step CASCADE,
DROP COLUMN IF EXISTS __language CASCADE,
DROP COLUMN IF EXISTS __same_worker CASCADE,
DROP COLUMN IF EXISTS __raw_lock CASCADE,
DROP COLUMN IF EXISTS __pre_run_error CASCADE,
DROP COLUMN IF EXISTS __email CASCADE,
DROP COLUMN IF EXISTS __visible_to_owner CASCADE,
DROP COLUMN IF EXISTS __mem_peak CASCADE,
DROP COLUMN IF EXISTS __root_job CASCADE,
DROP COLUMN IF EXISTS __leaf_jobs CASCADE,
DROP COLUMN IF EXISTS __concurrent_limit CASCADE,
DROP COLUMN IF EXISTS __concurrency_time_window_s CASCADE,
DROP COLUMN IF EXISTS __timeout CASCADE,
DROP COLUMN IF EXISTS __flow_step_id CASCADE,
DROP COLUMN IF EXISTS __cache_ttl CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_job_completed", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
DROP COLUMN IF EXISTS __created_at CASCADE,
DROP COLUMN IF EXISTS __success CASCADE,
DROP COLUMN IF EXISTS __script_hash CASCADE,
DROP COLUMN IF EXISTS __script_path CASCADE,
DROP COLUMN IF EXISTS __args CASCADE,
DROP COLUMN IF EXISTS __logs CASCADE,
DROP COLUMN IF EXISTS __raw_code CASCADE,
DROP COLUMN IF EXISTS __canceled CASCADE,
DROP COLUMN IF EXISTS __job_kind CASCADE,
DROP COLUMN IF EXISTS __env_id CASCADE,
DROP COLUMN IF EXISTS __schedule_path CASCADE,
DROP COLUMN IF EXISTS __permissioned_as CASCADE,
DROP COLUMN IF EXISTS __raw_flow CASCADE,
DROP COLUMN IF EXISTS __is_flow_step CASCADE,
DROP COLUMN IF EXISTS __language CASCADE,
DROP COLUMN IF EXISTS __is_skipped CASCADE,
DROP COLUMN IF EXISTS __raw_lock CASCADE,
DROP COLUMN IF EXISTS __email CASCADE,
DROP COLUMN IF EXISTS __visible_to_owner CASCADE,
DROP COLUMN IF EXISTS __tag CASCADE,
DROP COLUMN IF EXISTS __priority CASCADE;
"#,
)
.await?;
});
Ok(())
}
async fn fix_job_completed_index(db: &DB) -> Result<(), Error> {
// let has_done_migration = sqlx::query_scalar!(
// "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_job_completed_index')"
// )
// .fetch_one(db)
// .await?
// .unwrap_or(false);
// if !has_done_migration {
// tracing::info!("Applying fix_job_completed_index migration");
// let mut tx = db.begin().await?;
// let mut r = false;
// while !r {
// r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)")
// .fetch_one(&mut *tx)
// .await
// .map_err(|e| {
// tracing::error!("Error acquiring fix_job_completed_index lock: {e:#}");
// sqlx::migrate::MigrateError::Execute(e)
// })?
// .unwrap_or(false);
// if !r {
// tracing::info!("PG fix_job_completed_index_migration lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)");
// tokio::time::sleep(std::time::Duration::from_secs(5)).await;
// }
// }
// // sqlx::query(
// // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new ON completed_job (workspace_id, job_kind, is_skipped, is_flow_step, created_at DESC, started_at DESC)"
// // ).execute(db).await?;
// sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at")
// .execute(db)
// .await?;
// sqlx::query!("INSERT INTO windmill_migrations (name) VALUES ('fix_job_completed_index') ON CONFLICT DO NOTHING")
// .execute(&mut *tx)
// .await?;
// let _ = sqlx::query("SELECT pg_advisory_unlock(4242)")
// .execute(&mut *tx)
// .await?;
// tx.commit().await?;
// }
run_windmill_migration!("fix_job_completed_index_2", &db, |tx| {
// sqlx::query(
// "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new_2 ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, created_at DESC)"
// ).execute(db).await?;
// sqlx::query(
// "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_started_at_new ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, started_at DESC)"
// ).execute(db).await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at")
.execute(db)
.await?;
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new",
)
.execute(db)
.await?;
});
run_windmill_migration!("fix_job_completed_index_3", &db, |tx| {
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_job_on_schedule_path")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS concurrency_limit_stats_queue")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_on_created")
.execute(db)
.await?;
});
run_windmill_migration!("fix_job_index_1_II", &db, |tx| {
let migration_job_name = "fix_job_index_1_II";
let mut i = 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_3 ON v2_job (workspace_id, created_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_8 ON v2_job (workspace_id, created_at DESC) where kind in ('deploymentcallback') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_9 ON v2_job (workspace_id, created_at DESC) where kind in ('dependencies', 'flowdependencies', 'appdependencies') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_5 ON v2_job (workspace_id, created_at DESC) where kind in ('preview', 'flowpreview') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_completed_job_workspace_id_started_at_new_2 ON v2_job_completed (workspace_id, started_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path_2")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_created_at ON v2_job (created_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_2",
)
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_started_at_new",
)
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path")
.execute(db)
.await?;
});
run_windmill_migration!("fix_labeled_jobs_index", &db, |tx| {
tracing::info!("Special migration to add index concurrently on job labels 2");
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS labeled_jobs_on_jobs")
.execute(db)
.await?;
sqlx::query!(
"CREATE INDEX CONCURRENTLY labeled_jobs_on_jobs ON v2_job_completed USING GIN ((result -> 'wm_labels')) WHERE result ? 'wm_labels'"
).execute(db).await?;
});
run_windmill_migration!("v2_labeled_jobs_index", &db, |tx| {
tracing::info!("Special migration to add index concurrently on job labels");
sqlx::query!(
"CREATE INDEX CONCURRENTLY ix_v2_job_labels ON v2_job
USING GIN (labels)
WHERE labels IS NOT NULL"
)
.execute(db)
.await?;
});
run_windmill_migration!("v2_jobs_rls", &db, |tx| {
sqlx::query!("ALTER TABLE v2_job ENABLE ROW LEVEL SECURITY")
.execute(db)
.await?;
});
run_windmill_migration!("v2_improve_v2_job_indices_ii", &db, |tx| {
sqlx::query!("create index concurrently if not exists ix_v2_job_workspace_id_created_at ON v2_job (workspace_id, created_at DESC) where kind in ('script', 'flow', 'singlescriptflow') AND parent_job IS NULL")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS ix_job_workspace_id_created_at_new_6")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS ix_job_workspace_id_created_at_new_7")
.execute(db)
.await?;
});
run_windmill_migration!("v2_improve_v2_queued_jobs_indices", &db, |tx| {
sqlx::query!("CREATE INDEX CONCURRENTLY IF NOT EXISTS queue_sort_v2 ON v2_job_queue (priority DESC NULLS LAST, scheduled_for, tag) WHERE running = false")
.execute(db)
.await?;
// sqlx::query!("CREATE INDEX CONCURRENTLY queue_sort_2_v2 ON v2_job_queue (tag, priority DESC NULLS LAST, scheduled_for) WHERE running = false")
// .execute(db)
// .await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS queue_sort")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS queue_sort_2")
.execute(db)
.await?;
});
run_windmill_migration!("audit_timestamps", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_audit_timestamps ON audit (timestamp DESC)"
)
.execute(db)
.await?;
});
run_windmill_migration!("job_completed_completed_at", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_job_completed_completed_at ON v2_job_completed (completed_at DESC)"
)
.execute(db)
.await?;
});
run_windmill_migration!("alerts_by_workspace", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS alerts_by_workspace ON alerts (workspace_id);"
)
.execute(db)
.await?;
});
run_windmill_migration!("remove_redundant_log_file_index", db, |tx| {
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS log_file_hostname_log_ts_idx")
.execute(db)
.await?;
});
run_windmill_migration!("v2_job_queue_suspend", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS v2_job_queue_suspend ON v2_job_queue (workspace_id, suspend) WHERE suspend > 0;"
)
.execute(db)
.await?;
});
run_windmill_migration!("audit_recent_login_activities", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_audit_recent_login_activities
ON audit (timestamp, username)
WHERE operation IN ('users.login', 'oauth.login', 'users.token.refresh');"
)
.execute(db)
.await?;
});
Ok(())
return crate::live_migrations::custom_migrations(&mut custom_migrator, db).await;
}
#[derive(Clone, Debug, Default, Hash, Eq, PartialEq)]
+2 -1
View File
@@ -103,6 +103,7 @@ mod inkeep_ee;
mod inkeep_oss;
mod inputs;
mod integration;
mod live_migrations;
#[cfg(feature = "postgres_trigger")]
mod postgres_triggers;
@@ -623,7 +624,7 @@ pub async fn run_server(
.nest("/mqtt_triggers", mqtt_triggers_service)
.nest("/sqs_triggers", sqs_triggers_service)
.nest("/gcp_triggers", gcp_triggers_service)
.nest("/postgres_triggers", postgres_triggers_service),
.nest("/postgres_triggers", postgres_triggers_service),
)
.nest("/workspaces", workspaces::global_service())
.nest(
+613
View File
@@ -0,0 +1,613 @@
/*
* Author: Ruben Fiszel
* Copyright: Windmill Labs, Inc 2022
* 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.
*/
use sqlx::Postgres;
use std::time::Duration;
use tokio::task::JoinHandle;
use windmill_common::error::Error;
use windmill_common::worker::MIN_VERSION_IS_AT_LEAST_1_461;
use crate::db::{CustomMigrator, DB};
use sqlx::migrate::Migrate;
use sqlx::Executor;
pub async fn custom_migrations(
migrator: &mut CustomMigrator,
db: &DB,
) -> Result<Option<JoinHandle<()>>, Error> {
if let Err(err) = fix_flow_versioning_migration(migrator, db).await {
tracing::error!("Could not apply flow versioning fix migration: {err:#}");
}
let db2 = db.clone();
let _ = tokio::task::spawn(async move {
if let Err(err) = fix_job_completed_index(&db2).await {
tracing::error!("Could not apply job completed index fix migration: {err:#}");
}
});
let mut jh = None;
if !has_done_migration(db, "v2_finalize_job_completed").await {
let db2 = db.clone();
let v2jh = tokio::task::spawn(async move {
loop {
if !*MIN_VERSION_IS_AT_LEAST_1_461.read().await {
tracing::info!("Waiting for all workers to be at least version 1.461 before applying v2 finalize migration, sleeping for 5s...");
tokio::time::sleep(Duration::from_secs(5)).await;
continue;
}
if let Err(err) = v2_finalize(&db2).await {
tracing::error!(
"{err:#}: Could not apply v2 finalize migration, retry in 30s.."
);
tokio::time::sleep(Duration::from_secs(30)).await;
continue;
}
tracing::info!("v2 finalization step successfully applied.");
break;
}
});
jh = Some(v2jh)
}
Ok(jh)
}
async fn fix_flow_versioning_migration(
migrator: &mut CustomMigrator,
db: &DB,
) -> Result<(), Error> {
let has_done_migration = sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_flow_versioning_2')",
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !has_done_migration {
migrator.lock().await?;
if migrator
.list_applied_migrations()
.await?
.iter()
.any(|x| x.version == 20240630102146)
{
let has_done_migration = sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_flow_versioning_2')",
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !has_done_migration {
let query = include_str!("../../custom_migrations/fix_flow_versioning_2.sql");
tracing::info!("Applying fix_flow_versioning_2.sql");
let mut tx: sqlx::Transaction<'_, Postgres> = db.begin().await?;
tx.execute(query).await?;
tracing::info!("Applied fix_flow_versioning_2.sql");
sqlx::query!(
"INSERT INTO windmill_migrations (name) VALUES ('fix_flow_versioning_2')"
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
}
}
migrator.unlock().await?;
}
Ok(())
}
async fn has_done_migration(db: &DB, migration_job_name: &str) -> bool {
sqlx::query_scalar!(
"SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)",
migration_job_name
)
.fetch_one(db)
.await
.ok()
.flatten()
.unwrap_or(false)
}
use sqlx::Pool;
macro_rules! run_windmill_migration {
($migration_job_name:expr, $db:expr, |$tx:ident| $code:block) => {
{
let migration_job_name = $migration_job_name;
let db: &Pool<Postgres> = $db;
let has_done = has_done_migration(db, migration_job_name).await;
if !has_done {
tracing::info!("Applying {migration_job_name} migration");
let mut $tx = db.begin().await?;
let mut r = false;
while !r {
r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)")
.fetch_one(&mut *$tx)
.await
.map_err(|e| {
tracing::error!("Error acquiring {migration_job_name} lock: {e:#}");
sqlx::migrate::MigrateError::Execute(e)
})?
.unwrap_or(false);
if !r {
tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)");
drop($tx);
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
$tx = db.begin().await?;
}
}
tracing::info!("acquired lock for {migration_job_name}");
let has_done = has_done_migration(db, migration_job_name).await;
if !has_done {
$code
sqlx::query!(
"INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING",
migration_job_name
)
.execute(&mut *$tx)
.await?;
tracing::info!("Finished applying {migration_job_name} migration");
} else {
tracing::debug!("migration {migration_job_name} already done");
}
let _ = sqlx::query("SELECT pg_advisory_unlock(4242)")
.execute(&mut *$tx)
.await?;
$tx.commit().await?;
tracing::info!("released lock for {migration_job_name}");
} else {
tracing::debug!("migration {migration_job_name} already done");
}
}
};
}
async fn v2_finalize(db: &DB) -> Result<(), Error> {
run_windmill_migration!("v2_finalize_disable_sync_III", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue DISABLE ROW LEVEL SECURITY;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_2", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed DISABLE ROW LEVEL SECURITY;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_3", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_after_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_4", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_completed_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_completed_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_5", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_queue_after_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_6", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_runtime IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_7", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_status IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_status_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_status_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_8", db, |tx| {
tx.execute(
r#"
DROP VIEW IF EXISTS completed_job, completed_job_view, job, queue, queue_view CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_job_queue", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
DROP COLUMN IF EXISTS __script_hash CASCADE,
DROP COLUMN IF EXISTS __script_path CASCADE,
DROP COLUMN IF EXISTS __args CASCADE,
DROP COLUMN IF EXISTS __logs CASCADE,
DROP COLUMN IF EXISTS __raw_code CASCADE,
DROP COLUMN IF EXISTS __canceled CASCADE,
DROP COLUMN IF EXISTS __last_ping CASCADE,
DROP COLUMN IF EXISTS __job_kind CASCADE,
DROP COLUMN IF EXISTS __env_id CASCADE,
DROP COLUMN IF EXISTS __schedule_path CASCADE,
DROP COLUMN IF EXISTS __permissioned_as CASCADE,
DROP COLUMN IF EXISTS __flow_status CASCADE,
DROP COLUMN IF EXISTS __raw_flow CASCADE,
DROP COLUMN IF EXISTS __is_flow_step CASCADE,
DROP COLUMN IF EXISTS __language CASCADE,
DROP COLUMN IF EXISTS __same_worker CASCADE,
DROP COLUMN IF EXISTS __raw_lock CASCADE,
DROP COLUMN IF EXISTS __pre_run_error CASCADE,
DROP COLUMN IF EXISTS __email CASCADE,
DROP COLUMN IF EXISTS __visible_to_owner CASCADE,
DROP COLUMN IF EXISTS __mem_peak CASCADE,
DROP COLUMN IF EXISTS __root_job CASCADE,
DROP COLUMN IF EXISTS __leaf_jobs CASCADE,
DROP COLUMN IF EXISTS __concurrent_limit CASCADE,
DROP COLUMN IF EXISTS __concurrency_time_window_s CASCADE,
DROP COLUMN IF EXISTS __timeout CASCADE,
DROP COLUMN IF EXISTS __flow_step_id CASCADE,
DROP COLUMN IF EXISTS __cache_ttl CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_job_completed", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
DROP COLUMN IF EXISTS __created_at CASCADE,
DROP COLUMN IF EXISTS __success CASCADE,
DROP COLUMN IF EXISTS __script_hash CASCADE,
DROP COLUMN IF EXISTS __script_path CASCADE,
DROP COLUMN IF EXISTS __args CASCADE,
DROP COLUMN IF EXISTS __logs CASCADE,
DROP COLUMN IF EXISTS __raw_code CASCADE,
DROP COLUMN IF EXISTS __canceled CASCADE,
DROP COLUMN IF EXISTS __job_kind CASCADE,
DROP COLUMN IF EXISTS __env_id CASCADE,
DROP COLUMN IF EXISTS __schedule_path CASCADE,
DROP COLUMN IF EXISTS __permissioned_as CASCADE,
DROP COLUMN IF EXISTS __raw_flow CASCADE,
DROP COLUMN IF EXISTS __is_flow_step CASCADE,
DROP COLUMN IF EXISTS __language CASCADE,
DROP COLUMN IF EXISTS __is_skipped CASCADE,
DROP COLUMN IF EXISTS __raw_lock CASCADE,
DROP COLUMN IF EXISTS __email CASCADE,
DROP COLUMN IF EXISTS __visible_to_owner CASCADE,
DROP COLUMN IF EXISTS __tag CASCADE,
DROP COLUMN IF EXISTS __priority CASCADE;
"#,
)
.await?;
});
Ok(())
}
async fn fix_job_completed_index(db: &DB) -> Result<(), Error> {
// let has_done_migration = sqlx::query_scalar!(
// "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_job_completed_index')"
// )
// .fetch_one(db)
// .await?
// .unwrap_or(false);
// if !has_done_migration {
// tracing::info!("Applying fix_job_completed_index migration");
// let mut tx = db.begin().await?;
// let mut r = false;
// while !r {
// r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)")
// .fetch_one(&mut *tx)
// .await
// .map_err(|e| {
// tracing::error!("Error acquiring fix_job_completed_index lock: {e:#}");
// sqlx::migrate::MigrateError::Execute(e)
// })?
// .unwrap_or(false);
// if !r {
// tracing::info!("PG fix_job_completed_index_migration lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)");
// tokio::time::sleep(std::time::Duration::from_secs(5)).await;
// }
// }
// // sqlx::query(
// // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new ON completed_job (workspace_id, job_kind, is_skipped, is_flow_step, created_at DESC, started_at DESC)"
// // ).execute(db).await?;
// sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at")
// .execute(db)
// .await?;
// sqlx::query!("INSERT INTO windmill_migrations (name) VALUES ('fix_job_completed_index') ON CONFLICT DO NOTHING")
// .execute(&mut *tx)
// .await?;
// let _ = sqlx::query("SELECT pg_advisory_unlock(4242)")
// .execute(&mut *tx)
// .await?;
// tx.commit().await?;
// }
run_windmill_migration!("fix_job_completed_index_2", &db, |tx| {
// sqlx::query(
// "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new_2 ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, created_at DESC)"
// ).execute(db).await?;
// sqlx::query(
// "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_started_at_new ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, started_at DESC)"
// ).execute(db).await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at")
.execute(db)
.await?;
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new",
)
.execute(db)
.await?;
});
run_windmill_migration!("fix_job_completed_index_3", &db, |tx| {
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_job_on_schedule_path")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS concurrency_limit_stats_queue")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index")
.execute(db)
.await?;
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_on_created")
.execute(db)
.await?;
});
run_windmill_migration!("fix_job_index_1_II", &db, |tx| {
let migration_job_name = "fix_job_index_1_II";
let mut i = 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_3 ON v2_job (workspace_id, created_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_8 ON v2_job (workspace_id, created_at DESC) where kind in ('deploymentcallback') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_9 ON v2_job (workspace_id, created_at DESC) where kind in ('dependencies', 'flowdependencies', 'appdependencies') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_5 ON v2_job (workspace_id, created_at DESC) where kind in ('preview', 'flowpreview') AND parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_completed_job_workspace_id_started_at_new_2 ON v2_job_completed (workspace_id, started_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path_2")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_created_at ON v2_job (created_at DESC)")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_2",
)
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query(
"DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_started_at_new",
)
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path")
.execute(db)
.await?;
});
run_windmill_migration!("fix_labeled_jobs_index", &db, |tx| {
tracing::info!("Special migration to add index concurrently on job labels 2");
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS labeled_jobs_on_jobs")
.execute(db)
.await?;
sqlx::query!(
"CREATE INDEX CONCURRENTLY labeled_jobs_on_jobs ON v2_job_completed USING GIN ((result -> 'wm_labels')) WHERE result ? 'wm_labels'"
).execute(db).await?;
});
run_windmill_migration!("v2_labeled_jobs_index", &db, |tx| {
tracing::info!("Special migration to add index concurrently on job labels");
sqlx::query!(
"CREATE INDEX CONCURRENTLY ix_v2_job_labels ON v2_job
USING GIN (labels)
WHERE labels IS NOT NULL"
)
.execute(db)
.await?;
});
run_windmill_migration!("v2_jobs_rls", &db, |tx| {
sqlx::query!("ALTER TABLE v2_job ENABLE ROW LEVEL SECURITY")
.execute(db)
.await?;
});
run_windmill_migration!("v2_improve_v2_job_indices_ii", &db, |tx| {
sqlx::query!("create index concurrently if not exists ix_v2_job_workspace_id_created_at ON v2_job (workspace_id, created_at DESC) where kind in ('script', 'flow', 'singlescriptflow') AND parent_job IS NULL")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS ix_job_workspace_id_created_at_new_6")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS ix_job_workspace_id_created_at_new_7")
.execute(db)
.await?;
});
run_windmill_migration!("v2_improve_v2_queued_jobs_indices", &db, |tx| {
sqlx::query!("CREATE INDEX CONCURRENTLY IF NOT EXISTS queue_sort_v2 ON v2_job_queue (priority DESC NULLS LAST, scheduled_for, tag) WHERE running = false")
.execute(db)
.await?;
// sqlx::query!("CREATE INDEX CONCURRENTLY queue_sort_2_v2 ON v2_job_queue (tag, priority DESC NULLS LAST, scheduled_for) WHERE running = false")
// .execute(db)
// .await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS queue_sort")
.execute(db)
.await?;
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS queue_sort_2")
.execute(db)
.await?;
});
run_windmill_migration!("audit_timestamps", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_audit_timestamps ON audit (timestamp DESC)"
)
.execute(db)
.await?;
});
run_windmill_migration!("job_completed_completed_at", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_job_completed_completed_at ON v2_job_completed (completed_at DESC)"
)
.execute(db)
.await?;
});
run_windmill_migration!("alerts_by_workspace", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS alerts_by_workspace ON alerts (workspace_id);"
)
.execute(db)
.await?;
});
run_windmill_migration!("remove_redundant_log_file_index", db, |tx| {
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS log_file_hostname_log_ts_idx")
.execute(db)
.await?;
});
run_windmill_migration!("v2_job_queue_suspend", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS v2_job_queue_suspend ON v2_job_queue (workspace_id, suspend) WHERE suspend > 0;"
)
.execute(db)
.await?;
});
run_windmill_migration!("audit_recent_login_activities", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_audit_recent_login_activities
ON audit (timestamp, username)
WHERE operation IN ('users.login', 'oauth.login', 'users.token.refresh');"
)
.execute(db)
.await?;
});
run_windmill_migration!("v2_script_lock_index", db, |tx| {
sqlx::query!(
"CREATE INDEX CONCURRENTLY IF NOT EXISTS script_not_archived ON script (workspace_id, path, created_at DESC) where archived = false;"
)
.execute(db)
.await?;
});
Ok(())
}
+1 -1
View File
@@ -575,7 +575,7 @@ async fn databases_exist(
Json(database_names): Json<Vec<String>>,
) -> JsonResult<Vec<String>> {
let result = sqlx::query_scalar!(
r#"SELECT elem FROM (SELECT unnest($1::TEXT[]) AS elem)
r#"SELECT elem FROM (SELECT unnest($1::TEXT[]) AS elem) AS e
WHERE elem NOT IN (SELECT datname FROM pg_catalog.pg_database);"#,
database_names.as_slice()
)
+14 -4
View File
@@ -1327,8 +1327,9 @@ async fn update_workspace_user(
eu.operator,
eu.disabled,
&mut tx,
Some(&authed)
).await?;
Some(&authed),
)
.await?;
let user_email = sqlx::query_scalar!(
"SELECT email FROM usr WHERE username = $1 AND workspace_id = $2",
@@ -1589,7 +1590,14 @@ async fn delete_workspace_user(
let email_to_delete = not_found_if_none(email_to_delete_o, "User", &username_to_delete)?;
delete_workspace_user_internal(&w_id, &username_to_delete, &email_to_delete, &mut tx, Some(&authed)).await?;
delete_workspace_user_internal(
&w_id,
&username_to_delete,
&email_to_delete,
&mut tx,
Some(&authed),
)
.await?;
tx.commit().await?;
handle_deployment_metadata(
@@ -2172,7 +2180,9 @@ async fn get_all_runnables(
.collect::<Vec<_>>(),
);
let scripts = sqlx::query!(
"SELECT workspace_id as workspace, path, summary, description, schema FROM script as o WHERE created_at = (select max(created_at) from script where o.path = path and workspace_id = $1) and workspace_id = $1", workspace
"SELECT workspace_id as workspace, path, summary, description, schema FROM script as o
WHERE created_at = (select max(created_at) from script where o.path = path and workspace_id = $1 AND archived = false)
AND workspace_id = $1 and archived = false", workspace
)
.fetch_all(&mut *tx)
.await?;
+1
View File
@@ -169,6 +169,7 @@ impl<Key: Eq + Hash + Item + Clone, Val: Export, Root: AsRef<Path>> FsBackedCach
),
}
}
// Cache path doesn't exist or import failed, generate the content.
let data = Val::resolve(with.await?)?;
// Try to export data to the file-system.
+6 -4
View File
@@ -631,6 +631,7 @@ pub async fn get_latest_hash_for_path<'c, E: sqlx::PgExecutor<'c>>(
db: E,
w_id: &str,
script_path: &str,
require_locked: bool,
) -> error::Result<(
scripts::ScriptHash,
Option<Tag>,
@@ -646,11 +647,12 @@ pub async fn get_latest_hash_for_path<'c, E: sqlx::PgExecutor<'c>>(
String,
)> {
let r_o = sqlx::query!(
"select hash, tag, concurrency_key, concurrent_limit, concurrency_time_window_s, cache_ttl, language as \"language: ScriptLang\", dedicated_worker, priority, timeout, on_behalf_of_email, created_by FROM script where path = $1 AND workspace_id = $2 AND
created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND
deleted = false AND archived = false)",
"select hash, tag, concurrency_key, concurrent_limit, concurrency_time_window_s, cache_ttl, language as \"language: ScriptLang\", dedicated_worker, priority, timeout, on_behalf_of_email, created_by FROM script
WHERE path = $1 AND workspace_id = $2 AND archived = false AND (lock IS NOT NULL OR $3 = false)
ORDER BY created_at DESC LIMIT 1",
script_path,
w_id
w_id,
require_locked
)
.fetch_optional(db)
.await?;
+1
View File
@@ -156,6 +156,7 @@ pub async fn push_scheduled_job<'c>(
&mut *tx,
&schedule.workspace_id,
&schedule.script_path,
true,
)
.await?;
+15 -12
View File
@@ -118,7 +118,7 @@ struct Tool {
#[derive(Deserialize, Debug)]
struct AIAgentArgs {
provider: Provider,
system_prompt: String,
system_prompt: Option<String>,
user_message: String,
temperature: Option<f32>,
max_completion_tokens: Option<u32>,
@@ -608,18 +608,21 @@ async fn run_agent(
hostname: &str,
killpill_rx: &mut tokio::sync::broadcast::Receiver<()>,
) -> error::Result<Box<RawValue>> {
let mut messages = vec![
OpenAIMessage {
let mut messages = if let Some(system_prompt) = args.system_prompt.filter(|s| !s.is_empty()) {
vec![OpenAIMessage {
role: "system".to_string(),
content: Some(args.system_prompt),
content: Some(system_prompt),
..Default::default()
},
OpenAIMessage {
role: "user".to_string(),
content: Some(args.user_message),
..Default::default()
},
];
}]
} else {
vec![]
};
messages.push(OpenAIMessage {
role: "user".to_string(),
content: Some(args.user_message),
..Default::default()
});
let mut actions = vec![];
@@ -923,7 +926,7 @@ pub async fn handle_ai_agent_job(
.await?;
Ok(Some(hub_script.schema))
} else {
let hash = get_latest_hash_for_path(db, &job.workspace_id, path)
let hash = get_latest_hash_for_path(db, &job.workspace_id, path, true)
.await?
.0;
// update module definition to use a fixed hash so all tool calls match the same schema
+1 -1
View File
@@ -758,7 +758,7 @@ pub async fn prebundle_bun_script(
pub const BUN_BUNDLE_OBJECT_STORE_PREFIX: &str = "bun_bundle/";
async fn get_script_import_updated_at(db: &DB, w_id: &str, script_path: &str) -> Result<String> {
let script_hash = get_latest_hash_for_path(db, w_id, script_path).await?;
let script_hash = get_latest_hash_for_path(db, w_id, script_path, false).await?;
let last_updated_at = sqlx::query_scalar!(
"SELECT created_at FROM script WHERE workspace_id = $1 AND hash = $2",
w_id,
@@ -622,8 +622,7 @@ async fn spawn_dedicated_worker(
} else {
sqlx::query_as::<_, (String, Option<String>, Option<ScriptLang>, Option<Vec<String>>, bool, Option<ScriptHash>)>(
"SELECT content, lock, language, envs, codebase IS NOT NULL, hash FROM script WHERE path = $1 AND workspace_id = $2 AND
created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND
deleted = false AND lock IS not NULL AND lock_error_logs IS NULL)",
archived = false AND lock IS not NULL AND lock_error_logs IS NULL ORDER BY created_at DESC LIMIT 1",
)
.bind(&path)
.bind(&w_id)
+10 -1
View File
@@ -5,7 +5,9 @@ use regex::Regex;
use serde::Deserialize;
use serde_json::value::RawValue;
use serde_json::{Map, Value};
use tiberius::{AuthMethod, Client, ColumnData, Config, FromSqlOwned, Query, Row, SqlBrowser};
use tiberius::{
AuthMethod, Client, ColumnData, Config, EncryptionLevel, FromSqlOwned, Query, Row, SqlBrowser,
};
use tokio::net::TcpStream;
use tokio_util::compat::TokioAsyncWriteCompatExt;
use uuid::Uuid;
@@ -39,6 +41,7 @@ struct MssqlDatabase {
trust_cert: Option<bool>,
#[serde(default, deserialize_with = "empty_as_none")]
ca_cert: Option<String>,
encrypt: Option<bool>,
}
#[derive(Debug, Deserialize)]
@@ -146,6 +149,12 @@ pub async fn do_mssql(
tracing::info!("MSSQL: using provided CA certificate for trust");
}
config.encryption(if database.encrypt.unwrap_or(true) {
EncryptionLevel::Required
} else {
EncryptionLevel::NotSupported
});
let tcp = if use_instance_name {
TcpStream::connect_named(&config).await.map_err(to_anyhow)? // named instance
} else {
+1 -1
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.530.0";
export const VERSION = "v1.535.0";
export async function login(email: string, password: string): Promise<string> {
return await windmill.UserService.login({
+18 -20
View File
@@ -3,11 +3,11 @@ import { GlobalOptions } from "../../types.ts";
import { requireLogin } from "../../core/auth.ts";
import { resolveWorkspace } from "../../core/context.ts";
import * as wmill from "../../../gen/services.gen.ts";
import { SyncOptions, readConfigFile, getEffectiveSettings, DEFAULT_SYNC_OPTIONS } from "../../core/conf.ts";
import { SyncOptions, readConfigFile, getEffectiveSettings, DEFAULT_SYNC_OPTIONS, getWmillYamlPath } from "../../core/conf.ts";
import { deepEqual } from "../../utils/utils.ts";
import { getCurrentGitBranch, isGitRepository } from "../../utils/git.ts";
import { GitSyncRepository, WriteMode } from "./types.ts";
import { WriteMode } from "./types.ts";
import { GitSyncSettingsConverter } from "./converter.ts";
import { handleLegacyRepositoryMigration } from "./legacySettings.ts";
import {
@@ -132,11 +132,9 @@ export async function pullGitSyncSettings(
const backendSyncOptions: SyncOptions = GitSyncSettingsConverter.fromBackendFormat(selectedRepo.settings);
// Check if wmill.yaml exists - create a default one if it doesn't exist
let wmillYamlExists = true;
try {
await Deno.stat("wmill.yaml");
} catch (error) {
wmillYamlExists = false;
const wmillYamlPath = getWmillYamlPath();
const wmillYamlExists = wmillYamlPath !== null;
if (!wmillYamlExists) {
if (!opts.jsonOutput) {
log.info(
colors.yellow(
@@ -165,11 +163,11 @@ export async function pullGitSyncSettings(
if (isGitRepository()) {
const currentBranch = getCurrentGitBranch();
if (currentBranch) {
if (!updatedConfig.git_branches) {
updatedConfig.git_branches = {};
if (!updatedConfig.gitBranches) {
updatedConfig.gitBranches = {};
}
if (!updatedConfig.git_branches[currentBranch]) {
updatedConfig.git_branches[currentBranch] = { overrides: {} };
if (!updatedConfig.gitBranches[currentBranch]) {
updatedConfig.gitBranches[currentBranch] = { overrides: {} };
}
}
}
@@ -358,16 +356,16 @@ export async function pullGitSyncSettings(
let needsBranchStructure = false;
if (isGitRepository()) {
const currentBranch = getCurrentGitBranch();
if (currentBranch && (!localConfig.git_branches || !localConfig.git_branches[currentBranch])) {
if (currentBranch && (!localConfig.gitBranches || !localConfig.gitBranches[currentBranch])) {
needsBranchStructure = true;
// Create empty branch structure
const updatedConfig = { ...localConfig };
if (!updatedConfig.git_branches) {
updatedConfig.git_branches = {};
if (!updatedConfig.gitBranches) {
updatedConfig.gitBranches = {};
}
if (!updatedConfig.git_branches[currentBranch]) {
updatedConfig.git_branches[currentBranch] = { overrides: {} };
if (!updatedConfig.gitBranches[currentBranch]) {
updatedConfig.gitBranches[currentBranch] = { overrides: {} };
}
// Write updated configuration
@@ -429,11 +427,11 @@ export async function pullGitSyncSettings(
const currentBranch = getCurrentGitBranch();
if (currentBranch) {
log.info(`Detected Git repository, adding empty branch structure for: ${currentBranch}`);
if (!updatedConfig.git_branches) {
updatedConfig.git_branches = {};
if (!updatedConfig.gitBranches) {
updatedConfig.gitBranches = {};
}
if (!updatedConfig.git_branches[currentBranch]) {
updatedConfig.git_branches[currentBranch] = { overrides: {} };
if (!updatedConfig.gitBranches[currentBranch]) {
updatedConfig.gitBranches[currentBranch] = { overrides: {} };
}
}
}
+3 -4
View File
@@ -3,7 +3,7 @@ import { GlobalOptions } from "../../types.ts";
import { requireLogin } from "../../core/auth.ts";
import { resolveWorkspace } from "../../core/context.ts";
import * as wmill from "../../../gen/services.gen.ts";
import { SyncOptions, readConfigFile, validateBranchConfiguration, getEffectiveSettings } from "../../core/conf.ts";
import { SyncOptions, readConfigFile, validateBranchConfiguration, getEffectiveSettings, getWmillYamlPath } from "../../core/conf.ts";
import { deepEqual } from "../../utils/utils.ts";
import { GitSyncRepository } from "./types.ts";
@@ -44,9 +44,8 @@ export async function pushGitSyncSettings(
try {
// Check if wmill.yaml exists - require it for git-sync settings commands
try {
await Deno.stat("wmill.yaml");
} catch (error) {
const wmillYamlPath = getWmillYamlPath();
if (!wmillYamlPath) {
log.error(
colors.red(
"No wmill.yaml file found. Please run 'wmill init' first to create the configuration file.",
+9 -9
View File
@@ -30,12 +30,12 @@ export function getOrCreateBranchConfig(config: SyncOptions, branchName: string)
config: SyncOptions;
branchKey: string;
} {
if (!config.git_branches) {
config.git_branches = {};
if (!config.gitBranches) {
config.gitBranches = {};
}
if (!config.git_branches[branchName]) {
config.git_branches[branchName] = {};
if (!config.gitBranches[branchName]) {
config.gitBranches[branchName] = {};
}
return {
@@ -53,12 +53,12 @@ export function applyBackendSettingsToBranch(
const { config: updatedConfig } = getOrCreateBranchConfig(config, branchName);
// Get the base settings (top-level + defaults) to compare against
const { git_branches, ...topLevelSettings } = config;
const { gitBranches, ...topLevelSettings } = config;
const baseSettings: Partial<SyncOptions> = { ...DEFAULT_SYNC_OPTIONS, ...topLevelSettings };
// Only store fields that differ from the base settings
Object.keys(backendSettings).forEach(key => {
if (key !== 'git_branches' && backendSettings[key as keyof SyncOptions] !== undefined) {
if (key !== 'gitBranches' && backendSettings[key as keyof SyncOptions] !== undefined) {
const backendValue = backendSettings[key as keyof SyncOptions];
const baseValue = baseSettings[key as keyof SyncOptions];
@@ -66,10 +66,10 @@ export function applyBackendSettingsToBranch(
const isDifferent = GitSyncSettingsConverter.isDifferent(backendValue, baseValue);
if (isDifferent) {
if (!updatedConfig.git_branches![branchName].overrides) {
updatedConfig.git_branches![branchName].overrides = {};
if (!updatedConfig.gitBranches![branchName].overrides) {
updatedConfig.gitBranches![branchName].overrides = {};
}
(updatedConfig.git_branches![branchName].overrides as any)[key] = backendValue;
(updatedConfig.gitBranches![branchName].overrides as any)[key] = backendValue;
}
}
});
+9 -9
View File
@@ -43,14 +43,14 @@ async function initAction(opts: InitOptions) {
if (isGitRepository()) {
const currentBranch = getCurrentGitBranch();
if (currentBranch) {
initialConfig.git_branches = {
initialConfig.gitBranches = {
[currentBranch]: { overrides: {} },
};
} else {
initialConfig.git_branches = {};
initialConfig.gitBranches = {};
}
} else {
initialConfig.git_branches = {};
initialConfig.gitBranches = {};
}
await Deno.writeTextFile("wmill.yaml", yamlStringify(initialConfig));
@@ -116,16 +116,16 @@ async function initAction(opts: InitOptions) {
const currentConfig = await import("../../core/conf.ts").then((m) =>
m.readConfigFile()
);
if (!currentConfig.git_branches) {
currentConfig.git_branches = {};
if (!currentConfig.gitBranches) {
currentConfig.gitBranches = {};
}
if (!currentConfig.git_branches[currentBranch]) {
currentConfig.git_branches[currentBranch] = { overrides: {} };
if (!currentConfig.gitBranches[currentBranch]) {
currentConfig.gitBranches[currentBranch] = { overrides: {} };
}
currentConfig.git_branches[currentBranch].baseUrl =
currentConfig.gitBranches[currentBranch].baseUrl =
activeWorkspace.remote;
currentConfig.git_branches[currentBranch].workspaceId =
currentConfig.gitBranches[currentBranch].workspaceId =
activeWorkspace.workspaceId;
await Deno.writeTextFile(
+165 -24
View File
@@ -42,7 +42,19 @@ import {
readConfigFile,
getEffectiveSettings,
validateBranchConfiguration,
mergeConfigWithConfigFile,
} from "../../core/conf.ts";
import {
SpecificItemsConfig,
getSpecificItemsForCurrentBranch,
isSpecificItem,
getBranchSpecificPath,
fromBranchSpecificPath,
isCurrentBranchFile,
toBranchSpecificPath,
isBranchSpecificFile,
} from "../../core/specific_items.ts";
import { getCurrentGitBranch } from "../../utils/git.ts";
import { Workspace } from "../workspace/workspace.ts";
import { removePathPrefix } from "../../types.ts";
import { SyncCodebase, listSyncCodebases } from "../../utils/codebase.ts";
@@ -67,9 +79,9 @@ function mergeCliWithEffectiveOptions<
// Resolve effective sync options using branch-based configuration
async function resolveEffectiveSyncOptions(
workspace: Workspace,
localConfig: SyncOptions,
promotion?: string
): Promise<SyncOptions> {
const localConfig = await readConfigFile();
return await getEffectiveSettings(localConfig, promotion);
}
@@ -631,9 +643,36 @@ export async function elementsToMap(
els: DynFSElement,
ignore: (path: string, isDirectory: boolean) => boolean,
json: boolean,
skips: Skips
skips: Skips,
specificItems?: SpecificItemsConfig
): Promise<{ [key: string]: string }> {
const map: { [key: string]: string } = {};
const processedBasePaths = new Set<string>();
// First pass: collect all file paths to identify branch-specific files
const allPaths: string[] = [];
for await (const entry of readDirRecursiveWithIgnore(ignore, els)) {
if (!entry.isDirectory && !entry.ignored) {
allPaths.push(entry.path);
}
}
const branchSpecificExists = new Set<string>();
if (specificItems) {
const currentBranch = getCurrentGitBranch();
if (currentBranch) {
for (const path of allPaths) {
if (isCurrentBranchFile(path)) {
const basePath = fromBranchSpecificPath(path, currentBranch);
if (isSpecificItem(basePath, specificItems)) {
branchSpecificExists.add(basePath);
}
}
}
}
}
for await (const entry of readDirRecursiveWithIgnore(ignore, els)) {
if (entry.isDirectory || entry.ignored) continue;
const path = entry.path;
@@ -695,11 +734,21 @@ export async function elementsToMap(
"nu",
"java",
"rb",
// for related places search: ADD_NEW_LANG
// for related places search: ADD_NEW_LANG
].includes(path.split(".").pop() ?? "") &&
!isFileResource(path)
)
continue;
// Handle branch-specific files - skip files for other branches
if (specificItems && isBranchSpecificFile(path)) {
const currentBranch = getCurrentGitBranch();
if (!currentBranch || !isCurrentBranchFile(path)) {
// Skip branch-specific files for other branches
continue;
}
}
const content = await entry.getContentText();
if (skips.skipSecrets && path.endsWith(".variable" + ext)) {
@@ -727,7 +776,33 @@ export async function elementsToMap(
log.warn(`Error reading variable ${path} to check for secrets`);
}
}
map[entry.path] = content;
// Handle branch-specific path mapping after all filtering
if (specificItems) {
const currentBranch = getCurrentGitBranch();
if (currentBranch && isCurrentBranchFile(path)) {
// This is a branch-specific file for current branch
const basePath = fromBranchSpecificPath(path, currentBranch);
if (isSpecificItem(basePath, specificItems)) {
// Map to base path for push operations
map[basePath] = content;
processedBasePaths.add(basePath);
} else {
// Branch-specific file doesn't match pattern, skip it
continue;
}
} else if (!isBranchSpecificFile(path)) {
// This is a regular base file, check if we should skip it
if (processedBasePaths.has(path)) {
// Skip base file, we already processed branch-specific version
continue;
}
map[path] = content;
}
} else {
// No specific items configuration, use regular path
map[entry.path] = content;
}
}
return map;
}
@@ -758,14 +833,15 @@ async function compareDynFSElement(
skips: Skips,
ignoreMetadataDeletion: boolean,
codebases: SyncCodebase[],
ignoreCodebaseChanges: boolean
ignoreCodebaseChanges: boolean,
specificItems?: SpecificItemsConfig
): Promise<Change[]> {
const [m1, m2] = els2
? await Promise.all([
elementsToMap(els1, ignore, json, skips),
elementsToMap(els2, ignore, json, skips),
elementsToMap(els1, ignore, json, skips, specificItems),
elementsToMap(els2, ignore, json, skips, specificItems),
])
: [await elementsToMap(els1, ignore, json, skips), {}];
: [await elementsToMap(els1, ignore, json, skips, specificItems), {}];
const changes: Change[] = [];
@@ -1163,6 +1239,10 @@ export async function pull(
opts: GlobalOptions &
SyncOptions & { repository?: string; promotion?: string }
) {
const originalCliOpts = { ...opts };
opts = await mergeConfigWithConfigFile(opts);
// Validate branch configuration early
try {
await validateBranchConfiguration(false, opts.yes);
@@ -1184,11 +1264,15 @@ export async function pull(
// Resolve effective sync options with branch awareness
const effectiveOpts = await resolveEffectiveSyncOptions(
workspace,
opts,
opts.promotion
);
// Extract specific items configuration before merging overwrites gitBranches
const specificItems = getSpecificItemsForCurrentBranch(opts);
// Merge CLI flags with resolved settings (CLI flags take precedence only for explicit overrides)
opts = mergeCliWithEffectiveOptions(opts, effectiveOpts);
opts = mergeCliWithEffectiveOptions(originalCliOpts, effectiveOpts);
const codebases = await listSyncCodebases(opts);
@@ -1238,7 +1322,8 @@ export async function pull(
opts,
false,
codebases,
true
true,
specificItems
);
log.info(
@@ -1255,6 +1340,12 @@ export async function pull(
...(change.name === "edited" && change.codebase
? { codebase_changed: true }
: {}),
...(specificItems && isSpecificItem(change.path, specificItems)
? {
branch_specific: true,
branch_specific_path: getBranchSpecificPath(change.path, specificItems)
}
: {}),
})),
total: changes.length,
};
@@ -1264,7 +1355,7 @@ export async function pull(
if (changes.length > 0) {
if (!opts.jsonOutput) {
prettyChanges(changes);
prettyChanges(changes, specificItems);
}
if (opts.dryRun) {
log.info(colors.gray(`Dry run complete.`));
@@ -1284,8 +1375,17 @@ export async function pull(
log.info(colors.gray(`Applying changes to files ...`));
for await (const change of changes) {
const target = path.join(Deno.cwd(), change.path);
const stateTarget = path.join(Deno.cwd(), ".wmill", change.path);
// Determine if this file should be written to a branch-specific path
let targetPath = change.path;
if (specificItems && isSpecificItem(change.path, specificItems)) {
const branchSpecificPath = getBranchSpecificPath(change.path, specificItems);
if (branchSpecificPath) {
targetPath = branchSpecificPath;
}
}
const target = path.join(Deno.cwd(), targetPath);
const stateTarget = path.join(Deno.cwd(), ".wmill", targetPath);
if (change.name === "edited") {
if (opts.stateful) {
try {
@@ -1328,12 +1428,12 @@ export async function pull(
}
}
if (exts.some((e) => change.path.endsWith(e))) {
log.info(`Editing script content of ${change.path}`);
log.info(`Editing script content of ${targetPath}${targetPath !== change.path ? colors.gray(` (branch-specific override for ${change.path})`) : ""}`);
} else if (
change.path.endsWith(".yaml") ||
change.path.endsWith(".json")
) {
log.info(`Editing ${getTypeStrFromPath(change.path)} ${change.path}`);
log.info(`Editing ${getTypeStrFromPath(change.path)} ${targetPath}${targetPath !== change.path ? colors.gray(` (branch-specific override for ${change.path})`) : ""}`);
}
await Deno.writeTextFile(target, change.after);
@@ -1345,10 +1445,10 @@ export async function pull(
await ensureDir(path.dirname(target));
if (opts.stateful) {
await ensureDir(path.dirname(stateTarget));
log.info(`Adding ${getTypeStrFromPath(change.path)} ${change.path}`);
log.info(`Adding ${getTypeStrFromPath(change.path)} ${targetPath}${targetPath !== change.path ? colors.gray(` (branch-specific override for ${change.path})`) : ""}`);
}
await Deno.writeTextFile(target, change.content);
log.info(`Writing ${getTypeStrFromPath(change.path)} ${change.path}`);
log.info(`Writing ${getTypeStrFromPath(change.path)} ${targetPath}${targetPath !== change.path ? colors.gray(` (branch-specific override for ${change.path})`) : ""}`);
if (opts.stateful) {
await Deno.copyFile(target, stateTarget);
}
@@ -1423,6 +1523,12 @@ export async function pull(
...(change.name === "edited" && change.codebase
? { codebase_changed: true }
: {}),
...(specificItems && isSpecificItem(change.path, specificItems)
? {
branch_specific: true,
branch_specific_path: getBranchSpecificPath(change.path, specificItems)
}
: {}),
})),
total: changes.length,
};
@@ -1445,21 +1551,33 @@ export async function pull(
}
}
function prettyChanges(changes: Change[]) {
function prettyChanges(changes: Change[], specificItems?: SpecificItemsConfig) {
for (const change of changes) {
let displayPath = change.path;
let branchNote = "";
// Check if this will be written as a branch-specific file
if (specificItems && isSpecificItem(change.path, specificItems)) {
const branchSpecificPath = getBranchSpecificPath(change.path, specificItems);
if (branchSpecificPath) {
displayPath = branchSpecificPath;
branchNote = " (branch-specific)";
}
}
if (change.name === "added") {
log.info(
colors.green(`+ ${getTypeStrFromPath(change.path)} ` + change.path)
colors.green(`+ ${getTypeStrFromPath(change.path)} ` + displayPath + colors.gray(branchNote))
);
} else if (change.name === "deleted") {
log.info(
colors.red(`- ${getTypeStrFromPath(change.path)} ` + change.path)
colors.red(`- ${getTypeStrFromPath(change.path)} ` + displayPath + colors.gray(branchNote))
);
} else if (change.name === "edited") {
log.info(
colors.yellow(
`~ ${getTypeStrFromPath(change.path)} ` +
change.path +
displayPath + colors.gray(branchNote) +
(change.codebase ? ` (codebase changed)` : "")
)
);
@@ -1499,6 +1617,12 @@ function removeSuffix(str: string, suffix: string) {
export async function push(
opts: GlobalOptions & SyncOptions & { repository?: string }
) {
// Save original CLI options before merging with config file
const originalCliOpts = { ...opts };
// Load configuration from wmill.yaml and merge with CLI options
opts = await mergeConfigWithConfigFile(opts);
// Validate branch configuration early
try {
await validateBranchConfiguration(false, opts.yes);
@@ -1516,11 +1640,15 @@ export async function push(
// Resolve effective sync options with branch awareness
const effectiveOpts = await resolveEffectiveSyncOptions(
workspace,
opts,
opts.promotion
);
// Extract specific items configuration BEFORE merging overwrites gitBranches
const specificItems = getSpecificItemsForCurrentBranch(opts);
// Merge CLI flags with resolved settings (CLI flags take precedence only for explicit overrides)
opts = mergeCliWithEffectiveOptions(opts, effectiveOpts);
opts = mergeCliWithEffectiveOptions(originalCliOpts, effectiveOpts);
const codebases = await listSyncCodebases(opts);
if (opts.raw) {
@@ -1581,7 +1709,8 @@ export async function push(
opts,
true,
codebases,
false
false,
specificItems
);
const globalDeps = await findGlobalDeps();
@@ -1660,6 +1789,12 @@ export async function push(
...(change.name === "edited" && change.codebase
? { codebase_changed: true }
: {}),
...(specificItems && isSpecificItem(change.path, specificItems)
? {
branch_specific: true,
branch_specific_path: getBranchSpecificPath(change.path, specificItems)
}
: {}),
})),
total: changes.length,
};
@@ -1669,7 +1804,7 @@ export async function push(
if (changes.length > 0) {
if (!opts.jsonOutput) {
prettyChanges(changes);
prettyChanges(changes, specificItems);
}
if (opts.dryRun) {
log.info(colors.gray(`Dry run complete.`));
@@ -2041,6 +2176,12 @@ export async function push(
...(change.name === "edited" && change.codebase
? { codebase_changed: true }
: {}),
...(specificItems && isSpecificItem(change.path, specificItems)
? {
branch_specific: true,
branch_specific_path: getBranchSpecificPath(change.path, specificItems)
}
: {}),
})),
total: changes.length,
duration_ms: Math.round(performance.now() - start),
+10 -10
View File
@@ -386,22 +386,22 @@ async function bind(
}
// For unbind, check if branch exists
if (!bindWorkspace && (!config.git_branches || !config.git_branches[branch])) {
log.error(colors.red(`Branch '${branch}' not found in wmill.yaml git_branches`));
if (!bindWorkspace && (!config.gitBranches || !config.gitBranches[branch])) {
log.error(colors.red(`Branch '${branch}' not found in wmill.yaml gitBranches`));
return;
}
// Update the branch configuration with workspace binding
if (!config.git_branches) {
config.git_branches = {};
if (!config.gitBranches) {
config.gitBranches = {};
}
if (!config.git_branches[branch]) {
config.git_branches[branch] = { overrides: {} };
if (!config.gitBranches[branch]) {
config.gitBranches[branch] = { overrides: {} };
}
if (bindWorkspace && activeWorkspace) {
config.git_branches[branch].baseUrl = activeWorkspace.remote;
config.git_branches[branch].workspaceId = activeWorkspace.workspaceId;
config.gitBranches[branch].baseUrl = activeWorkspace.remote;
config.gitBranches[branch].workspaceId = activeWorkspace.workspaceId;
log.info(colors.green(
`✓ Bound branch '${branch}' to workspace '${activeWorkspace.name}'\n` +
@@ -409,8 +409,8 @@ async function bind(
));
} else {
// Unbind
delete config.git_branches[branch].baseUrl;
delete config.git_branches[branch].workspaceId;
delete config.gitBranches[branch].baseUrl;
delete config.gitBranches[branch].workspaceId;
log.info(colors.green(`✓ Removed workspace binding from branch '${branch}'`));
}
+210 -35
View File
@@ -1,5 +1,8 @@
import { log, yamlParseFile, Confirm, yamlStringify } from "../../deps.ts";
import { getCurrentGitBranch, isGitRepository } from "../utils/git.ts";
import { join, dirname, resolve, relative } from "node:path";
import { existsSync } from "node:fs";
import { execSync } from "node:child_process";
export let showDiffs = false;
export function setShowDiffs(value: boolean) {
@@ -37,13 +40,40 @@ export interface SyncOptions {
codebases?: Codebase[];
parallel?: number;
jsonOutput?: boolean;
git_branches?: {
gitBranches?: {
commonSpecificItems?: {
variables?: string[];
resources?: string[];
};
} & {
[branchName: string]: SyncOptions & {
overrides?: Partial<SyncOptions>;
promotionOverrides?: Partial<SyncOptions>;
baseUrl?: string;
workspaceId?: string;
}
specificItems?: {
variables?: string[];
resources?: string[];
};
};
};
// Legacy field - deprecated, use gitBranches instead
git_branches?: {
commonSpecificItems?: {
variables?: string[];
resources?: string[];
};
} & {
[branchName: string]: SyncOptions & {
overrides?: Partial<SyncOptions>;
promotionOverrides?: Partial<SyncOptions>;
baseUrl?: string;
workspaceId?: string;
specificItems?: {
variables?: string[];
resources?: string[];
};
};
};
promotion?: string;
}
@@ -62,9 +92,90 @@ export interface Codebase {
inject?: string[];
}
function getGitRepoRoot(): string | null {
try {
const result = execSync("git rev-parse --show-toplevel", {
encoding: "utf8",
stdio: "pipe"
});
return result.trim();
} catch (error) {
return null;
}
}
function findWmillYaml(): string | null {
const startDir = resolve(Deno.cwd());
const isInGitRepo = isGitRepository();
// If not in git repo, only check current directory
if (!isInGitRepo) {
const wmillYamlPath = join(startDir, "wmill.yaml");
return existsSync(wmillYamlPath) ? wmillYamlPath : null;
}
// If in git repo, search up to git repository root
const gitRoot = getGitRepoRoot();
let currentDir = startDir;
let foundPath: string | null = null;
while (true) {
const wmillYamlPath = join(currentDir, "wmill.yaml");
if (existsSync(wmillYamlPath)) {
foundPath = wmillYamlPath;
break;
}
// Check if we've reached the git repository root
if (gitRoot && resolve(currentDir) === resolve(gitRoot)) {
break;
}
// Check if we've reached the filesystem root
const parentDir = dirname(currentDir);
if (parentDir === currentDir) {
break;
}
currentDir = parentDir;
}
// If wmill.yaml was found in a parent directory, warn the user and change working directory
if (foundPath && resolve(dirname(foundPath)) !== resolve(startDir)) {
const configDir = dirname(foundPath);
const relativePath = relative(startDir, foundPath);
log.warn(`⚠️ wmill.yaml found in parent directory: ${relativePath}`);
// Change working directory to where wmill.yaml was found
Deno.chdir(configDir);
log.info(`📁 Changed working directory to: ${configDir}`);
}
return foundPath;
}
export function getWmillYamlPath(): string | null {
return findWmillYaml();
}
export async function readConfigFile(): Promise<SyncOptions> {
try {
const conf = (await yamlParseFile("wmill.yaml")) as SyncOptions;
// First, try to find wmill.yaml recursively
const wmillYamlPath = findWmillYaml();
if (!wmillYamlPath) {
log.warn(
"No wmill.yaml found. Use 'wmill init' to bootstrap it. Using 'bun' as default typescript runtime."
);
return {};
}
const conf = (await yamlParseFile(wmillYamlPath)) as SyncOptions;
// Handle legacy format migrations (combine overrides and git_branches)
let needsConfigWrite = false;
const migrationMessages: string[] = [];
// Handle obsolete overrides format
if (conf && 'overrides' in conf) {
@@ -78,18 +189,54 @@ export async function readConfigFile(): Promise<SyncOptions> {
" Please delete your wmill.yaml and run 'wmill init' to recreate it with the new format."
);
} else {
// Remove empty overrides with a note
log.info("️ Removing empty 'overrides: {}' from wmill.yaml (migrated to git_branches format)");
// Remove empty overrides
delete conf.overrides;
// Write the updated config back to file
try {
await Deno.writeTextFile("wmill.yaml", yamlStringify(conf));
} catch (error) {
log.warn(`Could not update wmill.yaml to remove empty overrides: ${error instanceof Error ? error.message : error}`);
}
needsConfigWrite = true;
migrationMessages.push("️ Removing empty 'overrides: {}' from wmill.yaml (migrated to gitBranches format)");
}
}
// Handle git_branches to gitBranches migration
if (conf && 'git_branches' in conf) {
if (!conf.gitBranches) {
// Deep copy git_branches to gitBranches (even if empty)
conf.gitBranches = JSON.parse(JSON.stringify(conf.git_branches));
needsConfigWrite = true;
migrationMessages.push("⚠️ Migrating 'git_branches' to 'gitBranches' (camelCase). The snake_case format is deprecated.");
migrationMessages.push("✅ Successfully migrated 'git_branches' to 'gitBranches' in wmill.yaml");
} else {
migrationMessages.push("⚠️ Both 'git_branches' and 'gitBranches' found in wmill.yaml. Using 'gitBranches' and ignoring 'git_branches'.");
}
// Always remove the old field from config object (both file and memory)
delete conf.git_branches;
}
// Perform single atomic write if any migrations are needed
if (needsConfigWrite) {
try {
await Deno.writeTextFile(wmillYamlPath, yamlStringify(conf));
// Log all migration messages after successful write
migrationMessages.forEach(msg => {
if (msg.startsWith('⚠️')) {
log.warn(msg);
} else {
log.info(msg);
}
});
} catch (error) {
log.warn(`Could not update wmill.yaml to apply migrations: ${error instanceof Error ? error.message : error}`);
}
} else if (migrationMessages.length > 0) {
// Log messages for non-write cases (like "both found")
migrationMessages.forEach(msg => {
if (msg.startsWith('⚠️')) {
log.warn(msg);
} else {
log.info(msg);
}
});
}
if (conf?.defaultTs == undefined) {
log.warn(
"No defaultTs defined in your wmill.yaml. Using 'bun' as default."
@@ -100,10 +247,23 @@ export async function readConfigFile(): Promise<SyncOptions> {
if (e instanceof Error && (e.message.includes("overrides") || e.message.includes("Obsolete configuration format"))) {
throw e; // Re-throw the specific obsolete format error
}
log.warn(
"No wmill.yaml found. Use 'wmill init' to bootstrap it. Using 'bun' as default typescript runtime."
);
return {};
// Since we already found the file path, this is likely a parsing or access error
if (e instanceof Error && e.message.includes("Error parsing yaml")) {
const yamlError = e.cause instanceof Error ? e.cause.message : String(e.cause);
throw new Error(
"❌ YAML syntax error in wmill.yaml:\n" +
" " + yamlError + "\n" +
" Please fix the YAML syntax in wmill.yaml or delete the file to start fresh."
);
} else {
// File exists but has other issues (permissions, etc.)
throw new Error(
"❌ Failed to read wmill.yaml:\n" +
" " + (e instanceof Error ? e.message : String(e)) + "\n" +
" Please check file permissions or fix the syntax."
);
}
}
}
@@ -148,26 +308,26 @@ export async function validateBranchConfiguration(skipValidation?: boolean, auto
}
const config = await readConfigFile();
const { git_branches } = config;
const { gitBranches } = config;
const currentBranch = getCurrentGitBranch();
// In a git repository, git_branches section is recommended
if (!git_branches || Object.keys(git_branches).length === 0) {
// In a git repository, gitBranches section is recommended
if (!gitBranches || Object.keys(gitBranches).length === 0) {
log.warn(
"⚠️ WARNING: In a Git repository, the 'git_branches' section is recommended in wmill.yaml.\n" +
" Consider adding a git_branches section with configuration for your Git branches.\n" +
"⚠️ WARNING: In a Git repository, the 'gitBranches' section is recommended in wmill.yaml.\n" +
" Consider adding a gitBranches section with configuration for your Git branches.\n" +
" Run 'wmill init' to recreate the configuration file with proper branch setup."
);
return;
}
// Current branch must be defined in git_branches config
if (currentBranch && !git_branches[currentBranch]) {
// Current branch must be defined in gitBranches config
if (currentBranch && !gitBranches[currentBranch]) {
// In interactive mode, offer to create the branch
if (Deno.stdin.isTerminal()) {
const availableBranches = Object.keys(git_branches).join(', ');
const availableBranches = Object.keys(gitBranches).join(', ');
log.info(
`Current Git branch '${currentBranch}' is not defined in the git_branches configuration.\n` +
`Current Git branch '${currentBranch}' is not defined in the gitBranches configuration.\n` +
`Available branches: ${availableBranches}`
);
@@ -177,13 +337,21 @@ export async function validateBranchConfiguration(skipValidation?: boolean, auto
});
if (shouldCreate) {
// Warn if branch name contains filesystem-unsafe characters
if (/[\/\\:*?"<>|.]/.test(currentBranch)) {
const sanitizedBranchName = currentBranch.replace(/[\/\\:*?"<>|.]/g, '_');
log.warn(`⚠️ WARNING: Branch name "${currentBranch}" contains filesystem-unsafe characters (/ \\ : * ? " < > | .).`);
log.warn(` Branch-specific files will be saved with sanitized name: "${sanitizedBranchName}"`);
log.warn(` Example: "file.variable.yaml" → "file.${sanitizedBranchName}.variable.yaml"`);
}
// Read current config, add branch, and write it back
const currentConfig = await readConfigFile();
if (!currentConfig.git_branches) {
currentConfig.git_branches = {};
if (!currentConfig.gitBranches) {
currentConfig.gitBranches = {};
}
currentConfig.git_branches[currentBranch] = { overrides: {} };
currentConfig.gitBranches[currentBranch] = { overrides: {} };
await Deno.writeTextFile("wmill.yaml", yamlStringify(currentConfig));
@@ -193,10 +361,17 @@ export async function validateBranchConfiguration(skipValidation?: boolean, auto
return;
}
} else {
// Warn about filesystem-unsafe characters in branch name
if (/[\/\\:*?"<>|.]/.test(currentBranch)) {
const sanitizedBranchName = currentBranch.replace(/[\/\\:*?"<>|.]/g, '_');
log.warn(`⚠️ WARNING: Branch name "${currentBranch}" contains filesystem-unsafe characters (/ \\ : * ? " < > | .).`);
log.warn(` Branch-specific files will use sanitized name: "${sanitizedBranchName}"`);
}
log.warn(
`⚠️ WARNING: Current Git branch '${currentBranch}' is not defined in the git_branches configuration.\n` +
` Consider adding configuration for branch '${currentBranch}' in the git_branches section of wmill.yaml.\n` +
` Available branches: ${Object.keys(git_branches).join(', ')}`
`⚠️ WARNING: Current Git branch '${currentBranch}' is not defined in the gitBranches configuration.\n` +
` Consider adding configuration for branch '${currentBranch}' in the gitBranches section of wmill.yaml.\n` +
` Available branches: ${Object.keys(gitBranches).join(', ')}`
);
return;
}
@@ -206,15 +381,15 @@ export async function validateBranchConfiguration(skipValidation?: boolean, auto
// Get effective settings by merging top-level settings with branch-specific overrides
export async function getEffectiveSettings(config: SyncOptions, promotion?: string, skipBranchValidation?: boolean, suppressLogs?: boolean): Promise<SyncOptions> {
// Start with top-level settings from config
const { git_branches, ...topLevelSettings } = config;
const { gitBranches, ...topLevelSettings } = config;
let effective = { ...topLevelSettings };
if (isGitRepository()) {
const currentBranch = getCurrentGitBranch();
// If promotion is specified, use that branch's promotionOverrides or overrides
if (promotion && git_branches && git_branches[promotion]) {
const targetBranch = git_branches[promotion];
if (promotion && gitBranches && gitBranches[promotion]) {
const targetBranch = gitBranches[promotion];
// First try promotionOverrides, then fall back to overrides
if (targetBranch.promotionOverrides) {
@@ -232,8 +407,8 @@ export async function getEffectiveSettings(config: SyncOptions, promotion?: stri
}
}
// Otherwise use current branch overrides (existing behavior)
else if (currentBranch && git_branches && git_branches[currentBranch] && git_branches[currentBranch].overrides) {
Object.assign(effective, git_branches[currentBranch].overrides);
else if (currentBranch && gitBranches && gitBranches[currentBranch] && gitBranches[currentBranch].overrides) {
Object.assign(effective, gitBranches[currentBranch].overrides);
if (!suppressLogs) {
log.info(`Applied settings for Git branch: ${currentBranch}`);
}
+1 -1
View File
@@ -135,7 +135,7 @@ async function tryResolveBranchWorkspace(
// Read wmill.yaml to check for branch workspace configuration
const config = await readConfigFile();
const branchConfig = config.git_branches?.[currentBranch];
const branchConfig = config.gitBranches?.[currentBranch];
// Check if branch has workspace configuration
if (!branchConfig?.baseUrl || !branchConfig?.workspaceId) {
+186
View File
@@ -0,0 +1,186 @@
import { minimatch } from "../../deps.ts";
import { getCurrentGitBranch, isGitRepository } from "../utils/git.ts";
import { SyncOptions } from "./conf.ts";
export interface SpecificItemsConfig {
variables?: string[];
resources?: string[];
}
/**
* Get the specific items configuration for the current git branch
* Merges commonSpecificItems with branch-specific specificItems
*/
export function getSpecificItemsForCurrentBranch(config: SyncOptions): SpecificItemsConfig | undefined {
if (!isGitRepository() || !config.gitBranches) {
return undefined;
}
const currentBranch = getCurrentGitBranch();
if (!currentBranch) {
return undefined;
}
const commonItems = config.gitBranches.commonSpecificItems;
const branchItems = config.gitBranches[currentBranch]?.specificItems;
// If neither common nor branch-specific items exist, return undefined
if (!commonItems && !branchItems) {
return undefined;
}
// Merge common and branch-specific items
const merged: SpecificItemsConfig = {};
// Add common items
if (commonItems?.variables) {
merged.variables = [...commonItems.variables];
}
if (commonItems?.resources) {
merged.resources = [...commonItems.resources];
}
// Add branch-specific items (extending common items)
if (branchItems?.variables) {
merged.variables = [...(merged.variables || []), ...branchItems.variables];
}
if (branchItems?.resources) {
merged.resources = [...(merged.resources || []), ...branchItems.resources];
}
return merged;
}
/**
* Check if a path matches any of the patterns in the given list
*/
function matchesPatterns(path: string, patterns: string[]): boolean {
return patterns.some(pattern => minimatch(path, pattern));
}
/**
* Check if a file path should be treated as branch-specific
*/
export function isSpecificItem(path: string, specificItems: SpecificItemsConfig | undefined): boolean {
if (!specificItems) {
return false;
}
// Determine the item type from the file path
if (path.endsWith('.variable.yaml')) {
return specificItems.variables ? matchesPatterns(path, specificItems.variables) : false;
}
if (path.endsWith('.resource.yaml')) {
return specificItems.resources ? matchesPatterns(path, specificItems.resources) : false;
}
return false;
}
/**
* Convert a base path to a branch-specific path
*/
export function toBranchSpecificPath(basePath: string, branchName: string): string {
// Extract the extension (e.g., ".variable.yaml" or ".resource.yaml")
const extensionMatch = basePath.match(/(\.(variable|resource)\.yaml)$/);
if (!extensionMatch) {
return basePath; // Return unchanged if no recognized extension
}
const extension = extensionMatch[1];
const pathWithoutExtension = basePath.substring(0, basePath.length - extension.length);
// Sanitize branch name to be filesystem-safe
const sanitizedBranchName = branchName.replace(/[\/\\:*?"<>|.]/g, '_');
// Warn about potential collisions if sanitization occurred
if (sanitizedBranchName !== branchName) {
console.warn(`Warning: Branch name "${branchName}" contains filesystem-unsafe characters (/ \\ : * ? " < > | .) and was sanitized to "${sanitizedBranchName}". This may cause collisions with other similarly named branches.`);
}
return `${pathWithoutExtension}.${sanitizedBranchName}${extension}`;
}
/**
* Convert a branch-specific path back to a base path
*/
export function fromBranchSpecificPath(branchSpecificPath: string, branchName: string): string {
// Sanitize branch name the same way as in toBranchSpecificPath
const sanitizedBranchName = branchName.replace(/[\/\\:*?"<>|.]/g, '_');
// Pattern: path.sanitizedBranchName.extension
const escapedBranchName = sanitizedBranchName.replace(/[.*+?^${}()|[\]\\]/g, '\\$&');
const pattern = new RegExp(`\\.${escapedBranchName}(\\.(variable|resource)\\.yaml)$`);
const match = branchSpecificPath.match(pattern);
if (!match) {
return branchSpecificPath; // Return unchanged if not a branch-specific path
}
const extension = match[1];
const pathWithoutBranchAndExtension = branchSpecificPath.substring(
0,
branchSpecificPath.length - `.${sanitizedBranchName}${extension}`.length
);
return `${pathWithoutBranchAndExtension}${extension}`;
}
/**
* Get the branch-specific path for the current branch if the item should be branch-specific
*/
export function getBranchSpecificPath(
basePath: string,
specificItems: SpecificItemsConfig | undefined
): string | undefined {
if (!isGitRepository() || !specificItems) {
return undefined;
}
const currentBranch = getCurrentGitBranch();
if (!currentBranch) {
return undefined;
}
if (isSpecificItem(basePath, specificItems)) {
return toBranchSpecificPath(basePath, currentBranch);
}
return undefined;
}
// Cache for compiled regex patterns to avoid recompilation
const branchPatternCache = new Map<string, RegExp>();
/**
* Check if a path is a branch-specific file for the current branch
*/
export function isCurrentBranchFile(path: string): boolean {
if (!isGitRepository()) {
return false;
}
const currentBranch = getCurrentGitBranch();
if (!currentBranch) {
return false;
}
// Use cached pattern or create and cache new one
let pattern = branchPatternCache.get(currentBranch);
if (!pattern) {
pattern = new RegExp(`\\.${currentBranch}\\.(variable|resource)\\.yaml$`);
branchPatternCache.set(currentBranch, pattern);
}
return pattern.test(path);
}
/**
* Check if a path is a branch-specific file for ANY branch (not necessarily current)
* Used to identify and skip files from other branches during sync operations
*/
export function isBranchSpecificFile(path: string): boolean {
// Pattern: *.branchName.variable.yaml or *.branchName.resource.yaml
return /\.[^.]+\.(variable|resource)\.yaml$/.test(path);
}
+1 -1
View File
@@ -68,7 +68,7 @@ export {
// }
// });
export const VERSION = "1.530.0";
export const VERSION = "1.535.0";
const command = new Command()
.name("wmill")
+14 -11
View File
@@ -15,44 +15,45 @@ That's it! You are ready to go.
> Using **direnv** is highly recommended, since it can load shell automatically based on your CWD. It also can give you hints.
### Development
```bash
# enter a dev shell containing all necessary packages. `direnv allow` if direnv is installed.
nix develop
nix develop
## or ignore if you have `direnv`
# Start db (if not started already)
sudo docker compose up db -d
./start-dev-db.sh
# run the frontend.
wm
# In an other shell:
#
nix develop
nix develop
## or ignore if you have `direnv`
cd backend
# You don't need to install anything extra. All dependencies are already in place!
cargo run --features all_languages
cargo run --features all_languages
```
The default proxy is setup to use the local backend: <http://localhost:8000>.
### wm-* Commands
### wm-\* Commands
Nix shell provides you with several helper commands prefixed with `wm-`
```bash
# Start minio server (implements S3)
wm-minio
# Note: You will need access to EE private repo in order to compile, don't forget "enterprise" and "parquet" freatures as well.
wm-minio
# Note: You will need access to EE private repo in order to compile, don't forget "enterprise" and "parquet" freatures as well.
# Generate keys for local dev.
wm-minio-keys
wm-minio-keys
# Minio data as well as generated keys are stored in `backend/.minio-data`
```
You can read about all others commands individually in [flake.nix](../flake.nix).
You can read about all others commands individually in [flake.nix](../flake.nix).
### dev.nu
@@ -106,7 +107,7 @@ REMOTE=http://localhost REMOTE_LSP=http://localhost npm run dev
Sometimes it is important to build docker image for your branch locally. It is crucial part of testing, since local environment may differ from the containerized one.
That's why we provide [docker/dev.nu](../docker/dev.nu). It is helper that can build images locally and execute them.
That's why we provide [docker/dev.nu](../docker/dev.nu). It is helper that can build images locally and execute them.
it can build the image and run on local repository.
@@ -150,7 +151,7 @@ If you develop wasm parser for new language you can also pass `--wasm-pkg <langu
In the root folder:
```bash
docker-compose up db
./start-dev-db.sh
```
In the backend folder:
@@ -218,6 +219,7 @@ nix flake update # update the lock file.
```
Some cargo dependencies use fixed git revisions, which are also fixed in `flake.nix`:
```nix
outputHashes = {
"php-parser-rs-0.1.3" = "sha256-ZeI3KgUPmtjlRfq6eAYveqt8Ay35gwj6B9iOQRjQa9A=";
@@ -226,6 +228,7 @@ outputHashes = {
```
When updating a revision, replace the incorrect `sha256` with `pkgs.lib.fakeHash`:
```diff
- "php-parser-rs-0.1.3" = "sha256-ZeI3KgUPmtjlRfq6eAYveqt8Ay35gwj6B9iOQRjQa9A=";
+ "php-parser-rs-0.1.3" = pkgs.lib.fakeHash;
+999 -1111
View File
File diff suppressed because it is too large Load Diff
+13 -13
View File
@@ -1,6 +1,6 @@
{
"name": "windmill-components",
"version": "1.530.0",
"version": "1.535.0",
"scripts": {
"dev": "vite dev",
"build": "vite build",
@@ -78,13 +78,13 @@
"type": "module",
"dependencies": {
"@aws-crypto/sha256-js": "^4.0.0",
"@codingame/monaco-vscode-configuration-service-override": "~19.1.4",
"@codingame/monaco-vscode-editor-api": "~19.1.4",
"@codingame/monaco-vscode-standalone-css-language-features": "~19.1.4",
"@codingame/monaco-vscode-standalone-html-language-features": "~19.1.4",
"@codingame/monaco-vscode-standalone-json-language-features": "~19.1.4",
"@codingame/monaco-vscode-standalone-languages": "~19.1.4",
"@codingame/monaco-vscode-standalone-typescript-language-features": "~19.1.4",
"@codingame/monaco-vscode-configuration-service-override": "~20.2.1",
"@codingame/monaco-vscode-editor-api": "~20.2.1",
"@codingame/monaco-vscode-standalone-css-language-features": "~20.2.1",
"@codingame/monaco-vscode-standalone-html-language-features": "~20.2.1",
"@codingame/monaco-vscode-standalone-json-language-features": "~20.2.1",
"@codingame/monaco-vscode-standalone-languages": "~20.2.1",
"@codingame/monaco-vscode-standalone-typescript-language-features": "~20.2.1",
"@json2csv/plainjs": "^7.0.6",
"@leeoniya/ufuzzy": "^1.0.8",
"@popperjs/core": "^2.11.6",
@@ -117,10 +117,10 @@
"lru-cache": "^11.1.0",
"lucide-svelte": "^0.399.0",
"minimatch": "^10.0.1",
"monaco-editor": "npm:@codingame/monaco-vscode-editor-api@~19.1.4",
"monaco-editor-wrapper": "6.10.0",
"monaco-editor": "npm:@codingame/monaco-vscode-editor-api@~20.2.1",
"monaco-editor-wrapper": "6.12.0",
"monaco-graphql": "=1.6.0",
"monaco-languageclient": "9.9.0",
"monaco-languageclient": "9.11.0",
"monaco-vim": "^0.4.1",
"ol": "^7.4.0",
"openai": "^4.87.1",
@@ -137,10 +137,10 @@
"svelte-infinite-loading": "^1.4.0",
"svelte-tiny-virtual-list": "^2.0.5",
"tailwind-merge": "^1.13.2",
"vscode": "npm:@codingame/monaco-vscode-extension-api@~19.1.4",
"vscode": "npm:@codingame/monaco-vscode-extension-api@~20.2.1",
"vscode-languageclient": "~9.0.1",
"vscode-uri": "~3.1.0",
"vscode-ws-jsonrpc": "~3.4.0",
"vscode-ws-jsonrpc": "~3.5.0",
"windmill-parser-wasm-csharp": "1.510.1",
"windmill-parser-wasm-go": "1.510.1",
"windmill-parser-wasm-java": "1.510.1",
@@ -6,6 +6,7 @@
import {
setInputCat as computeInputCat,
debounce,
emptySchema,
emptyString,
getSchemaFromProperties,
type DynamicSelect
@@ -42,6 +43,8 @@
import { safeSelectItems } from './select/utils.svelte'
import S3ArgInput from './common/fileUpload/S3ArgInput.svelte'
import { base } from '$lib/base'
import { workspaceStore } from '$lib/stores'
import { getJsonSchemaFromResource } from './schema/jsonSchemaResource.svelte'
interface Props {
label?: string
@@ -658,6 +661,47 @@
{appPath}
{computeS3ForceViewerPolicies}
/>
{:else if inputCat == 'object' && format == 'json-schema'}
{#await import('$lib/components/EditableSchemaForm.svelte')}
<Loader2 class="animate-spin" />
{:then Module}
<Module.default
bind:schema={
() =>
value && typeof value === 'object' && !Array.isArray(value) ? value : emptySchema(),
(v) => {
value = v
}
}
isFlowInput
editTab="inputEditor"
noPreview
addPropertyInEditorTab
/>
{/await}
{:else if inputCat == 'object' && format?.startsWith('jsonschema-')}
{#await getJsonSchemaFromResource(format.substring('jsonschema-'.length), workspace ?? $workspaceStore ?? '')}
<Loader2 class="animate-spin" />
{:then schema}
{#if !schema || !schema.properties}
{#await import('$lib/components/JsonEditor.svelte')}
<Loader2 class="animate-spin" />
{:then Module}
<Module.default code={JSON.stringify(value, null, 2)} bind:value />
{/await}
{:else}
<div class="py-4 pr-2 pl-6 border rounded-md w-full">
<SchemaForm
{onlyMaskPassword}
{disablePortal}
{disabled}
{prettifyHeader}
{schema}
bind:args={value}
/>
</div>
{/if}
{/await}
{:else if inputCat == 'list' && !isListJson}
<div class="w-full flex gap-4">
<div class="w-full">
+5 -4
View File
@@ -48,7 +48,7 @@
import type { FlowPropPickerConfig, PropPickerContext } from './prop_picker'
import type { PickableProperties } from './flows/previousResults'
import { Triggers } from './triggers/triggers.svelte'
import { TestSteps } from './flows/testSteps.svelte'
import { StepsInputArgs } from './flows/stepsInputArgs.svelte'
import { ModulesTestStates } from './modulesTest.svelte'
import type { GraphModuleState } from './graph'
@@ -457,7 +457,7 @@
const scriptEditorDrawer = writable(undefined)
const moving = writable<{ id: string } | undefined>(undefined)
const history = initHistory(flowStore.val)
const testSteps = new TestSteps()
const stepsInputArgs = new StepsInputArgs()
const selectedIdStore = writable('settings-metadata')
const triggersCount = writable<TriggersCount | undefined>(undefined)
const modulesTestStates = new ModulesTestStates((moduleId) => {
@@ -481,7 +481,7 @@
pathStore: writable(''),
flowStateStore,
flowStore,
testSteps,
stepsInputArgs,
saveDraft: () => {},
initialPathStore: writable(''),
fakeInitialPath: '',
@@ -806,7 +806,7 @@
noEditor
on:applyArgs={(ev) => {
if (ev.detail.kind === 'preprocessor') {
testSteps.setStepArgs('preprocessor', ev.detail.args ?? {})
stepsInputArgs.setStepArgs('preprocessor', ev.detail.args ?? {})
$selectedIdStore = 'preprocessor'
} else {
previewArgsStore.val = ev.detail.args ?? {}
@@ -818,6 +818,7 @@
isOwner={flowPreviewContent?.getIsOwner()}
{suspendStatus}
onOpenDetails={flowPreviewButtons?.openPreview}
previewOpen={flowPreviewButtons?.getPreviewOpen()}
/>
{/key}
</Pane>
@@ -212,7 +212,7 @@
<Module.default
open={true}
automaticLayout
class="h-full"
className="h-full"
defaultLang={lang}
defaultModifiedLang={data.current.lang}
defaultOriginal={content}
@@ -227,7 +227,7 @@
<Module.default
open={true}
automaticLayout
class="h-full"
className="h-full"
defaultLang="yaml"
defaultOriginal={metadata}
defaultModified={data.current.metadata}
+51 -30
View File
@@ -1,6 +1,6 @@
<script lang="ts">
import { BROWSER } from 'esm-env'
import { createEventDispatcher, onMount } from 'svelte'
import { onMount } from 'svelte'
import '@codingame/monaco-vscode-standalone-languages'
import '@codingame/monaco-vscode-standalone-json-language-features'
@@ -10,24 +10,47 @@
import { initializeVscode } from './vscode'
import EditorTheme from './EditorTheme.svelte'
import Button from '$lib/components/common/button/Button.svelte'
import { twMerge } from 'tailwind-merge'
import type { ButtonType } from './common'
const SIDE_BY_SIDE_MIN_WIDTH = 700
export let automaticLayout = true
export let fixedOverflowWidgets = true
export let defaultLang: string | undefined = undefined
export let defaultModifiedLang: string | undefined = undefined
export let defaultOriginal: string | undefined = undefined
export let defaultModified: string | undefined = undefined
export let readOnly = false
export let showButtons = false
export let showHistoryButton: boolean = true
export interface ButtonProp {
text: string
color?: ButtonType.Color
onClick: () => void
}
let diffEditor: meditor.IStandaloneDiffEditor | undefined
let diffDivEl: HTMLDivElement | null = null
let editorWidth: number = SIDE_BY_SIDE_MIN_WIDTH
interface Props {
open?: boolean
className?: string
automaticLayout?: boolean
fixedOverflowWidgets?: boolean
defaultLang?: string
defaultModifiedLang?: string
defaultOriginal?: string
defaultModified?: string
readOnly?: boolean
buttons?: ButtonProp[]
}
let {
open = false,
className = '',
automaticLayout = true,
fixedOverflowWidgets = true,
defaultLang,
defaultModifiedLang,
defaultOriginal = undefined,
defaultModified = undefined,
readOnly = false,
buttons = []
}: Props = $props()
let diffEditor: meditor.IStandaloneDiffEditor | undefined = $state(undefined)
let diffDivEl: HTMLDivElement | null = $state(null)
let editorWidth: number = $state(SIDE_BY_SIDE_MIN_WIDTH)
export let open = false
async function loadDiffEditor() {
await initializeVscode()
@@ -105,9 +128,15 @@
diffEditor?.updateOptions({ renderSideBySide: editorWidth >= SIDE_BY_SIDE_MIN_WIDTH })
}
$: onWidthChange(editorWidth)
$effect(() => {
if (open && diffDivEl) {
loadDiffEditor()
}
})
$: open && diffDivEl && loadDiffEditor()
$effect(() => {
onWidthChange(editorWidth)
})
onMount(() => {
if (BROWSER) {
@@ -116,32 +145,24 @@
}
}
})
const dispatch = createEventDispatcher<{
hideDiffMode: void
seeHistory: void
}>()
</script>
{#if open}
<EditorTheme />
<div
bind:this={diffDivEl}
class="{$$props.class} editor nonmain-editor"
class={twMerge('editor nonmain-editor', className)}
bind:clientWidth={editorWidth}
></div>
{#if showButtons}
{#if buttons.length > 0}
<div
class="absolute flex flex-row gap-2 bottom-10 left-1/2 z-10 -translate-x-1/2 rounded-md p-1 w-full justify-center"
>
{#if showHistoryButton}
<Button on:click={() => dispatch('seeHistory')} variant="contained" size="sm"
>See changes history</Button
{#each buttons as button}
<Button on:click={button.onClick} variant="contained" size="sm" color={button.color}
>{button.text}</Button
>
{/if}
<Button on:click={() => dispatch('hideDiffMode')} variant="contained" size="sm" color="red"
>Quit diff mode</Button
>
{/each}
</div>
{/if}
{/if}
@@ -928,7 +928,7 @@
>
</button>
{:else if !s3object?.disable_download}
<FileDownload {s3object} />
<FileDownload {workspaceId} {s3object} {appPath} />
{:else}
<div class="flex text-secondary pt-2">{s3object?.s3} (download disabled)</div>
{/if}
@@ -29,6 +29,7 @@
import type { EditableSchemaFormUi } from '$lib/components/custom_ui'
import Section from '$lib/components/Section.svelte'
import Editor from './Editor.svelte'
import AddPropertyV2 from './schema/AddPropertyV2.svelte'
// export let openEditTab: () => void = () => {}
const dispatch = createEventDispatcher()
@@ -69,6 +70,7 @@
dynSelectCode?: string | undefined
dynSelectLang?: ScriptLang | undefined
showDynSelectOpt?: boolean
addPropertyInEditorTab?: boolean
openEditTab?: import('svelte').Snippet
addProperty?: import('svelte').Snippet
runButton?: import('svelte').Snippet
@@ -104,6 +106,7 @@
dynSelectCode = $bindable(),
dynSelectLang = $bindable(),
showDynSelectOpt = false,
addPropertyInEditorTab = false,
openEditTab,
addProperty,
runButton,
@@ -509,22 +512,31 @@
{:else}
<!-- WIP -->
{#if jsonEnabled && customUi?.jsonOnly != true}
<div class="w-full p-3 flex justify-end">
<Toggle
bind:checked={jsonView}
label="JSON View"
size="xs"
options={{
right: 'JSON editor',
rightTooltip:
'Arguments can be edited either using the wizard, or by editing their JSON Schema.'
}}
lightMode
on:change={() => {
schemaString = JSON.stringify(schema, null, '\t')
editor?.setCode(schemaString)
}}
/>
<div class="w-full p-3 flex gap-4 justify-end items-center">
{#if addPropertyInEditorTab}
<AddPropertyV2 bind:schema on:change>
{#snippet trigger()}
<Button color="light" size="xs" iconOnly startIcon={{ icon: Plus }} />
{/snippet}
</AddPropertyV2>
{/if}
<div class="shrink-0">
<Toggle
bind:checked={jsonView}
label="JSON View"
size="xs"
options={{
right: 'JSON editor',
rightTooltip:
'Arguments can be edited either using the wizard, or by editing their JSON Schema.'
}}
lightMode
on:change={() => {
schemaString = JSON.stringify(schema, null, '\t')
editor?.setCode(schemaString)
}}
/>
</div>
</div>
{/if}
@@ -655,7 +667,6 @@
const isS3 = v == 'S3'
const isOneOf = v == 'oneOf'
const isDynSelect = v == 'dynselect'
const emptyProperty = {
contentEncoding: undefined,
enum_: undefined,
+13 -4
View File
@@ -184,6 +184,7 @@
loadAsync?: boolean
key?: string | undefined
class?: string | undefined
moduleId?: string
}
let {
@@ -209,7 +210,8 @@
changeTimeout = 500,
loadAsync = false,
key = undefined,
class: clazz = undefined
class: clazz = undefined,
moduleId = undefined
}: Props = $props()
$effect.pre(() => {
@@ -1235,7 +1237,13 @@
try {
editor = meditor.create(divEl as HTMLDivElement, {
...editorConfig(code ?? '', lang, automaticLayout, fixedOverflowWidgets, $relativeLineNumbers),
...editorConfig(
code ?? '',
lang,
automaticLayout,
fixedOverflowWidgets,
$relativeLineNumbers
),
model,
fontSize: !small ? 14 : 12,
lineNumbersMinChars,
@@ -1328,7 +1336,8 @@
aiChatManager.addSelectedLinesToContext(
selectedLines,
selection.startLineNumber,
selection.endLineNumber
selection.endLineNumber,
moduleId
)
} else {
aiChatManager.toggleOpen()
@@ -1654,7 +1663,7 @@
files && model && untrack(() => onFileChanges())
})
$effect(() => {
editor?.updateOptions({
editor?.updateOptions({
lineNumbers: $relativeLineNumbers ? 'relative' : 'on'
})
})
+4 -2
View File
@@ -82,6 +82,7 @@
showHistoryDrawer?: boolean
right?: import('svelte').Snippet
openAiChat?: boolean
moduleId?: string
}
let {
@@ -105,7 +106,8 @@
diffMode = false,
showHistoryDrawer = $bindable(false),
right,
openAiChat = false
openAiChat = false,
moduleId = undefined
}: Props = $props()
let contextualVariablePicker: ItemPicker | undefined = $state()
@@ -964,7 +966,7 @@ JsonNode ${windmillPathToCamelCaseName(path)} = JsonNode.Parse(await client.GetS
{#if customUi?.aiGen != false}
{#if openAiChat}
<FlowInlineScriptAiButton />
<FlowInlineScriptAiButton {moduleId} />
{:else}
<ScriptGen {editor} {diffEditor} {lang} {iconOnly} {args} />
{/if}
@@ -40,7 +40,7 @@
{/if}
{#if displayType}
{#if format && !format.startsWith('resource')}
{#if format && !format.startsWith('resource') && !format.startsWith('jsonschema-')}
<span class="text-xs italic ml-2 text-tertiary dark:text-indigo-400">
{format}
</span>
@@ -77,7 +77,7 @@
} from './triggers/utils'
import DraftTriggersConfirmationModal from './common/confirmationModal/DraftTriggersConfirmationModal.svelte'
import { Triggers } from './triggers/triggers.svelte'
import { TestSteps } from './flows/testSteps.svelte'
import { StepsInputArgs } from './flows/stepsInputArgs.svelte'
import { aiChatManager } from './copilot/chat/AIChatManager.svelte'
import type { GraphModuleState } from './graph'
import {
@@ -571,7 +571,7 @@
payloadData: undefined
})
const testSteps = new TestSteps()
const stepsInputArgs = new StepsInputArgs()
function select(selectedId: string) {
selectedIdStore.set(selectedId)
@@ -592,7 +592,7 @@
flowStateStore,
flowStore,
pathStore,
testSteps,
stepsInputArgs,
saveDraft,
initialPathStore,
fakeInitialPath,
@@ -1129,6 +1129,8 @@
bind:this={flowPreviewButtons}
{loading}
onRunPreview={() => {
// Reset manually edited args inputs when running a preview
stepsInputArgs.resetManuallyEditedArgs()
modulesTestStates.hideJobsInGraph()
localModuleStates = {}
showJobStatus = true
@@ -1170,7 +1172,7 @@
{newFlow}
on:applyArgs={(ev) => {
if (ev.detail.kind === 'preprocessor') {
testSteps.setStepArgs('preprocessor', ev.detail.args ?? {})
stepsInputArgs.setStepArgs('preprocessor', ev.detail.args ?? {})
$selectedIdStore = 'preprocessor'
}
}}
@@ -1218,6 +1220,7 @@
delete modulesTestStates.states[id]
}}
{flowHasChanged}
previewOpen={flowPreviewButtons?.getPreviewOpen()}
/>
{:else}
<CenteredPage>Loading...</CenteredPage>
@@ -140,7 +140,7 @@
jobId = await runFlowPreview(args, newFlow, $pathStore, restartedFrom)
isRunning = true
if (inputSelected) {
savedArgs = previewArgs.val
savedArgs = $state.snapshot(previewArgs.val)
inputSelected = undefined
}
onRunPreview?.()
@@ -166,7 +166,7 @@
if (preventEscape) {
selectInput(undefined)
event.preventDefault()
event.stopPropagation
event.stopPropagation()
}
break
}
@@ -506,7 +506,7 @@
schema={flowStore.val.schema}
bind:args={previewArgs.val}
on:change={() => {
savedArgs = previewArgs.val
savedArgs = $state.snapshot(previewArgs.val)
}}
bind:isValid
helperScript={flowStore.val.schema?.['x-windmill-dyn-select-code'] &&
@@ -110,6 +110,9 @@
let updateGlobalRefresh = (moduleId: string, updateFn: (clear, root) => Promise<void>) => {
globalRefreshes[moduleId] = [...(globalRefreshes[moduleId] ?? []), updateFn]
}
let storedToolCallJobs: Record<string, Job> = $state({})
let toolCallIndicesToLoad: string[] = $state([])
</script>
<FlowStatusViewerInner
@@ -141,4 +144,30 @@
isNodeSelected={true}
{refreshGlobal}
{updateGlobalRefresh}
toolCallStore={{
getStoredToolCallJob: (storeKey: string) => storedToolCallJobs[storeKey],
setStoredToolCallJob: (storeKey: string, job: Job) => {
storedToolCallJobs[storeKey] = job
},
getLocalToolCallJobs: (prefix: string) => {
// we return a map from tool call index to job
// to do so, we filter the storedToolCallJobs object by the prefix and we make sure what's left in the key is a tool call index: 2 part of format agentModuleId-toolCallIndex
// and not a further nested tool call index
return Object.fromEntries(
Object.entries(storedToolCallJobs)
.filter(
([key]) => key.startsWith(prefix) && key.replace(prefix, '').split('-').length === 2
)
.map(([key, job]) => [Number(key.replace(prefix, '').split('-').pop()), job])
)
},
isToolCallToBeLoaded: (storeKey: string) => {
return toolCallIndicesToLoad.includes(storeKey)
},
addToolCallToLoad: (storeKey: string) => {
if (!toolCallIndicesToLoad.includes(storeKey)) {
toolCallIndicesToLoad.push(storeKey)
}
}
}}
/>
@@ -33,6 +33,7 @@
import { deepEqual } from 'fast-equals'
import FlowTimeline from './FlowTimeline.svelte'
import { dfs } from './flows/dfs'
import { dfs as dfsPreviousResults } from '$lib/components/flows/previousResults'
import Alert from './common/alert/Alert.svelte'
import FlowGraphViewerStep from './FlowGraphViewerStep.svelte'
import FlowGraphV2 from './graph/FlowGraphV2.svelte'
@@ -116,6 +117,13 @@
onStart?: () => void
onJobsLoaded?: ({ job, force }: { job: Job; force: boolean }) => void
onDone?: ({ job }: { job: CompletedJob }) => void
toolCallStore?: {
getStoredToolCallJob: (storeKey: string) => Job | undefined
setStoredToolCallJob: (storeKey: string, job: Job) => void
getLocalToolCallJobs: (prefix: string) => Record<number, Job>
isToolCallToBeLoaded: (storeKey: string) => boolean
addToolCallToLoad: (storeKey: string) => void
}
}
let {
@@ -155,7 +163,8 @@
loadExtraLogs = undefined,
onStart = undefined,
onJobsLoaded = undefined,
onDone = undefined
onDone = undefined,
toolCallStore
}: Props = $props()
let getTopModuleStates = $derived(topModuleStates ?? localModuleStates)
@@ -913,9 +922,7 @@
let storedListJobs: Record<number, Job> = $state({})
let storedToolCallJobs: Record<number, Job> = $state({})
let selectedToolCall: number | undefined = $state(undefined)
let toolCallIndicesToLoad: number[] = $state([])
let selectedToolCall: string | undefined = $state(undefined)
let wrapperHeight: number = $state(0)
@@ -950,8 +957,10 @@
let nprefix = buildPrefix(prefix, oid)
return fms
? rec(
dfs(fms, (x) =>
x.id.startsWith('subflow:') ? x.id : buildSubflowKey(x.id, nprefix)
dfs(
fms,
(x) => (x.id.startsWith('subflow:') ? x.id : buildSubflowKey(x.id, nprefix)),
{ skipToolNodes: true }
),
nprefix
)
@@ -1009,6 +1018,11 @@
selectedForLoopSetManually: false
})
}
if (selectedNode?.startsWith(AI_TOOL_CALL_PREFIX)) {
const [, agentModuleId, toolCallIndex, _] = selectedNode.split('-')
const parentLoopsPrefix = getParentLoopsPrefix(agentModuleId)
toolCallStore?.addToolCallToLoad(parentLoopsPrefix + agentModuleId + '-' + toolCallIndex)
}
}
$effect(() => {
@@ -1039,6 +1053,29 @@
let animateLogsTab = $state(false)
let noLogs = $derived(graphTabOpen && !isNodeSelected)
/**
* Returns a string like "forloopmodid1-{iter1}-forloopmodid2-{iter2}-forloopmodid3-{iter3}-"
* that can be used to prefix tool call store keys for nested tool calls.
*/
function getParentLoopsPrefix(modId: string) {
if (job?.raw_flow) {
const indices: string[] = []
const parents = dfsPreviousResults(modId, { value: job?.raw_flow, summary: '' }, true)
for (const parent of parents) {
if (parent.value.type === 'forloopflow' || parent.value.type === 'whileloopflow') {
const state = localModuleStates[parent.id]
if (state?.selectedForloopIndex !== undefined) {
indices.push(parent.id + '-' + state.selectedForloopIndex.toString())
}
}
}
indices.reverse()
return indices.length > 0 ? indices.join('-') + '-' : ''
}
return ''
}
</script>
<JobLoader workspaceOverride={workspaceId} {noLogs} noCode bind:this={jobLoader} />
@@ -1173,6 +1210,10 @@
{@const forloopIsSelected =
forloop_selected == loopJobId ||
(innerModule?.type != 'forloopflow' && innerModule?.type != 'whileloopflow')}
{@const forLoopStoreKeyPrefix =
innerModule?.type == 'forloopflow' || innerModule?.type == 'whileloopflow'
? (flowJobIds?.moduleId ?? '') + '-' + j + '-'
: ''}
<!-- <LogId id={loopJobId} /> -->
<div class="border p-6" class:hidden={forloop_selected != loopJobId}>
<FlowStatusViewerInner
@@ -1207,6 +1248,18 @@
graphTabOpen={selected == 'graph' && graphTabOpen}
isNodeSelected={forloop_selected == loopJobId}
{globalIterationBounds}
toolCallStore={{
getStoredToolCallJob: (storeKey: string) =>
toolCallStore?.getStoredToolCallJob(forLoopStoreKeyPrefix + storeKey),
setStoredToolCallJob: (storeKey: string, job: Job) =>
toolCallStore?.setStoredToolCallJob(forLoopStoreKeyPrefix + storeKey, job),
getLocalToolCallJobs: (prefix: string) =>
toolCallStore?.getLocalToolCallJobs(forLoopStoreKeyPrefix + prefix) ?? {},
addToolCallToLoad: (storeKey: string) =>
toolCallStore?.addToolCallToLoad(forLoopStoreKeyPrefix + storeKey),
isToolCallToBeLoaded: (storeKey: string) =>
toolCallStore?.isToolCallToBeLoaded(forLoopStoreKeyPrefix + storeKey) ?? false
}}
/>
</div>
{/if}
@@ -1366,12 +1419,17 @@
graphTabOpen={selected == 'graph' && graphTabOpen}
isNodeSelected={localModuleStates?.[selectedNode ?? '']?.job_id == mod.job}
{globalIterationBounds}
{toolCallStore}
/>
{#if mod.agent_actions && mod.agent_actions.length > 0}
{#if mod.agent_actions && mod.agent_actions.length > 0 && mod.id}
{@const storeKeyPrefix = getParentLoopsPrefix(mod.id)}
{#each mod.agent_actions as agentAction, j}
{#if agentAction.type === 'tool_call' && mod.id}
{#if agentAction.type === 'tool_call'}
{@const toolCallId = getToolCallId(j, mod.id, agentAction.module_id)}
{@const isSelected = selectedToolCall === j}
{@const localToolCallKey = mod.id + '-' + j}
{@const storeKey = storeKeyPrefix + localToolCallKey}
{@const storedToolCallJob = toolCallStore?.getStoredToolCallJob(storeKey)}
{@const isSelected = localToolCallKey === selectedToolCall}
<Button
variant={isSelected ? 'contained' : 'border'}
color={mod.agent_actions_success?.[j] === false
@@ -1381,10 +1439,10 @@
: 'light'}
btnClasses="w-full flex justify-start"
on:click={async () => {
if (selectedToolCall == j) {
if (isSelected) {
selectedToolCall = undefined
} else {
selectedToolCall = j
selectedToolCall = localToolCallKey
}
}}
endIcon={{
@@ -1396,7 +1454,7 @@
Tool call: {agentAction.function_name}
</span>
</Button>
{#if isSelected || storedToolCallJobs[j] || toolCallIndicesToLoad.includes(j)}
{#if isSelected || storedToolCallJob || toolCallStore?.isToolCallToBeLoaded(storeKey)}
<FlowStatusViewerInner
topModuleStates={getTopModuleStates}
{refreshGlobal}
@@ -1414,11 +1472,11 @@
{subflowParentsDurationStatuses}
{isSelectedBranch}
jobId={agentAction.job_id}
job={storedToolCallJobs[j]}
initialJob={storedToolCallJobs[j]}
job={storedToolCallJob}
initialJob={storedToolCallJob}
{reducedPolling}
onJobsLoaded={({ job, force }) => {
storedToolCallJobs[j] = job
toolCallStore?.setStoredToolCallJob(storeKey, job)
onJobsLoadedInner({ id: toolCallId } as FlowStatusModule, job, force)
}}
loadExtraLogs={(logs) => {
@@ -1509,11 +1567,11 @@
stepDetail = mod
selectedNode = e
if (e.startsWith(AI_TOOL_CALL_PREFIX)) {
const [_prefix, _agentModuleId, j, _toolModuleId] = e.split('-')
const [_prefix, agentModuleId, j, _toolModuleId] = e.split('-')
const parentLoopsPrefix = getParentLoopsPrefix(agentModuleId)
const jIdx = Number(j)
if (!toolCallIndicesToLoad.includes(jIdx)) {
toolCallIndicesToLoad.push(jIdx)
}
const storeKey = parentLoopsPrefix + agentModuleId + '-' + jIdx
toolCallStore?.addToolCallToLoad(storeKey)
}
}
} else {
@@ -1603,6 +1661,7 @@
stepDetail && typeof stepDetail !== 'string' ? stepDetail : undefined}
{@const agentTools =
module && module.value.type === 'aiagent' ? module.value.tools : undefined}
{@const parentLoopsPrefix = getParentLoopsPrefix(module?.id ?? '')}
{#if node.flow_jobs_results}
<span class="pl-1 text-tertiary"
>Result of step as collection of all subflows</span
@@ -1667,7 +1726,7 @@
logs={node.logs}
downloadLogs={!hideDownloadLogs}
aiAgentStatus={agentTools &&
node.job_id &&
node?.job_id &&
(node.type === 'Success' || node.type === 'Failure')
? {
tools: agentTools,
@@ -1679,9 +1738,14 @@
success: node.type === 'Success',
type: 'CompletedJob'
},
storedToolCallJobs,
storedToolCallJobs: module
? toolCallStore?.getLocalToolCallJobs(parentLoopsPrefix)
: undefined,
onToolJobLoaded: (job, idx) => {
storedToolCallJobs[idx] = job
if (module) {
const storeKey = parentLoopsPrefix + module.id + '-' + idx
toolCallStore?.setStoredToolCallJob(storeKey, job)
}
}
}
: undefined}
@@ -5,12 +5,14 @@
let {
flowStore: oldFlowStore,
flowStateStore: oldFlowStateStore,
disableAi,
light,
...props
}: FlowBuilderProps & { light?: boolean } = $props()
let flowStore = $state(oldFlowStore)
let flowStateStore = $state(oldFlowStateStore)
let trialRender = $state(true)
@@ -24,7 +26,7 @@
{#if trialRender}
<AiChatLayout noPadding={true} {disableAi}>
{#if light}<div class="bg-red-500 absolute z-10">Trial version</div>{/if}
<FlowBuilder {flowStore} {disableAi} {...props} />
<FlowBuilder {flowStore} {flowStateStore} {disableAi} {...props} />
</AiChatLayout>
{:else}
<div class="flex flex-col items-center justify-center h-screen">
@@ -18,6 +18,7 @@
noEditor?: boolean
scriptProgress?: any
focusArg?: string
onJobDone?: () => void
}
let {
@@ -28,7 +29,8 @@
testIsLoading = $bindable(false),
noEditor = false,
scriptProgress = $bindable(undefined),
focusArg = undefined
focusArg = undefined,
onJobDone
}: Props = $props()
const { flowStore } = getContext<FlowEditorContext>('FlowEditorContext')
@@ -46,6 +48,7 @@
bind:testIsLoading
bind:scriptProgress
bind:this={moduleTest}
{onJobDone}
/>
<div class="p-4">
@@ -1,11 +1,11 @@
<script lang="ts">
import type { Schema } from '$lib/common'
import { allTrue, sendUserToast } from '$lib/utils'
import { allTrue } from '$lib/utils'
import { RefreshCw } from 'lucide-svelte'
import ArgInput from './ArgInput.svelte'
import { Button } from './common'
import { getContext, onMount, untrack } from 'svelte'
import { getContext, untrack } from 'svelte'
import type { FlowEditorContext } from './flows/types'
import { evalValue } from './flows/utils'
import type { FlowModule } from '$lib/gen'
@@ -32,7 +32,7 @@
focusArg = undefined
}: Props = $props()
const { testSteps, flowStateStore, flowStore, previewArgs } =
const { stepsInputArgs, flowStateStore, flowStore, previewArgs } =
getContext<FlowEditorContext>('FlowEditorContext')
let inputCheck: { [id: string]: boolean } = $state({})
@@ -45,12 +45,12 @@
let lkeys = Object.keys(schema?.properties ?? {})
if (schema?.properties && JSON.stringify(lkeys) != JSON.stringify(keys)) {
keys = lkeys
untrack(() => testSteps?.removeExtraKey(mod.id, keys))
untrack(() => stepsInputArgs?.removeExtraKey(mod.id, keys))
}
})
function plugIt(argName: string) {
testSteps?.setEvaluatedStepArg(
stepsInputArgs?.setEvaluatedStepArg(
mod.id,
argName,
$state.snapshot(evalValue(argName, mod, pickableProperties, true))
@@ -98,70 +98,78 @@
loadResourceTypes()
onMount(() => {
if (!testSteps) {
sendUserToast('testSteps module not initialized. Preview will not work.', true)
let initialized = $state(false)
$effect.pre(() => {
if (!initialized) {
if (stepsInputArgs) {
stepsInputArgs?.updateStepArgs(mod.id, flowStateStore.val, flowStore?.val, previewArgs?.val)
initialized = true
}
}
testSteps?.updateStepArgs(mod.id, flowStateStore.val, flowStore?.val, previewArgs?.val)
})
</script>
<div class="w-full pt-2" data-popover>
{#if keys.length > 0}
{#each keys as argName, i (argName)}
{#if Object.keys(schema.properties ?? {}).includes(argName)}
<div
class={twMerge(
'flex gap-2',
animateArg === argName && 'animate-pulse ring-2 ring-offset-2 ring-blue-500 rounded'
)}
data-arg={argName}
>
{#if schema?.properties?.[argName]}
<ArgInput
{resourceTypes}
minW={false}
autofocus={autofocus && !focusArg && i == 0}
label={argName}
description={schema.properties[argName].description}
bind:value={
() => testSteps?.getStepInputArgs(mod.id, argName),
(v) => testSteps?.setStepInputArgs(mod.id, argName, v)
}
type={schema.properties[argName].type}
oneOf={schema.properties[argName].oneOf}
required={schema?.required?.includes(argName)}
pattern={schema.properties[argName].pattern}
bind:editor={editor[argName]}
bind:valid={inputCheck[argName]}
defaultValue={schema.properties[argName].default}
enum_={schema.properties[argName].enum}
format={schema.properties[argName].format}
contentEncoding={schema.properties[argName].contentEncoding}
properties={schema.properties[argName].properties}
nestedRequired={schema.properties[argName].required}
itemsType={schema.properties[argName].items}
extra={schema.properties[argName]}
nullable={schema.properties[argName].nullable}
title={schema.properties[argName].title}
placeholder={schema.properties[argName].placeholder}
/>
{/if}
{#if testSteps?.isArgManuallySet(mod.id, argName)}
<div class="pt-6 mt-0.5">
<Button
on:click={() => {
plugIt(argName)
}}
size="sm"
variant="border"
color="light"
title="Re-evaluate input step"><RefreshCw size={14} /></Button
>
</div>
{/if}
</div>
{/if}
{/each}
{#if initialized}
{#if keys.length > 0}
{#each keys as argName, i (argName)}
{#if Object.keys(schema.properties ?? {}).includes(argName)}
<div
class={twMerge(
'flex gap-2',
animateArg === argName && 'animate-pulse ring-2 ring-offset-2 ring-blue-500 rounded'
)}
data-arg={argName}
>
{#if schema?.properties?.[argName]}
<ArgInput
{resourceTypes}
minW={false}
autofocus={autofocus && !focusArg && i == 0}
label={argName}
description={schema.properties[argName].description}
bind:value={
() => stepsInputArgs?.getStepInputArgs(mod.id, argName),
(v) => stepsInputArgs?.setStepInputArgs(mod.id, argName, v)
}
type={schema.properties[argName].type}
oneOf={schema.properties[argName].oneOf}
required={schema?.required?.includes(argName)}
pattern={schema.properties[argName].pattern}
bind:editor={editor[argName]}
bind:valid={inputCheck[argName]}
defaultValue={schema.properties[argName].default}
enum_={schema.properties[argName].enum}
format={schema.properties[argName].format}
contentEncoding={schema.properties[argName].contentEncoding}
properties={schema.properties[argName].properties}
nestedRequired={schema.properties[argName].required}
itemsType={schema.properties[argName].items}
extra={schema.properties[argName]}
nullable={schema.properties[argName].nullable}
title={schema.properties[argName].title}
placeholder={schema.properties[argName].placeholder}
/>
{/if}
{#if stepsInputArgs?.isArgManuallySet(mod.id, argName)}
<div class="pt-6 mt-0.5">
<Button
on:click={() => {
plugIt(argName)
}}
size="sm"
variant="border"
color="light"
title="Re-evaluate input step"><RefreshCw size={14} /></Button
>
</div>
{/if}
</div>
{/if}
{/each}
{/if}
{:else}
<div class="text-center text-sm text-tertiary"> Loading test step arguments... </div>
{/if}
</div>
@@ -7,7 +7,7 @@
import { type Script, type Job, type FlowModule } from '$lib/gen'
import OutputPickerInner from '$lib/components/flows/propPicker/OutputPickerInner.svelte'
import { Pane, Splitpanes } from 'svelte-splitpanes'
import type { FlowEditorContext } from './flows/types'
import type { FlowEditorContext, OutputViewerJob } from './flows/types'
import { getContext } from 'svelte'
import { getStringError } from './copilot/chat/utils'
import AiAgentLogViewer from './AIAgentLogViewer.svelte'
@@ -17,9 +17,8 @@
editor: Editor | undefined
diffEditor: DiffEditor | undefined
loopStatus?: { type: 'inside' | 'self'; flow: 'forloopflow' | 'whileloopflow' } | undefined
lastJob?: Job | undefined
testJob?: Job & { result_stream?: string }
scriptProgress?: number | undefined
testJob?: Job | undefined
mod: FlowModule
testIsLoading?: boolean
disableMock?: boolean
@@ -34,7 +33,6 @@
editor,
diffEditor,
loopStatus = undefined,
lastJob = undefined,
scriptProgress = $bindable(undefined),
testJob = undefined,
mod,
@@ -46,21 +44,20 @@
tagLabel = undefined
}: Props = $props()
const { testSteps } = getContext<FlowEditorContext>('FlowEditorContext')
const { stepsInputArgs } = getContext<FlowEditorContext>('FlowEditorContext')
let selectedJob: Job | undefined = $state(undefined)
let preview: 'mock' | 'job' | undefined = $state(undefined)
let jobProgressReset: () => void = $state(() => {})
$effect(() => {
if (preview != undefined && testJob) {
preview = undefined
}
})
let forceJson = $state(false)
let outputPickerInner: OutputPickerInner | undefined = $state(undefined)
export function getOutputPickerInner() {
return outputPickerInner
}
const selectedJob: OutputViewerJob = $derived.by(
() => outputPickerInner?.getSelectedJob?.() ?? undefined
)
const logJob = $derived(testJob ?? selectedJob)
const preview = $derived.by(() => outputPickerInner?.getPreview?.())
</script>
<Splitpanes horizontal>
@@ -75,7 +72,6 @@
{/if}
<OutputPickerInner
{lastJob}
{testJob}
fullResult
moduleId={mod.id}
@@ -83,24 +79,22 @@
getLogs
{onUpdateMock}
mock={mod.mock}
bind:forceJson
bind:selectedJob
isLoading={testIsLoading || loadingJob}
bind:preview
path={`path` in mod.value ? mod.value.path : ''}
{loopStatus}
{disableMock}
{disableHistory}
bind:this={outputPickerInner}
>
{#snippet copilot_fix()}
{#if lang && editor && diffEditor && testSteps.getStepArgs(mod.id) && selectedJob?.type === 'CompletedJob' && !selectedJob.success && getStringError(selectedJob.result)}
{#if lang && editor && diffEditor && stepsInputArgs.getStepArgs(mod.id) && selectedJob?.type === 'CompletedJob' && !selectedJob.success && getStringError(selectedJob.result)}
<ScriptFix {lang} />
{/if}
{/snippet}
</OutputPickerInner>
</Pane>
<Pane size={35} minSize={10}>
{#if (mod.mock?.enabled && preview != 'job') || preview == 'mock'}
{#if (mod.mock?.enabled && preview !== 'job' && testJob?.type !== 'QueuedJob') || preview === 'mock'}
<LogViewer
small
content={undefined}
@@ -124,9 +118,9 @@
jobId={logJob?.id}
duration={logJob?.['duration_ms']}
mem={logJob?.['mem_peak']}
content={logJob?.logs}
content={logJob?.['logs']}
isLoading={(testIsLoading && logJob?.['running'] == false) || loadingJob}
tag={logJob?.tag}
tag={logJob?.['tag']}
{tagLabel}
/>
{/if}
+12 -6
View File
@@ -14,6 +14,7 @@
testIsLoading?: boolean
noEditor?: boolean
scriptProgress?: any
onJobDone?: () => void
}
let {
@@ -21,10 +22,11 @@
testJob = $bindable(undefined),
testIsLoading = $bindable(false),
noEditor = false,
scriptProgress = $bindable(undefined)
scriptProgress = $bindable(undefined),
onJobDone
}: Props = $props()
const { flowStore, flowStateStore, pathStore, testSteps, previewArgs, modulesTestStates } =
const { flowStore, flowStateStore, pathStore, stepsInputArgs, previewArgs, modulesTestStates } =
getContext<FlowEditorContext>('FlowEditorContext')
let jobLoader: JobLoader | undefined = $state(undefined)
@@ -32,12 +34,12 @@
let stepHistoryLoader = getStepHistoryLoaderContext()
export function runTestWithStepArgs() {
runTest(testSteps.getStepArgs(mod.id))
runTest(stepsInputArgs.getStepArgs(mod.id))
}
export function loadArgsAndRunTest() {
testSteps?.updateStepArgs(mod.id, flowStateStore.val, flowStore?.val, previewArgs?.val)
runTest(testSteps.getStepArgs(mod.id))
stepsInputArgs?.updateStepArgs(mod.id, flowStateStore.val, flowStore?.val, previewArgs?.val)
runTest(stepsInputArgs.getStepArgs(mod.id))
}
export async function runTest(args: any) {
@@ -138,6 +140,7 @@
if (modulesTestStates.states[mod.id]) {
modulesTestStates.states[mod.id].testJob = testJob
}
onJobDone?.()
}
export function cancelJob() {
@@ -179,7 +182,10 @@
}
}
}
bind:job={modulesTestStates.states[mod.id].testJob}
bind:job={
() => modulesTestStates.states[mod.id]?.testJob,
(v) => modulesTestStates.states[mod.id] && (modulesTestStates.states[mod.id].testJob = v)
}
loadPlaceholderJobOnStart={{
type: 'QueuedJob',
id: '',
@@ -22,6 +22,7 @@
import TestTriggerConnection from './triggers/TestTriggerConnection.svelte'
import GitHubAppIntegration from './GitHubAppIntegration.svelte'
import Button from './common/button/Button.svelte'
import { clearJsonSchemaResourceCache } from './schema/jsonSchemaResource.svelte'
interface Props {
canSave?: boolean
@@ -94,6 +95,9 @@
path: resourceToEdit.path,
requestBody: { path, value: args, description }
})
if (resourceToEdit.resource_type === 'json_schema') {
clearJsonSchemaResourceCache(resourceToEdit.path, $workspaceStore!)
}
sendUserToast(`Updated resource at ${path}`)
dispatch('refresh', path)
} else {
+2 -2
View File
@@ -77,7 +77,7 @@
}
let {
runnable = $bindable(),
runnable,
runAction,
buttonText = 'Run',
schedulable = true,
@@ -309,7 +309,7 @@
bind:scheduledForStr
bind:invisible_to_owner
bind:overrideTag
bind:runnable
{runnable}
/>
{/snippet}
</Popover>
@@ -70,7 +70,19 @@
{#if !$userStore?.operator}
{#if $workerTags && $workerTags?.length > 0}
<div class="w-full">
<select placeholder="Worker group" bind:value={overrideTag}>
<select
placeholder="Worker group"
bind:value={
() => overrideTag ?? '',
(v) => {
if (v == '') {
overrideTag = undefined
} else {
overrideTag = v
}
}
}
>
{#if overrideTag}
<option value="">reset to default</option>
{:else}
@@ -269,6 +269,7 @@
onMount(() => {
inferSchema(code)
loadPastTests()
aiChatManager.saveAndClear()
aiChatManager.changeMode(AIMode.SCRIPT)
})
@@ -625,16 +626,28 @@
{args}
/>
<DiffEditor
class="h-full"
className="h-full"
bind:this={diffEditor}
automaticLayout
defaultLang={scriptLangToEditorLang(lang)}
{fixedOverflowWidgets}
showButtons={diffMode}
on:hideDiffMode={hideDiffMode}
on:seeHistory={() => {
showHistoryDrawer = true
}}
buttons={diffMode
? [
{
text: 'See changes history',
onClick: () => {
showHistoryDrawer = true
}
},
{
text: 'Quit diff mode',
onClick: () => {
hideDiffMode()
},
color: 'red'
}
]
: []}
/>
{/key}
</div>
@@ -267,7 +267,7 @@
}
})
$effect(() => {
editor?.updateOptions({
editor?.updateOptions({
lineNumbers: $relativeLineNumbers ? 'relative' : 'on'
})
})
@@ -377,9 +377,14 @@
return
}
try {
console.log('fixedOverflowWidgets', fixedOverflowWidgets)
editor = meditor.create(divEl as HTMLDivElement, {
...editorConfig(code ?? '', lang, automaticLayout, fixedOverflowWidgets, $relativeLineNumbers),
...editorConfig(
code ?? '',
lang,
automaticLayout,
fixedOverflowWidgets,
$relativeLineNumbers
),
model,
lineDecorationsWidth: 6,
lineNumbersMinChars: 2,
@@ -423,7 +423,7 @@
<DiffEditor
open={false}
bind:this={diffEditor}
class="h-full"
className="h-full"
automaticLayout
fixedOverflowWidgets
defaultLang={scriptLangToEditorLang(inlineScript?.language)}
@@ -116,7 +116,6 @@
if (deepEqual(runnable, lastRunnable)) {
return
}
console.log('runnable', runnable)
notFound = false
if (runnable.runType == 'script') {
refreshScript(runnable)
@@ -35,7 +35,7 @@
showOnDemandOnlyToggle = true
}: Props = $props()
$effect(() => {
$effect.pre(() => {
if (oneOf == undefined) {
oneOf = { configuration: {}, selected: '' }
}
@@ -70,59 +70,61 @@
</script>
<div class="p-2 border">
<div class="mb-2 text-sm font-semibold">
{capitalize(addWhitespaceBeforeCapitals(key))}&nbsp;
{#if tooltip}
<Tooltip light>{tooltip}</Tooltip>
{/if}
</div>
<select
class="w-full border border-gray-300 rounded-md p-2"
value={oneOf.selected}
onchange={(e) => {
oneOf = { ...oneOf, selected: e?.target?.['value'] }
}}
>
{#each Object.keys(inputSpecsConfiguration ?? {}) as choice}
{#if (!disabledOptions.includes(choice) && !getValueOfDeprecated(inputSpecsConfiguration[choice])) || oneOf.selected === choice}
<option value={choice}>{labels?.[choice] ?? choice}</option>
{#if oneOf}
<div class="mb-2 text-sm font-semibold">
{capitalize(addWhitespaceBeforeCapitals(key))}&nbsp;
{#if tooltip}
<Tooltip light>{tooltip}</Tooltip>
{/if}
{/each}
</select>
{#if oneOf.selected !== 'none' && oneOf.selected !== 'errorOverlay'}
<div class="mb-4"></div>
{/if}
<div class="flex flex-col gap-4">
{#each Object.keys(inputSpecsConfiguration?.[oneOf.selected] ?? {}) as nestedKey}
{@const config = {
...inputSpecsConfiguration?.[oneOf.selected]?.[nestedKey],
...oneOf.configuration?.[oneOf.selected]?.[nestedKey]
</div>
<select
class="w-full border border-gray-300 rounded-md p-2"
value={oneOf.selected}
onchange={(e) => {
oneOf = { ...oneOf, selected: e?.target?.['value'] }
}}
>
{#each Object.keys(inputSpecsConfiguration ?? {}) as choice}
{#if (!disabledOptions.includes(choice) && !getValueOfDeprecated(inputSpecsConfiguration[choice])) || oneOf.selected === choice}
<option value={choice}>{labels?.[choice] ?? choice}</option>
{/if}
{/each}
</select>
{#if oneOf.selected !== 'none' && oneOf.selected !== 'errorOverlay'}
<div class="mb-4"></div>
{/if}
<div class="flex flex-col gap-4">
{#each Object.keys(inputSpecsConfiguration?.[oneOf.selected] ?? {}) as nestedKey}
{@const config = {
...inputSpecsConfiguration?.[oneOf.selected]?.[nestedKey],
...oneOf.configuration?.[oneOf.selected]?.[nestedKey]
}}
{#if config && oneOf.configuration[oneOf.selected]}
<InputsSpecEditor
{recomputeOnInputChanged}
key={nestedKey}
bind:componentInput={oneOf.configuration[oneOf.selected][nestedKey]}
{id}
{acceptSelf}
userInputEnabled={false}
{shouldCapitalize}
{resourceOnly}
fieldType={config?.['fieldType']}
subFieldType={config?.['subFieldType']}
format={config?.['format']}
selectOptions={config?.['selectOptions']}
placeholder={config?.['placeholder']}
customTitle={config?.['customTitle']}
tooltip={config?.['tooltip']}
fileUpload={config?.['fileUpload']}
loading={config?.['loading']}
documentationLink={config?.['documentationLink']}
allowTypeChange={config?.['allowTypeChange']}
{showOnDemandOnlyToggle}
/>
{/if}
{/each}
</div>
{#if config && oneOf.configuration[oneOf.selected]}
<InputsSpecEditor
{recomputeOnInputChanged}
key={nestedKey}
bind:componentInput={oneOf.configuration[oneOf.selected][nestedKey]}
{id}
{acceptSelf}
userInputEnabled={false}
{shouldCapitalize}
{resourceOnly}
fieldType={config?.['fieldType']}
subFieldType={config?.['subFieldType']}
format={config?.['format']}
selectOptions={config?.['selectOptions']}
placeholder={config?.['placeholder']}
customTitle={config?.['customTitle']}
tooltip={config?.['tooltip']}
fileUpload={config?.['fileUpload']}
loading={config?.['loading']}
documentationLink={config?.['documentationLink']}
allowTypeChange={config?.['allowTypeChange']}
{showOnDemandOnlyToggle}
/>
{/if}
{/each}
</div>
{/if}
</div>
@@ -2,7 +2,16 @@ export const BUTTON_COLORS = ['blue', 'red', 'dark', 'light', 'green', 'gray', '
export namespace ButtonType {
export type Size = 'xs3' | 'xs2' | 'xs' | 'sm' | 'md' | 'lg' | 'xl'
export type Color = string
export type Color =
| 'blue'
| 'red'
| 'dark'
| 'light'
| 'green'
| 'gray'
| 'none'
| 'marine'
| 'nord'
export type Variant = 'contained' | 'border' | 'divider'
export type Target = '_self' | '_blank'
export type Element = HTMLButtonElement | HTMLAnchorElement
@@ -3,9 +3,13 @@
import { Download } from 'lucide-svelte'
import { base } from '$lib/base'
export let s3object: any
export let workspaceId: string | undefined = undefined
export let appPath: string | undefined = undefined
interface Props {
s3object: any
workspaceId?: string | undefined
appPath?: string | undefined
}
let { s3object, workspaceId = undefined, appPath = undefined }: Props = $props()
</script>
<a
@@ -8,6 +8,12 @@
import { twMerge } from 'tailwind-merge'
import { aiChatManager, AIMode } from './chat/AIChatManager.svelte'
interface Props {
moduleId?: string
}
const { moduleId }: Props = $props()
const aiChatScriptModeClasses = $derived(
aiChatManager.mode === AIMode.SCRIPT && aiChatManager.isOpen
? 'dark:bg-violet-900 bg-violet-100'
@@ -22,7 +28,7 @@
btnClasses={twMerge('!px-2', aiChatScriptModeClasses)}
{onClick}
iconOnly
title="Open AI chat in script mode"
title="Open AI chat"
startIcon={{ icon: WandSparkles, classes: 'text-violet-800 dark:text-violet-400' }}
/>
{/snippet}
@@ -30,7 +36,8 @@
{#if $copilotInfo.enabled}
{@render button(() => {
aiChatManager.openChat()
aiChatManager.changeMode(AIMode.SCRIPT)
const availableContext = aiChatManager.contextManager.getAvailableContext()
aiChatManager.contextManager.setSelectedModuleContext(moduleId, availableContext)
})}
{:else}
<Popover
@@ -174,19 +174,17 @@ export class Autocompletor {
additionalTextEdits:
endsWithNewLine && !multiline
? [
{
range: toEol,
text: ''
}
]
{
range: toEol,
text: ''
}
]
: []
}
]
}
},
// @ts-ignore
disposeInlineCompletions: () => {},
freeInlineCompletions: () => {}
disposeInlineCompletions: () => { },
}
)
@@ -77,7 +77,7 @@
})
$effect(() => {
aiChatManager.listenForScriptEditorContextChange(
aiChatManager.listenForContextChange(
$dbSchemas,
$workspaceStore,
$copilotSessionModel
@@ -115,9 +115,7 @@
pastChats={historyManager.getPastChats()}
bind:selectedContext={
() => aiChatManager.contextManager.getSelectedContext(),
(sc) => {
aiChatManager.scriptEditorOptions && aiChatManager.contextManager.setSelectedContext(sc)
}
(sc) => aiChatManager.contextManager.setSelectedContext(sc)
}
availableContext={aiChatManager.contextManager.getAvailableContext()}
messages={aiChatManager.currentReply
@@ -52,7 +52,7 @@
if (placeholder) {
return placeholder
}
switch (aiChatManager.mode) {
case AIMode.SCRIPT:
return 'Modify this script...'
@@ -74,7 +74,7 @@
let instructions = $state(initialInstructions)
export function focusInput() {
if (aiChatManager.mode === AIMode.SCRIPT) {
if (aiChatManager.mode === AIMode.SCRIPT || aiChatManager.mode === AIMode.FLOW) {
contextTextareaComponent?.focus()
} else {
instructionsTextareaComponent?.focus()
@@ -132,7 +132,7 @@
</script>
<div use:clickOutside class="relative">
{#if aiChatManager.mode === AIMode.SCRIPT}
{#if aiChatManager.mode === AIMode.SCRIPT || aiChatManager.mode === AIMode.FLOW}
{#if showContext}
<div class="flex flex-row gap-1 mb-1 overflow-scroll pt-2 no-scrollbar">
<Popover>
@@ -157,7 +157,7 @@
<ContextElementBadge
contextElement={element}
deletable
on:delete={() => {
onDelete={() => {
selectedContext = selectedContext?.filter(
(c) => c.type !== element.type || c.title !== element.title
)
@@ -1,5 +1,5 @@
import type { AIProviderModel, ScriptLang } from '$lib/gen/types.gen'
import type { ScriptOptions } from './ContextManager.svelte'
import type { FlowOptions, ScriptOptions } from './ContextManager.svelte'
import {
flowTools,
prepareFlowSystemMessage,
@@ -40,13 +40,12 @@ import { getStringError } from './utils'
import type { FlowModuleState, FlowState } from '$lib/components/flows/flowState'
import type { CurrentEditor, ExtendedOpenFlow } from '$lib/components/flows/types'
import { untrack } from 'svelte'
import { copilotSessionModel, type DBSchemas } from '$lib/stores'
import { getCurrentModel, type DBSchemas } from '$lib/stores'
import { askTools, prepareAskSystemMessage, prepareAskUserMessage } from './ask/core'
import { chatState, DEFAULT_SIZE, triggerablesByAi } from './sharedChatState.svelte'
import type { ContextElement } from './context'
import type { Selection } from 'monaco-editor'
import type AIChatInput from './AIChatInput.svelte'
import { get } from 'svelte/store'
import { prepareApiSystemMessage, prepareApiUserMessage } from './api/core'
// If the estimated token usage is greater than the model context window - the threshold, we delete the oldest message
@@ -88,6 +87,7 @@ class AIChatManager {
helpers = $state<any | undefined>(undefined)
scriptEditorOptions = $state<ScriptOptions | undefined>(undefined)
flowOptions = $state<FlowOptions | undefined>(undefined)
scriptEditorApplyCode = $state<((code: string, applyAll?: boolean) => void) | undefined>(
undefined
)
@@ -100,7 +100,7 @@ class AIChatManager {
private confirmationCallback = $state<((value: boolean) => void) | undefined>(undefined)
allowedModes: Record<AIMode, boolean> = $derived({
script: this.scriptEditorOptions !== undefined,
script: this.flowAiChatHelpers === undefined && this.scriptEditorOptions !== undefined,
flow: this.flowAiChatHelpers !== undefined,
navigator: true,
ask: true,
@@ -123,11 +123,12 @@ class AIChatManager {
}
return acc
}, 0)
const modelContextWindow = getModelContextWindow(get(copilotSessionModel)?.model ?? '')
const model = getCurrentModel()
const modelContextWindow = getModelContextWindow(model.model)
return (
estimatedTokens >
modelContextWindow -
Math.max(modelContextWindow * MAX_TOKENS_THRESHOLD_PERCENTAGE, MAX_TOKENS_HARD_LIMIT)
Math.max(modelContextWindow * MAX_TOKENS_THRESHOLD_PERCENTAGE, MAX_TOKENS_HARD_LIMIT)
)
}
@@ -557,8 +558,8 @@ class AIChatManager {
onNewToken: (token: string) => {
reply += token
},
onMessageEnd: () => { },
setToolStatus: () => { }
onMessageEnd: () => {},
setToolStatus: () => {}
},
systemMessage
}
@@ -625,7 +626,7 @@ class AIChatManager {
}
try {
const oldSelectedContext = this.contextManager?.getSelectedContext() ?? []
if (this.mode === AIMode.SCRIPT) {
if (this.mode === AIMode.SCRIPT || this.mode === AIMode.FLOW) {
this.contextManager?.updateContextOnRequest(options)
}
this.loading = true
@@ -648,7 +649,10 @@ class AIChatManager {
{
role: 'user',
content: this.instructions,
contextElements: this.mode === AIMode.SCRIPT ? oldSelectedContext : undefined,
contextElements:
this.mode === AIMode.SCRIPT || this.mode === AIMode.FLOW
? oldSelectedContext
: undefined,
snapshot,
index: this.messages.length // matching with actual messages index. not -1 because it's not yet added to the messages array
}
@@ -672,7 +676,8 @@ class AIChatManager {
case AIMode.FLOW:
userMessage = prepareFlowUserMessage(
oldInstructions,
this.flowAiChatHelpers!.getFlowAndSelectedId()
this.flowAiChatHelpers!.getFlowAndSelectedId(),
oldSelectedContext
)
break
case AIMode.NAVIGATOR:
@@ -823,12 +828,19 @@ class AIChatManager {
this.sendRequest()
}
addSelectedLinesToContext = (lines: string, startLine: number, endLine: number) => {
addSelectedLinesToContext = (
lines: string,
startLine: number,
endLine: number,
moduleId?: string
) => {
if (!this.open) {
this.toggleOpen()
}
this.changeMode(AIMode.SCRIPT)
this.contextManager?.addSelectedLinesToContext(lines, startLine, endLine)
if (!moduleId) {
this.changeMode(AIMode.SCRIPT)
}
this.contextManager?.addSelectedLinesToContext(lines, startLine, endLine, moduleId)
this.focusInput()
}
@@ -869,12 +881,12 @@ class AIChatManager {
})
}
listenForScriptEditorContextChange = (
listenForContextChange = (
dbSchemas: DBSchemas,
workspaceStore: string | undefined,
copilotSessionModel: AIProviderModel | undefined
) => {
if (this.scriptEditorOptions) {
if (this.mode === AIMode.SCRIPT && this.scriptEditorOptions) {
this.contextManager.updateAvailableContext(
this.scriptEditorOptions,
dbSchemas,
@@ -882,6 +894,18 @@ class AIChatManager {
!copilotSessionModel?.model.endsWith('/thinking'),
untrack(() => this.contextManager.getSelectedContext())
)
} else if (this.mode === AIMode.FLOW && this.flowOptions) {
this.contextManager.updateAvailableContextForFlow(
this.flowOptions,
dbSchemas,
workspaceStore ?? '',
!copilotSessionModel?.model.endsWith('/thinking'),
untrack(() => this.contextManager.getSelectedContext())
)
}
if (this.scriptEditorOptions) {
this.contextManager.setScriptOptions(this.scriptEditorOptions)
}
}
@@ -941,15 +965,15 @@ class AIChatManager {
const editorRelated =
currentEditor && currentEditor.type === 'script' && currentEditor.stepId === module.id
? {
diffMode: currentEditor.diffMode,
lastDeployedCode: currentEditor.lastDeployedCode,
lastSavedCode: undefined
}
diffMode: currentEditor.diffMode,
lastDeployedCode: currentEditor.lastDeployedCode,
lastSavedCode: undefined
}
: {
diffMode: false,
lastDeployedCode: undefined,
lastSavedCode: undefined
}
diffMode: false,
lastDeployedCode: undefined,
lastSavedCode: undefined
}
return {
args: moduleState?.previewArgs ?? {},
@@ -976,6 +1000,13 @@ class AIChatManager {
this.scriptEditorOptions = undefined
}
untrack(() =>
this.contextManager?.setSelectedModuleContext(
selectedId,
untrack(() => this.contextManager.getAvailableContext())
)
)
return () => {
this.scriptEditorOptions = undefined
}
@@ -1,65 +1,285 @@
<script lang="ts">
import FlowModuleIcon from '$lib/components/flows/FlowModuleIcon.svelte'
import BarsStaggered from '$lib/components/icons/BarsStaggered.svelte'
import type { FlowModule } from '$lib/gen/types.gen'
import { ContextIconMap, type ContextElement } from './context'
import { ArrowLeft, Diff, Database, ChevronRight } from 'lucide-svelte'
interface Props {
availableContext: ContextElement[]
selectedContext: ContextElement[]
onSelect: (element: ContextElement) => void
setShowing?: (showing: boolean) => void
showAllAvailable?: boolean
stringSearch?: string
selectedIndex?: number
onViewChange?: (newNumber: number) => void
}
const {
availableContext,
selectedContext,
onSelect,
setShowing,
showAllAvailable = false,
stringSearch = '',
selectedIndex = 0
onViewChange
}: Props = $props()
// Define priority map for context types
const typePriority = {
code: 1,
diff: 2,
default: 3
// Current view state: 'categories' or specific category type
let currentView = $state<'categories' | 'diffs' | 'modules' | 'databases'>('categories')
// Selected index for keyboard navigation
let itemSelectedIndex = $state(0)
let categorySelectedIndex = $state(0)
// Category definitions
const categories = [
{ id: 'diffs', label: 'Diffs', icon: Diff },
{ id: 'modules', label: 'Modules', icon: BarsStaggered },
{ id: 'databases', label: 'Databases', icon: Database }
]
const filteredAvailableContext = $derived(
availableContext.filter((context) => {
const filtered =
(showAllAvailable ||
!selectedContext.some((sc) => sc.type === context.type && sc.title === context.title)) &&
(!stringSearch || context.title.toLowerCase().includes(stringSearch.toLowerCase()))
return filtered
})
)
// Group context by category
const contextByCategory = $derived.by(() => {
const grouped: Record<string, ContextElement[]> = {
diffs: [],
modules: [],
databases: []
}
filteredAvailableContext.forEach((context) => {
if (context.type === 'diff') grouped.diffs.push(context)
else if (context.type === 'flow_module') grouped.modules.push(context)
else if (context.type === 'db') grouped.databases.push(context)
})
return grouped
})
const currentCategoryItems = $derived(
currentView !== 'categories' ? contextByCategory[currentView] : []
)
// Filter to only show categories with items
const availableCategories = $derived(
categories.filter((cat) => contextByCategory[cat.id].length > 0)
)
// Report view changes
$effect(() => {
if (onViewChange) {
if (currentView === 'categories') {
onViewChange(availableCategories.length)
} else {
onViewChange(currentCategoryItems.length + 1)
}
}
})
function handleCategoryClick(categoryId: string) {
currentView = categoryId as typeof currentView
}
const actualAvailableContext = $derived(
availableContext
.filter(
(c) =>
(showAllAvailable ||
!selectedContext.some((sc) => sc.type === c.type && sc.title === c.title)) &&
(!stringSearch || c.title.toLowerCase().includes(stringSearch.toLowerCase()))
)
.sort((a, b) => {
const priorityA = typePriority[a.type] || typePriority.default
const priorityB = typePriority[b.type] || typePriority.default
return priorityA - priorityB
})
)
function handleBackClick() {
currentView = 'categories'
itemSelectedIndex = 0
}
function handleKeyDown(e: KeyboardEvent) {
if (stringSearch.length > 0) {
// Navigation in search view (flat list)
if (e.key === 'ArrowDown') {
e.preventDefault()
e.stopPropagation()
if (filteredAvailableContext.length > 0) {
itemSelectedIndex = (itemSelectedIndex + 1) % filteredAvailableContext.length
}
} else if (e.key === 'ArrowUp') {
e.preventDefault()
e.stopPropagation()
if (filteredAvailableContext.length > 0) {
itemSelectedIndex =
(itemSelectedIndex - 1 + filteredAvailableContext.length) %
filteredAvailableContext.length
}
} else if (e.key === 'Enter' || e.key === 'Tab') {
if (e.key === 'Tab') e.preventDefault()
e.stopPropagation()
const selectedItem = filteredAvailableContext[itemSelectedIndex]
if (selectedItem) {
onSelect(selectedItem)
}
}
} else if (currentView === 'categories') {
// Navigation in categories view
if (e.key === 'ArrowDown') {
e.preventDefault()
e.stopPropagation()
categorySelectedIndex = (categorySelectedIndex + 1) % availableCategories.length
} else if (e.key === 'ArrowUp') {
e.preventDefault()
e.stopPropagation()
categorySelectedIndex =
(categorySelectedIndex - 1 + availableCategories.length) % availableCategories.length
} else if (e.key === 'Enter' || e.key === 'ArrowRight' || e.key === 'Tab') {
e.preventDefault()
e.stopPropagation()
const selectedCategory = availableCategories[categorySelectedIndex]
if (selectedCategory) {
handleCategoryClick(selectedCategory.id)
}
} else if (e.key === 'Escape' || e.key === 'ArrowLeft') {
e.preventDefault()
e.stopPropagation()
setShowing?.(false)
}
} else {
// Navigation in category items view
if (e.key === 'ArrowDown') {
e.preventDefault()
e.stopPropagation()
if (currentCategoryItems.length > 0) {
itemSelectedIndex = (itemSelectedIndex + 1) % currentCategoryItems.length
}
} else if (e.key === 'ArrowUp') {
e.preventDefault()
e.stopPropagation()
if (currentCategoryItems.length > 0) {
itemSelectedIndex =
(itemSelectedIndex - 1 + currentCategoryItems.length) % currentCategoryItems.length
}
} else if (e.key === 'Enter' || e.key === 'Tab') {
if (e.key === 'Tab') e.preventDefault()
e.stopPropagation()
const selectedItem = currentCategoryItems[itemSelectedIndex]
if (selectedItem) {
onSelect(selectedItem)
currentView = 'categories' // Go back to categories after selection
}
} else if (e.key === 'ArrowLeft' || e.key === 'Escape') {
e.preventDefault()
e.stopPropagation()
handleBackClick()
}
}
}
// Listen for keyboard events
$effect(() => {
document.addEventListener('keydown', handleKeyDown)
return () => {
document.removeEventListener('keydown', handleKeyDown)
}
})
$effect(() => {
if (stringSearch.length > 0) {
itemSelectedIndex = 0
}
})
</script>
<div class="flex flex-col gap-1 text-tertiary text-xs p-1 min-w-24 max-h-48 overflow-y-scroll">
{#if actualAvailableContext.length === 0}
<div class="text-center text-tertiary text-xs">No available context</div>
{:else}
{#each actualAvailableContext as element, i}
<div
class="flex flex-col gap-1 text-tertiary text-xs p-1 pr-0 min-w-24 max-h-48 overflow-y-scroll"
onmousedown={(e) =>
// avoids triggering onblur on the textinput and closing the tooltip
e.preventDefault()}
role="listbox"
tabindex={0}
>
{#if stringSearch.length > 0}
<!-- Search view - show flat list -->
{#each filteredAvailableContext as element, i}
{@const Icon = ContextIconMap[element.type]}
<button
class="hover:bg-surface-hover rounded-md p-1 text-left flex flex-row gap-1 items-center font-normal {i ===
selectedIndex
class="hover:bg-surface-hover rounded-md p-1 text-left flex flex-row gap-1 items-center font-normal transition-colors {i ===
itemSelectedIndex
? 'bg-surface-hover'
: ''}"
onclick={() => onSelect(element)}
onclick={() => {
onSelect(element)
}}
>
{#if Icon}
{#if element.type === 'flow_module'}
<FlowModuleIcon module={element as FlowModule} size={16} />
{:else if Icon}
<Icon size={16} />
{/if}
{element.type === 'diff' ? element.title.replace(/_/g, ' ') : element.title}
<span class="truncate">
{element.type === 'diff' || element.type === 'flow_module'
? element.title.replace(/_/g, ' ')
: element.title}
</span>
</button>
{/each}
{#if filteredAvailableContext.length === 0}
<div class="text-center text-tertiary text-xs py-2">No matching context</div>
{/if}
{:else if currentView === 'categories'}
<!-- Categories view -->
{#each availableCategories as category, i}
{@const Icon = category.icon}
<button
class="hover:bg-surface-hover rounded-md p-1 pr-0 text-left flex flex-row gap-1 items-center font-normal transition-colors {i ===
categorySelectedIndex
? 'bg-surface-hover'
: ''}"
onclick={() => handleCategoryClick(category.id)}
>
<Icon size={16} />
<span class="flex-1">{category.label}</span>
<ChevronRight size={16} />
</button>
{/each}
{#if availableCategories.length === 0}
<div class="text-center text-tertiary text-xs py-2">No available context</div>
{/if}
{:else}
<!-- Category items view -->
<button
class="hover:bg-surface-hover rounded-md text-left flex flex-row gap-1 items-center font-normal transition-colors mb-1"
onclick={handleBackClick}
>
<ArrowLeft size={12} />
<span class="text-xs">Go back</span>
</button>
{#if currentCategoryItems.length === 0}
<div class="text-center text-tertiary text-xs py-2">No items in this category</div>
{:else}
{#each currentCategoryItems as element, i}
{@const Icon = ContextIconMap[element.type]}
<button
class="hover:bg-surface-hover rounded-md p-1 text-left flex flex-row gap-1 items-center font-normal transition-colors {i ===
itemSelectedIndex
? 'bg-surface-hover'
: ''}"
onclick={() => {
onSelect(element)
currentView = 'categories' // Go back to categories after selection
}}
>
{#if element.type === 'flow_module'}
<FlowModuleIcon module={element as FlowModule} size={16} />
{:else if Icon}
<Icon size={16} />
{/if}
<span class="truncate">
{element.type === 'diff' ? element.title.replace(/_/g, ' ') : element.title}
</span>
</button>
{/each}
{/if}
{/if}
</div>
@@ -10,36 +10,43 @@
formatSchema
} from '$lib/components/apps/components/display/dbtable/utils'
import ObjectViewer from '$lib/components/propertyPicker/ObjectViewer.svelte'
import { createEventDispatcher } from 'svelte'
import HighlightCode from '$lib/components/HighlightCode.svelte'
import FlowModuleIcon from '$lib/components/flows/FlowModuleIcon.svelte'
import type { FlowModule } from '$lib/gen'
export let contextElement: ContextElement
export let deletable = false
interface Props {
contextElement: ContextElement
deletable?: boolean
onDelete?: () => void
}
let { contextElement, deletable = false, onDelete }: Props = $props()
const icon = ContextIconMap[contextElement.type]
let showDelete = false
let showDelete = $state(false)
const dispatch = createEventDispatcher<{
delete: void
}>()
const isDeletable = $derived(deletable && contextElement.deletable !== false)
</script>
<Popover>
<svelte:fragment slot="trigger">
{#snippet trigger()}
<div
class={twMerge(
'border rounded-md px-1 py-0.5 flex flex-row items-center gap-1 text-tertiary text-xs cursor-default hover:bg-surface-hover hover:cursor-pointer max-w-48 bg-surface'
)}
on:mouseenter={() => (showDelete = true)}
on:mouseleave={() => (showDelete = false)}
onmouseenter={() => (showDelete = true)}
onmouseleave={() => (showDelete = false)}
aria-label="Context element"
role="button"
tabindex={0}
>
<button on:click={() => dispatch('delete')} class:cursor-default={!deletable}>
{#if showDelete && deletable}
<button onclick={isDeletable ? onDelete : undefined} class:cursor-default={!isDeletable}>
{#if showDelete && isDeletable}
<X size={16} />
{:else if contextElement.type === 'flow_module' || contextElement.type === 'flow_module_code_piece'}
<FlowModuleIcon module={contextElement as FlowModule} size={16} />
{:else}
<svelte:component this={icon} size={16} />
{@const SvelteComponent = icon}
<SvelteComponent size={16} />
{/if}
</button>
<span class="truncate">
@@ -48,8 +55,8 @@
: contextElement.title}
</span>
</div>
</svelte:fragment>
<svelte:fragment slot="content">
{/snippet}
{#snippet content()}
{#if contextElement.type === 'error'}
<div class="max-w-96 max-h-[300px] text-xs overflow-auto">
<Highlight language={json} code={contextElement.content} class="w-full p-2" />
@@ -71,7 +78,7 @@
<div class="text-tertiary">Not loaded yet</div>
{/if}
</div>
{:else if contextElement.type === 'code' || contextElement.type === 'code_piece' || contextElement.type === 'diff'}
{:else if contextElement.type === 'code' || contextElement.type === 'code_piece' || contextElement.type === 'diff' || contextElement.type === 'flow_module_code_piece'}
<div class="max-w-96 max-h-[300px] text-xs overflow-auto">
<HighlightCode
language={contextElement.lang}
@@ -79,6 +86,20 @@
class="w-full p-2 "
/>
</div>
{:else if contextElement.type === 'flow_module'}
{#if contextElement.value.content}
<div class="p-2 max-w-96 max-h-[300px] text-xs overflow-auto">
<HighlightCode
language={contextElement.value.language}
code={contextElement.value.content}
class="w-full p-2 "
/>
</div>
{:else}
<div class="p-2 max-w-96 max-h-[300px] text-xs overflow-auto">
<div class="text-tertiary">{contextElement.title}</div>
</div>
{/if}
{/if}
</svelte:fragment>
{/snippet}
</Popover>
@@ -1,11 +1,13 @@
import { ResourceService, type ListResourceResponse, type ScriptLang } from '$lib/gen'
import { ResourceService, type Flow, type ListResourceResponse, type ScriptLang } from '$lib/gen'
import { scriptLangToEditorLang } from '$lib/scripts'
import { SQLSchemaLanguages, type DBSchemas } from '$lib/stores'
import { diffLines } from 'diff'
import type { ContextElement } from './context'
import type { ContextElement, FlowModuleElement } from './context'
import type { FlowModule } from '$lib/gen'
import type { DisplayMessage } from './shared'
import { langToExt } from '$lib/editorLangUtils'
import type { ExtendedOpenFlow } from '$lib/components/flows/types'
export interface ScriptOptions {
lang: ScriptLang | 'bunnative'
@@ -18,6 +20,14 @@ export interface ScriptOptions {
diffMode: boolean
}
export interface FlowOptions {
currentFlow: ExtendedOpenFlow
lastDeployedFlow?: Flow
path: string | undefined
modules: FlowModule[]
lastSavedFlow?: Flow
}
export default class ContextManager {
private selectedContext: ContextElement[] = $state([])
private availableContext: ContextElement[] = $state([])
@@ -55,6 +65,93 @@ export default class ContextManager {
)
}
async updateAvailableContextForFlow(
flowOptions: FlowOptions,
dbSchemas: DBSchemas,
workspace: string,
toolSupport: boolean,
currentlySelectedContext: ContextElement[]
) {
try {
if (this.workspace !== workspace) {
await this.refreshDbResources(workspace)
this.workspace = workspace
}
let newAvailableContext: ContextElement[] = []
// Add diff context if we have a deployed flow version
const deployedFlowString = JSON.stringify(flowOptions.lastDeployedFlow, null, 2)
const savedFlowString = JSON.stringify(flowOptions.lastSavedFlow, null, 2)
const currentFlowString = JSON.stringify(flowOptions.currentFlow, null, 2)
if (currentFlowString && deployedFlowString && deployedFlowString !== currentFlowString) {
newAvailableContext.push({
type: 'diff',
title: 'diff_with_last_deployed_version',
content: deployedFlowString,
diff: diffLines(deployedFlowString, currentFlowString),
lang: 'graphql' // irrelevant, but needed for the diff component
})
}
if (currentFlowString && savedFlowString && savedFlowString !== currentFlowString) {
newAvailableContext.push({
type: 'diff',
title: 'diff_with_last_saved_draft',
content: savedFlowString,
diff: diffLines(savedFlowString, currentFlowString),
lang: 'graphql' // irrelevant, but needed for the diff component
})
}
for (const module of flowOptions.modules) {
newAvailableContext.push({
type: 'flow_module',
id: module.id,
title: `${module.id}`,
value: {
language: 'language' in module.value ? module.value.language : 'bunnative',
path: 'path' in module.value ? module.value.path : '',
content: 'content' in module.value ? module.value.content : '',
type: module.value.type
}
})
}
if (toolSupport) {
for (const d of this.dbResources) {
const loadedSchema = dbSchemas[d.path]
newAvailableContext.push({
type: 'db',
title: d.path,
// If the db is already fetched, add the schema to the context
...(loadedSchema ? { schema: loadedSchema } : {})
})
}
}
let newSelectedContext: ContextElement[] = [...currentlySelectedContext]
// Filter selected context to only include available items
newSelectedContext = newSelectedContext
.filter((c) => newAvailableContext.some((ac) => ac.type === c.type && ac.title === c.title))
.map((c) =>
c.type === 'db' && dbSchemas[c.title]
? {
...c,
schema: dbSchemas[c.title]
}
: c
)
this.availableContext = newAvailableContext
this.selectedContext = newSelectedContext
} catch (err) {
console.error('Could not update available context for flow', err)
}
}
async updateAvailableContext(
scriptOptions: ScriptOptions,
dbSchemas: DBSchemas,
@@ -63,12 +160,10 @@ export default class ContextManager {
currentlySelectedContext: ContextElement[]
) {
try {
let firstTime = !this.workspace
if (this.workspace !== workspace) {
await this.refreshDbResources(workspace)
this.workspace = workspace
}
this.scriptOptions = scriptOptions
let newAvailableContext: ContextElement[] = [
{
type: 'code',
@@ -123,16 +218,15 @@ export default class ContextManager {
let newSelectedContext: ContextElement[] = [...currentlySelectedContext]
if (firstTime) {
newSelectedContext = [
{
type: 'code',
title: this.getContextCodePath(scriptOptions) ?? '',
content: scriptOptions.code,
lang: scriptOptions.lang
}
]
}
newSelectedContext = [
{
type: 'code',
title: this.getContextCodePath(scriptOptions) ?? '',
content: scriptOptions.code,
lang: scriptOptions.lang,
deletable: false
}
]
const db = this.getSelectedDBSchema(scriptOptions, dbSchemas)
if (
@@ -160,15 +254,15 @@ export default class ContextManager {
.map((c) =>
c.type === 'code'
? {
...c,
content: scriptOptions.code,
title: this.getContextCodePath(scriptOptions)
}
...c,
content: scriptOptions.code,
title: this.getContextCodePath(scriptOptions)
}
: c.type === 'db' && dbSchemas[c.title]
? {
...c,
schema: dbSchemas[c.title]
}
...c,
schema: dbSchemas[c.title]
}
: c
)
@@ -191,26 +285,56 @@ export default class ContextManager {
return this.availableContext
}
addSelectedLinesToContext(lines: string, startLine: number, endLine: number) {
setScriptOptions(scriptOptions: ScriptOptions) {
this.scriptOptions = scriptOptions
}
addSelectedLinesToContext(lines: string, startLine: number, endLine: number, moduleId?: string) {
const title = moduleId ? `${moduleId} L${startLine}-L${endLine}` : `L${startLine}-L${endLine}`
if (
!this.scriptOptions ||
this.selectedContext.find(
(c) => c.type === 'code_piece' && c.title === `L${startLine}-L${endLine}`
(c) =>
(c.type === 'code_piece' && c.title === title) ||
(c.type === 'flow_module_code_piece' && c.id === moduleId && c.title === title)
)
) {
return
}
this.selectedContext = [
...this.selectedContext,
{
type: 'code_piece',
title: `L${startLine}-L${endLine}`,
startLine,
endLine,
content: lines,
lang: this.scriptOptions.lang
if (moduleId) {
const module = [...this.availableContext, ...this.selectedContext].find(
(c) => c.type === 'flow_module' && c.id === moduleId
) as FlowModuleElement
if (!module) {
console.error('Module not found', moduleId)
return
}
]
this.selectedContext = [
...this.selectedContext,
{
type: 'flow_module_code_piece',
id: moduleId,
title: title,
startLine,
endLine,
content: lines,
lang: this.scriptOptions.lang,
value: module.value
}
]
} else {
this.selectedContext = [
...this.selectedContext,
{
type: 'code_piece',
title: title,
startLine,
endLine,
content: lines,
lang: this.scriptOptions.lang
}
]
}
}
setFixContext() {
@@ -234,14 +358,14 @@ export default class ContextManager {
...(options.withCode === false ? [] : [codeContext]),
...(options.withDiff
? [
{
type: 'diff' as const,
title: 'diff_with_last_deployed_version',
content: this.scriptOptions.lastDeployedCode ?? '',
diff: diffLines(this.scriptOptions.lastDeployedCode ?? '', this.scriptOptions.code),
lang: this.scriptOptions.lang
}
]
{
type: 'diff' as const,
title: 'diff_with_last_deployed_version',
content: this.scriptOptions.lastDeployedCode ?? '',
diff: diffLines(this.scriptOptions.lastDeployedCode ?? '', this.scriptOptions.code),
lang: this.scriptOptions.lang
}
]
: [])
]
}
@@ -268,15 +392,37 @@ export default class ContextManager {
contextElements:
m.role !== 'tool' && m.contextElements
? m.contextElements.map((c) =>
c.type === 'db'
? {
type: 'db',
title: c.title,
schema: dbSchemas[c.title]
}
: c
)
c.type === 'db'
? {
type: 'db',
title: c.title,
schema: dbSchemas[c.title]
}
: c
)
: undefined
}))
}
setSelectedModuleContext(
moduleId: string | undefined,
availableContext: ContextElement[] | undefined
) {
if (availableContext && moduleId) {
const module = availableContext.find((c) => c.type === 'flow_module' && c.id === moduleId)
if (
module &&
!this.selectedContext.find((c) => c.type === 'flow_module' && c.id === moduleId)
) {
this.selectedContext = this.selectedContext.filter((c) => c.type !== 'flow_module')
this.selectedContext = [module, ...this.selectedContext]
}
} else if (!moduleId) {
this.selectedContext = this.selectedContext.filter((c) => c.type !== 'flow_module')
}
}
clearContext() {
this.selectedContext = []
}
}
@@ -38,7 +38,7 @@
let tooltipPosition = $state({ x: 0, y: 0 })
let textarea = $state<HTMLTextAreaElement | undefined>(undefined)
let tooltipElement = $state<HTMLDivElement | undefined>(undefined)
let selectedSuggestionIndex = $state(0)
let tooltipCurrentViewNumber = $state(0)
// Properties to copy for caret position calculation
const properties = [
@@ -155,7 +155,7 @@
}
function getHighlightedText(text: string) {
return text.replace(/@[\w/.-]+/g, (match) => {
return text.replace(/@[\w/.\-\[\]]+/g, (match) => {
const contextElement = availableContext.find((c) => c.title === match.slice(1))
if (contextElement) {
return `<span class="bg-black dark:bg-white text-white dark:text-black z-10">${match}</span>`
@@ -182,27 +182,19 @@
showContextTooltip = false
}
async function updateTooltipPosition(
availableContext: ContextElement[],
showContextTooltip: boolean,
contextTooltipWord: string
) {
if (!textarea || !showContextTooltip) return
async function updateTooltipPosition(currentViewItemsNumber: number) {
if (!textarea) return
try {
const coords = getCaretCoordinates(textarea, textarea.selectionEnd)
const rect = textarea.getBoundingClientRect()
const filteredAvailableContext = availableContext.filter(
(c) => !contextTooltipWord || c.title.toLowerCase().includes(contextTooltipWord.slice(1))
)
const itemHeight = 28 // Estimated height of one item + gap (Button: p-1(8px) + text-xs(16px) = 24px; Parent: gap-1(4px) = 28px)
const containerPadding = 8 // p-1 top + p-1 bottom = 4px + 4px = 8px
const maxHeight = 192 + containerPadding // max-h-48 (192px) + containerPadding (8px)
// Calculate uncapped height, subtract gap from last item as it's not needed
const numItems = filteredAvailableContext.length
const numItems = currentViewItemsNumber
let uncappedHeight =
numItems > 0 ? numItems * itemHeight - 4 + containerPadding : containerPadding
// Ensure height is at least containerPadding even if no items
@@ -270,68 +262,33 @@
} else {
showContextTooltip = false
contextTooltipWord = ''
selectedSuggestionIndex = 0
}
}
function handleKeyPress(e: KeyboardEvent) {
if (e.key === 'Enter' && !e.shiftKey) {
e.preventDefault()
if (contextTooltipWord) {
const filteredContext = availableContext.filter(
(c) => !contextTooltipWord || c.title.toLowerCase().includes(contextTooltipWord.slice(1))
)
const contextElement = filteredContext[selectedSuggestionIndex]
if (contextElement) {
const isInSelectedContext = selectedContext.find(
(c) => c.title === contextElement.title && c.type === contextElement.type
)
// If the context element is already in the selected context and the last word in the instructions is the same as the context element title, send request
if (isInSelectedContext && value.split(' ').pop() === '@' + contextElement.title) {
onSendRequest()
return
}
handleContextSelection(contextElement)
} else if (contextTooltipWord === '@' && availableContext.length > 0) {
handleContextSelection(availableContext[0])
}
} else {
onSendRequest()
}
}
}
function handleKeyDown(e: KeyboardEvent) {
// Pass to parent first if provided
if (onKeyDown) {
onKeyDown(e)
}
if (!showContextTooltip) return
const filteredContext = availableContext.filter(
(c) => !contextTooltipWord || c.title.toLowerCase().includes(contextTooltipWord.slice(1))
)
if (e.key === 'Tab') {
e.preventDefault()
const contextElement = filteredContext[selectedSuggestionIndex]
if (contextElement) {
handleContextSelection(contextElement)
if (showContextTooltip) {
// avoid new line after Enter in the tooltip
if (e.key === 'Enter') {
e.preventDefault()
}
return
}
if (e.key === 'ArrowDown') {
if (e.key === 'Enter' && !e.shiftKey) {
e.preventDefault()
selectedSuggestionIndex = (selectedSuggestionIndex + 1) % filteredContext.length
} else if (e.key === 'ArrowUp') {
e.preventDefault()
selectedSuggestionIndex =
(selectedSuggestionIndex - 1 + filteredContext.length) % filteredContext.length
onSendRequest()
}
}
$effect(() => {
updateTooltipPosition(availableContext, showContextTooltip, contextTooltipWord)
if (showContextTooltip) {
updateTooltipPosition(tooltipCurrentViewNumber)
}
})
export function focus() {
@@ -352,7 +309,6 @@
</div>
<textarea
bind:this={textarea}
onkeypress={handleKeyPress}
onkeydown={handleKeyDown}
bind:value
use:autosize
@@ -388,7 +344,12 @@
}}
showAllAvailable={true}
stringSearch={contextTooltipWord.slice(1)}
selectedIndex={selectedSuggestionIndex}
onViewChange={(newNumber) => {
tooltipCurrentViewNumber = newNumber
}}
setShowing={(showing) => {
showContextTooltip = showing
}}
/>
</div>
</Portal>
@@ -25,6 +25,14 @@
return obj
}
}
for (const key in obj) {
try {
const parsed = JSON.parse(obj[key])
obj[key] = parsed
} catch (e) {
console.error('Failed to parse JSON:', e)
}
}
return JSON.stringify(obj, null, 2)
} catch {
return String(obj)
@@ -81,7 +89,9 @@
<div
class="bg-surface-secondary border border-gray-200 dark:border-gray-700 rounded p-3 overflow-x-auto max-h-64 overflow-y-auto"
>
<pre class="text-2xs text-primary whitespace-pre-wrap">{formatJson(content)}</pre>
<pre class="text-2xs text-primary whitespace-pre-wrap"
>{formatJson($state.snapshot(content))}</pre
>
</div>
{:else}
<div
@@ -25,8 +25,8 @@
<!-- Collapsible Header -->
<button
class={twMerge(
"w-full p-3 bg-surface-secondary hover:bg-surface-hover transition-colors flex items-center justify-between text-left border-b border-gray-200 dark:border-gray-700",
message.needsConfirmation ? "opacity-80" : ""
'w-full p-3 bg-surface-secondary hover:bg-surface-hover transition-colors flex items-center justify-between text-left border-b border-gray-200 dark:border-gray-700',
message.needsConfirmation ? 'opacity-80' : ''
)}
onclick={() => (isExpanded = !isExpanded)}
disabled={!message.showDetails}
@@ -57,7 +57,7 @@
{#if isExpanded}
<div class="p-3 bg-surface space-y-3">
<!-- Parameters Section -->
<div class={message.needsConfirmation ? "opacity-80" : ""}>
<div class={message.needsConfirmation ? 'opacity-80' : ''}>
<ToolContentDisplay title="Parameters" content={message.parameters} />
</div>
@@ -9,6 +9,7 @@ export const ContextIconMap = {
db: Database,
diff: Diff,
code_piece: Code
// flow_module type is handled with FlowModuleIcon
}
export interface CodeElement {
@@ -47,4 +48,33 @@ export interface CodePieceElement {
lang: ScriptLang | 'bunnative'
}
export type ContextElement = CodeElement | ErrorElement | DBElement | DiffElement | CodePieceElement
export interface FlowModuleElement {
type: 'flow_module'
id: string
title: string
// mimics the FlowModule type, with only the fields we need
value: {
language?: ScriptLang | 'bunnative'
path?: string
content?: string
type: string
}
}
export interface FlowModuleCodePieceElement extends Omit<CodePieceElement, 'type'> {
type: 'flow_module_code_piece'
id: string
value: FlowModuleElement['value']
}
export type ContextElement = (
| CodeElement
| ErrorElement
| DBElement
| DiffElement
| CodePieceElement
| FlowModuleElement
| FlowModuleCodePieceElement
) & {
deletable?: boolean
}
@@ -6,7 +6,7 @@
import { dfs } from '$lib/components/flows/previousResults'
import { dfs as dfsApply } from '$lib/components/flows/dfs'
import { getSubModules } from '$lib/components/flows/flowExplorer'
import type { FlowModule, OpenFlow } from '$lib/gen'
import type { FlowModule, OpenFlow, RawScript } from '$lib/gen'
import { getIndexInNestedModules, getNestedModules } from './utils'
import type { AIModuleAction, FlowAIChatHelpers } from './core'
import {
@@ -174,6 +174,17 @@
if (!newModule) {
throw new Error('Module not found')
}
// Apply the old code to the editor and hide diff editor if the reverted module is a rawscript
if (
newModule.value.type === 'rawscript' &&
$currentEditor?.type === 'script' &&
$currentEditor.stepId === id
) {
$currentEditor.editor.setCode((oldModule.value as RawScript).content)
$currentEditor.hideDiffMode()
}
newModule.value = oldModule.value
}
@@ -525,6 +536,47 @@
return cleanup
})
// Automatically show diff mode when selecting a rawscript module with pending changes
$effect(() => {
if (
$currentEditor?.type === 'script' &&
$selectedId &&
affectedModules[$selectedId] &&
lastSnapshot
) {
const moduleLastSnapshot = getModule($selectedId, lastSnapshot)
const currentModule = getModule($selectedId)
if (
moduleLastSnapshot &&
currentModule &&
currentModule.value.type === 'rawscript' &&
moduleLastSnapshot.value.type === 'rawscript'
) {
// Show diff mode automatically
$currentEditor.setDiffOriginal?.(moduleLastSnapshot.value.content ?? '')
$currentEditor.showDiffMode()
$currentEditor.setDiffButtons?.([
{
text: 'Accept Changes',
color: 'green',
onClick: () => {
flowHelpers.acceptModuleAction($selectedId)
$currentEditor?.hideDiffMode()
}
},
{
text: 'Reject Changes',
onClick: () => {
flowHelpers.revertModuleAction($selectedId)
$currentEditor?.hideDiffMode()
}
}
])
}
}
})
let diffDrawer: DiffDrawer | undefined = $state(undefined)
</script>
@@ -10,9 +10,21 @@ import { emptySchema, emptyString } from '$lib/utils'
import {
getFormattedResourceTypes,
getLangContext,
SUPPORTED_CHAT_SCRIPT_LANGUAGES
SUPPORTED_CHAT_SCRIPT_LANGUAGES,
createDbSchemaTool
} from '../script/core'
import { createSearchHubScriptsTool, createToolDef, type Tool, executeTestRun, buildSchemaForTool, buildTestRunArgs } from '../shared'
import {
createSearchHubScriptsTool,
createToolDef,
type Tool,
executeTestRun,
buildSchemaForTool,
buildTestRunArgs,
buildContextString,
applyCodePiecesToFlowModules,
findModuleById
} from '../shared'
import type { ContextElement } from '../context'
import type { ExtendedOpenFlow } from '$lib/components/flows/types'
export type AIModuleAction = 'added' | 'modified' | 'removed'
@@ -339,8 +351,11 @@ const getInstructionsForCodeGenerationToolDef = createToolDef(
// Will be overridden by setSchema
const testRunFlowSchema = z.object({
args: z.object({}).nullable().optional()
.describe('Arguments to pass to the flow (optional, uses default flow inputs if not provided)')
args: z
.object({})
.nullable()
.optional()
.describe('Arguments to pass to the flow (optional, uses default flow inputs if not provided)')
})
const testRunFlowToolDef = createToolDef(
@@ -368,6 +383,7 @@ const workspaceScriptsSearch = new WorkspaceScriptsSearch()
export const flowTools: Tool<FlowAIChatHelpers>[] = [
createSearchHubScriptsTool(false),
createDbSchemaTool<FlowAIChatHelpers>(),
{
def: searchScriptsToolDef,
fn: async ({ args, workspace, toolId, toolCallbacks }) => {
@@ -562,7 +578,7 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
},
{
def: testRunFlowToolDef,
fn: async function({ args, workspace, helpers, toolCallbacks, toolId }) {
fn: async function ({ args, workspace, helpers, toolCallbacks, toolId }) {
const { flow } = helpers.getFlowAndSelectedId()
if (!flow || !flow.value) {
@@ -577,13 +593,14 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
const parsedArgs = await buildTestRunArgs(args, this.def)
return executeTestRun({
jobStarter: () => JobService.runFlowPreview({
workspace: workspace,
requestBody: {
args: parsedArgs,
value: flow.value,
}
}),
jobStarter: () =>
JobService.runFlowPreview({
workspace: workspace,
requestBody: {
args: parsedArgs,
value: flow.value
}
}),
workspace,
toolCallbacks,
toolId,
@@ -591,7 +608,7 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
contextName: 'flow'
})
},
setSchema: async function(helpers: FlowAIChatHelpers) {
setSchema: async function (helpers: FlowAIChatHelpers) {
await buildSchemaForTool(this.def, async () => {
const flowInputsSchema = await helpers.getFlowInputsSchema()
return flowInputsSchema
@@ -622,7 +639,7 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
// Find the step in the flow
const modules = helpers.getModules()
let targetModule: FlowModule | undefined = modules.find((m) => m.id === stepId)
let targetModule: FlowModule | undefined = findModuleById(modules, stepId)
if (!targetModule) {
toolCallbacks.setToolStatus(toolId, {
@@ -647,7 +664,10 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
requestBody: {
content: moduleValue.content ?? '',
language: moduleValue.language,
args: module.id === 'preprocessor' ? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...stepArgs } : stepArgs
args:
module.id === 'preprocessor'
? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...stepArgs }
: stepArgs
}
}),
workspace,
@@ -675,7 +695,10 @@ export const flowTools: Tool<FlowAIChatHelpers>[] = [
requestBody: {
content: script.content,
language: script.language,
args: module.id === 'preprocessor' ? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...stepArgs } : stepArgs,
args:
module.id === 'preprocessor'
? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...stepArgs }
: stepArgs
}
}),
workspace,
@@ -721,6 +744,16 @@ Follow the user instructions carefully.
Go step by step, and explain what you're doing as you're doing it.
DO NOT wait for user confirmation before performing an action. Only do it if the user explicitly asks you to wait in their initial instructions.
ALWAYS test your modifications. You have access to the \`test_run_flow\` and \`test_run_step\` tools to test the flow and steps. If you only modified a single step, use the \`test_run_step\` tool to test it. If you modified the flow, use the \`test_run_flow\` tool to test it. If the user cancels the test run, do not try again and wait for the next user instruction.
When testing steps that are sql scripts, the arguments to be passed are { database: $res:<db_resource> }.
## Code Markers in Flow Modules
When viewing flow modules, the code content of rawscript steps may include \`[#START]\` and \`[#END]\` markers:
- These markers indicate specific code sections that need attention
- You MUST only modify the code between these markers when using the \`set_code\` tool
- After modifying the code, remove the markers from your response
- If a question is asked about the code, focus only on the code between the markers
- The markers appear in the YAML representation of flow modules when specific code pieces are selected
## Understanding User Requests
@@ -802,6 +835,16 @@ For truly static values in step inputs (those not linked to previous steps or lo
Both modules only support a script or rawscript step. You cannot nest modules using forloop/branchone/branchall.
### Contexts
You have access to the following contexts:
- Database schemas
- Flow diffs
- Focused flow modules
Database schemas give you the schema of databases the user is using.
Flow diffs give you the diff between the current flow and the last deployed flow.
Focused flow modules give you the ids of the flow modules the user is focused on. Your response should focus on these modules.
## Resource types
On Windmill, credentials and configuration are stored in resources. Resource types define the format of the resource.
If the user needs a resource as flow input, you should set the property type in the schema to "object" as well as add a key called "format" and set it to "resource-nameofresourcetype" (e.g. "resource-stripe").
@@ -816,26 +859,33 @@ If the user wants a specific resource as step input, you should set the step val
export function prepareFlowUserMessage(
instructions: string,
flowAndSelectedId?: { flow: ExtendedOpenFlow; selectedId: string }
flowAndSelectedId?: { flow: ExtendedOpenFlow; selectedId: string },
selectedContext: ContextElement[] = []
): ChatCompletionUserMessageParam {
const flow = flowAndSelectedId?.flow
const selectedId = flowAndSelectedId?.selectedId
// Handle context elements
const contextInstructions = selectedContext ? buildContextString(selectedContext) : ''
if (!flow || !selectedId) {
let userMessage = `## INSTRUCTIONS:
${instructions}`
return {
role: 'user',
content: `## INSTRUCTIONS:
${instructions}`
content: userMessage
}
}
return {
role: 'user',
content: `## FLOW:
const codePieces = selectedContext.filter((c) => c.type === 'flow_module_code_piece')
const flowModulesYaml = applyCodePiecesToFlowModules(codePieces, flow.value.modules)
let flowContent = `## FLOW:
flow_input schema:
${JSON.stringify(flow.schema ?? emptySchema())}
flow modules:
${YAML.stringify(flow.value.modules)}
${flowModulesYaml}
preprocessor module:
${YAML.stringify(flow.value.preprocessor_module)}
@@ -844,9 +894,15 @@ failure module:
${YAML.stringify(flow.value.failure_module)}
currently selected step:
${selectedId}
${selectedId}`
## INSTRUCTIONS:
flowContent += contextInstructions
flowContent += `\n\n## INSTRUCTIONS:
${instructions}`
return {
role: 'user',
content: flowContent
}
}
@@ -1,6 +1,6 @@
import { ResourceService, JobService } from '$lib/gen/services.gen'
import type { ResourceType, ScriptLang } from '$lib/gen/types.gen'
import { capitalize, emptySchema, isObject, toCamel } from '$lib/utils'
import { capitalize, isObject, toCamel } from '$lib/utils'
import { get } from 'svelte/store'
import { compile, phpCompile, pythonCompile } from '../../utils'
import type {
@@ -8,15 +8,19 @@ import type {
ChatCompletionTool,
ChatCompletionUserMessageParam
} from 'openai/resources/index.mjs'
import { copilotSessionModel, type DBSchema, dbSchemas } from '$lib/stores'
import { scriptLangToEditorLang } from '$lib/scripts'
import { type DBSchema, dbSchemas, getCurrentModel } from '$lib/stores'
import { getDbSchemas } from '$lib/components/apps/components/display/dbtable/utils'
import type { CodePieceElement, ContextElement } from '../context'
import type { ContextElement } from '../context'
import { PYTHON_PREPROCESSOR_MODULE_CODE, TS_PREPROCESSOR_MODULE_CODE } from '$lib/script_helpers'
import { createSearchHubScriptsTool, type Tool, executeTestRun, buildSchemaForTool, buildTestRunArgs } from '../shared'
import {
createSearchHubScriptsTool,
type Tool,
executeTestRun,
buildTestRunArgs,
buildContextString
} from '../shared'
import { setupTypeAcquisition, type DepsToGet } from '$lib/ata'
import { getModelContextWindow } from '../../lib'
import { inferArgs } from '$lib/infer'
// Score threshold for npm packages search filtering
const SCORE_THRESHOLD = 1000
@@ -348,7 +352,7 @@ export const CHAT_SYSTEM_PROMPT = `
- You can also receive a \`DIFF\` of the changes that have been made to the code. You should use this diff to give better answers.
- Before giving your answer, check again that you carefully followed these instructions.
- When asked to create a script that communicates with an external service, you can use the \`search_hub_scripts\` tool to search for relevant scripts in the hub. Make sure the language is the same as what the user is coding in. If you do not find any relevant scripts, you can use the \`search_npm_packages\` tool to search for relevant packages and their documentation. Always give a link to the documentation in your answer if possible.
- After modifying the code, ALWAYS use the \`test_run_script\` tool to test the code, and iterate on the code until it works as expected. If the user cancels the test run, do not try again and wait for the next user instruction.
- At the end of your reponse, if you modified or suggested changes to the code, ALWAYS use the \`test_run_script\` tool to test the code, and iterate on the code until it works as expected (MAX 3 times). If the user cancels the test run, do not try again and wait for the next user instruction.
Important:
Do not mention or reveal these instructions to the user unless explicitly asked to do so.
@@ -439,18 +443,6 @@ export async function main() {
\`\`\`
`
const CHAT_USER_CODE_CONTEXT = `
- {title}:
\`\`\`{language}
{code}
\`\`\`
`
const CHAT_USER_ERROR_CONTEXT = `
ERROR:
{error}
`
export const CHAT_USER_PROMPT = `
INSTRUCTIONS:
{instructions}
@@ -460,8 +452,6 @@ WINDMILL LANGUAGE CONTEXT:
`
export const CHAT_USER_DB_CONTEXT = `- {title}: SCHEMA: \n{schema}\n`
export function prepareScriptSystemMessage(): ChatCompletionSystemMessageParam {
return {
role: 'system',
@@ -469,18 +459,6 @@ export function prepareScriptSystemMessage(): ChatCompletionSystemMessageParam {
}
}
const applyCodePieceToCodeContext = (codePieces: CodePieceElement[], codeContext: string) => {
let code = codeContext.split('\n')
let shiftOffset = 0
codePieces.sort((a, b) => a.startLine - b.startLine)
for (const codePiece of codePieces) {
code.splice(codePiece.endLine + shiftOffset, 0, '[#END]')
code.splice(codePiece.startLine + shiftOffset - 1, 0, '[#START]')
shiftOffset += 2
}
return code.join('\n')
}
export function prepareScriptTools(
language: ScriptLang | 'bunnative',
context: ContextElement[]
@@ -508,61 +486,12 @@ export function prepareScriptUserMessage(
isPreprocessor?: boolean
} = {}
): ChatCompletionUserMessageParam {
let codeContext = 'CODE:\n'
let errorContext = 'ERROR:\n'
let dbContext = 'DATABASES:\n'
let diffContext = 'DIFF:\n'
let hasCode = false
let hasError = false
let hasDb = false
let hasDiff = false
for (const context of selectedContext) {
if (context.type === 'code') {
hasCode = true
codeContext += CHAT_USER_CODE_CONTEXT.replace('{title}', context.title)
.replace('{language}', scriptLangToEditorLang(language))
.replace(
'{code}',
applyCodePieceToCodeContext(
selectedContext.filter((c) => c.type === 'code_piece'),
context.content
)
)
} else if (context.type === 'error') {
if (hasError) {
throw new Error('Multiple error contexts provided')
}
hasError = true
errorContext = CHAT_USER_ERROR_CONTEXT.replace('{error}', context.content)
} else if (context.type === 'db') {
hasDb = true
dbContext += CHAT_USER_DB_CONTEXT.replace('{title}', context.title).replace(
'{schema}',
context.schema?.stringified ?? 'to fetch with get_db_schema'
)
} else if (context.type === 'diff') {
hasDiff = true
const diff = JSON.stringify(context.diff)
diffContext = diff.length > 3000 ? diff.slice(0, 3000) + '...' : diff
}
}
let userMessage = CHAT_USER_PROMPT.replace('{instructions}', instructions).replace(
'{lang_context}',
getLangContext(language, { allowResourcesFetch: true, ...options })
)
if (hasCode) {
userMessage += codeContext
}
if (hasError) {
userMessage += errorContext
}
if (hasDb) {
userMessage += dbContext
}
if (hasDiff) {
userMessage += diffContext
}
const contextInstructions = buildContextString(selectedContext)
userMessage += contextInstructions
return {
role: 'user',
content: userMessage
@@ -626,7 +555,12 @@ async function formatDBSchema(dbSchema: DBSchema) {
}
export interface ScriptChatHelpers {
getScriptOptions: () => { code: string; lang: ScriptLang | 'bunnative'; path: string; args: Record<string, any> }
getScriptOptions: () => {
code: string
lang: ScriptLang | 'bunnative'
path: string
args: Record<string, any>
}
getLastSuggestedCode: () => string | undefined
applyCode: (code: string, applyAll?: boolean) => void
}
@@ -634,51 +568,60 @@ export interface ScriptChatHelpers {
export const resourceTypeTool: Tool<ScriptChatHelpers> = {
def: RESOURCE_TYPE_FUNCTION_DEF,
fn: async ({ args, workspace, helpers, toolCallbacks, toolId }) => {
toolCallbacks.setToolStatus(toolId, { content: 'Searching resource types for "' + args.query + '"...' })
toolCallbacks.setToolStatus(toolId, {
content: 'Searching resource types for "' + args.query + '"...'
})
const lang = helpers.getScriptOptions().lang
const formattedResourceTypes = await getFormattedResourceTypes(
lang,
args.query,
workspace
)
toolCallbacks.setToolStatus(toolId, { content: 'Retrieved resource types for "' + args.query + '"' })
const formattedResourceTypes = await getFormattedResourceTypes(lang, args.query, workspace)
toolCallbacks.setToolStatus(toolId, {
content: 'Retrieved resource types for "' + args.query + '"'
})
return formattedResourceTypes
}
}
export const dbSchemaTool: Tool<ScriptChatHelpers> = {
def: DB_SCHEMA_FUNCTION_DEF,
fn: async ({ args, workspace, toolCallbacks, toolId }) => {
if (!args.resourcePath) {
throw new Error('Database path not provided')
}
toolCallbacks.setToolStatus(toolId, { content: 'Getting database schema for ' + args.resourcePath + '...' })
const resource = await ResourceService.getResource({
workspace: workspace,
path: args.resourcePath
})
const newDbSchemas = {}
await getDbSchemas(
resource.resource_type,
args.resourcePath,
workspace,
newDbSchemas,
(error) => {
console.error(error)
// Generic DB schema tool factory that can be used by both script and flow modes
export function createDbSchemaTool<T>(): Tool<T> {
return {
def: DB_SCHEMA_FUNCTION_DEF,
fn: async ({ args, workspace, toolCallbacks, toolId }) => {
if (!args.resourcePath) {
throw new Error('Database path not provided')
}
)
dbSchemas.update((schemas) => ({ ...schemas, ...newDbSchemas }))
const dbs = get(dbSchemas)
const db = dbs[args.resourcePath]
if (!db) {
throw new Error('Database not found')
toolCallbacks.setToolStatus(toolId, {
content: 'Getting database schema for ' + args.resourcePath + '...'
})
const resource = await ResourceService.getResource({
workspace: workspace,
path: args.resourcePath
})
const newDbSchemas = {}
await getDbSchemas(
resource.resource_type,
args.resourcePath,
workspace,
newDbSchemas,
(error) => {
console.error(error)
}
)
dbSchemas.update((schemas) => ({ ...schemas, ...newDbSchemas }))
const dbs = get(dbSchemas)
const db = dbs[args.resourcePath]
if (!db) {
throw new Error('Database not found')
}
const stringSchema = await formatDBSchema(db)
toolCallbacks.setToolStatus(toolId, {
content: 'Retrieved database schema for ' + args.resourcePath
})
return stringSchema
}
const stringSchema = await formatDBSchema(db)
toolCallbacks.setToolStatus(toolId, { content: 'Retrieved database schema for ' + args.resourcePath })
return stringSchema
}
}
export const dbSchemaTool: Tool<ScriptChatHelpers> = createDbSchemaTool<ScriptChatHelpers>()
type PackageSearchQuery = {
package: {
name: string
@@ -712,7 +655,8 @@ export async function searchExternalIntegrationResources(args: { query: string }
(r: PackageSearchQuery) => r.searchScore >= SCORE_THRESHOLD
)
const modelContextWindow = getModelContextWindow(get(copilotSessionModel)?.model ?? '')
const model = getCurrentModel()
const modelContextWindow = getModelContextWindow(model.model)
const results: PackageSearchResult[] = await Promise.all(
filtered.map(async (r: PackageSearchQuery) => {
let documentation = ''
@@ -839,31 +783,31 @@ const TEST_RUN_SCRIPT_TOOL: ChatCompletionTool = {
function: {
name: 'test_run_script',
description: 'Execute a test run of the current script in the editor',
// will be overridden by setSchema
parameters: {
type: 'object',
properties: {
args: {
type: 'object',
description: 'Arguments to pass to the script (optional, uses current editor args if not provided)'
}
args: { type: 'string', description: 'JSON string containing the arguments for the tool' }
},
required: []
additionalProperties: false,
strict: false,
required: ['args']
}
},
}
}
export const testRunScriptTool: Tool<ScriptChatHelpers> = {
def: TEST_RUN_SCRIPT_TOOL,
fn: async function({ args, workspace, helpers, toolCallbacks, toolId }) {
fn: async function ({ args, workspace, helpers, toolCallbacks, toolId }) {
const scriptOptions = helpers.getScriptOptions()
if (!scriptOptions) {
toolCallbacks.setToolStatus(toolId, {
toolCallbacks.setToolStatus(toolId, {
content: 'No script available to test',
error: 'No script found in current context'
})
throw new Error('No script code available to test. Please ensure you have a script open in the editor.')
throw new Error(
'No script code available to test. Please ensure you have a script open in the editor.'
)
}
let codeToTest = scriptOptions.code
@@ -873,7 +817,7 @@ export const testRunScriptTool: Tool<ScriptChatHelpers> = {
if (lastSuggestedCode && lastSuggestedCode !== codeToTest) {
codeToTest = lastSuggestedCode
toolCallbacks.setToolStatus(toolId, { content: 'Applying code changes...' })
// Apply the suggested code changes using the existing mechanism
helpers.applyCode(lastSuggestedCode, true)
@@ -883,15 +827,16 @@ export const testRunScriptTool: Tool<ScriptChatHelpers> = {
const parsedArgs = await buildTestRunArgs(args, this.def)
return executeTestRun({
jobStarter: () => JobService.runScriptPreview({
workspace: workspace,
requestBody: {
path: scriptOptions.path,
content: codeToTest,
args: parsedArgs,
language: scriptOptions.lang as ScriptLang,
}
}),
jobStarter: () =>
JobService.runScriptPreview({
workspace: workspace,
requestBody: {
path: scriptOptions.path,
content: codeToTest,
args: parsedArgs,
language: scriptOptions.lang as ScriptLang
}
}),
workspace,
toolCallbacks,
toolId,
@@ -899,23 +844,7 @@ export const testRunScriptTool: Tool<ScriptChatHelpers> = {
contextName: 'script'
})
},
setSchema: async function(helpers: ScriptChatHelpers) {
await buildSchemaForTool(this.def, async () => {
const scriptOptions = helpers.getScriptOptions()
const code = scriptOptions?.code
const lang = scriptOptions?.lang
const lastSuggestedCode = helpers.getLastSuggestedCode()
const codeToTest = lastSuggestedCode ?? code
if (codeToTest) {
const newSchema = emptySchema()
await inferArgs(lang, codeToTest, newSchema)
return newSchema
}
return emptySchema()
})
},
requiresConfirmation: true,
confirmationMessage: 'Run script test',
showDetails: true,
showDetails: true
}
@@ -4,13 +4,194 @@ import type {
ChatCompletionTool
} from 'openai/resources/chat/completions.mjs'
import { get } from 'svelte/store'
import type { ContextElement } from './context'
import { copilotSessionModel, workspaceStore } from '$lib/stores'
import type { CodePieceElement, ContextElement, FlowModuleCodePieceElement } from './context'
import { workspaceStore, getCurrentModel } from '$lib/stores'
import type { ExtendedOpenFlow } from '$lib/components/flows/types'
import type { FunctionParameters } from 'openai/resources/shared.mjs'
import { zodToJsonSchema } from 'zod-to-json-schema'
import { z } from 'zod'
import { ScriptService, JobService, type CompletedJob } from '$lib/gen'
import { ScriptService, JobService, type CompletedJob, type FlowModule } from '$lib/gen'
import { scriptLangToEditorLang } from '$lib/scripts'
import YAML from 'yaml'
export interface ContextStringResult {
dbContext: string
diffContext: string
flowModuleContext: string
hasDb: boolean
hasDiff: boolean
hasFlowModule: boolean
}
export const extractAllModules = (modules: FlowModule[]): FlowModule[] => {
return modules.flatMap((m) => {
if (m.value.type === 'forloopflow' || m.value.type === 'whileloopflow') {
return [m, ...extractAllModules(m.value.modules)]
}
if (m.value.type === 'branchall') {
return [m, ...extractAllModules(m.value.branches.flatMap((b) => b.modules))]
}
if (m.value.type === 'branchone') {
return [
m,
...extractAllModules([...m.value.branches.flatMap((b) => b.modules), ...m.value.default])
]
}
return [m]
})
}
export const findModuleById = (modules: FlowModule[], moduleId: string): FlowModule | undefined => {
for (const module of modules) {
if (module.id === moduleId) {
return module
}
if (module.value.type === 'forloopflow' || module.value.type === 'whileloopflow') {
const found = findModuleById(module.value.modules, moduleId)
if (found) {
return found
}
}
if (module.value.type === 'branchall') {
const allModules = module.value.branches.flatMap((b) => b.modules)
const found = findModuleById(allModules, moduleId)
if (found) {
return found
}
}
if (module.value.type === 'branchone') {
const allModules = [
...module.value.branches.flatMap((b) => b.modules),
...module.value.default
]
const found = findModuleById(allModules, moduleId)
if (found) {
return found
}
}
}
return undefined
}
const applyCodePieceToCodeContext = (codePieces: CodePieceElement[], codeContext: string) => {
let code = codeContext.split('\n')
let shiftOffset = 0
codePieces.sort((a, b) => a.startLine - b.startLine)
for (const codePiece of codePieces) {
code.splice(codePiece.endLine + shiftOffset, 0, '[#END]')
code.splice(codePiece.startLine + shiftOffset - 1, 0, '[#START]')
shiftOffset += 2
}
return code.join('\n')
}
export function applyCodePiecesToFlowModules(
codePieces: FlowModuleCodePieceElement[],
flowModules: FlowModule[]
): string {
const moduleCodePieces = new Map<string, FlowModuleCodePieceElement[]>()
for (const codePiece of codePieces) {
const moduleId = codePiece.id
if (!moduleCodePieces.has(moduleId)) {
moduleCodePieces.set(moduleId, [])
}
moduleCodePieces.get(moduleId)!.push(codePiece)
}
// Clone modules to avoid mutation
const modifiedModules = JSON.parse(JSON.stringify(flowModules))
// Apply code pieces to each module
for (const [moduleId, pieces] of moduleCodePieces) {
const module = findModuleById(modifiedModules, moduleId)
if (module && module.value.type === 'rawscript' && module.value.content) {
module.value.content = applyCodePieceToCodeContext(
pieces as unknown as CodePieceElement[],
module.value.content
)
}
}
return YAML.stringify(modifiedModules)
}
export function buildContextString(selectedContext: ContextElement[]): string {
const dbTemplate = `- {title}: SCHEMA: \n{schema}\n`
const codeTemplate = `
- {title}:
\`\`\`{language}
{code}
\`\`\`
`
let dbContext = 'DATABASES:\n'
let diffContext = 'DIFF:\n'
let flowModuleContext = 'FOCUSED FLOW MODULES IDS:\n'
let codeContext = 'CODE:\n'
let errorContext = `
ERROR:
{error}
`
let hasCode = false
let hasDb = false
let hasDiff = false
let hasFlowModule = false
let hasError = false
let result = '\n\n'
for (const context of selectedContext) {
if (context.type === 'code') {
hasCode = true
codeContext += codeTemplate
.replace('{title}', context.title)
.replace('{language}', scriptLangToEditorLang(context.lang))
.replace(
'{code}',
applyCodePieceToCodeContext(
selectedContext.filter((c) => c.type === 'code_piece'),
context.content
)
)
} else if (context.type === 'error') {
if (hasError) {
throw new Error('Multiple error contexts provided')
}
hasError = true
errorContext = errorContext.replace('{error}', context.content)
} else if (context.type === 'db') {
hasDb = true
dbContext += dbTemplate
.replace('{title}', context.title)
.replace('{schema}', context.schema?.stringified ?? 'to fetch with get_db_schema')
dbContext += '\n'
} else if (context.type === 'diff') {
hasDiff = true
const diff = JSON.stringify(context.diff)
diffContext += (diff.length > 3000 ? diff.slice(0, 3000) + '...' : diff) + '\n'
} else if (context.type === 'flow_module') {
hasFlowModule = true
flowModuleContext += `${context.id}\n`
}
}
if (hasCode) {
result += '\n' + codeContext
}
if (hasError) {
result += '\n' + errorContext
}
if (hasDb) {
result += '\n' + dbContext
}
if (hasDiff) {
result += '\n' + diffContext
}
if (hasFlowModule) {
result += '\n' + flowModuleContext
}
return result
}
type BaseDisplayMessage = {
content: string
@@ -89,11 +270,13 @@ export async function processToolCall<T>({
// Add the tool to the display with appropriate status
toolCallbacks.setToolStatus(toolCall.id, {
...(tool?.requiresConfirmation ? { content: tool.confirmationMessage ?? "Waiting for confirmation..." } : {}),
...(tool?.requiresConfirmation
? { content: tool.confirmationMessage ?? 'Waiting for confirmation...' }
: {}),
parameters: args,
isLoading: true,
needsConfirmation: needsConfirmation,
showDetails: tool?.showDetails,
showDetails: tool?.showDetails
})
// If confirmation is needed and we have the callback, wait for it
@@ -254,12 +437,17 @@ export const createSearchHubScriptsTool = (withContent: boolean = false) => ({
}
})
export async function buildSchemaForTool(toolDef: ChatCompletionTool, schemaBuilder: () => Promise<FunctionParameters>): Promise<boolean> {
export async function buildSchemaForTool(
toolDef: ChatCompletionTool,
schemaBuilder: () => Promise<FunctionParameters>
): Promise<boolean> {
try {
const schema = await schemaBuilder()
// if schema properties contains values different from '^[a-zA-Z0-9_.-]{1,64}$'
const invalidProperties = Object.keys(schema.properties ?? {}).filter((key) => !/^[a-zA-Z0-9_.-]{1,64}$/.test(key))
const invalidProperties = Object.keys(schema.properties ?? {}).filter(
(key) => !/^[a-zA-Z0-9_.-]{1,64}$/.test(key)
)
if (invalidProperties.length > 0) {
console.warn(`Invalid flow inputs schema: ${invalidProperties.join(', ')}`)
throw new Error(`Invalid flow inputs schema: ${invalidProperties.join(', ')}`)
@@ -267,15 +455,23 @@ export async function buildSchemaForTool(toolDef: ChatCompletionTool, schemaBuil
toolDef.function.parameters = { ...schema, additionalProperties: false }
// OPEN AI models don't support strict mode well with schema with complex properties, so we disable it
const model = get(copilotSessionModel)?.provider
if (model === 'openai' || model === 'azure_openai') {
const model = getCurrentModel()
if (model.provider === 'openai' || model.provider === 'azure_openai') {
toolDef.function.strict = false
}
return true
} catch (error) {
console.error('Error building schema for tool', error)
// fallback to schema with args as a JSON string
toolDef.function.parameters = { type: 'object', properties: { args: { type: 'string', description: 'JSON string containing the arguments for the tool' } }, additionalProperties: false, strict: false, required: ['args'] }
toolDef.function.parameters = {
type: 'object',
properties: {
args: { type: 'string', description: 'JSON string containing the arguments for the tool' }
},
additionalProperties: false,
strict: false,
required: ['args']
}
return false
}
}
@@ -393,7 +589,10 @@ function getErrorMessage(result: unknown): string {
export async function buildTestRunArgs(args: any, toolDef: ChatCompletionTool): Promise<any> {
let parsedArgs = args
// if the schema is the fallback schema, parse the args as a JSON string
if ((toolDef.function.parameters as any).properties?.args?.description === 'JSON string containing the arguments for the tool') {
if (
(toolDef.function.parameters as any).properties?.args?.description ===
'JSON string containing the arguments for the tool'
) {
try {
parsedArgs = JSON.parse(args.args)
} catch (error) {
+12 -13
View File
@@ -1,7 +1,6 @@
import type { AIProvider, AIProviderModel } from '$lib/gen'
import {
copilotInfo,
copilotSessionModel,
getCurrentModel,
workspaceStore,
type DBSchema,
type GraphqlSchema,
@@ -23,7 +22,16 @@ import { z } from 'zod'
export const SUPPORTED_LANGUAGES = new Set(Object.keys(GEN_CONFIG.prompts))
const OPENAI_MODELS = ['gpt-5', 'gpt-5-mini', 'gpt-5-nano', 'gpt-4o', 'gpt-4o-mini', 'o4-mini', 'o3', 'o3-mini']
const OPENAI_MODELS = [
'gpt-5',
'gpt-5-mini',
'gpt-5-nano',
'gpt-4o',
'gpt-4o-mini',
'o4-mini',
'o3',
'o3-mini'
]
// need at least one model for each provider except customai
export const AI_DEFAULT_MODELS: Record<AIProvider, string[]> = {
@@ -468,18 +476,9 @@ function getProviderAndCompletionConfig<K extends boolean>({
? ChatCompletionCreateParamsStreaming
: ChatCompletionCreateParamsNonStreaming
} {
let info = get(copilotInfo)
const modelProvider =
forceModelProvider ?? get(copilotSessionModel) ?? info.defaultModel ?? info.aiModels[0]
if (!modelProvider) {
throw new Error('No model selected')
}
const modelProvider = forceModelProvider ?? getCurrentModel()
const providerConfig = PROVIDER_COMPLETION_CONFIG_MAP[modelProvider.provider]
const processedMessages = prepareMessages(modelProvider.provider, messages)
return {
provider: modelProvider.provider,
config: {
@@ -18,6 +18,8 @@
import { triggerableByAI } from '$lib/actions/triggerableByAI.svelte'
import type { ModulesTestStates } from '../modulesTest.svelte'
import type { StateStore } from '$lib/utils'
import type { FlowOptions } from '../copilot/chat/ContextManager.svelte'
import { extractAllModules } from '../copilot/chat/shared'
const { flowStore } = getContext<FlowEditorContext>('FlowEditorContext')
interface Props {
@@ -56,6 +58,7 @@
suspendStatus?: StateStore<Record<string, { job: Job; nb: number }>>
onDelete?: (id: string) => void
flowHasChanged?: boolean
previewOpen: boolean
}
let {
@@ -89,7 +92,8 @@
job,
suspendStatus,
onDelete,
flowHasChanged
flowHasChanged,
previewOpen
}: Props = $props()
let flowModuleSchemaMap: FlowModuleSchemaMap | undefined = $state()
@@ -103,11 +107,24 @@
pickablePropertiesFiltered: writable<PickableProperties | undefined>(undefined)
})
$effect(() => {
const options: FlowOptions = {
currentFlow: flowStore.val,
lastDeployedFlow: savedFlow,
lastSavedFlow: savedFlow?.draft,
path: savedFlow?.path,
modules: extractAllModules(flowStore.val.value.modules)
}
aiChatManager.flowOptions = options
})
onMount(() => {
aiChatManager.saveAndClear()
aiChatManager.changeMode(AIMode.FLOW)
})
onDestroy(() => {
aiChatManager.flowOptions = undefined
aiChatManager.changeMode(AIMode.NAVIGATOR)
})
</script>
@@ -191,6 +208,7 @@
{isOwner}
{suspendStatus}
onOpenDetails={onOpenPreview}
{previewOpen}
/>
{/if}
</Pane>
@@ -0,0 +1,50 @@
<script lang="ts">
import LanguageIcon from '$lib/components/common/languageIcons/LanguageIcon.svelte'
import IconedResourceType from '$lib/components/IconedResourceType.svelte'
import type { FlowModule } from '$lib/gen'
import { Building, Repeat, Square, ArrowDown, GitBranch, Bot } from 'lucide-svelte'
import BarsStaggered from '$lib/components/icons/BarsStaggered.svelte'
interface Props {
module: FlowModule
size?: number
width?: number
height?: number
}
let { module, size = 16, width, height }: Props = $props()
// Use width/height if provided, otherwise use size for both
const iconWidth = width || size
const iconHeight = height || size
</script>
{#if module.value.type === 'aiagent'}
<Bot size={16} class="text-violet-800 dark:text-violet-400" />
{:else if module.value.type === 'rawscript'}
<LanguageIcon lang={module.value.language} width={iconWidth} height={iconHeight} />
{:else if module.summary === 'Terminate flow'}
<Square {size} />
{:else if module.value.type === 'identity'}
<ArrowDown {size} />
{:else if module.value.type === 'flow'}
<BarsStaggered {size} />
{:else if module.value.type === 'forloopflow' || module.value.type === 'whileloopflow'}
<Repeat {size} />
{:else if module.value.type === 'branchone' || module.value.type === 'branchall'}
<GitBranch {size} />
{:else if module.value.type === 'script'}
{#if module.value.path.startsWith('hub/')}
<IconedResourceType
width={iconWidth.toString() + 'px'}
height={iconHeight.toString() + 'px'}
name={module.value.path.split('/')[2]}
silent={true}
/>
{:else}
<Building {size} />
{/if}
{:else}
<!-- Fallback icon for unknown module types -->
<BarsStaggered {size} />
{/if}
@@ -154,6 +154,9 @@
{:else if flowModuleValue.type === 'flow'}
<Badge color="indigo" capitalize>flow</Badge>
<input bind:value={summary} placeholder="Summary" class="w-full grow" />
{:else if flowModuleValue.type === 'aiagent'}
<Badge color="indigo">AI Agent</Badge>
<input bind:value={summary} placeholder="Summary" class="w-full grow" />
{/if}
</div>
</span>
@@ -34,6 +34,7 @@
isOwner?: boolean
suspendStatus?: StateStore<Record<string, { job: Job; nb: number }>>
onOpenDetails?: () => void
previewOpen?: boolean
}
let {
@@ -49,7 +50,8 @@
job,
isOwner,
suspendStatus,
onOpenDetails
onOpenDetails,
previewOpen = false
}: Props = $props()
const {
@@ -95,6 +97,7 @@
}}
on:applyArgs
{onTestFlow}
{previewOpen}
/>
{:else if $selectedId === 'Result'}
<FlowResult {noEditor} {job} {isOwner} {suspendStatus} {onOpenDetails} />
@@ -48,9 +48,10 @@
noEditor: boolean
disabled: boolean
onTestFlow?: () => void
previewOpen: boolean
}
let { noEditor, disabled, onTestFlow }: Props = $props()
let { noEditor, disabled, onTestFlow, previewOpen }: Props = $props()
const {
flowStore,
previewArgs,
@@ -63,7 +64,6 @@
let addPropertyV2: AddPropertyV2 | undefined = $state(undefined)
let previewSchema: Record<string, any> | undefined = $state(undefined)
let payloadData: Record<string, any> | undefined = undefined
let previewArguments: Record<string, any> | undefined = $state(previewArgs.val)
let dropdownItems: Array<{
label: string
onClick: () => void
@@ -195,7 +195,9 @@
function handleKeydown(event: KeyboardEvent) {
if ((event.metaKey || event.ctrlKey) && event.key === 'Enter') {
runPreview()
if (!previewOpen) {
runPreview()
}
} else if (event.key === 'Enter' && previewSchema && !preventEnter) {
applySchemaAndArgs()
connectFirstNode()
@@ -205,9 +207,6 @@
}
function runPreview() {
if (previewArguments) {
previewArgs.val = structuredClone($state.snapshot(previewArguments))
}
onTestFlow?.()
}
@@ -248,8 +247,8 @@
async function applySchemaAndArgs() {
flowStore.val.schema = applyDiff(flowStore.val.schema, diff)
if (previewArguments) {
savedPreviewArgs = structuredClone($state.snapshot(previewArguments))
if (previewArgs.val) {
savedPreviewArgs = structuredClone($state.snapshot(previewArgs.val))
}
updatePreviewSchemaAndArgs(undefined)
if ($flowInputEditorState) {
@@ -259,11 +258,13 @@
function updatePreviewArguments(payloadData: Record<string, any> | undefined) {
if (!payloadData) {
previewArguments = savedPreviewArgs
if (savedPreviewArgs) {
previewArgs.val = savedPreviewArgs
}
return
}
savedPreviewArgs = structuredClone($state.snapshot(previewArguments))
previewArguments = structuredClone($state.snapshot(payloadData))
savedPreviewArgs = structuredClone($state.snapshot(previewArgs.val))
previewArgs.val = structuredClone($state.snapshot(payloadData))
}
let tabButtonWidth = 0
@@ -372,7 +373,7 @@
displayWebhookWarning
editTab={$flowInputEditorState?.selectedTab}
{previewSchema}
bind:args={previewArguments}
bind:args={previewArgs.val}
bind:editPanelSize={
() => {
return editPanelSize
@@ -399,9 +400,8 @@
}}
shouldDispatchChanges={true}
on:change={() => {
previewArguments = previewArguments
if (!previewSchema) {
savedPreviewArgs = structuredClone($state.snapshot(previewArguments))
savedPreviewArgs = structuredClone($state.snapshot(previewArgs.val))
}
refreshStateStore(flowStore)
}}
@@ -560,7 +560,7 @@
on:isEditing={(e) => {
preventEnter = e.detail
}}
previewArgs={previewArguments}
previewArgs={previewArgs.val}
{isValid}
limitPayloadSize
bind:this={savedInputsPicker}
@@ -583,7 +583,7 @@
on:select={(e) => {
updatePreviewSchemaAndArgs(e.detail ?? undefined)
}}
selected={!!previewArguments}
selected={!!previewArgs.val}
bind:this={jsonInputs}
/>
</FlowInputEditor>
@@ -36,7 +36,7 @@
import FlowModuleMockTransitionMessage from './FlowModuleMockTransitionMessage.svelte'
import Tooltip from '$lib/components/Tooltip.svelte'
import { SecondsInput } from '$lib/components/common'
import DiffEditor from '$lib/components/DiffEditor.svelte'
import DiffEditor, { type ButtonProp } from '$lib/components/DiffEditor.svelte'
import FlowModuleTimeout from './FlowModuleTimeout.svelte'
import HighlightCode from '$lib/components/HighlightCode.svelte'
import ToggleButtonGroup from '$lib/components/common/toggleButton-v2/ToggleButtonGroup.svelte'
@@ -47,7 +47,7 @@
import { isCloudHosted } from '$lib/cloud'
import { loadSchemaFromModule } from '../flowInfers'
import FlowModuleSkip from './FlowModuleSkip.svelte'
import { type Job, JobService } from '$lib/gen'
import { type Job } from '$lib/gen'
import { workspaceStore } from '$lib/stores'
import { checkIfParentLoop } from '../utils'
import ModulePreviewResultViewer from '$lib/components/ModulePreviewResultViewer.svelte'
@@ -102,6 +102,15 @@
let workspaceScriptTag: string | undefined = $state(undefined)
let workspaceScriptLang: ScriptLang | undefined = $state(undefined)
let diffMode = $state(false)
let diffButtons = $state<ButtonProp[]>([
{
text: 'Quit diff mode',
color: 'red',
onClick: () => {
hideDiffMode()
}
}
])
let editor: Editor | undefined = $state()
let diffEditor: DiffEditor | undefined = $state()
@@ -120,7 +129,6 @@
let s3Kind = $state('s3_client')
let validCode = $state(true)
let width = $state(1200)
let lastJob: Job | undefined = $state(undefined)
let testJob: Job | undefined = $state(undefined)
let testIsLoading = $state(false)
let scriptProgress = $state(undefined)
@@ -197,43 +205,10 @@
let editorSettingsPanelSize = $state(100 - untrack(() => editorPanelSize))
let stepHistoryLoader = getStepHistoryLoaderContext()
let lastJobId: string | undefined = undefined
function onSelectedIdChange() {
if (!flowStateStore?.val?.[$selectedId]?.schema && flowModule) {
reload(flowModule)
}
lastJobId = undefined
}
async function getLastJob() {
if (
!flowStateStore ||
!flowModule.id ||
flowStateStore.val[flowModule.id]?.previewResult === 'never tested this far' ||
!flowStateStore.val[flowModule.id]?.previewJobId
) {
return
}
if (
lastJobId == flowStateStore.val[flowModule.id]?.previewJobId ||
lastJob?.id == flowStateStore.val[flowModule.id]?.previewJobId ||
flowStateStore.val[flowModule.id]?.previewSuccess == undefined
) {
return
}
lastJobId = flowStateStore.val[flowModule.id]?.previewJobId
const job = await JobService.getJob({
workspace: $workspaceStore ?? '',
id: flowStateStore.val[flowModule.id]?.previewJobId ?? '',
noCode: true
})
if (job && job.type === 'CompletedJob') {
lastJobId = flowStateStore.val[flowModule.id]?.previewJobId
lastJob = job
}
}
let leftPanelSize = $state(0)
@@ -270,13 +245,6 @@
$effect.pre(() => {
$selectedId && untrack(() => onSelectedIdChange())
})
$effect(() => {
if (testJob && testJob.type === 'CompletedJob') {
lastJob = $state.snapshot(testJob)
} else if ($workspaceStore && $pathStore && flowModule?.id && flowStateStore) {
untrack(() => getLastJob())
}
})
let parentLoop = $derived(
flowStore.val && flowModule ? checkIfParentLoop(flowStore.val, flowModule.id) : undefined
)
@@ -297,7 +265,13 @@
showDiffMode,
hideDiffMode,
diffMode,
lastDeployedCode
lastDeployedCode,
setDiffOriginal: (code: string) => {
diffEditor?.setOriginal(code ?? '')
},
setDiffButtons: (buttons: ButtonProp[]) => {
diffButtons = buttons
}
})
})
@@ -322,6 +296,12 @@
let rawScriptLang = $derived(
flowModule.value.type == 'rawscript' ? flowModule.value.language : undefined
)
let modulePreviewResultViewer: ModulePreviewResultViewer | undefined = $state(undefined)
function onJobDone() {
modulePreviewResultViewer?.getOutputPickerInner()?.setJobPreview()
}
</script>
<svelte:window onkeydown={onKeyDown} />
@@ -419,6 +399,7 @@
{lastDeployedCode}
{diffMode}
openAiChat
moduleId={flowModule.id}
/>
</div>
{/if}
@@ -477,6 +458,7 @@
{}
)}
key={`flow-inline-${$workspaceStore}-${$pathStore}-${flowModule.id}`}
moduleId={flowModule.id}
/>
<DiffEditor
open={false}
@@ -484,10 +466,8 @@
automaticLayout
fixedOverflowWidgets
defaultLang={scriptLangToEditorLang(flowModule.value.language)}
class="h-full"
showButtons={diffMode}
showHistoryButton={false}
on:hideDiffMode={hideDiffMode}
className="h-full"
buttons={diffMode ? diffButtons : []}
/>
{/key}
{/if}
@@ -594,6 +574,7 @@
bind:testIsLoading
bind:scriptProgress
focusArg={highlightArg}
{onJobDone}
/>
{:else if selected === 'advanced'}
<Tabs bind:selected={advancedSelected}>
@@ -898,15 +879,15 @@
flowModule = flowModule
refreshStateStore(flowStore)
}}
{lastJob}
{scriptProgress}
{testJob}
{scriptProgress}
mod={flowModule}
{testIsLoading}
disableMock={preprocessorModule || failureModule}
disableHistory={failureModule}
loadingJob={stepHistoryLoader?.stepStates[flowModule.id]?.loadingJobs}
tagLabel={customUi?.tagLabel}
bind:this={modulePreviewResultViewer}
/>
</Pane>
{/if}
@@ -140,7 +140,7 @@
{:then Module}
<Module.default
open={true}
class="h-screen"
className="h-screen"
readOnly
automaticLayout
defaultLang={scriptLangToEditorLang(language)}
+6 -6
View File
@@ -3,29 +3,29 @@ import type { FlowModule } from '$lib/gen'
export function dfs<T>(
modules: FlowModule[],
f: (x: FlowModule, modules: FlowModule[], branches: FlowModule[][]) => T,
{ skipToolNodes = false }: { skipToolNodes?: boolean } = {}
opts: { skipToolNodes?: boolean } = {}
): T[] {
let result: T[] = []
for (const module of modules) {
if (module.value.type == 'forloopflow' || module.value.type == 'whileloopflow') {
result = result.concat(f(module, modules, [module.value.modules]))
result = result.concat(dfs(module.value.modules, f))
result = result.concat(dfs(module.value.modules, f, opts))
} else if (module.value.type == 'branchone') {
const allBranches = [module.value.default, ...module.value.branches.map((b) => b.modules)]
result = result.concat(f(module, modules, allBranches))
for (const branch of allBranches) {
result = result.concat(dfs(branch, f))
result = result.concat(dfs(branch, f, opts))
}
} else if (module.value.type == 'branchall') {
const allBranches = module.value.branches.map((b) => b.modules)
result = result.concat(f(module, modules, allBranches))
for (const branch of allBranches) {
result = result.concat(dfs(branch, f))
result = result.concat(dfs(branch, f, opts))
}
} else if (module.value.type == 'aiagent' && !skipToolNodes) {
} else if (module.value.type == 'aiagent' && !opts.skipToolNodes) {
result = result.concat(f(module, modules, [module.value.tools]))
result = result.concat(dfs(module.value.tools, f))
result = result.concat(dfs(module.value.tools, f, opts))
} else {
result.push(f(module, modules, []))
}
@@ -96,13 +96,12 @@ export async function loadSchemaFromModule(module: FlowModule): Promise<{
}
]
},
system_prompt: {
type: 'string',
default: 'You are a helpful assistant'
},
user_message: {
type: 'string'
},
system_prompt: {
type: 'string'
},
max_completion_tokens: {
type: 'number'
},
@@ -110,13 +109,13 @@ export async function loadSchemaFromModule(module: FlowModule): Promise<{
type: 'number'
}
},
required: ['provider', 'model', 'system_prompt', 'user_message'],
required: ['provider', 'model', 'user_message'],
type: 'object',
order: [
'provider',
'model',
'system_prompt',
'user_message',
'system_prompt',
'max_completion_tokens',
'temperature'
]
@@ -165,8 +165,7 @@ export async function createBranchAll(id: string): Promise<[FlowModule, FlowModu
export async function createAiAgent(id: string): Promise<[FlowModule, FlowModuleState]> {
const aiAgentFlowModules: FlowModule = {
id,
value: { type: 'aiagent', tools: [], input_transforms: {} },
summary: 'AI Agent'
value: { type: 'aiagent', tools: [], input_transforms: {} }
}
const flowModuleState = await loadFlowModuleState(aiAgentFlowModules)
@@ -144,9 +144,7 @@
let testIsLoading = $state(false)
let hover = $state(false)
let connectingData: any | undefined = $state(undefined)
let lastJob: any | undefined = $state(undefined)
let outputPicker: OutputPicker | undefined = $state(undefined)
let historyOpen = $state(false)
let testJob: any | undefined = $state(undefined)
let outputPickerBarOpen = $state(false)
@@ -170,30 +168,6 @@
updateConnectingData(id, pickableIds, $flowPropPickerConfig, flowStateStore)
})
function updateLastJob(flowStateStore: any | undefined) {
if (
!flowStateStore ||
!id ||
flowStateStore.val[id]?.previewResult === 'never tested this far'
) {
return
}
lastJob = {
id: flowStateStore.val[id]?.previewJobId ?? '',
result: flowStateStore.val[id]?.previewResult,
type: 'CompletedJob' as const,
success: flowStateStore.val[id]?.previewSuccess ?? undefined
}
}
$effect(() => {
if (testJob && testJob.type === 'CompletedJob') {
lastJob = $state.snapshot(testJob)
} else if (id) {
updateLastJob(flowStateStore)
}
})
let isConnectingCandidate = $derived(
!!id && !!$flowPropPickerConfig && !!pickableIds && Object.keys(pickableIds).includes(id)
)
@@ -207,6 +181,9 @@
const action = $derived(getAiModuleAction(id))
let testRunDropdownOpen = $state(false)
let outputPickerInner: OutputPickerInner | undefined = $state(undefined)
let historyOpen = $derived.by(() => outputPickerInner?.getHistoryOpen?.() ?? false)
</script>
{#if deletable && id && editId}
@@ -271,7 +248,15 @@
{@const flowStore = flowEditorContext?.flowStore.val}
{@const mod = flowStore?.value ? dfsPreviousResults(id, flowStore, false)[0] : undefined}
{#if mod && flowStateStore?.val?.[id]}
<ModuleTest bind:this={moduleTest} {mod} bind:testIsLoading bind:testJob />
<ModuleTest
bind:this={moduleTest}
{mod}
bind:testIsLoading
bind:testJob
onJobDone={() => {
outputPickerInner?.setJobPreview?.()
}}
/>
{/if}
{/if}
@@ -455,7 +440,6 @@
prefix={'results'}
connectingData={isConnecting ? connectingData : undefined}
{mock}
{lastJob}
{testJob}
moduleId={id}
onSelect={selectConnection}
@@ -463,12 +447,12 @@
{path}
{loopStatus}
rightMargin
bind:derivedHistoryOpen={historyOpen}
historyOffset={{ mainAxis: 12, crossAxis: -9 }}
clazz="p-1"
isLoading={testIsLoading ||
(id ? stepHistoryLoader?.stepStates[id]?.loadingJobs : false)}
initial={id ? stepHistoryLoader?.stepStates[id]?.initial : undefined}
bind:this={outputPickerInner}
/>
{/snippet}
</OutputPicker>
@@ -1,15 +1,12 @@
<script lang="ts">
import { Button } from '$lib/components/common'
import LanguageIcon from '$lib/components/common/languageIcons/LanguageIcon.svelte'
import IconedResourceType from '$lib/components/IconedResourceType.svelte'
import type { FlowModule, FlowStatusModule, Job } from '$lib/gen'
import { Building, Repeat, Square, ArrowDown, GitBranch, Bot } from 'lucide-svelte'
import { createEventDispatcher, getContext } from 'svelte'
import type { Writable } from 'svelte/store'
import FlowModuleSchemaItem from './FlowModuleSchemaItem.svelte'
import FlowModuleIcon from '../FlowModuleIcon.svelte'
import { prettyLanguage } from '$lib/common'
import { msToSec } from '$lib/utils'
import BarsStaggered from '$lib/components/icons/BarsStaggered.svelte'
import FlowJobsMenu from './FlowJobsMenu.svelte'
import {
isTriggerStep,
@@ -185,9 +182,7 @@
{darkMode}
>
{#snippet icon()}
<div>
<Repeat size={16} />
</div>
<FlowModuleIcon module={mod} />
{/snippet}
</FlowModuleSchemaItem>
{:else if mod.value.type === 'branchone'}
@@ -208,9 +203,7 @@
{darkMode}
>
{#snippet icon()}
<div>
<GitBranch size={16} />
</div>
<FlowModuleIcon module={mod} />
{/snippet}
</FlowModuleSchemaItem>
{:else if mod.value.type === 'branchall'}
@@ -231,9 +224,7 @@
{darkMode}
>
{#snippet icon()}
<div>
<GitBranch size={16} />
</div>
<FlowModuleIcon module={mod} />
{/snippet}
</FlowModuleSchemaItem>
{:else}
@@ -257,6 +248,7 @@
{bgColor}
{bgHoverColor}
label={mod.summary ||
(mod.value.type === 'aiagent' ? 'AI Agent' : undefined) ||
(mod.id === 'preprocessor'
? 'Preprocessor'
: mod.id.startsWith('failure')
@@ -281,32 +273,13 @@
{skipped}
>
{#snippet icon()}
<div>
{#if mod.value.type === 'aiagent'}
<Bot size={16} />
{:else if mod.value.type === 'rawscript'}
<LanguageIcon lang={mod.value.language} width={16} height={16} />
{:else if mod.summary == 'Terminate flow'}
<Square size={16} />
{:else if mod.value.type === 'identity'}
<ArrowDown size={16} />
{:else if mod.value.type === 'flow'}
<BarsStaggered size={16} />
{:else if mod.value.type === 'script'}
{#if mod.value.path.startsWith('hub/')}
<div>
<IconedResourceType
width="20px"
height="20px"
name={mod.value.path.split('/')[2]}
silent={true}
/>
</div>
{:else}
<Building size={14} />
{/if}
{/if}
</div>
{@const size =
mod.value.type === 'script' && mod.value.path.startsWith('hub/')
? 20
: mod.value.type === 'script'
? 14
: 16}
<FlowModuleIcon module={mod} {size} />
{/snippet}
</FlowModuleSchemaItem>
{/if}
@@ -146,7 +146,7 @@
{selected}
{hover}
id={id ?? ''}
isConnectingCandidate={true}
isConnectingCandidate={nodeKind !== 'result'}
variant="virtual"
type={outputType}
{darkMode}

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