mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 08:02:19 +00:00
feat: schedule hub scripts (#11330)
* feat: schedule hub scripts Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * test: pin hub script schedule push Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * feat: retry scheduled hub scripts natively Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * chore: address review nits on hub schedules Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * feat: open the hub picker from the script select Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * chore: use a subtle button for the hub row Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
421fdab3d8
commit
733c119fd9
@@ -1007,7 +1007,7 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
));
|
||||
}
|
||||
|
||||
// Native script retry: a failed `Script` job that carries a retry policy and
|
||||
// Native script retry: a failed `Script` or `Script_Hub` job that carries a retry policy and
|
||||
// has attempts left gets its next attempt enqueued here — before the queue
|
||||
// row (which holds the attempt counter) is removed by commit. The failed
|
||||
// attempt is still recorded as a completed job below. `maybe_enqueue_…`
|
||||
@@ -1077,9 +1077,9 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
// Auto-resolve a retry chain that ultimately worked, from whichever of the two
|
||||
// completions lands last (see resolve_retry_chain_if_succeeded): a success that has a
|
||||
// parent (so is a possible retry attempt), or a failure that just enqueued a retry.
|
||||
// `retry_pending` already implies a non-flow-step `Script`.
|
||||
// `retry_pending` already implies a non-flow-step `Script` or `Script_Hub`.
|
||||
let resolve_root = if success && !skipped && !completed_job.is_flow_step() {
|
||||
matches!(completed_job.kind, JobKind::Script)
|
||||
matches!(completed_job.kind, JobKind::Script | JobKind::Script_Hub)
|
||||
.then(|| completed_job.parent_job)
|
||||
.flatten()
|
||||
} else if !success && !skipped && retry_pending {
|
||||
@@ -1843,10 +1843,10 @@ async fn eval_retry_if(
|
||||
false
|
||||
}
|
||||
|
||||
/// Native script retry. When a failed `Script` job carries a retry policy (via
|
||||
/// `runnable_settings_handle`) and has attempts left, enqueue a fresh attempt of
|
||||
/// the same script after the policy's backoff delay — instead of having wrapped
|
||||
/// it in a one-step flow. Each attempt is a real `Script` job; the attempt
|
||||
/// Native script retry. When a failed `Script` or `Script_Hub` job carries a retry
|
||||
/// policy (via `runnable_settings_handle`) and has attempts left, enqueue a fresh
|
||||
/// attempt of the same script after the policy's backoff delay — instead of having
|
||||
/// wrapped it in a one-step flow. Each attempt is a job of the same kind; the attempt
|
||||
/// counter lives in the `native_retry_attempt` marker, written here and read only
|
||||
/// on the next failure (never on the hot job-pull path).
|
||||
///
|
||||
@@ -1867,7 +1867,10 @@ pub async fn maybe_enqueue_native_script_retry(
|
||||
result_fn: &(dyn Fn() -> Option<Box<serde_json::value::RawValue>> + Sync),
|
||||
) -> Result<bool, Error> {
|
||||
// Only plain top-level scripts retry natively; cancellation always wins.
|
||||
if canceled_by.is_some() || !matches!(job.kind, JobKind::Script) || job.is_flow_step() {
|
||||
if canceled_by.is_some()
|
||||
|| !matches!(job.kind, JobKind::Script | JobKind::Script_Hub)
|
||||
|| job.is_flow_step()
|
||||
{
|
||||
return Ok(false);
|
||||
}
|
||||
|
||||
@@ -6154,10 +6157,12 @@ async fn push_inner<'c, 'd>(
|
||||
// `quickjs` feature it cannot be evaluated and fails closed (no retry);
|
||||
// the flow path is not a fallback, since the flow runtime needs quickjs
|
||||
// too.
|
||||
// A hub script has no hash and runs as a `Script_Hub` job.
|
||||
let is_hub = hash.is_none() && path.starts_with("hub/");
|
||||
let native_retry = !is_flow
|
||||
&& skip_handler.is_none()
|
||||
&& error_handler_path.is_none()
|
||||
&& hash.is_some()
|
||||
&& (hash.is_some() || is_hub)
|
||||
&& language.is_some()
|
||||
&& windmill_common::runnable_settings::min_version_supports_runnable_settings_v0()
|
||||
.await;
|
||||
@@ -6187,7 +6192,11 @@ async fn push_inner<'c, 'd>(
|
||||
break 'ssf JobPayloadUntagged {
|
||||
runnable_id: hash.map(|h| h.0),
|
||||
runnable_path: Some(path),
|
||||
job_kind: JobKind::Script,
|
||||
job_kind: if is_hub {
|
||||
JobKind::Script_Hub
|
||||
} else {
|
||||
JobKind::Script
|
||||
},
|
||||
language,
|
||||
dedicated_worker,
|
||||
concurrency_settings,
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
* LICENSE-AGPL for a copy of the license.
|
||||
*/
|
||||
|
||||
use crate::jobs::HTTP_CLIENT;
|
||||
use crate::push;
|
||||
use crate::PushIsolationLevel;
|
||||
use anyhow::Context;
|
||||
@@ -25,6 +26,7 @@ use windmill_common::jobs::OnBehalfOf;
|
||||
use windmill_common::runnable_settings::ConcurrencySettings;
|
||||
use windmill_common::runnable_settings::DebouncingSettings;
|
||||
use windmill_common::schedule::schedule_to_user;
|
||||
use windmill_common::scripts::get_full_hub_script_by_path;
|
||||
use windmill_common::scripts::ScriptHash;
|
||||
use windmill_common::triggers::TriggerMetadata;
|
||||
use windmill_common::utils::WarnAfterExt;
|
||||
@@ -81,6 +83,8 @@ async fn get_schedule_metadata<'c>(
|
||||
Some(version),
|
||||
parsed_retry,
|
||||
))
|
||||
} else if schedule.script_path.starts_with("hub/") {
|
||||
Ok((None, None, None, None, None, parsed_retry))
|
||||
} else {
|
||||
let (
|
||||
hash,
|
||||
@@ -322,6 +326,48 @@ pub async fn push_scheduled_job<'c>(
|
||||
None,
|
||||
on_behalf_of,
|
||||
)
|
||||
} else if schedule.script_path.starts_with("hub/") {
|
||||
let tag = schedule.tag.clone().filter(|t| !t.is_empty());
|
||||
let payload = match &schedule.retry {
|
||||
// The language is what lets `push` materialize this as a native retry
|
||||
// instead of a flow wrapper, which would queue the next tick at start.
|
||||
Some(retry) => JobPayload::SingleStepFlow {
|
||||
path: schedule.script_path.clone(),
|
||||
hash: None,
|
||||
flow_version: None,
|
||||
language: Some(
|
||||
get_full_hub_script_by_path(
|
||||
StripPath(schedule.script_path.clone()),
|
||||
&HTTP_CLIENT,
|
||||
Some(db),
|
||||
)
|
||||
.await?
|
||||
.language,
|
||||
),
|
||||
retry: Some(serde_json::from_value::<Retry>(retry.clone()).map_err(|e| {
|
||||
error::Error::internal_err(format!(
|
||||
"Unable to parse retry information from schedule: {e}"
|
||||
))
|
||||
})?),
|
||||
error_handler_path: None,
|
||||
error_handler_args: None,
|
||||
skip_handler: None,
|
||||
args: args.clone(),
|
||||
cache_ttl: None,
|
||||
cache_ignore_s3_path: None,
|
||||
priority: None,
|
||||
tag_override: tag.clone(),
|
||||
trigger_path: None,
|
||||
apply_preprocessor: false,
|
||||
concurrency_settings: ConcurrencySettings::default(),
|
||||
debouncing_settings: DebouncingSettings::default(),
|
||||
},
|
||||
None => JobPayload::ScriptHub {
|
||||
path: schedule.script_path.clone(),
|
||||
apply_preprocessor: false,
|
||||
},
|
||||
};
|
||||
(payload, tag, None, None)
|
||||
} else {
|
||||
let (
|
||||
hash,
|
||||
|
||||
@@ -1757,4 +1757,59 @@ mod schedule_push {
|
||||
assert_eq!(count_queued_jobs(&db).await, 1);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
|
||||
async fn test_push_hub_script_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
// Seed the hub cache so the push resolves the script without the network.
|
||||
let version = "990000001";
|
||||
let hub_dir = &*windmill_common::worker::HUB_CACHE_DIR;
|
||||
tokio::fs::create_dir_all(hub_dir).await?;
|
||||
tokio::fs::write(
|
||||
format!("{hub_dir}/{version}"),
|
||||
r#"{"content":"echo hi","lockfile":null,"language":"bash","schema":{},"summary":null}"#,
|
||||
)
|
||||
.await?;
|
||||
let hub_path = format!("hub/{version}/test/echo");
|
||||
|
||||
// A retry must stay a native `script_hub` job: a flow wrapper would queue the
|
||||
// next tick at start and let slow runs overlap.
|
||||
let retry = serde_json::json!({ "constant": { "attempts": 2, "seconds": 1 } });
|
||||
for (path, retry, dynamic_skip) in [
|
||||
("f/system/hub_plain", None, None),
|
||||
("f/system/hub_retry", Some(retry), None),
|
||||
// A skip handler needs a flow wrapper, so its metadata lookup must not
|
||||
// go through the `script` table.
|
||||
("f/system/hub_skip", None, Some("f/system/skip".to_string())),
|
||||
] {
|
||||
let has_retry = retry.is_some();
|
||||
let has_skip = dynamic_skip.is_some();
|
||||
let schedule = make_schedule(|s| {
|
||||
s.path = path.to_string();
|
||||
s.script_path = hub_path.clone();
|
||||
s.retry = retry;
|
||||
s.dynamic_skip = dynamic_skip;
|
||||
});
|
||||
let tx = db.begin().await?;
|
||||
let tx = push_scheduled_job(&db, tx, &schedule, Some(&make_authed()), None).await?;
|
||||
tx.commit().await?;
|
||||
|
||||
let (job_kind, runnable_path, language, handle) =
|
||||
sqlx::query_as::<_, (String, Option<String>, Option<String>, Option<i64>)>(
|
||||
"SELECT j.kind::text, j.runnable_path, j.script_lang::text, q.runnable_settings_handle
|
||||
FROM v2_job j JOIN v2_job_queue q ON j.id = q.id WHERE j.trigger = $1",
|
||||
)
|
||||
.bind(path)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
assert_eq!(runnable_path.as_deref(), Some(hub_path.as_str()));
|
||||
if has_skip {
|
||||
assert_eq!(job_kind, "singlestepflow");
|
||||
} else {
|
||||
assert_eq!(job_kind, "script_hub");
|
||||
assert_eq!(language.as_deref(), Some("bash"));
|
||||
assert_eq!(handle.is_some(), has_retry);
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,7 +12,9 @@
|
||||
import FlowPathViewer from './flows/content/FlowPathViewer.svelte'
|
||||
import ToggleButton from './common/toggleButton-v2/ToggleButton.svelte'
|
||||
import ToggleButtonGroup from './common/toggleButton-v2/ToggleButtonGroup.svelte'
|
||||
import { Code, Code2, ExternalLink, Pen, RefreshCw } from 'lucide-svelte'
|
||||
import { Code, Code2, ExternalLink, Globe2, Pen, RefreshCw } from 'lucide-svelte'
|
||||
import PickHubScript from './flows/pickers/PickHubScript.svelte'
|
||||
import { disableHubStore } from '$lib/stores'
|
||||
import type { SupportedLanguage } from '$lib/common'
|
||||
import FlowIcon from './home/FlowIcon.svelte'
|
||||
import DarkModeObserver from './DarkModeObserver.svelte'
|
||||
@@ -32,6 +34,9 @@
|
||||
allowEdit?: boolean
|
||||
allowView?: boolean
|
||||
clearable?: boolean
|
||||
/** Offer picking a script from the Hub; the picked path is `hub/...`. Browses plain
|
||||
* scripts only, whatever `kinds` is. */
|
||||
allowHub?: boolean
|
||||
/** Workspace to list runnables from. Defaults to the operating workspace (see
|
||||
* `useOperatingWorkspace`). */
|
||||
workspace?: string
|
||||
@@ -48,9 +53,15 @@
|
||||
allowEdit = true,
|
||||
allowView = true,
|
||||
clearable = false,
|
||||
allowHub = false,
|
||||
workspace = undefined
|
||||
}: Props = $props()
|
||||
|
||||
let isHubPath = $derived(itemKind == 'script' && !!scriptPath?.startsWith('hub/'))
|
||||
let drawerHub: Drawer | undefined = $state()
|
||||
let filterText = $state('')
|
||||
let hubFilter = $state('')
|
||||
|
||||
let effectiveWorkspace = $derived(workspace ?? $operatingWorkspace)
|
||||
// Edit/View routes open in the workspace listed here, not wherever the tab lands.
|
||||
let wsParam = $derived(
|
||||
@@ -115,6 +126,21 @@
|
||||
</DrawerContent>
|
||||
</Drawer>
|
||||
|
||||
{#if allowHub}
|
||||
<Drawer bind:this={drawerHub} size="900px">
|
||||
<DrawerContent title="Pick a Hub script" on:close={drawerHub.closeDrawer}>
|
||||
<PickHubScript
|
||||
bind:filter={hubFilter}
|
||||
on:pick={(e) => {
|
||||
scriptPath = e.detail.path
|
||||
dispatch('select', { path: e.detail.path, itemKind })
|
||||
drawerHub?.closeDrawer()
|
||||
}}
|
||||
/>
|
||||
</DrawerContent>
|
||||
</Drawer>
|
||||
{/if}
|
||||
|
||||
<div class="flex flex-row items-center gap-1 w-full">
|
||||
{#if options.length > 1}
|
||||
<div>
|
||||
@@ -145,9 +171,11 @@
|
||||
}
|
||||
}
|
||||
class="grow shrink max-w-full"
|
||||
{items}
|
||||
items={isHubPath ? [{ value: scriptPath!, label: scriptPath! }, ...items] : items}
|
||||
{clearable}
|
||||
bind:filterText
|
||||
placeholder="Pick {itemKind === 'app' ? 'an' : 'a'} {itemKind}"
|
||||
bottomSnippet={allowHub && itemKind == 'script' && !$disableHubStore ? hubHint : undefined}
|
||||
/>
|
||||
{/if}
|
||||
|
||||
@@ -212,7 +240,7 @@
|
||||
</div>
|
||||
{:else}
|
||||
<div class="flex gap-2">
|
||||
{#if allowEdit}
|
||||
{#if allowEdit && !isHubPath}
|
||||
<Button
|
||||
startIcon={{ icon: Pen }}
|
||||
target="_blank"
|
||||
@@ -245,3 +273,21 @@
|
||||
{/if}
|
||||
{/if}
|
||||
</div>
|
||||
|
||||
{#snippet hubHint({ close }: { close: () => void })}
|
||||
<Button
|
||||
variant="subtle"
|
||||
size="xs2"
|
||||
startIcon={{ icon: Globe2 }}
|
||||
wrapperClasses="w-full border-t border-border-light"
|
||||
btnClasses="w-full rounded-none font-normal"
|
||||
onClick={() => {
|
||||
// Read before close(): Select clears its filter text when the list closes.
|
||||
hubFilter = filterText
|
||||
close()
|
||||
drawerHub?.openDrawer()
|
||||
}}
|
||||
>
|
||||
Browse Hub scripts
|
||||
</Button>
|
||||
{/snippet}
|
||||
|
||||
@@ -13,6 +13,7 @@
|
||||
import LabelsInput from '$lib/components/LabelsInput.svelte'
|
||||
import Required from '$lib/components/Required.svelte'
|
||||
import ScriptPicker from '$lib/components/ScriptPicker.svelte'
|
||||
import { loadSchema } from '$lib/infer'
|
||||
import PipelineLockedRunnableInfo from '$lib/components/triggers/PipelineLockedRunnableInfo.svelte'
|
||||
import ErrorOrRecoveryHandler, {
|
||||
handlerFullPath
|
||||
@@ -112,7 +113,7 @@
|
||||
// already-bound script. We swap the runnable ScriptPicker for a read-only
|
||||
// viewer so the trigger can't be silently reassigned off the pipeline.
|
||||
let fixedScriptPath = $state('')
|
||||
let runnable: Script | Flow | undefined = $state()
|
||||
let runnable: Pick<Script | Flow, 'schema'> | undefined = $state()
|
||||
let args: Record<string, any> = $state({})
|
||||
let loading = $state(false)
|
||||
let drawerLoading = $state(true)
|
||||
@@ -418,6 +419,8 @@
|
||||
try {
|
||||
if (is_flow) {
|
||||
runnable = await FlowService.getFlowByPath({ workspace: wsId!, path: p })
|
||||
} else if (p.startsWith('hub/')) {
|
||||
runnable = await loadSchema(wsId!, p, 'hubscript')
|
||||
} else {
|
||||
runnable = await ScriptService.getScriptByPath({ workspace: wsId!, path: p })
|
||||
}
|
||||
@@ -980,6 +983,7 @@
|
||||
initialPath={initialScriptPath}
|
||||
kinds={['script']}
|
||||
allowFlow={true}
|
||||
allowHub={true}
|
||||
allowRefresh={can_write}
|
||||
bind:itemKind
|
||||
bind:scriptPath={script_path}
|
||||
|
||||
Reference in New Issue
Block a user