This is an automated email from the ASF dual-hosted git repository.

Sxnan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-agents-demos.git


The following commit(s) were added to refs/heads/main by this push:
     new b0cda72  Update Flink Agents to 0.3.1 (#8)
b0cda72 is described below

commit b0cda72ee378db0314beef882d8d00ee4c82d752
Author: yunfengzhou-hub <[email protected]>
AuthorDate: Thu Aug 6 17:46:40 2026 +0800

    Update Flink Agents to 0.3.1 (#8)
    
    AI-Contributed/Feature: 0/133
    AI-Contributed/UT: 0/0
---
 flink-operations-agent-demo/README.md              |  2 +-
 .../bin/internal/setup_flink.sh                    | 26 +++++++-
 .../custom_types_and_prompts.py                    | 25 ++++++-
 .../operations-agent-job/operations_agent.py       | 78 +++++++++++++---------
 .../operations-agent-job/pyproject.toml            |  2 +-
 5 files changed, 98 insertions(+), 35 deletions(-)

diff --git a/flink-operations-agent-demo/README.md 
b/flink-operations-agent-demo/README.md
index 7f0ecac..dc75119 100644
--- a/flink-operations-agent-demo/README.md
+++ b/flink-operations-agent-demo/README.md
@@ -62,7 +62,7 @@ Finally, the system outputs a complete **operations record** 
and feeds it back t
 - **DashScope API Key**: For AI-powered diagnosis
     - Sign up at https://dashscope.aliyun.com
     - Set environment variable: `export DASHSCOPE_API_KEY=your_api_key_here`
-    - If you prefer not to use Tongyi, you can modify the chat model section 
in `operations-agent-job/operations_agent.py`. We didn't use Ollama by default 
because smaller models have poor demo performance, while larger models run 
slowly on personal computers. For more chat model options, see the [Flink 
Agents 
Documentation](https://nightlies.apache.org/flink/flink-agents-docs-release-0.2/docs/development/chat_models/).
+    - If you prefer not to use Tongyi, you can modify the chat model section 
in `operations-agent-job/operations_agent.py`. We didn't use Ollama by default 
because smaller models have poor demo performance, while larger models run 
slowly on personal computers. For more chat model options, see the [Flink 
Agents 
Documentation](https://nightlies.apache.org/flink/flink-agents-docs-release-0.3/docs/development/chat_models/).
 
 
 ## Quick Start
diff --git a/flink-operations-agent-demo/bin/internal/setup_flink.sh 
b/flink-operations-agent-demo/bin/internal/setup_flink.sh
index ff70ca7..0adf933 100755
--- a/flink-operations-agent-demo/bin/internal/setup_flink.sh
+++ b/flink-operations-agent-demo/bin/internal/setup_flink.sh
@@ -71,15 +71,37 @@ else
 fi
 
 # Create a virtual environment in a directory named 'venv'
+# flink-agents 0.3.x requires Python >=3.10,<3.13; apache-flink 1.20 supports 
up to 3.11
 if [ ! -d "venv" ]; then
-    python3 -m venv venv
+    PYTHON_BIN=""
+    for candidate in python3.11 python3.10 python3; do
+        if command -v "$candidate" >/dev/null 2>&1; then
+            if "$candidate" -c 'import sys; sys.exit(0 if (3,10) <= 
sys.version_info[:2] <= (3,11) else 1)'; then
+                PYTHON_BIN="$candidate"
+                break
+            fi
+        fi
+    done
+    if [ -z "$PYTHON_BIN" ]; then
+        echo "Error: No suitable Python interpreter found (requires Python 
3.10 or 3.11)"
+        exit 1
+    fi
+    echo "Creating virtual environment with $PYTHON_BIN ($($PYTHON_BIN 
--version))"
+    "$PYTHON_BIN" -m venv venv
 fi
 
 # Activate the virtual environment
 # On Linux/macOS:
 source venv/bin/activate
 
