mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-05 16:02:24 +00:00
fix: python results with a KeyError-raising __getattr__ fail the job (#11517)
* fix: python results with a KeyError-raising __getattr__ no longer fail Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: keep the stream probe on the instance so enum results still serialize Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: probe python results through class dicts, never their __getattr__ 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
289d3c3941
commit
3fff1bd3dc
@@ -496,6 +496,28 @@ async def main(x: int):
|
||||
);
|
||||
}
|
||||
|
||||
// `hasattr` does not swallow the KeyError this `__getattr__` raises; released
|
||||
// `wmill` clients ship an `S3Object` built this way.
|
||||
#[test]
|
||||
fn test_python_execd_result_with_keyerror_getattr() {
|
||||
let script = r#"
|
||||
class AttrDict(dict):
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
|
||||
def main(x: str):
|
||||
return AttrDict(s3=x)
|
||||
"#;
|
||||
let results = run_py_raw_protocol_test(
|
||||
&[("f/test/attr_dict", script)],
|
||||
vec![ProtocolCmd::Execd { args: serde_json::json!({"x": "a/b.txt"}) }],
|
||||
);
|
||||
assert_eq!(
|
||||
results,
|
||||
vec![DedicatedWorkerResult::Success(serde_json::json!({"s3": "a/b.txt"}))]
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== Argument Transformation Tests ====================
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -860,7 +860,7 @@ pub async fn handle_python_job(
|
||||
if v == '<function call>':
|
||||
del pre_args[k]
|
||||
_pre_result = inner_script.preprocessor(**pre_args)
|
||||
if hasattr(_pre_result, '__await__'):
|
||||
if any('__await__' in vars(c) for c in type(_pre_result).__mro__):
|
||||
import asyncio
|
||||
_pre_result = asyncio.run(_pre_result)
|
||||
kwargs = _pre_result if _pre_result is not None else {{}}
|
||||
@@ -981,7 +981,7 @@ try:
|
||||
raise ValueError("{main_override} function is missing")
|
||||
res = _await(inner_script.{main_override}(**args))
|
||||
typ = type(res)
|
||||
if hasattr(res, '__iter__') and not isinstance(res, (str, dict, list, bytes, tuple, set, frozenset, range, memoryview, bytearray)) and typ.__name__ != 'DataFrame':
|
||||
if _defines(res, '__iter__') and not isinstance(res, (str, dict, list, bytes, tuple, set, frozenset, range, memoryview, bytearray)) and typ.__name__ != 'DataFrame':
|
||||
for chunk in res:
|
||||
print("WM_STREAM: " + chunk.replace('\n', '\\n'))
|
||||
res = None
|
||||
@@ -3340,10 +3340,16 @@ This is not normal behavior, please make sure all workers have enough memory.\n
|
||||
/// loop-per-call would break them from the second job on. A loop the script
|
||||
/// set while importing (`asyncio.set_event_loop`) is reused. asyncio is
|
||||
/// imported lazily so sync jobs don't pay its import time.
|
||||
/// `_defines` reads the class dicts instead of calling `hasattr`: a result is
|
||||
/// arbitrary user data, and `hasattr` runs its `__getattr__` (or its
|
||||
/// metaclass's), which may raise something other than AttributeError or
|
||||
/// answer every name.
|
||||
const PY_AWAIT_HELPER: &str = r#"_loop = None
|
||||
def _defines(r, name):
|
||||
return any(name in vars(c) for c in type(r).__mro__)
|
||||
def _await(r):
|
||||
global _loop
|
||||
if not hasattr(r, '__await__'):
|
||||
if not _defines(r, '__await__'):
|
||||
return r
|
||||
import asyncio
|
||||
if _loop is None or _loop.is_closed():
|
||||
|
||||
@@ -1,6 +1,15 @@
|
||||
from typing import Optional
|
||||
|
||||
|
||||
def _item_as_attr(d: dict, attr: str):
|
||||
# AttributeError, not KeyError: hasattr, copy and pickle probe optional
|
||||
# attributes and only treat AttributeError as "absent".
|
||||
try:
|
||||
return d[attr]
|
||||
except KeyError:
|
||||
raise AttributeError(attr) from None
|
||||
|
||||
|
||||
class S3Object(dict):
|
||||
"""S3 file reference with file key, optional storage identifier, and presigned token."""
|
||||
s3: str
|
||||
@@ -8,7 +17,7 @@ class S3Object(dict):
|
||||
presigned: Optional[str]
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class S3FsClientKwargs(dict):
|
||||
@@ -16,7 +25,7 @@ class S3FsClientKwargs(dict):
|
||||
region_name: str
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class S3FsArgs(dict):
|
||||
@@ -29,7 +38,7 @@ class S3FsArgs(dict):
|
||||
client_kwargs: S3FsClientKwargs
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class StorageOptions(dict):
|
||||
@@ -41,7 +50,7 @@ class StorageOptions(dict):
|
||||
aws_allow_http: str
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class PolarsConnectionSettings(dict):
|
||||
@@ -50,7 +59,7 @@ class PolarsConnectionSettings(dict):
|
||||
storage_options: StorageOptions
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class Boto3ConnectionSettings(dict):
|
||||
@@ -63,7 +72,7 @@ class Boto3ConnectionSettings(dict):
|
||||
aws_session_token: Optional[str]
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
|
||||
class DuckDbConnectionSettings(dict):
|
||||
@@ -71,4 +80,4 @@ class DuckDbConnectionSettings(dict):
|
||||
connection_settings_str: str
|
||||
|
||||
def __getattr__(self, attr):
|
||||
return self[attr]
|
||||
return _item_as_attr(self, attr)
|
||||
|
||||
Reference in New Issue
Block a user