From 105a5fc08f8bf934ce3a3a33c967c30c372c043c Mon Sep 17 00:00:00 2001 From: Steven Enamakel <31011319+senamakel@users.noreply.github.com> Date: Wed, 11 Feb 2026 02:50:50 +0530 Subject: [PATCH] Fix/tuesday (#97) * Update submodule URL for skills and refresh commit reference - Changed the submodule URL for the "skills" repository to use SSH instead of HTTPS. - Updated the subproject commit reference in the skills directory to the latest commit. * Update subproject commit reference for skills to latest version b5d138f * updated skills * Update ping interval and refactor context usage in QjsSkillInstance - Reduced the PING_INTERVAL from 60 seconds to 5 seconds for more frequent pings. - Refactored the context parameter in lifecycle calls and message handling functions to use the runtime context, improving consistency and performance. - Updated subproject commit reference for skills to the latest version. --- .gitmodules | 2 +- skills | 2 +- src-tauri/src/runtime/ping_scheduler.rs | 2 +- src-tauri/src/runtime/qjs_skill_instance.rs | 480 +++++++++++++++++--- 4 files changed, 417 insertions(+), 69 deletions(-) diff --git a/.gitmodules b/.gitmodules index f0987d0b8..3c583d071 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,3 +1,3 @@ [submodule "skills"] path = skills - url = https://github.com/alphahumanxyz/skills + url = git@github.com:vezuresdotxyz/alphahuman-skills.git diff --git a/skills b/skills index 91b4f1f7e..640518e8c 160000 --- a/skills +++ b/skills @@ -1 +1 @@ -Subproject commit 91b4f1f7e7918c3eded06a53bcb87a354c972b9f +Subproject commit 640518e8cc3792c8f0a5565294ff791d75f375d5 diff --git a/src-tauri/src/runtime/ping_scheduler.rs b/src-tauri/src/runtime/ping_scheduler.rs index a6fa681ca..fb63646db 100644 --- a/src-tauri/src/runtime/ping_scheduler.rs +++ b/src-tauri/src/runtime/ping_scheduler.rs @@ -23,7 +23,7 @@ use crate::runtime::skill_registry::SkillRegistry; use crate::runtime::types::{events, SkillMessage, SkillStatus}; /// Interval between ping sweeps. -const PING_INTERVAL: Duration = Duration::from_secs(60); +const PING_INTERVAL: Duration = Duration::from_secs(5); /// Per-skill timeout when waiting for a ping reply. const PING_TIMEOUT: Duration = Duration::from_secs(30); diff --git a/src-tauri/src/runtime/qjs_skill_instance.rs b/src-tauri/src/runtime/qjs_skill_instance.rs index 2fa13af5a..773001f09 100644 --- a/src-tauri/src/runtime/qjs_skill_instance.rs +++ b/src-tauri/src/runtime/qjs_skill_instance.rs @@ -230,7 +230,7 @@ impl QjsSkillInstance { restore_oauth_credential(&ctx, &config.skill_id).await; // Call init() lifecycle - if let Err(e) = call_lifecycle(&ctx, "init").await { + if let Err(e) = call_lifecycle(&rt, &ctx, "init").await { let mut s = state.write(); s.status = SkillStatus::Error; s.error = Some(format!("init() failed: {e}")); @@ -242,7 +242,7 @@ impl QjsSkillInstance { drive_jobs(&rt).await; // Call start() lifecycle - if let Err(e) = call_lifecycle(&ctx, "start").await { + if let Err(e) = call_lifecycle(&rt, &ctx, "start").await { let mut s = state.write(); s.status = SkillStatus::Error; s.error = Some(format!("start() failed: {e}")); @@ -258,7 +258,7 @@ impl QjsSkillInstance { log::info!("[skill:{}] Running (QuickJS)", config.skill_id); // Immediate ping to verify the connection is healthy - match handle_js_call(&ctx, "onPing", "{}").await { + match handle_js_call(&rt, &ctx, "onPing", "{}").await { Ok(value) => { log::info!("[skill:{}] Initial ping result: {}", config.skill_id, value); } @@ -329,7 +329,7 @@ async fn run_event_loop( // messages (events, pings, etc.) but queue any new CallTool. match rx.try_recv() { Ok(msg) => { - let should_stop = handle_message(ctx, msg, state, skill_id, &mut pending_tool).await; + let should_stop = handle_message(rt, ctx, msg, state, skill_id, &mut pending_tool).await; if should_stop { break; } @@ -439,6 +439,7 @@ async fn fire_timer_callback(ctx: &rquickjs::AsyncContext, timer_id: u32) { /// Handle a single message from the channel. /// Returns true if the skill should stop. async fn handle_message( + rt: &rquickjs::AsyncRuntime, ctx: &rquickjs::AsyncContext, msg: SkillMessage, state: &Arc>, @@ -470,13 +471,13 @@ async fn handle_message( } } SkillMessage::ServerEvent { event, data } => { - let _ = handle_server_event(ctx, &event, data).await; + let _ = handle_server_event(rt, ctx, &event, data).await; } SkillMessage::CronTrigger { schedule_id } => { - let _ = handle_cron_trigger(ctx, &schedule_id).await; + let _ = handle_cron_trigger(rt, ctx, &schedule_id).await; } SkillMessage::Stop { reply } => { - let _ = call_lifecycle(ctx, "stop").await; + let _ = call_lifecycle(rt, ctx, "stop").await; // Clear OAuth credential from memory and mark as disconnected in store let clear_code = r#"(function() { @@ -497,39 +498,39 @@ async fn handle_message( return true; // Signal to stop } SkillMessage::SetupStart { reply } => { - let result = handle_js_call(ctx, "onSetupStart", "{}").await; + let result = handle_js_call(rt, ctx, "onSetupStart", "{}").await; let _ = reply.send(result); } SkillMessage::SetupSubmit { step_id, values, reply } => { let args = serde_json::json!({ "stepId": step_id, "values": values }); - let result = handle_js_call(ctx, "onSetupSubmit", &args.to_string()).await; + let result = handle_js_call(rt, ctx, "onSetupSubmit", &args.to_string()).await; let _ = reply.send(result); } SkillMessage::SetupCancel { reply } => { - let result = handle_js_void_call(ctx, "onSetupCancel", "{}").await; + let result = handle_js_void_call(rt, ctx, "onSetupCancel", "{}").await; let _ = reply.send(result); } SkillMessage::ListOptions { reply } => { - let result = handle_js_call(ctx, "onListOptions", "{}").await; + let result = handle_js_call(rt, ctx, "onListOptions", "{}").await; let _ = reply.send(result); } SkillMessage::SetOption { name, value, reply } => { let args = serde_json::json!({ "name": name, "value": value }); - let result = handle_js_void_call(ctx, "onSetOption", &args.to_string()).await; + let result = handle_js_void_call(rt, ctx, "onSetOption", &args.to_string()).await; let _ = reply.send(result); } SkillMessage::SessionStart { session_id, reply } => { let args = serde_json::json!({ "sessionId": session_id }); - let result = handle_js_void_call(ctx, "onSessionStart", &args.to_string()).await; + let result = handle_js_void_call(rt, ctx, "onSessionStart", &args.to_string()).await; let _ = reply.send(result); } SkillMessage::SessionEnd { session_id, reply } => { let args = serde_json::json!({ "sessionId": session_id }); - let result = handle_js_void_call(ctx, "onSessionEnd", &args.to_string()).await; + let result = handle_js_void_call(rt, ctx, "onSessionEnd", &args.to_string()).await; let _ = reply.send(result); } SkillMessage::Tick { reply } => { - let result = handle_js_void_call(ctx, "onTick", "{}").await; + let result = handle_js_void_call(rt, ctx, "onTick", "{}").await; let _ = reply.send(result); } SkillMessage::Error { error_type, message, source, recoverable } => { @@ -539,7 +540,7 @@ async fn handle_message( "source": source, "recoverable": recoverable, }); - if let Err(e) = handle_js_void_call(ctx, "onError", &args.to_string()).await { + if let Err(e) = handle_js_void_call(rt, ctx, "onError", &args.to_string()).await { log::warn!("[skill:{}] onError() handler failed: {e}", skill_id); } } @@ -564,10 +565,10 @@ async fn handle_message( }).await; log::info!("[skill:{}] OAuth credential set and persisted to store", skill_id); let params_str = serde_json::to_string(¶ms).unwrap_or_else(|_| "{}".to_string()); - handle_js_call(ctx, "onOAuthComplete", ¶ms_str).await + handle_js_call(rt, ctx, "onOAuthComplete", ¶ms_str).await } "skill/ping" => { - handle_js_call(ctx, "onPing", "{}").await + handle_js_call(rt, ctx, "onPing", "{}").await } "oauth/revoked" => { // Clear credential: set to empty string so it's clearly "disconnected" @@ -584,12 +585,12 @@ async fn handle_message( }).await; log::info!("[skill:{}] OAuth credential cleared from store", skill_id); let params_str = serde_json::to_string(¶ms).unwrap_or_else(|_| "{}".to_string()); - handle_js_void_call(ctx, "onOAuthRevoked", ¶ms_str).await + handle_js_void_call(rt, ctx, "onOAuthRevoked", ¶ms_str).await .map(|()| serde_json::json!({ "ok": true })) } _ => { let args = serde_json::json!({ "method": method, "params": params }); - handle_js_call(ctx, "onRpc", &args.to_string()).await + handle_js_call(rt, ctx, "onRpc", &args.to_string()).await } }; let _ = reply.send(result); @@ -638,29 +639,97 @@ fn format_js_exception(js_ctx: &rquickjs::Ctx<'_>, err: &rquickjs::Error) -> Str } /// Call a lifecycle function on the skill object. -/// Looks for the skill at globalThis.__skill.default first, then falls back to globalThis. -async fn call_lifecycle(ctx: &rquickjs::AsyncContext, name: &str) -> Result<(), String> { +/// Handles async (Promise-returning) lifecycle methods (init, start, stop). +async fn call_lifecycle( + rt: &rquickjs::AsyncRuntime, + ctx: &rquickjs::AsyncContext, + name: &str, +) -> Result<(), String> { let name = name.to_string(); - ctx.with(|js_ctx| { + let is_promise = ctx.with(|js_ctx| { let code = format!( r#"(function() {{ var skill = globalThis.__skill && globalThis.__skill.default ? globalThis.__skill.default : (globalThis.__skill || globalThis); - if (typeof skill.{name} === 'function') {{ - skill.{name}(); - }} else if (typeof globalThis.{name} === 'function') {{ - globalThis.{name}(); + var fn = skill.{name} || globalThis.{name}; + if (typeof fn === 'function') {{ + var result = fn.call(skill); + if (result && typeof result.then === 'function') {{ + globalThis.__pendingLifecycleDone = false; + globalThis.__pendingLifecycleError = undefined; + result.then( + function() {{ globalThis.__pendingLifecycleDone = true; }}, + function(e) {{ + globalThis.__pendingLifecycleError = e && e.message ? e.message : String(e); + globalThis.__pendingLifecycleDone = true; + }} + ); + return "1"; + }} }} + return "0"; }})()"# ); - js_ctx.eval::(code.as_bytes()) - .map_err(|e| { + match js_ctx.eval::(code.as_bytes()) { + Ok(s) => Ok(s == "1"), + Err(e) => { let detail = format_js_exception(&js_ctx, &e); - format!("{name}() failed: {detail}") - })?; - Ok(()) - }).await + Err(format!("{name}() failed: {detail}")) + } + } + }).await?; + + if is_promise { + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + + loop { + drive_jobs(rt).await; + + let done = ctx.with(|js_ctx| { + js_ctx + .eval::(b"globalThis.__pendingLifecycleDone === true") + .unwrap_or(false) + }).await; + + if done { + let error = ctx.with(|js_ctx| { + let has_error = js_ctx + .eval::(b"globalThis.__pendingLifecycleError !== undefined") + .unwrap_or(false); + let err = if has_error { + Some(js_ctx + .eval::(b"String(globalThis.__pendingLifecycleError)") + .unwrap_or_else(|_| "Unknown error".to_string())) + } else { + None + }; + let _ = js_ctx.eval::( + b"delete globalThis.__pendingLifecycleDone; delete globalThis.__pendingLifecycleError;" + ); + err + }).await; + + if let Some(err_msg) = error { + return Err(format!("{name}() rejected: {err_msg}")); + } + return Ok(()); + } + + if tokio::time::Instant::now() >= deadline { + ctx.with(|js_ctx| { + let _ = js_ctx.eval::( + b"delete globalThis.__pendingLifecycleDone; delete globalThis.__pendingLifecycleError;" + ); + }).await; + return Err(format!("{name}() timed out after 30s")); + } + + tokio::time::sleep(Duration::from_millis(5)).await; + } + } + + Ok(()) } /// Extract tool definitions from the skill. @@ -821,7 +890,9 @@ async fn read_pending_tool_result( } /// Handle a server event. +/// Handles async (Promise-returning) onServerEvent handlers. async fn handle_server_event( + rt: &rquickjs::AsyncRuntime, ctx: &rquickjs::AsyncContext, event: &str, data: serde_json::Value, @@ -829,56 +900,189 @@ async fn handle_server_event( let data_str = serde_json::to_string(&data).unwrap_or_else(|_| "null".to_string()); let event = event.to_string(); - ctx.with(|js_ctx| { + let is_promise = ctx.with(|js_ctx| { let code = format!( r#"(function() {{ var skill = globalThis.__skill && globalThis.__skill.default ? globalThis.__skill.default : (globalThis.__skill || globalThis); - if (typeof skill.onServerEvent === 'function') {{ - skill.onServerEvent("{}", {}); - }} else if (typeof globalThis.onServerEvent === 'function') {{ - globalThis.onServerEvent("{}", {}); + var fn = skill.onServerEvent || globalThis.onServerEvent; + if (typeof fn === 'function') {{ + var result = fn.call(skill, "{}", {}); + if (result && typeof result.then === 'function') {{ + globalThis.__pendingEventDone = false; + globalThis.__pendingEventError = undefined; + result.then( + function() {{ globalThis.__pendingEventDone = true; }}, + function(e) {{ + globalThis.__pendingEventError = e && e.message ? e.message : String(e); + globalThis.__pendingEventDone = true; + }} + ); + return "1"; + }} }} + return "0"; }})()"#, event.replace('"', r#"\""#), data_str, - event.replace('"', r#"\""#), - data_str, ); - js_ctx.eval::(code.as_bytes()) - .map_err(|e| format!("Event handler failed: {e}"))?; - Ok(()) - }).await + match js_ctx.eval::(code.as_bytes()) { + Ok(s) => Ok(s == "1"), + Err(e) => Err(format!("Event handler failed: {e}")) + } + }).await?; + + if is_promise { + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + + loop { + drive_jobs(rt).await; + + let done = ctx.with(|js_ctx| { + js_ctx + .eval::(b"globalThis.__pendingEventDone === true") + .unwrap_or(false) + }).await; + + if done { + let error = ctx.with(|js_ctx| { + let has_error = js_ctx + .eval::(b"globalThis.__pendingEventError !== undefined") + .unwrap_or(false); + let err = if has_error { + Some(js_ctx + .eval::(b"String(globalThis.__pendingEventError)") + .unwrap_or_else(|_| "Unknown error".to_string())) + } else { + None + }; + let _ = js_ctx.eval::( + b"delete globalThis.__pendingEventDone; delete globalThis.__pendingEventError;" + ); + err + }).await; + + if let Some(err_msg) = error { + return Err(format!("onServerEvent() rejected: {err_msg}")); + } + return Ok(()); + } + + if tokio::time::Instant::now() >= deadline { + ctx.with(|js_ctx| { + let _ = js_ctx.eval::( + b"delete globalThis.__pendingEventDone; delete globalThis.__pendingEventError;" + ); + }).await; + return Err("onServerEvent() timed out after 30s".to_string()); + } + + tokio::time::sleep(Duration::from_millis(5)).await; + } + } + + Ok(()) } /// Handle a cron trigger. -async fn handle_cron_trigger(ctx: &rquickjs::AsyncContext, schedule_id: &str) -> Result<(), String> { +/// Handles async (Promise-returning) onCronTrigger handlers. +async fn handle_cron_trigger( + rt: &rquickjs::AsyncRuntime, + ctx: &rquickjs::AsyncContext, + schedule_id: &str, +) -> Result<(), String> { let schedule_id = schedule_id.to_string(); - ctx.with(|js_ctx| { + + let is_promise = ctx.with(|js_ctx| { let code = format!( r#"(function() {{ var skill = globalThis.__skill && globalThis.__skill.default ? globalThis.__skill.default : (globalThis.__skill || globalThis); - if (typeof skill.onCronTrigger === 'function') {{ - skill.onCronTrigger("{}"); - }} else if (typeof globalThis.onCronTrigger === 'function') {{ - globalThis.onCronTrigger("{}"); + var fn = skill.onCronTrigger || globalThis.onCronTrigger; + if (typeof fn === 'function') {{ + var result = fn.call(skill, "{}"); + if (result && typeof result.then === 'function') {{ + globalThis.__pendingCronDone = false; + globalThis.__pendingCronError = undefined; + result.then( + function() {{ globalThis.__pendingCronDone = true; }}, + function(e) {{ + globalThis.__pendingCronError = e && e.message ? e.message : String(e); + globalThis.__pendingCronDone = true; + }} + ); + return "1"; + }} }} + return "0"; }})()"#, schedule_id.replace('"', r#"\""#), - schedule_id.replace('"', r#"\""#), ); - js_ctx.eval::(code.as_bytes()) - .map_err(|e| format!("Cron trigger failed: {e}")) - .map(|_| ()) - }).await + match js_ctx.eval::(code.as_bytes()) { + Ok(s) => Ok(s == "1"), + Err(e) => Err(format!("Cron trigger failed: {e}")) + } + }).await?; + + if is_promise { + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + + loop { + drive_jobs(rt).await; + + let done = ctx.with(|js_ctx| { + js_ctx + .eval::(b"globalThis.__pendingCronDone === true") + .unwrap_or(false) + }).await; + + if done { + let error = ctx.with(|js_ctx| { + let has_error = js_ctx + .eval::(b"globalThis.__pendingCronError !== undefined") + .unwrap_or(false); + let err = if has_error { + Some(js_ctx + .eval::(b"String(globalThis.__pendingCronError)") + .unwrap_or_else(|_| "Unknown error".to_string())) + } else { + None + }; + let _ = js_ctx.eval::( + b"delete globalThis.__pendingCronDone; delete globalThis.__pendingCronError;" + ); + err + }).await; + + if let Some(err_msg) = error { + return Err(format!("onCronTrigger() rejected: {err_msg}")); + } + return Ok(()); + } + + if tokio::time::Instant::now() >= deadline { + ctx.with(|js_ctx| { + let _ = js_ctx.eval::( + b"delete globalThis.__pendingCronDone; delete globalThis.__pendingCronError;" + ); + }).await; + return Err("onCronTrigger() timed out after 30s".to_string()); + } + + tokio::time::sleep(Duration::from_millis(5)).await; + } + } + + Ok(()) } /// Call a JS function on the skill object that returns a JSON value. +/// Handles both sync and async (Promise-returning) functions. async fn handle_js_call( + rt: &rquickjs::AsyncRuntime, ctx: &rquickjs::AsyncContext, fn_name: &str, args_json: &str, @@ -894,8 +1098,23 @@ async fn handle_js_call( : (globalThis.__skill || globalThis); var fn = skill.{fn_name} || globalThis.{fn_name}; if (typeof fn === 'function') {{ - var args = {args_json}; - var result = fn.call(skill, args); + var result = fn.call(skill, {args_json}); + if (result && typeof result.then === 'function') {{ + globalThis.__pendingRpcResult = undefined; + globalThis.__pendingRpcError = undefined; + globalThis.__pendingRpcDone = false; + result.then( + function(v) {{ + globalThis.__pendingRpcResult = v; + globalThis.__pendingRpcDone = true; + }}, + function(e) {{ + globalThis.__pendingRpcError = e && e.message ? e.message : String(e); + globalThis.__pendingRpcDone = true; + }} + ); + return "__PROMISE__"; + }} return JSON.stringify(result === undefined ? null : result); }} return "null"; @@ -911,12 +1130,74 @@ async fn handle_js_call( } }).await?; - serde_json::from_str(&result_text) - .map_err(|e| format!("{fn_name}() returned invalid JSON: {e}")) + if result_text == "__PROMISE__" { + // Async — drive the QuickJS job queue until the promise resolves + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + + loop { + drive_jobs(rt).await; + + let done = ctx.with(|js_ctx| { + js_ctx + .eval::(b"globalThis.__pendingRpcDone === true") + .unwrap_or(false) + }).await; + + if done { + let result = ctx.with(|js_ctx| { + let has_error = js_ctx + .eval::(b"globalThis.__pendingRpcError !== undefined") + .unwrap_or(false); + + let val = if has_error { + let err_msg = js_ctx + .eval::(b"String(globalThis.__pendingRpcError)") + .unwrap_or_else(|_| "Unknown error".to_string()); + Err(format!("{fn_name}() rejected: {err_msg}")) + } else { + let json_str = js_ctx + .eval::( + b"JSON.stringify(globalThis.__pendingRpcResult === undefined ? null : globalThis.__pendingRpcResult)" + ) + .unwrap_or_else(|_| "null".to_string()); + Ok(json_str) + }; + + // Clean up globals + let _ = js_ctx.eval::( + b"delete globalThis.__pendingRpcDone; delete globalThis.__pendingRpcResult; delete globalThis.__pendingRpcError;" + ); + val + }).await; + + return match result { + Ok(json_str) => serde_json::from_str(&json_str) + .map_err(|e| format!("{fn_name}() returned invalid JSON: {e}")), + Err(e) => Err(e), + }; + } + + if tokio::time::Instant::now() >= deadline { + ctx.with(|js_ctx| { + let _ = js_ctx.eval::( + b"delete globalThis.__pendingRpcDone; delete globalThis.__pendingRpcResult; delete globalThis.__pendingRpcError;" + ); + }).await; + return Err(format!("{fn_name}() timed out after 30s")); + } + + tokio::time::sleep(Duration::from_millis(5)).await; + } + } else { + serde_json::from_str(&result_text) + .map_err(|e| format!("{fn_name}() returned invalid JSON: {e}")) + } } /// Call a JS function on the skill object that returns void. +/// Handles both sync and async (Promise-returning) functions. async fn handle_js_void_call( + rt: &rquickjs::AsyncRuntime, ctx: &rquickjs::AsyncContext, fn_name: &str, args_json: &str, @@ -924,7 +1205,7 @@ async fn handle_js_void_call( let fn_name = fn_name.to_string(); let args_json = args_json.to_string(); - ctx.with(|js_ctx| { + let is_promise = ctx.with(|js_ctx| { let code = format!( r#"(function() {{ var skill = globalThis.__skill && globalThis.__skill.default @@ -932,16 +1213,83 @@ async fn handle_js_void_call( : (globalThis.__skill || globalThis); var fn = skill.{fn_name} || globalThis.{fn_name}; if (typeof fn === 'function') {{ - var args = {args_json}; - fn.call(skill, args); + var result = fn.call(skill, {args_json}); + if (result && typeof result.then === 'function') {{ + globalThis.__pendingRpcVoidDone = false; + globalThis.__pendingRpcVoidError = undefined; + result.then( + function() {{ globalThis.__pendingRpcVoidDone = true; }}, + function(e) {{ + globalThis.__pendingRpcVoidError = e && e.message ? e.message : String(e); + globalThis.__pendingRpcVoidDone = true; + }} + ); + return "1"; + }} }} + return "0"; }})()"# ); - js_ctx.eval::(code.as_bytes()) - .map_err(|e| format!("{fn_name}() failed: {e}")) - .map(|_| ()) - }).await + match js_ctx.eval::(code.as_bytes()) { + Ok(s) => Ok(s == "1"), + Err(e) => { + let detail = format_js_exception(&js_ctx, &e); + Err(format!("{fn_name}() failed: {detail}")) + } + } + }).await?; + + if is_promise { + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + + loop { + drive_jobs(rt).await; + + let done = ctx.with(|js_ctx| { + js_ctx + .eval::(b"globalThis.__pendingRpcVoidDone === true") + .unwrap_or(false) + }).await; + + if done { + let error = ctx.with(|js_ctx| { + let has_error = js_ctx + .eval::(b"globalThis.__pendingRpcVoidError !== undefined") + .unwrap_or(false); + let err = if has_error { + Some(js_ctx + .eval::(b"String(globalThis.__pendingRpcVoidError)") + .unwrap_or_else(|_| "Unknown error".to_string())) + } else { + None + }; + let _ = js_ctx.eval::( + b"delete globalThis.__pendingRpcVoidDone; delete globalThis.__pendingRpcVoidError;" + ); + err + }).await; + + if let Some(err_msg) = error { + return Err(format!("{fn_name}() rejected: {err_msg}")); + } + return Ok(()); + } + + if tokio::time::Instant::now() >= deadline { + ctx.with(|js_ctx| { + let _ = js_ctx.eval::( + b"delete globalThis.__pendingRpcVoidDone; delete globalThis.__pendingRpcVoidError;" + ); + }).await; + return Err(format!("{fn_name}() timed out after 30s")); + } + + tokio::time::sleep(Duration::from_millis(5)).await; + } + } + + Ok(()) } /// Load a persisted OAuth credential from the skill's store and inject it