-pip install "flink-agents==0.2.1" "apache-flink==1.20.3" "elasticsearch~=8.19" 
"setuptools>=75.3,<82"
+# Constrain setuptools in PEP 517 build environments (apache-beam sdist build
+# imports pkg_resources, which was removed in setuptools>=81)
+BUILD_CONSTRAINT_FILE="$(mktemp)"
+echo "setuptools<81" > "$BUILD_CONSTRAINT_FILE"
+
+PIP_CONSTRAINT="$BUILD_CONSTRAINT_FILE" pip install "flink-agents==0.3.1" 
"apache-flink==1.20.3" "elasticsearch~=8.19" "setuptools>=75.3,<82"
+
+rm -f "$BUILD_CONSTRAINT_FILE"
 
 # Set PYTHONPATH to your Python site-packages directory
 export PYTHONPATH=$(python -c 'import sysconfig; 
print(sysconfig.get_paths()["purelib"])')
diff --git 
a/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py 
b/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
index c7c6240..fee81cf 100644
--- 
a/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
+++ 
b/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
@@ -16,6 +16,8 @@
 # limitations under the License.
 
#################################################################################
 
+from typing import ClassVar
+
 from flink_agents.api.chat_message import ChatMessage, MessageRole
 from flink_agents.api.events.event import Event
 from flink_agents.api.prompts.prompt import Prompt
@@ -257,4 +259,25 @@ class ProblemRemedyResult(BaseModel):
 
 class ProblemRemedyRequestEvent(Event):
     """Job Adjust Request Event."""
-    diagnoses_result: str
+
+    EVENT_TYPE: ClassVar[str] = "problem_remedy_request_event"
+
+    def __init__(self, diagnoses_result: str) -> None:
+        """Create a ProblemRemedyRequestEvent with the diagnosis result."""
+        super().__init__(
+            type=ProblemRemedyRequestEvent.EVENT_TYPE,
+            attributes={"diagnoses_result": diagnoses_result},
+        )
+
+    @classmethod
+    def from_event(cls, event: Event) -> "ProblemRemedyRequestEvent":
+        """Reconstruct a typed ProblemRemedyRequestEvent from a base Event."""
+        assert "diagnoses_result" in event.attributes
+        result = 
ProblemRemedyRequestEvent(diagnoses_result=event.attributes["diagnoses_result"])
+        result.id = event.id
+        return result
+
+    @property
+    def diagnoses_result(self) -> str:
+        """Return the diagnosis result."""
+        return self.get_attr("diagnoses_result")
diff --git 
a/flink-operations-agent-demo/operations-agent-job/operations_agent.py 
b/flink-operations-agent-demo/operations-agent-job/operations_agent.py
index ffa3aff..554dfb0 100644
--- a/flink-operations-agent-demo/operations-agent-job/operations_agent.py
+++ b/flink-operations-agent-demo/operations-agent-job/operations_agent.py
@@ -132,7 +132,7 @@ class FlinkJobOperationsAgent(Agent):
         """EmbeddingModelConnection responsible for ollama model service 
connection."""
         return ResourceDescriptor(
             clazz=ResourceName.EmbeddingModel.OLLAMA_CONNECTION,
-            host="http://localhost:11434";,
+            base_url="http://localhost:11434";,
         )
 
     @embedding_model_setup
@@ -158,7 +158,7 @@ class FlinkJobOperationsAgent(Agent):
             dims=768,
         )
 
-    @action(InputEvent, ChatResponseEvent)
+    @action(InputEvent.EVENT_TYPE, ChatResponseEvent.EVENT_TYPE)
     @staticmethod
     def simple_problem_identification(event: Event, ctx: RunnerContext) -> 
None:
         """Perform simple problem identification for job diagnosis.
@@ -168,11 +168,13 @@ class FlinkJobOperationsAgent(Agent):
             ctx: Runner context for sending events
         """
         # Extract job information from input
-        if isinstance(event, InputEvent):
-            job_info: JobInfo = event.input
-            if job_info is None:
+        if event.get_type() == InputEvent.EVENT_TYPE:
+            raw_input = InputEvent.from_event(event).input
+            if raw_input is None:
                 logging.error("job_info is None, skipping")
                 return
+            # Events are JSON-serialized in transit, so the input may arrive 
as a plain dict
+            job_info: JobInfo = JobInfo.model_validate(raw_input) if 
isinstance(raw_input, dict) else raw_input
             logging.info(f"🔍 Starting diagnosis: {job_info}")
 
             # Store job information in short-term memory
