Skip to content

Commit d1cc9f8

Browse files
committed
update
1 parent 64c47a4 commit d1cc9f8

9 files changed

Lines changed: 1386 additions & 17 deletions

File tree

‎src/server/handler.rs‎

Lines changed: 25 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,25 @@ impl FunctionStreamServiceImpl {
9191
}
9292
}
9393

94+
fn classify_error(message: &str) -> StatusCode {
95+
let lower = message.to_lowercase();
96+
if lower.contains("not found") || lower.contains("not exist") {
97+
StatusCode::NotFound
98+
} else if lower.contains("uniqueness violation")
99+
|| lower.contains("already exists")
100+
|| lower.contains("duplicate")
101+
{
102+
StatusCode::Conflict
103+
} else if lower.contains("invalid")
104+
|| lower.contains("unsupported")
105+
|| lower.contains("missing")
106+
{
107+
StatusCode::BadRequest
108+
} else {
109+
StatusCode::InternalServerError
110+
}
111+
}
112+
94113
async fn execute_statement(
95114
&self,
96115
stmt: &dyn Statement,
@@ -101,7 +120,8 @@ impl FunctionStreamServiceImpl {
101120
if result.success {
102121
Self::build_success_response(success_status, result.message, result.data)
103122
} else {
104-
Self::build_error_response(StatusCode::InternalServerError, result.message)
123+
let status = Self::classify_error(&result.message);
124+
Self::build_error_response(status, result.message)
105125
}
106126
}
107127
}
@@ -139,8 +159,9 @@ impl FunctionStreamService for FunctionStreamServiceImpl {
139159

140160
if !result.success {
141161
error!("SQL execution aborted: {}", result.message);
162+
let status = Self::classify_error(&result.message);
142163
return Ok(TonicResponse::new(Self::build_error_response(
143-
StatusCode::InternalServerError,
164+
status,
144165
result.message,
145166
)));
146167
}
@@ -235,8 +256,9 @@ impl FunctionStreamService for FunctionStreamServiceImpl {
235256

236257
if !result.success {
237258
error!("show_functions execution failed: {}", result.message);
259+
let status = Self::classify_error(&result.message);
238260
return Ok(TonicResponse::new(ShowFunctionsResponse {
239-
status_code: StatusCode::InternalServerError as i32,
261+
status_code: status as i32,
240262
message: result.message,
241263
functions: vec![],
242264
}));

‎tests/integration/Makefile‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,13 +43,11 @@ help:
4343
@echo " test Setup Python env + run pytest (PYTEST_ARGS=...)"
4444
@echo " clean Remove .venv and target/tests output"
4545

46-
install $(VENV)/.installed: requirements.txt $(PYTHON_ROOT)/functionstream-api/pyproject.toml $(PYTHON_ROOT)/functionstream-client/pyproject.toml
46+
install: requirements.txt $(PYTHON_ROOT)/functionstream-api/pyproject.toml $(PYTHON_ROOT)/functionstream-client/pyproject.toml
4747
$(call log,ENV,Setting up Python virtual environment)
4848
@test -d $(VENV) || python3 -m venv $(VENV)
4949
@$(PIP) install --quiet --upgrade pip
5050
@$(PIP) install --quiet -r requirements.txt
51-
@$(PIP) install --quiet -e $(PYTHON_ROOT)/functionstream-api
52-
@$(PIP) install --quiet -e $(PYTHON_ROOT)/functionstream-client
5351
@touch $@
5452
$(call success,Python environment ready)
5553

@@ -62,4 +60,5 @@ clean:
6260
$(call log,CLEAN,Removing test artifacts)
6361
@rm -rf $(VENV)
6462
@rm -rf $(CURDIR)/target
63+
@rm -rf $(CURDIR)/install
6564
$(call success,Clean complete)

‎tests/integration/framework/config.py‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,21 @@
2222

2323
from .workspace import InstanceWorkspace
2424

25+
_INTEGRATION_DIR = Path(__file__).resolve().parents[1]
26+
_PROJECT_ROOT = _INTEGRATION_DIR.parents[1]
27+
28+
29+
def _find_python_wasm() -> str:
30+
"""Locate the Python WASM runtime for server initialisation."""
31+
candidates = [
32+
_PROJECT_ROOT / "python" / "functionstream-runtime" / "target" / "functionstream-python-runtime.wasm",
33+
_PROJECT_ROOT / "dist" / "function-stream" / "data" / "cache" / "python-runner" / "functionstream-python-runtime.wasm",
34+
]
35+
for c in candidates:
36+
if c.exists():
37+
return str(c.resolve())
38+
return str(candidates[0])
39+
2540

2641
class InstanceConfig:
2742
"""Generates and persists config.yaml for one FunctionStream instance."""
@@ -44,6 +59,11 @@ def __init__(self, host: str, port: int, workspace: InstanceWorkspace):
4459
"max_file_size": 50,
4560
"max_files": 3,
4661
},
62+
"python": {
63+
"wasm_path": _find_python_wasm(),
64+
"cache_dir": str(workspace.data_dir / "cache" / "python-runner"),
65+
"enable_cache": True,
66+
},
4767
"state_storage": {
4868
"storage_type": "memory",
4969
},

‎tests/integration/requirements.txt‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,12 @@
1010
# See the License for the specific language governing permissions and
1111
# limitations under the License.
1212

13+
# Third-party dependencies
1314
pytest>=7.0
1415
pyyaml>=6.0
1516
grpcio>=1.60.0
1617
protobuf>=4.25.0
18+
19+
# FunctionStream Python packages (local editable installs)
20+
-e ../../python/functionstream-api
21+
-e ../../python/functionstream-client

‎tests/integration/test/wasm/go_sdk/__init__.py‎

Lines changed: 0 additions & 11 deletions
This file was deleted.
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
# Licensed under the Apache License, Version 2.0 (the "License");
2+
# you may not use this file except in compliance with the License.
3+
# You may obtain a copy of the License at
4+
#
5+
# http://www.apache.org/licenses/LICENSE-2.0
6+
#
7+
# Unless required by applicable law or agreed to in writing, software
8+
# distributed under the License is distributed on an "AS IS" BASIS,
9+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
10+
# See the License for the specific language governing permissions and
11+
# limitations under the License.
12+
13+
"""
14+
Fixtures for Python SDK integration tests.
15+
A single FunctionStreamInstance is shared across the entire module.
16+
"""
17+
18+
import sys
19+
from pathlib import Path
20+
21+
import pytest
22+
23+
sys.path.insert(0, str(Path(__file__).resolve().parents[3]))
24+
25+
from framework import FunctionStreamInstance
26+
27+
PROJECT_ROOT = Path(__file__).resolve().parents[5]
28+
PYTHON_EXAMPLE_DIR = PROJECT_ROOT / "examples" / "python-processor"
29+
30+
31+
@pytest.fixture(scope="session")
32+
def fs_server():
33+
"""Start a FunctionStream server once for all Python SDK tests."""
34+
instance = FunctionStreamInstance(test_name="wasm_python_sdk")
35+
instance.start()
36+
yield instance
37+
instance.kill()
38+
39+
40+
@pytest.fixture(scope="session")
41+
def python_example_dir():
42+
"""Path to the Python processor example directory."""
43+
return PYTHON_EXAMPLE_DIR

0 commit comments

Comments
 (0)