apache/flink-agents
55.4
Adequate · 22 September 2026
81.5k
lines of production code
Java
with Python
6
measurements over time
What this system is
This system is a Java and Python API for building, executing, and managing AI agents within the Apache Flink streaming runtime. It provides a declarative framework for defining agent workflows, tools, and memory, while integrating with various LLM providers, vector stores, and external protocols like MCP. The system ensures resilient execution through durable state persistence, cross-language interoperability, and comprehensive observability via structured event logging and metrics.
How it got here
2025 — Flink Agents initial release and cross-language support
88 changes.
This period established the foundational architecture for Apache Flink Agents, introducing the core Java and Python APIs for defining, planning, and executing agent workflows. It implemented critical cross-language interoperability, allowing Java and Python components to seamlessly share resources like chat models, vector stores, and tools. The release also delivered essential runtime features including durable execution, memory management, and observability, alongside comprehensive test coverage and integration examples.
2026 — Agent integrations and declarative configuration
44 changes.
This period focused on expanding Flink Agents with extensive integrations for chat models (OpenAI, Azure, Bedrock, Gemini, Watsonx), vector stores, and external protocols like MCP. It also introduced a declarative YAML API for defining agents and resources, alongside new capabilities for long-term memory, agent skills, and secure shell execution.
Features
Add Amazon Bedrock embedding model integration
Users can now connect to Amazon Bedrock to generate text embeddings using the Amazon Titan Text Embeddings V2 model. This integration allows configuring the AWS region, model ID, embedding dimensions, and concurrency settings, and automatically tracks token usage metrics for billing and monitoring.
integrations/embedding-models/bedrock · high confidence
Add Amazon OpenSearch and S3 Vectors vector store integrations
Users can now persist and query vector embeddings using Amazon OpenSearch (including OpenSearch Serverless) and Amazon S3 Vectors. The OpenSearch integration supports both IAM (SigV4) and basic authentication, allows configuration of index names, vector/content field names, and AWS regions, and maps vector store collections to OpenSearch indices. The S3 Vectors integration connects to an S3 Vectors bucket and index, supporting document upserts, metadata filtering, and similarity search with relevance scores.
integrations/vector-stores/opensearch · high confidence
Add ChromaDB vector store integration with Cloud support
Introduces a new ChromaVectorStore implementation for the Flink Agents vector store API, enabling semantic search capabilities using ChromaDB. This integration supports multiple connection modes including in-memory, persistent storage, server-based connections, and Chroma Cloud (via API key). The implementation includes thread-safe client initialization to handle concurrent access and automatically manages collection creation, providing a seamless way to store and query vector embeddings within Flink Agents applications.
_python/flink\_agents/integrations/vector\stores/chroma · high confidence
Add Elasticsearch vector store integration for Java agents
Introduces the ElasticsearchVectorStore implementation for Java, enabling agents to store and retrieve documents using Elasticsearch's approximate nearest neighbor (KNN) search capabilities. Users can configure the store via ResourceDescriptor arguments to specify the target index, vector field, dimensionality, and connection details (including basic auth or API key authentication). The implementation supports filtering KNN queries with raw JSON Elasticsearch DSL filters and allows overriding default neighbor counts (k) and candidate set sizes (num\_candidates) on a per-query basis.
integrations/vector-stores/elasticsearch/src/main · high confidence
Add Google Gemini chat model integration
This change introduces a new built-in integration for the Google Gemini chat model, allowing agents to connect to the Gemini API via the official google-genai Java SDK. The integration includes a connection class that handles API authentication (including Vertex AI support), request configuration, and protocol translation for system messages and tool calls, as well as a setup class for configuring model parameters like temperature and max output tokens. It also adds comprehensive unit tests for both the connection and setup components to verify configuration handling and schema generation.
integrations/chat-models/gemini · high confidence
Add Milvus vector store integration
Users can now store and query vector embeddings using Milvus as a backend for Flink Agents. This new integration provides a Milvus-backed implementation of the vector store interface, supporting dense-vector similarity search with configurable collection schemas (id, content, metadata, embedding), metric types (defaulting to COSINE), and index types. It allows configuration of connection details (URI, authentication), collection parameters (dimensionality, shard count, consistency level), and metadata indexing (including JSON path indexes for filtering). The change includes the main implementation class and corresponding unit tests.
_integrations/vector-stores/milvus, python/flink\_agents/integrations/vector\stores/mem0 · high confidence
Add Ollama and Tongyi chat model integrations with native structured output
New Ollama and Tongyi (DashScope) chat model integrations are now available in the Python SDK. These models support native structured output for Pydantic BaseModel schemas, allowing the Ollama server and DashScope API to enforce response formats directly rather than relying on prompt engineering. The integrations also include tool support, with utility functions to convert internal tool metadata into the specific formats required by each provider.
_python/flink\_agents/integrations/chat\models · high confidence
Add Ollama embedding model integration for Java
Users can now use Ollama as an embedding model provider in Java-based Flink Agents. This change introduces the OllamaEmbeddingModelConnection and OllamaEmbeddingModelSetup classes, which leverage the ollama4j client to generate embeddings. The integration supports configuring the Ollama host (defaulting to localhost:11434) and the model name (defaulting to nomic-embed-text), and automatically pulls the required model if it is not already present. Tests are included to verify the connection setup and parameter handling.
integrations/embedding-models/ollama · high confidence
Add OpenAI and Tongyi embedding model integrations
Users can now integrate with OpenAI and Tongyi (DashScope) embedding models within Flink Agents. This change introduces new connection and setup classes for both providers, allowing users to generate embeddings via the OpenAI API and DashScope TextEmbedding API, including support for token usage tracking and provider-specific configuration options.
_python/flink\_agents/integrations/embedding\models · high confidence
Add Python-based Vector Store implementations with collection management
Introduced Java wrappers for Python vector stores, enabling Flink Agents to interact with Python-based vector storage backends. The new \PythonVectorStore\ class bridges Java and Python, handling document storage, retrieval, and querying while managing embedding generation to avoid cross-language deadlocks. Additionally, \PythonCollectionManageableVectorStore\ extends this functionality to support creating and deleting vector collections, providing a more complete integration for users leveraging Python vector store libraries.
api/src/main/java/org/apache/flink/agents/api/vectorstores/python · high confidence
Add built-in Ollama embedding model integration
Users can now generate text embeddings using a local Ollama server via the new \OllamaEmbeddingModelConnection\ and \OllamaEmbeddingModelSetup\ classes. This integration allows connecting to an Ollama instance (defaulting to localhost:11434) and supports configuration options such as model selection, truncation behavior, and request timeouts, enabling local embedding generation without external cloud dependencies.
_python/flink\_agents/integrations/embedding\models/local · high confidence
Add built-in support for Azure OpenAI Chat Model
Users can now connect Flink Agents to Azure OpenAI via a new \AzureOpenAIChatModelConnection\. This integration supports native structured output for specific GPT and o-series models by leveraging Azure's JSON schema response format, and includes configuration options for API keys, endpoints, timeouts, and retries.
_python/flink\_agents/integrations/chat\models/azure · high confidence
Add built-in support for Azure OpenAI and the OpenAI Responses API
This release introduces two new chat model integrations for the OpenAI ecosystem. First, it adds a dedicated Azure OpenAI connection and setup, allowing users to connect to Azure deployments using the official OpenAI Java SDK with native structured output support for specific models (e.g., gpt-4o, o1) and configurable URL path modes. Second, it adds a new OpenAI Responses API integration, providing a separate connection and setup class that leverages OpenAI's newer Responses API for features like reasoning effort and response storage, distinct from the existing Chat Completions integration.
integrations/chat-models/openai/src/main · high confidence
Add vLLM chat model integration for Python
Users can now connect Flink Agents to vLLM servers using a new \VLLMChatModelConnection\ and \VLLMChatModelSetup\. The connection defaults to a local vLLM instance (\http://localhost:8000/v1\) and uses a placeholder API key, ignoring standard OpenAI environment variables to match Java behavior. The setup requires an explicit model name and supports native structured outputs (JSON schema) for served models like Qwen and Llama via guided decoding.
_python/flink\_agents/integrations/chat\models/vllm · high confidence
Added utility to configure operator chaining strategy in Flink 1.20 agents runtime
The Flink 1.20 agents runtime now includes an \OperatorUtils\ helper that allows setting the chaining strategy (e.g., ALWAYS, NEVER, HEAD) on stream operators. This utility enables more granular control over how operators are chained together during execution, which can impact performance and resource usage in streaming jobs. A corresponding test suite verifies the correct application of these strategies across different operator types.
dist/flink-1.20 · high confidence
Anthropic Chat Model integration with structured output and tool support
This change introduces the Python integration for the Anthropic chat model within Flink Agents. It enables users to connect to Anthropic's API, supporting tool use via the standard tool metadata conversion and handling of multi-modal content blocks. The integration also implements native structured output support for specific Claude model generations (such as Claude 4.5+ and Mythos) and manages model-specific constraints like prefilling restrictions, ensuring robust interaction with the Anthropic API.
_python/flink\_agents/integrations/chat\models/anthropic · high confidence
CLI tool for visualizing agent execution traces
A new command-line interface has been added to Flink Agents, allowing users to reconstruct and view trace trees from Event Log files. The \trace\_tree\ module reads event logs (supporting both current and legacy formats) and outputs a structured representation of the agent's execution lineage in either text or JSON format, enabling users to inspect the sequence of actions and events that occurred during an agent run.
_python/flink\agents/cli · high confidence
Embedding models now report token usage metrics
Embedding operations now return token usage metadata (prompt tokens and total tokens) alongside the generated embeddings. This is implemented via new \EmbeddingResult\ and \EmbeddingTokenUsage\ classes, and utility methods in \EmbeddingModelUtils\ that parse usage data from cross-language (Python) responses, allowing users to track consumption when calling embedding models.
api/src/main/java/org/apache/flink/agents/api/embedding/model · high confidence
Framework-managed LLM-as-judge routing strategy
The plan-side routing layer now includes a built-in LLM-as-judge strategy (Strategies.llm) that automatically selects the target chat model by sending the request to a dedicated judge model and parsing its JSON verdict. This strategy is executed via the new LlmJudgeRoutingExecutor, which integrates with the standard durable, metered, and observable chat invocation path, and is registered in the plan routing dispatch alongside existing rule-based and custom executors.
plan/src/main/java/org/apache/flink/agents/plan/routing · high confidence
Initial release of the Python Flink Agents package and examples
The Python Flink Agents library is now available, introducing the core \flink\_agents\ package and a companion \examples\ subpackage. This initial release establishes the foundational structure for building and running Flink-based agents in Python, providing the necessary entry points and licensing headers to begin development.
_python/flink\_agents, python/flink\agents/examples · high confidence
Introduce AGENT resource type and sub-agent invocation API
The API now supports defining and invoking sub-agents as first-class resources. A new AGENT type is added to the ResourceType enum, and the SubagentSetup class allows agents to declare nested sub-agents. Invocations are managed through SubagentFuture and SubagentFutures, which provide asynchronous execution, cancellation, and grouping capabilities. Results are captured in SubagentResult, designed to be serializable for durable execution and failover recovery, allowing callers to inspect success or error states without relying on exception propagation.
api/src/main/java/org/apache/flink/agents/api/resource · high confidence
Introduce Anthropic chat model integration
Adds a new integration for the Anthropic Messages API, allowing agents to connect to Anthropic models via the \AnthropicChatModelConnection\ and \AnthropicChatModelSetup\ classes. The connection supports configuration of API keys, timeouts, and retry limits, while the setup layer manages per-chat parameters such as model selection (defaulting to \claude-sonnet-4-6\), temperature, max tokens, and tool bindings. The integration also implements native structured output support for specific Claude 4.5+ and Claude 5 models, and includes logic to handle JSON prefilling and strict tool adherence.
_integrations/chat-models/anthropic/src/main, python/flink\_agents/integrations/chat\models/openai · high confidence
Introduce Flink Agents metric group and updatable gauge interfaces
The API now exposes a new metrics subsystem for monitoring agent actions. Users can access metric groups via the new FlinkAgentsMetricGroup interface, which supports hierarchical grouping (including key-value pairs for model and action dimensions) and provides access to counters, meters, histograms, and gauges. A new UpdatableGauge interface extends the standard Gauge to allow runtime value updates, standardizing on String values to ensure consistency across Python and Java interactions.
api/src/main/java/org/apache/flink/agents/api/metrics · high confidence
Introduce FunctionTool for executable function definitions
The \python/flink\_agents/plan/tools\ module now provides the \FunctionTool\ class, which wraps Python or Java functions to make them executable tools. This component handles the automatic derivation of tool metadata (such as name, description, and argument schemas) from function signatures and docstrings, and supports parameter injection via \InjectedArg\ for both local Python functions and remote Java functions bridged through a resource adapter.
_python/flink\agents/plan/tools · high confidence
Introduce Java API for Long-Term Memory
This change adds the Java-side API for long-term memory management, introducing the BaseLongTermMemory interface and the MemorySet class to allow agents to store, retrieve, and search persistent memory items. It includes configuration options for the Mem0 backend (chat model, embedding model, and vector store) and defines the MemorySetItem data structure, enabling users to integrate semantic memory capabilities into their Flink agents.
api/src/main/java/org/apache/flink/agents/api/memory · high confidence
Introduce Java API for building and executing Flink Agents
This change adds the core Java API classes for defining and running Flink Agents, including the AgentsExecutionEnvironment for configuring inputs from DataStreams, Tables, or lists, and the AgentBuilder for chaining agent configurations and producing outputs. It introduces a unified Event model with type-based routing, attachments, and lineage tracking, along with specific event types (InputEvent, OutputEvent) and a centralized EventType constant registry. Additionally, it provides a RetryExecutor utility with configurable backoff for remote calls and an EventContext for runtime event observation.
api/src/main/java/org/apache/flink/agents/api · high confidence
Introduce Java AgentPlan and configuration model with cross-language function support
The plan module now includes the core Java model for defining and serializing agent workflows. AgentPlan serves as the central structure, holding actions, resource providers, and configuration, and supports JSON serialization via custom Jackson serializers. AgentConfiguration provides a key-value store for agent settings. To enable cross-language execution, the plan layer introduces Function, JavaFunction, and PythonFunction classes, allowing the plan to represent and serialize both Java methods and Python functions for later invocation. Additionally, trigger conditions for actions can now be defined using CEL expressions, validated at plan time, and parsed by a shared dialect.
plan/src/main/java/org/apache/flink/agents/plan · high confidence
Introduce Java ChatMessage and MessageRole API classes
Added the ChatMessage class and MessageRole enumeration to the Flink Agents API, providing a structured way to represent chat interactions with support for user, system, assistant, and tool roles. The ChatMessage class includes fields for content, tool calls, and extra arguments, along with static factory methods for creating messages of each role type, enabling developers to build and manage conversation histories programmatically in Java.
api/src/main/java/org/apache/flink/agents/api/chat/messages · high confidence
Introduce Java Prompt API for template-based message generation
The Flink Agents Java API now includes a new \Prompt\ class and its \LocalPrompt\ implementation, allowing users to define templates using either plain text strings or sequences of chat messages. This API supports placeholder substitution (e.g., \{variable}\) in a single pass, enabling the dynamic generation of formatted text or structured chat message lists for use with language models.
api/src/main/java/org/apache/flink/agents/api/prompt · high confidence
Introduce Java and Python resource providers for Flink Agents
This change adds the resource provider implementation classes in the plan module, enabling the agent framework to instantiate and manage resources at runtime. It introduces JavaResourceProvider and JavaSerializableResourceProvider to create Java-based resources (including serializable ones via Jackson), and PythonResourceProvider and PythonSerializableResourceProvider to bridge Python resources (such as ChatModels, EmbeddingModels, VectorStores, MCP servers, Prompts, and Tools) into the Java agent execution context. These providers carry the necessary metadata (module, class, arguments) to dynamically load and construct the appropriate resource objects when an AgentPlan is executed.
plan/src/main/java/org/apache/flink/agents/plan/resourceprovider · high confidence
Introduce Java integration for Model Context Protocol (MCP) servers
Added a new Java integration for connecting to Model Context Protocol (MCP) servers, enabling Flink Agents to discover and execute remote tools and prompts. This update introduces the \MCPServer\ resource class, which manages server connections, authentication, and timeouts, alongside \MCPTool\ and \MCPPrompt\ classes that expose remote capabilities as native agent tools and prompts. The integration supports multiple authentication methods—including Bearer tokens, Basic auth, and API keys—via a pluggable \Auth\ interface, and includes utilities for normalizing diverse content types (text, images, embedded resources) returned by MCP servers.
integrations/mcp/src/main · high confidence
Introduce Java integration for Ollama chat models
Users can now connect Flink Agents to local Ollama instances via a new Java integration. This adds the OllamaChatModelConnection and OllamaChatModelSetup classes, enabling tool calling with native structured output support and configurable reasoning parameters (think/extract\_reasoning). The integration also tolerates tool schemas that omit the 'required' key, ensuring compatibility with various tool definitions.
integrations/chat-models/ollama · high confidence
Introduce Mem0-based long-term memory with event attachment support
This change adds a new long-term memory implementation backed by the Mem0 library, allowing agents to store and retrieve persistent information using vector-based storage and LLM-driven fact extraction. It includes adapters that map Flink Agents' chat, embedding, and vector store resources to Mem0's internal components, and enforces scoping so that long-term memory operations are isolated to the specific action that obtained the memory set. Additionally, it introduces event attachment handling in the runtime memory layer, enabling agents to pass large or structured data across actions via sensory memory references instead of embedding them directly in events.
_python/flink\agents/runtime/memory · high confidence
Introduce Prompt abstraction for templating chat messages
This change adds a new \Prompt\ API under \python/flink\_agents/api/prompts\ that allows users to define and format chat templates. It includes a \LocalPrompt\ implementation that supports both plain text and sequence-of-messages templates, utilizing a custom \SafeFormatter\ (based on \re\ rather than Python's \Formatter\) to safely substitute variables without raising errors on missing keys. This enables users to create reusable prompt structures that can be formatted with dynamic arguments to generate \ChatMessage\ objects or raw strings.
_python/flink\agents/api/prompts · high confidence
Introduce Python AgentPlan for serializing and executing agent workflows
The \python/flink\_agents/plan\ module now provides the core planning layer for Python agents, introducing the \AgentPlan\ class to compile user-defined agents into a structured representation of actions and resource providers. This change adds support for serializing and deserializing agent plans (including cross-language Java resources) via Pydantic models, defines \AgentConfiguration\ for managing nested configuration data, and implements \PythonFunction\ with a global cache to reuse identical function instances. It also introduces specific resource providers (\PythonResourceProvider\, \PythonSerializableResourceProvider\, \JavaResourceProvider\) to handle resource creation and deserialization at runtime, effectively establishing the mechanism for defining and executing agent logic in Python.
_python/flink\agents/plan · high confidence
Introduce Python Bash tool with AST-based command validation
A new standalone Bash execution tool is added to the Python Flink Agents plan, allowing LLMs to run shell commands via a \command\, \timeout\, and \cwd\ interface. To ensure safety, the tool uses an AST-based validator (tree-sitter) to enforce strict allowlists for commands and script directories, while explicitly blocking dangerous shell constructs such as command substitution, process substitution, subshells, loops, and file redirects.
_python/flink\agents/plan/tools/bash · high confidence
Introduce Python ChatModel API with structured output and token metrics
The Python API now includes a new \chat\_models\ module providing \BaseChatModelConnection\ and \BaseChatModelSetup\ abstractions, enabling developers to integrate with chat providers. This release adds explicit support for structured output via the \StructuredOutputStrategy\ enum (AUTO, NATIVE, PROMPT) and an \output\_schema\ parameter, allowing users to enforce response formats. It also introduces token usage metrics tracking (prompt and completion tokens) scoped to specific models and actions. Additionally, the API supports bridging to Java-based chat models through \JavaChatModelConnection\ and \JavaChatModelSetup\, facilitating cross-language resource usage.
_python/flink\_agents/api/chat\models · high confidence
Introduce Python Embedding Model abstraction and cross-language support
This change adds the core Python API for embedding models, introducing abstract base classes for connections and setups that allow users to define and manage text embedding services. It includes a built-in implementation for bridging Java embedding models into the Python environment via the \@java\_resource\ decorator, enabling cross-language usage. The API also tracks token usage metrics, returning provider-specific prompt and total token counts alongside embedding results when available.
_python/flink\_agents/api/embedding\models · high confidence
Introduce Python Tool API with parameter injection and declarative function tools
The \python/flink\_agents/api/tools\ module now provides a new declarative Tool API for Python agents. Users can wrap Python callables into \FunctionTool\ instances via \Tool.from\_callable()\, which automatically generates Pydantic argument schemas from function signatures and docstrings. The API supports declarative parameter injection, allowing tool arguments to be automatically populated from agent configuration, sensory memory, or short-term memory using \InjectedArg\ and \ToolParameterSource\. Additionally, a \ToolResponse\ dataclass is introduced to allow tools to explicitly report success or failure states, and a \ToolExecutionMetadataProvider\ abstraction is added to support optional stable metadata for tool executions.
_python/flink\agents/api/tools · high confidence
Introduce Python agent skill loading and management
This change adds the Python runtime support for agent skills, enabling agents to load and use skills defined in SKILL.md files. The implementation includes a SkillManager to orchestrate loading, a registry for different source types (local filesystem, URL, and Python package), and parsers to extract skill metadata and content. Skills can now be discovered and activated by agents, with support for lazy resource loading and secure URL fetching with optional SHA-256 verification.
_python/flink\agents/runtime/skill · high confidence
Introduce Python resource lifecycle and interoperability APIs
This change adds the core Java-side infrastructure for bridging Java and Python resources in Flink Agents. It introduces PythonObjectScope to manage the lifecycle of temporary Pemja references, ensuring Python objects are properly closed and preventing memory leaks when crossing the language boundary. It also defines the PythonResourceAdapter interface, which provides methods for initializing Python resources, converting data types (such as ChatMessage, Document, and VectorStoreQuery) between Java and Python, and invoking Python functions as tools. Additionally, the PythonResourceWrapper interface allows Java objects to expose their underlying Python resources and bind metric groups, enabling metrics to flow from Python components back to the Java runtime.
api/src/main/java/org/apache/flink/agents/api/resource/python · high confidence
Introduce YAML API for declaring agents and resources
Users can now define agents, actions, tools, prompts, and resources (such as chat models, embedding models, and vector stores) using YAML configuration files. This change adds a new YAML loader that parses these declarations, resolving short aliases for providers (e.g., 'ollama', 'openai') to their full class implementations and supporting both Java and Python resource types. The loader also handles cross-language wrapping, allowing Python implementations to be used within the Java agent runtime, and enforces security requirements for URL-based skill sources by requiring HTTPS and optional digest pinning.
api/src/main/java/org/apache/flink/agents/api/yaml · high confidence
Introduce agents configuration API
The agents API now exposes a new configuration mechanism in the \org.apache.flink.agents.api.configuration\ package, allowing users to control runtime behavior through typed configuration options. This includes settings for event logging (such as logger type, log level, trace enablement, and payload size limits), action state store backends (with specific configuration keys for Kafka and Fluss, including bootstrap servers, topics, and security protocols), and memory event observation (master switches and per-operation toggles). The API also provides options for handling trigger condition evaluation failures and defining custom event listeners.
_api/src/main/java/org/apache/flink/agents/api/configuration, python/flink\_agents/api, python/flink\agents/runtime · high confidence
Introduce built-in Java agent actions with durable execution and parallel tool calls
This change adds the core built-in action implementations for the Flink Agents plan in Java, including Action, ChatModelAction, ContextRetrievalAction, and ToolCallAction. These actions provide the execution logic for chat requests, context retrieval (RAG), and tool calls, leveraging the engine's durable execution machinery to ensure resilience. A key behavioral addition is support for parallel tool call execution, allowing multiple tool invocations to run concurrently when configured, while also handling tool parameter injection and error strategies. The implementation includes model routing capabilities within ChatModelAction and ensures that interrupted or cancelled calls are not retried or persisted, improving reliability during agent cancellation.
plan/src/main/java/org/apache/flink/agents/plan/actions · high confidence
Introduce built-in ReAct agent with structured output support
The Python API now includes a built-in ReAct agent that leverages LLM function-calling capabilities. Users can define agents by specifying a chat model, optional prompts, and an output schema (Pydantic models or Flink RowTypeInfo). The agent automatically generates system prompts to enforce JSON output matching the schema and handles input formatting. Clear errors are raised if the output schema cannot be rendered to JSON Schema.
_python/flink\agents/api/agents · high confidence
Introduce built-in actions for chat, tool calls, and context retrieval
The Python agent plan now includes built-in actions for processing chat model requests, executing tool calls, and retrieving context. The \ChatModelAction\ handles LLM interactions with support for retry intervals, error handling strategies, and checkpoint-stable tool call context management. The \ToolCallAction\ processes tool requests, supporting parallel execution when enabled, and injects framework-owned arguments to prevent spoofing. The \ContextRetrievalAction\ enables RAG-style context retrieval via vector stores, with async execution support for cross-language resources where available. These actions provide the core execution logic for agent workflows.
_python/flink\agents/plan/actions · high confidence
Introduce declarative YAML API for declaring agents and resources
Users can now define Flink Agents, their actions, tools, prompts, and connected resources (such as chat models, embedding models, and vector stores) using a structured YAML format. This new API provides a declarative way to configure agent behavior, supporting both Python and Java implementations through a unified loader that resolves aliases and handles cross-language resource descriptors.
_python/flink\agents/api/yaml · high confidence
Introduce dedicated event types for chat, tools, and memory operations
The \python/flink\_agents/api/events\ module now exposes specific event classes to provide structured observability for agent internals. \ChatRequestEvent\ and \ChatResponseEvent\ track LLM interactions, including prompt arguments and retry metrics. \ToolRequestEvent\ and \ToolResponseEvent\ expose tool call details and execution results. \ContextRetrievalRequestEvent\ and \ContextRetrievalResponseEvent\ record vector store queries and retrieved documents. Additionally, \MemoryEvent\ subclasses (e.g., \ShortTermWriteEvent\, \LongTermSearchEvent\) and \AgentRunBeginEvent\ provide granular visibility into memory operations and run lifecycle, all accessible via the unified \EventType\ constants.
_python/flink\agents/api/events · high confidence
Introduce durable action state persistence with Kafka and Fluss backends
The runtime now persists agent action state to durable external backends (Kafka or Fluss) to enable fine-grained recovery and exactly-once semantics. This change introduces the \ActionStateStore\ interface and its implementations, along with \ActionStateSerde\ for serializing state (including memory updates via a versioned Kryo envelope) and \ActionStateKeyEncoder\ for generating stable, type-preserving keys. It also adds \ResourceCache\ for lazy, thread-safe resolution of agent resources (including Python MCP tools and prompts) and \CompileUtils\ to bridge Flink DataStreams with the agent runtime, including a safety check that rejects incompatible Flink batch configurations.
runtime/src/main · high confidence
Introduce long-term memory interface and Mem0 support
This change introduces a new long-term memory API for Flink Agents, allowing actions to persist and retrieve context across sessions. It provides a \MemorySet\ interface scoped to individual actions to ensure correct partitioning, supports vector-store-based implementations, and adds specific configuration options for the Mem0 backend (chat model, embedding model, and vector store). Users can now add, search, get, and delete memory items within the context of their agent actions.
_python/flink\agents/api/memory · high confidence
Introduce plan-level tool execution and schema generation for Java and Python functions
This change adds the plan-module implementation for executing tools defined as static Java methods or Python functions. It introduces FunctionTool to handle invocation logic for both Java and Python backends, SchemaUtils to generate JSON schemas from method signatures (correctly mapping numeric types to integer/number), and ToolMetadataFactory to create tool metadata from @Tool-annotated methods. This enables the agent planning layer to properly describe and execute user-defined tools.
plan/src/main/java/org/apache/flink/agents/plan/tools · high confidence
Introduce unified Java Vector Store API
This change introduces a new, unified API for vector store integrations in Java, centered around the new \BaseVectorStore\ abstract class. It provides a consistent interface for adding, updating, and querying documents, including automatic embedding generation via configured models. The API supports a unified filter DSL for metadata equality matching, allows specifying target collections, and exposes store-specific arguments through \extraArgs\. Supporting types include \Document\ for content and metadata, \VectorStoreQuery\ for query parameters (mode, limit, filters), and \VectorStoreQueryResult\ for results. An optional \CollectionManageableVectorStore\ interface is also added for stores that support collection lifecycle management.
api/src/main/java/org/apache/flink/agents/api/vectorstores · high confidence
Introduce unified Python Vector Store API with Java bridge support
This change establishes the core Python API for vector stores, introducing abstract base classes (BaseVectorStore, CollectionManageableVectorStore) and structured query models (VectorStoreQuery, Document) that define a unified interface for semantic search and metadata filtering. It also adds a JavaVectorStore implementation that bridges Python code to Java-based vector store backends, enabling cross-language resource usage within the Flink Agents runtime.
_python/flink\_agents/api/vector\stores · high confidence
Introduces durable execution and memory reference APIs for Flink Agents
The API context package now exposes the core interfaces for durable action execution and in-memory data passing. Users can perform resilient, restart-safe operations via the new \RunnerContext.durableExecute\ and \durableExecuteAsync\ methods, which accept \DurableCallable\ instances that support optional reconciliation logic for in-flight tasks. Asynchronous calls can be composed into parallel batches using \RunnerContext.gather\, which leverages \DurableFuture\ handles. Additionally, the package introduces \MemoryObject\ and \MemoryRef\ to manage sensory and short-term memory, allowing actions to pass large data structures efficiently via lightweight, serializable references rather than copying values.
api/src/main/java/org/apache/flink/agents/api/context · high confidence
Introduces pluggable in-chat model routing with LLM-as-judge and custom strategies
The API now supports dynamic model selection within chat sessions via the new \ModelRouter\ resource and \RoutingStrategy\ system. Users can configure routing strategies to select from multiple candidate models based on request content, including a framework-managed \Strategies.llm()\ (LLM-as-judge) that delegates the decision to a specified judge model, a \Strategies.rules()\ for keyword/regex matching, and \Strategies.custom()\ for user-defined \CustomRoutingExecutor\ logic. This allows agents to adaptively choose the most appropriate model for each request while keeping the routing decision separate from the model execution path.
api/src/main/java/org/apache/flink/agents/api/chat/model/routing · high confidence
JSON serialization support for Flink Agent plans, actions, and resource providers
The plan serializer module now provides custom Jackson serializers and deserializers for AgentPlan, Action, and ResourceProvider objects. This enables the complete serialization and deserialization of agent plans to and from JSON, including the preservation of Java and Python function details, trigger conditions, configuration data, and resource provider implementations (Java and Python variants). Users can now persist and restore agent plans in a structured JSON format, facilitating cross-language compatibility and plan portability.
plan/src/main/java/org/apache/flink/agents/plan/serializer · high confidence
Java integration with Python embedding models
Users can now use Python-based embedding models directly from Java code. This change introduces PythonEmbeddingModelConnection and PythonEmbeddingModelSetup, which bridge Java and Python to allow embedding operations and usage metrics tracking across languages.
api/src/main/java/org/apache/flink/agents/api/embedding/model/python · high confidence
New API for defining and injecting tool parameters
The API now introduces a declarative tool parameter injection system, allowing tools to automatically receive values from framework sources like configuration, sensory memory, or short-term memory. This is enabled by new classes in the tools package: ToolParameterInjection and ToolParameterSource define the binding, while ToolParameterInjectionValidator ensures that injected arguments match the underlying Java method signatures. Additionally, the API layer now includes pure-data descriptors for cross-language function execution (JavaFunction and PythonFunction) and a FunctionTool class that wraps these descriptors, enabling users to define tools backed by Java methods or Python callables that are serialized for execution in the plan layer.
api/src/main/java/org/apache/flink/agents/api/tools · high confidence
New AWS Bedrock chat model integration with native structured output
Users can now connect to Amazon Bedrock via the Converse API to invoke models such as Claude, Mistral, and Llama. This integration supports tool calling and, for a specific set of models, native structured output (JSON schema) to ensure responses match expected types. It includes a connection resource for configuring the AWS region and default model, a setup resource for specifying model parameters like temperature and max tokens, and a shared POJO-to-JSON-Schema generator to handle output schema generation.
integrations/chat-models/bedrock · high confidence
New Flink Agents quickstart examples and YAML API support
The quickstart directory now includes a comprehensive set of end-to-end examples demonstrating Flink Agents capabilities. Users can explore parallel LLM fan-out patterns, ReAct agents with tool use (such as notifying a shipping manager), and agent skills (like a math calculator). The examples also cover multi-agent workflows, including a single-agent review analysis pipeline and a multi-stage workflow that aggregates review scores and generates improvement suggestions using the Flink Table API. Additionally, a new YAML-based agent declaration approach is introduced, allowing users to define agents, prompts, tools, and chat model connections via a YAML file and load them into the execution environment, providing an alternative to code-defined agents.
_python/flink\agents/examples/quickstart · high confidence
New IBM watsonx.ai chat model integration
This release adds a new integration for IBM watsonx.ai chat models, introducing \WatsonxChatModelConnection\ and \WatsonxChatModelSetup\ classes that allow users to connect to and interact with watsonx.ai services. The integration supports native structured output by translating Pydantic \BaseModel\ schemas into watsonx's \json\_schema\ response format, while preserving finish reasons and handling tool calls. It includes comprehensive test coverage for both mocked and integration scenarios, ensuring correct parameter passing and response parsing.
_integrations/chat-models/watsonx, python/flink\_agents/integrations/chat\models/watsonx · high confidence
New Java API for configurable event logging
The Flink Agents API now includes a new event logging system in the \org.apache.flink.agents.api.logger\ package, allowing users to capture and persist events from the agent execution pipeline. This change introduces the \EventLogger\ interface and an \EventLoggerFactory\ that supports built-in logging backends: \SLF4J\ (which outputs to the Flink Web UI via log4j2) and \FILE\ (which writes to per-subtask log files). Configuration is managed via \EventLoggerConfig\, which uses a fluent builder to select the logger type and pass properties like the base log directory. The system also supports \EventLogLevel\ (OFF, STANDARD, VERBOSE) to control verbosity, enabling users to choose between truncated summaries or full, untruncated event payloads.
api/src/main/java/org/apache/flink/agents/api/logger · high confidence
New Java Bash execution tool with strict security validation
A new Java implementation of the Bash tool has been added to the Flink Agents plan module, mirroring the existing Python version. This tool allows agents to execute shell commands while enforcing strict security policies: it uses an AST-based validator (via tree-sitter) to reject disallowed shell constructs, blocks file redirects (allowing only file-descriptor duplication/closure), and prevents assignments to dangerous environment variables like PATH or LD\_. Execution is sandboxed to allowed commands and directories, with configurable timeouts and working directories.
plan/src/main/java/org/apache/flink/agents/plan/tools/bash · high confidence
New Java EventListener interface for business event monitoring
A new EventListener interface has been added to the Java API, allowing users to register custom listeners that are notified synchronously when business events are received. This enables monitoring, metrics collection, and debugging by inspecting events before they are processed by actions, while explicitly excluding runtime observability trace records from this callback mechanism.
api/src/main/java/org/apache/flink/agents/api/listener · high confidence
New Java annotations for agent resources, actions, and cross-language support
The Java API now provides a set of annotations to declaratively define agent components and behavior. Developers can use @Action to mark methods as triggers for events or conditions, including cross-language targets via @PythonFunction. Resource management is simplified with annotations like @ChatModelSetup, @EmbeddingModelSetup, @Tool, @Prompt, @VectorStore, and @MCPServer, which allow the agent plan to automatically discover and manage these resources. Additionally, @ToolParam enables fine-grained control over tool arguments, including framework-injected parameters, and @Skills marks methods that provide agent skill definitions.
api/src/main/java/org/apache/flink/agents/api/annotation · high confidence
New Java examples for Flink Agents workflows, model routing, and skills
Added a suite of Java examples in the examples module demonstrating key Flink Agents capabilities: rule-based and LLM-as-judge model routing (ModelRoutingExample, ModelRoutingJudgeExample), parallel LLM fan-out for sentiment analysis (ParallelChatRequestExample), ReAct agent workflows for product review analysis (ReActAgentExample), agent skills integration (SkillsAgentExample), single and multi-agent streaming pipelines (WorkflowSingleAgentExample, WorkflowMultipleAgentExample), and YAML-declared agents (YamlWorkflowAgentExample). These examples also include shared helper classes and resource definitions (CustomTypesAndResources) to support the demonstrations.
examples · high confidence
New RAG example agent with ChromaDB and Ollama integration
A new Retrieval-Augmented Generation (RAG) example has been added to the Python Flink Agents library. This example demonstrates how to build an agent that retrieves context from a ChromaDB vector store using Ollama embeddings and generates responses via an Ollama chat model. The package includes the agent definition, a utility to populate the knowledge base, a standalone execution script, and tests verifying that the agent's actions are correctly resolved and that the vector store uses cross-process persistence.
_python/flink\agents/examples/rag · high confidence
New build, test, and validation tooling for Flink Agents
The tools directory now includes a comprehensive suite of scripts to streamline development and CI workflows. The new build.sh script automates the Java build and the packaging of Flink distribution JARs into the Python wheel structure, utilizing uv for Python dependency management. A new install.sh script provides an interactive, cross-platform installation wizard with a bundled UI helper (gum). Test execution is standardized via ut.sh (unit tests) and e2e.sh (end-to-end tests), which now support running against specific Flink versions and include cross-language resource consistency checks. Additionally, new validation tools ensure project hygiene: check-agents-md.py verifies that documentation version facts match source-of-truth files, check-license.sh runs Apache RAT license checks, and check-skill-schema.py ensures the coding-agent's bundled YAML schema remains in sync with the repository. The lint.sh script has been updated to use uv for Python linting and supports a check-only mode that no longer modifies files.
tools · high confidence
New execution tracing API for agent lifecycle events
The API now exposes a new \trace\ package (in both Java and Python) that enables recording agent execution lifecycles. This includes an \ExecutionReporter\ interface and helper utilities to report nested execution events (created, started, succeeded, failed) for entities like LLMs, tools, and actions, along with \ExecutionTraceContext\ to propagate lineage and run IDs. Standardized metadata keys for LLM and tool executions are also provided to structure the reported data.
_api/src/main/java/org/apache/flink/agents/api/trace, python/flink\agents/api/trace · high confidence
New observability event types for agent lifecycle, memory, and model interactions
The API now exposes a comprehensive set of typed event classes to track agent execution and internal state changes. This includes lifecycle tracking via AgentRunBeginEvent, chat interaction visibility through ChatRequestEvent and ChatResponseEvent, and context retrieval observability with ContextRetrievalRequestEvent and ContextRetrievalResponseEvent. Memory operations are now observable via ShortTermWriteEvent, ShortTermReadEvent, SensoryWriteEvent, SensoryReadEvent, LongTermGetEvent, LongTermSearchEvent, and LongTermUpdateEvent, all unified under the MemoryEvent base class. Additionally, model routing decisions are captured in ModelRoutingEvent, and tool interactions are tracked via ToolRequestEvent and ToolResponseEvent. These events enable users to monitor agent behavior, debug memory usage, and trace model selection logic.
api/src/main/java/org/apache/flink/agents/api/event · high confidence
New quickstart agent examples for Flink Agents
The quickstart examples now include a set of new agent implementations demonstrating various Flink Agents capabilities. These include a MathAgent that uses skills and tools for arithmetic, a ParallelChatAgent that fans out sentiment analysis across multiple LLM calls, and agents for product review analysis (ReviewAnalysisAgent, TableReviewAnalysisAgent) and product suggestion generation. The examples also introduce custom Pydantic types for structured data handling and demonstrate prompt construction, tool usage, and parallel event processing patterns.
_python/flink\agents/examples/quickstart/agents · high confidence
New release tooling and multi-JDK support in release scripts
The release process in tools/releasing has been expanded with new scripts to automate branching (create\_release\_branch.sh, create\_snapshot\_branch.sh), version updates (update\_branch\_version.sh), and contributor listing (list\_contributors.sh). The binary release script (create\_binary\_release.sh) now builds and signs both binary JARs and Python sdists, while deploy\_staging\_jars.sh introduces a two-phase deployment that builds the API module with a JDK 11 classifier alongside the default JDK 17 build, enabling multi-JDK support for consumers.
tools/releasing · high confidence
Python package now downloads required JARs from Maven Central during installation
The Python build process has been updated to automatically fetch and bundle necessary JAR dependencies from Maven Central when installing the package via pip. This change introduces a custom PEP 517 build backend that intercepts the wheel-building process to download JAR files listed in a manifest, verify their integrity using SHA-256 checksums, and include them in the final distribution. Users can now install the package without manually managing these external binary dependencies, and the download step can be skipped or redirected using environment variables if needed.
_python/\_build\backend · high confidence
Support for Python ChatModel in Java via new bridge classes
Added PythonChatModelConnection and PythonChatModelSetup classes to enable Java applications to interact with Python-based chat models. These new components act as a bridge, wrapping Python chat model objects and delegating chat operations to the underlying Python implementation while maintaining Java interface compatibility. The implementation handles message conversion between Java and Python, manages Python resource lifecycles, and supports metric group propagation for cross-language resources.
api/src/main/java/org/apache/flink/agents/api/chat/model/python · high confidence
Support for Python Model Context Protocol (MCP) resources in Java agents
Flink Agents now allows Java-based agents to interact with Python MCP servers, tools, and prompts. This change introduces new resource wrappers (PythonMCPServer, PythonMCPTool, PythonMCPPrompt) that bridge Python MCP objects to the Java agent runtime, enabling Java code to list and invoke Python-defined tools and prompts while handling serialization, metric tracking, and resource lifecycle management.
plan/src/main/java/org/apache/flink/agents/plan/resource · high confidence
Support using Java ChatModel, EmbeddingModel, and VectorStore in Python
Python agents can now directly use Java-based ChatModel, EmbeddingModel, and VectorStore resources. This change adds Python-side bridge implementations (JavaChatModelConnectionImpl, JavaEmbeddingModelConnectionImpl, JavaVectorStoreImpl, etc.) that wrap Java objects via pemja, enabling cross-language resource usage. The implementation handles message and document serialization between Python and Java, preserves token usage metrics, and ensures thread-safe execution by performing embeddings on the Python side to avoid CPython state corruption when crossing the language boundary.
_python/flink\agents/runtime/java · high confidence
Behavioural changes
Added bundled dependency notices and license files
The distribution package now includes a comprehensive NOTICE file and individual license texts for all bundled third-party dependencies. This ensures compliance with open-source licensing requirements by explicitly listing dependencies such as Jackson, Netty, AWS SDK, and Anthropic Java client, along with their respective licenses (Apache 2.0, MIT, BSD, EPL2, etc.), and providing the full license text for each in the \META-INF/licenses\ directory.
dist · high confidence
Introduce secure, digest-pinned URL skill sources in Java API
The Java API for agent skills now includes a new \Skills\ configuration resource and \SkillSourceSpec\ model to define where skills are loaded from. A key behavioral change is the introduction of strict security requirements for URL-based skill sources: the new \SkillUrlUtils\ enforces HTTPS by default and adds support for SHA-256 digest pinning to verify the integrity of downloaded skill archives. Users can now explicitly allow insecure HTTP via \fromUrlUnsafe\ or use \fromUrlWithSha256\ to require a specific digest, ensuring that remote skill sources are both encrypted and tamper-proof.
api/src/main/java/org/apache/flink/agents/api/skills · high confidence
Java Agent API introduces execution options, error handling, and TTL controls
The Java agent API now exposes \AgentExecutionOptions\ to configure runtime behavior, including error-handling strategies (fail, retry, ignore), retry intervals, and async execution settings for chat, RAG, and tool calls. The \Agent\ base class adds an \ErrorHandlingStrategy\ enum and validates that chat models and model routers do not share names. A built-in \ReActAgent\ is provided, which supports structured output schemas and injects schema prompts into the system message. Short-term memory now supports configurable TTL, with options to define when the TTL is updated and whether expired state is visible before cleanup.
api/src/main/java/org/apache/flink/agents/api/agents · high confidence
Java chat model API refactored into separate connection and setup abstractions with structured output support
The Java chat model API has been restructured to separate connection management from request setup. A new BaseChatModelConnection class handles provider-specific connection details and native structured-output capabilities, while BaseChatModelSetup manages prompt application, tool binding, and request preparation. The API now supports explicit output schemas via a new StructuredOutputStrategy enum (AUTO, NATIVE, PROMPT) that controls whether native structured output is used or falls back to prompt engineering. Token usage metrics are now recorded through the setup layer with model-scoped metric groups.
api/src/main/java/org/apache/flink/agents/api/chat/model · high confidence
Move MCP integration to dedicated module
The Model Context Protocol (MCP) integration has been relocated from the API package to a new \python/flink\_agents/integrations/mcp\ module. This change introduces the \MCPServer\, \MCPTool\, and \MCPPrompt\ classes, along with utility functions for handling MCP content types, within the integrations namespace to better organize external protocol support.
_python/flink\agents/integrations/mcp · high confidence
Python package now downloads Flink JARs from Maven Central during installation
The Python package has been restructured to include a build backend and a manifest file (jar\_manifest.json) that specifies Apache Flink JAR artifacts (versions 1.20 through 2.3) to be downloaded from Maven Central during pip install. This change shifts the distribution model from bundling JARs directly to fetching them at install time, ensuring the package includes the necessary Flink runtime libraries for agent execution.
python · high confidence
Test coverage
Add end-to-end tests for MCP server integration with and without prompts; Added Java test coverage for agent plan configuration, cross-language execution, and resource declarations; Added Java-side e2e tests for cross-language resource integration; Added Python API tests for Flink Agents; Added Python runtime test coverage for memory, sub-agents, and durable execution; Added Python tests for Java/Python agent plan and config compatibility; Added comprehensive test coverage for the Flink Agents runtime operator; Added cross-language agent e2e test resources; Added e2e test resources for Python Flink Agents; Added end-to-end tests for the Java YAML agent API; Added integration tests for Elasticsearch vector store; Added integration tests for OpenAI and Tongyi embedding models; Added integration tests for the Chroma vector store implementation; Added integration tests for the local Ollama embedding model; Added test coverage for Flink Agents runtime components; Added test coverage for Java agent skill runtime components; Added test coverage for plan actions; Added test resources for CEL condition evaluation and agent skill definitions; Added tests for AsyncExecutorThreadFactory thread naming and lifecycle; Added tests for MCP integration server and tool serialization; Added tests for Python and Java FunctionTool serialization and metadata handling; Added tests for Python environment management and remote execution configuration; Added tests for Python runtime resource handling and sub-agent invocation lifecycles; Added tests for Python-Java bridge utilities; Added tests for chat model integrations and structured output contracts; Added tests for feedback runtime components; Added tests for the Agents event logging runtime components; Added tests for the Azure OpenAI chat model integration; Added tests for the Python OpenAI chat model integration; Added tests for the agents runtime metrics subsystem; Added unit and integration tests for the agents runtime memory subsystem; Added unit tests for ActionState runtime components; Added unit tests for Anthropic chat model connection and setup; Added unit tests for MCP integration components; Added unit tests for OpenAI chat model integrations; Added unit tests for Python agent skill loading and repository components; Added unit tests for durable execution, observability, and lifecycle in the Java runner context; Added unit tests for runtime condition evaluation and compilation; Added unit tests for the Flink Agents API; Added unit tests for the Flink Agents installation and recovery tooling; Added unit tests for tool schema conversion utilities; Initial test coverage for the Python Flink Agents plan module; Integration test suite for Flink Agents installer; Introduce YAML API for declaring agents; New cross-language e2e tests for Flink Agents resources; New e2e test scripts for agent plan compatibility, resource consistency, and checkpoint recovery; New end-to-end integration tests for Flink Agents; New integration tests for agent skills, async execution, and checkpoint recovery; Structured tool-invocation assertions for end-to-end tests.
Dependencies
Multi-version Flink distribution support
The distribution module now builds and ships Flink Agents packages for Apache Flink versions 1.20, 2.0, 2.1, 2.2, and 2.3. This is achieved through a new Maven multi-module structure that includes a shared 'common' module for dependencies and specific modules for each Flink version, allowing users to select the distribution artifact that matches their cluster's Flink release.
(dependencies) · high confidence
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
How this codebase got here
Score
- CAI 49 → 55 (+6.2)
- Rubric changed (rubric-2026.08.19 → rubric-2026.09.15) — scores are not directly comparable.
Lenses
- Code Health 84 → 88 (+3.9)
- Architecture 99 → 84 (-14.8)
- Maturity 65 → 65 (+0.3)
- Readiness 38 → 39 (+0.5)
- Security 42 → 67 (+25.5)
Resolved (111)
- ActionExecutionOperator.processActionTaskForKey (cognitive 17) (runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java)
- ActionJsonSerializer.serialize (cognitive 16) (plan/src/main/java/org/apache/flink/agents/plan/serializer/ActionJsonSerializer.java)
- AgentPlanJsonDeserializer.deserialize (cyclomatic 17) (plan/src/main/java/org/apache/flink/agents/plan/serializer/AgentPlanJsonDeserializer.java)
- AzureOpenAIChatModelConnection.chat (cognitive 17) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/AzureOpenAIChatModelConnection.java)
- AzureOpenAIChatModelConnection.chat (cyclomatic 16) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/AzureOpenAIChatModelConnection.java)
- BaseChatModelSetup.chat (cognitive 20) (api/src/main/java/org/apache/flink/agents/api/chat/model/BaseChatModelSetup.java)
- Boundary-crossing change coupling: ChatModelAction.java ↔ chat_model_action.py (plan/src/main/java/org/apache/flink/agents/plan/actions/ChatModelAction.java)
- Change coupling: ollama_chat_model.py ↔ tongyi_chat_model.py (python/flink_agents/integrations/chat_models/ollama_chat_model.py)
- Coverage not included — suite not readable by the collector
- Dependency hygiene not measured — dependency manifest found but not parsed for hygiene
- Duplicated block (10 lines × 2) (api/src/main/java/org/apache/flink/agents/api/vectorstores/python/PythonVectorStore.java)
- Duplicated block (10 lines × 2) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/AzureOpenAIChatModelConnection.java)
- Duplicated block (10 lines × 2) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/OpenAICompletionsConnection.java)
- Duplicated block (10 lines × 2) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/OpenAICompletionsSetup.java)
- Duplicated block (10 lines × 2) (integrations/vector-stores/elasticsearch/src/main/java/org/apache/flink/agents/integrations/vectorstores/elasticsearch/ElasticsearchVectorStore.java)
- Duplicated block (11 lines × 2) (api/src/main/java/org/apache/flink/agents/api/embedding/model/python/PythonEmbeddingModelConnection.java)
- Duplicated block (11 lines × 2) (api/src/main/java/org/apache/flink/agents/api/embedding/model/python/PythonEmbeddingModelSetup.java)
- Duplicated block (11 lines × 2) (integrations/chat-models/bedrock/src/main/java/org/apache/flink/agents/integrations/chatmodels/bedrock/BedrockChatModelConnection.java)
- Duplicated block (11 lines × 2) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/OpenAICompletionsConnection.java)
- Duplicated block (11 lines × 2) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/OpenAICompletionsSetup.java)
- …and 91 more
New (387)
- ActionExecutionOperator.processActionTask (cognitive 29) (runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java)
- ActionExecutionOperator.processActionTask (cyclomatic 16) (runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java)
- ActionExecutionOperator.processEvent (cognitive 21) (runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java)
- ActionJsonDeserializer.deserialize (cognitive 21) (plan/src/main/java/org/apache/flink/agents/plan/serializer/ActionJsonDeserializer.java)
- ActionJsonDeserializer.deserialize (cyclomatic 16) (plan/src/main/java/org/apache/flink/agents/plan/serializer/ActionJsonDeserializer.java)
- AgentPlan.__custom_deserialize (cognitive 21) (python/flink_agents/plan/agent_plan.py)
- AgentsExecutionEnvironment.get_execution_environment (cognitive 18) (python/flink_agents/api/execution_environment.py)
- AnthropicChatModelConnection.chat (cognitive 32) (python/flink_agents/integrations/chat_models/anthropic/anthropic_chat_model.py)
- AnthropicChatModelConnection.chat (cyclomatic 22) (python/flink_agents/integrations/chat_models/anthropic/anthropic_chat_model.py)
- AzureOpenAIChatModelConnection.buildRequest (cognitive 18) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/AzureOpenAIChatModelConnection.java)
- AzureOpenAIChatModelConnection.buildRequest (cyclomatic 16) (integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/AzureOpenAIChatModelConnection.java)
- AzureOpenAIChatModelConnection.chat (cognitive 16) (python/flink_agents/integrations/chat_models/azure/azure_openai_chat_model.py)
- BaseChatModelSetup.prepareRequestMessages (cognitive 19) (api/src/main/java/org/apache/flink/agents/api/chat/model/BaseChatModelSetup.java)
- BashValidator.walk (cognitive 22) (plan/src/main/java/org/apache/flink/agents/plan/tools/bash/BashValidator.java)
- BashValidator.walk (cyclomatic 17) (plan/src/main/java/org/apache/flink/agents/plan/tools/bash/BashValidator.java)
- ChatModelInvoker.chatWithRetries (cognitive 29) (plan/src/main/java/org/apache/flink/agents/plan/actions/ChatModelInvoker.java)
- ChatModelInvoker.chatWithRetries (cyclomatic 16) (plan/src/main/java/org/apache/flink/agents/plan/actions/ChatModelInvoker.java)
- ClassTooLong: ActionExecutionOperator (runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java)
- ClassTooLong: AgentPlan (plan/src/main/java/org/apache/flink/agents/plan/AgentPlan.java)
- ClassTooLong: AnthropicChatModelConnection (integrations/chat-models/anthropic/src/main/java/org/apache/flink/agents/integrations/chatmodels/anthropic/AnthropicChatModelConnection.java)
- …and 367 more
Changes since last survey
- 108 commits — 90 feature/other, 18 fixes
By area
- runtime/src — 23 commits
- python/flink_agents — 21 commits
- integrations/chat-models — 16 commits
- (root) — 8 commits
- plan/src — 7 commits
- .github/workflows — 6 commits
- api/src — 5 commits
- docs/content — 4 commits
- tools/test — 4 commits
- dist/src — 2 commits
- e2e-test/flink-agents-end-to-end-tests-resource-cross-language — 2 commits
- integrations/vector-stores — 2 commits
- tools/build.sh — 2 commits
- docs/README.md — 1 commit
- e2e-test/cross-language-event-snapshots — 1 commit
- python/pyproject.toml — 1 commit
- tools/check-license.sh — 1 commit
- tools/docker — 1 commit
- tools/e2e.sh — 1 commit
Notable commits
- fix: [Bug] Recovered Durable ActionState for keys owned by other subtasks is never pruned (#1024)
- fix: [docs] Fix two Python snippets that do not run and an inaccurate schema parity claim (#972)
- fix: [e2e][java][hotfix] Deep-copy Mem0 action event payload (#1123)
- fix: [hotfix] Add help handling to build script (#995)
- fix: [hotfix][api][java] Add tests for enum value parsers (#1135)
- fix: [hotfix][ci] Pull MinIO from Quay for Milvus tests (#1118)
- fix: [hotfix][e2e] Derive recovery test version from project (#1041)
- fix: [hotfix][java][python] Harden Bash arithmetic evaluation (#1083)
- fix: [hotfix][runtime] Aggregate close failures across the Exception/Error boundary (#974)
- fix: [hotfix][runtime] Close the Kafka consumer when the producer fails to close (#948)
- fix: [hotfix][runtime] Raise Kafka admin future timeout from 100ms to 30s (#1095)
- fix: [hotfix][runtime][python] Close every component when an earlier close fails (#987)
- fix: [hotfix][tools] Report RAT validation dependencies correctly (#996)
- fix: [hotfix][vector-store][java] Update Elasticsearch filters documentation (#997)
- fix: [infra] Guard uv in e2e.sh and fix the pytest log level in ut.sh (#986)
- fix: [python] Fix jar URI construction in execution_environment.py (#1025)
- fix: [python][test] Fix OpenAI E2E model configuration (#1053)
- fix: [tools][ci] Fix the bats suite's boundary with the scripts it tests (#1081)
- change: Update discussion page link in README
- change: Update version to 0.4-SNAPSHOT (#994)
- …and 88 more
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
apache/flink-agents was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.
About this page
- The score is its most recent published measurement, taken on 22 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
- Measured at commit ea5e5ea89c1e76f8ea18ad766f9ee14d67a3ec44 — the exact code this score is about.
- Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
- Measured by watchdog.canine.dev using codehealth-analyzer preprod-90d5d2fe38ee.