@@ -184,16 +186,23 @@ class FlinkJobOperationsAgent(Agent):
             # Send chat request for AI analysis - let LLM decide which tools 
to use
             logging.info("🤖 Requesting AI problem identification with tool 
selection...")
 
-            msg = ChatMessage(role=MessageRole.USER, extra_args={"job_info": 
job_info.model_dump_json()})
-            
ctx.send_event(ChatRequestEvent(model="problem_identification_chat_model", 
messages=[msg]))
-        elif isinstance(event, ChatResponseEvent):
+            msg = ChatMessage(role=MessageRole.USER)
+            ctx.send_event(
+                ChatRequestEvent(
+                    model="problem_identification_chat_model",
+                    messages=[msg],
+                    prompt_args={"job_info": job_info.model_dump_json()},
+                )
+            )
+        elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+            response = ChatResponseEvent.from_event(event).response
             sop = ctx.sensory_memory.get("sop")
             problem_diagnosis_res = 
ctx.sensory_memory.get("problem_diagnosis_res")
             if sop is None and problem_diagnosis_res is None:
-                ctx.sensory_memory.set("problem_identification_res", 
event.response.content)
-                logging.info(f"SOP is not set, retrieving SOP with query: 
{event.response.content}")
+                ctx.sensory_memory.set("problem_identification_res", 
response.content)
+                logging.info(f"SOP is not set, retrieving SOP with query: 
{response.content}")
                 try:
-                    result = 
ProblemIdentificationResult.model_validate_json(event.response.content)
+                    result = 
ProblemIdentificationResult.model_validate_json(response.content)
                     if result.has_issue:
                         ctx.send_event(
                             ContextRetrievalRequestEvent(
@@ -224,11 +233,11 @@ class FlinkJobOperationsAgent(Agent):
                     # In production, use high-performance models and improve 
error handling logic.
                     ctx.send_event(
                         ContextRetrievalRequestEvent(
-                            query=event.response.content, 
vector_store="vector_store", max_results=1
+                            query=response.content, 
vector_store="vector_store", max_results=1
                         )
                     )
 
-    @action(ContextRetrievalResponseEvent, ChatResponseEvent)
+    @action(ContextRetrievalResponseEvent.EVENT_TYPE, 
ChatResponseEvent.EVENT_TYPE)
     @staticmethod
     def deep_problem_analysis(event: Event, ctx: RunnerContext) -> None:
         """Perform deep problem diagnosis using retrieved SOP context.
@@ -237,8 +246,8 @@ class FlinkJobOperationsAgent(Agent):
             event: ContextRetrievalResponseEvent containing retrieved SOP, or 
ChatResponseEvent containing AI diagnosis result
             ctx: Runner context for sending output events
         """
-        if isinstance(event, ContextRetrievalResponseEvent):
-            sop = event.documents[0].content
+        if event.get_type() == ContextRetrievalResponseEvent.EVENT_TYPE:
+            sop = 
ContextRetrievalResponseEvent.from_event(event).documents[0].content
             logging.info(f"📚 SOP retrieved: {sop}")
             ctx.sensory_memory.set("sop", sop)
 
@@ -253,11 +262,14 @@ class FlinkJobOperationsAgent(Agent):
                 # Send chat request for AI analysis - let LLM decide which 
tools to use
                 logging.info("🤖 Requesting AI analysis with tool selection...")
 
-                msg = ChatMessage(
-                    role=MessageRole.USER,
-                    extra_args={"job_info": job_info, "sop": sop, 
"problem_identification_res": problem_identification_res},
+                msg = ChatMessage(role=MessageRole.USER)
+                ctx.send_event(
+                    ChatRequestEvent(
+                        model="diagnosis_chat_model",
+                        messages=[msg],
+                        prompt_args={"job_info": job_info, "sop": sop, 
"problem_identification_res": problem_identification_res},
+                    )
                 )
-                ctx.send_event(ChatRequestEvent(model="diagnosis_chat_model", 
messages=[msg]))
 
             except Exception as e:
                 error_msg = f"Diagnosis failed for job {job_info.job_id}: 
{e!s}"
@@ -276,7 +288,8 @@ class FlinkJobOperationsAgent(Agent):
                     )
                 )
 
-        elif isinstance(event, ChatResponseEvent):
+        elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+            response = ChatResponseEvent.from_event(event).response
             sop = ctx.sensory_memory.get("sop")
             problem_diagnosis_res = 
ctx.sensory_memory.get("problem_diagnosis_res")
             if sop is not None and problem_diagnosis_res is None:
@@ -288,7 +301,7 @@ class FlinkJobOperationsAgent(Agent):
 
                     logging.info(f"🤖 Received AI analysis for job: {job_id}")
 
-                    result = 
ProblemDiagnosisResult.model_validate_json(event.response.content)
+                    result = 
ProblemDiagnosisResult.model_validate_json(response.content)
 
                     if result.need_adjustment:
                         ctx.send_event(
@@ -328,7 +341,7 @@ class FlinkJobOperationsAgent(Agent):
                         )
                     )
 
-    @action(ProblemRemedyRequestEvent, ChatResponseEvent)
+    @action(ProblemRemedyRequestEvent.EVENT_TYPE, ChatResponseEvent.EVENT_TYPE)
     @staticmethod
     def try_problem_remedy(event: Event, ctx: RunnerContext) -> None:
         """Try to address the problem based on the diagnosis results.
@@ -337,8 +350,9 @@ class FlinkJobOperationsAgent(Agent):
             event: ProblemRemedyRequestEvent containing job adjustment request
             ctx: Runner context for sending output events
         """
-        if isinstance(event, ProblemRemedyRequestEvent):
-            ctx.sensory_memory.set("problem_diagnosis_res", 
event.diagnoses_result)
+        if event.get_type() == ProblemRemedyRequestEvent.EVENT_TYPE:
+            diagnoses_result = 
ProblemRemedyRequestEvent.from_event(event).diagnoses_result
+            ctx.sensory_memory.set("problem_diagnosis_res", diagnoses_result)
             # Send chat request for AI analysis - let LLM decide which tools 
to use
             logging.info("🤖 Requesting AI job adjustment...")
 
@@ -347,12 +361,16 @@ class FlinkJobOperationsAgent(Agent):
             job_name = ctx.sensory_memory.get("job_name")
             job_info = JobInfo(base_url=base_url, job_id=job_id, 
job_name=job_name)
 
-            msg = ChatMessage(
-                role=MessageRole.USER,
-                extra_args={"job_info": job_info, "problem_diagnosis_res": 
event.diagnoses_result},
+            msg = ChatMessage(role=MessageRole.USER)
+            ctx.send_event(
+                ChatRequestEvent(
+                    model="remedy_chat_model",
+                    messages=[msg],
+                    prompt_args={"job_info": job_info, 
"problem_diagnosis_res": diagnoses_result},
+                )
             )
-            ctx.send_event(ChatRequestEvent(model="remedy_chat_model", 
messages=[msg]))
-        elif isinstance(event, ChatResponseEvent):
+        elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+            response = ChatResponseEvent.from_event(event).response
             problem_diagnosis_res = 
ctx.sensory_memory.get("problem_diagnosis_res")
             if problem_diagnosis_res is not None:
                 try:
@@ -363,7 +381,7 @@ class FlinkJobOperationsAgent(Agent):
 
                     logging.info(f"🤖 Received AI job adjustment result for 
job: {job_id}")
 
-                    result = 
ProblemRemedyResult.model_validate_json(event.response.content)
+                    result = 
ProblemRemedyResult.model_validate_json(response.content)
 
                     # Create output event
                     output_data = {
diff --git a/flink-operations-agent-demo/operations-agent-job/pyproject.toml 
b/flink-operations-agent-demo/operations-agent-job/pyproject.toml
index 8e58c1f..c3e379d 100644
--- a/flink-operations-agent-demo/operations-agent-job/pyproject.toml
+++ b/flink-operations-agent-demo/operations-agent-job/pyproject.toml
@@ -24,7 +24,7 @@ requires-python = ">=3.10"
 
 dependencies = [
     "apache-flink==1.20.3",
-    "flink-agents",
+    "flink-agents>=0.3.1",
     "pydantic>=2.0.0",
     "pyyaml>=6.0.0",
 ]

Reply via email to