Skip to content

[Bug] Python resource bridges do not consistently release Pemja references #1057

Description

@joeyutong

Search before asking

  • I searched in the issues and found nothing similar.

Description

Flink Agents invokes Python resources through Pemja in several runtime paths. Values crossing the Java/Python boundary may be represented by PyObject handles, and each such handle owns a native Python reference that must be released explicitly. Java garbage collection is not an ownership mechanism for these references.

The current resource bridges do not apply one consistent ownership contract. Depending on the path, temporary arguments or return values can remain open, a value can outlive the bridge call without an explicit ownership transfer, or a partially initialized resource can be cached without rolling back its Python and native state. The affected paths include:

  • chat and embedding model setup/connections;
  • vector-store calls and returned document collections;
  • MCP servers, tools, prompts, and discovery;
  • Mem0 call-scoped values;
  • Python-backed FunctionTool calls;
  • Java and Python resource-cache initialization failures.

Nested Python values are also relevant: closing only a top-level collection is insufficient when its elements are separate PyObject handles.

The expected ownership contract is:

  1. one operation scope owns every temporary Python argument and result, including nested handles;
  2. values that intentionally outlive the operation are transferred to an explicit long-lived owner;
  3. results returned to ordinary Java callers are materialized into detached Java values before Python handles are released;
  4. resource creation and MCP discovery are transactional, so any failure closes all resources created before cache registration.

This is related to, but distinct from, #942/#944 and #1048/#1049. Those pairs cover executor-lifetime handles across task restarts and per-Action handles respectively. This issue covers the remaining resource-bridge operations and partial-initialization paths.

Independent local Flink A/B validation

Six bounded jobs ran in a local Flink 2.2 session cluster. Every fixed and leak variant used a fresh TaskManager, one probe subtask, Pemja 0.5.7 in MULTI_THREAD mode, embedded CPython 3.11.15, 200 measured iterations, and 256 KiB per tracked Python object.

  • Chat bridge: fixed 0 live messages/tools and 0 retained payload; leak variant retained 600 messages, 200 tool values, and 200 MiB, with RSS increasing by 160.3 MiB.
  • Vector bridge: fixed 0 live documents/normalized handles and 0 retained payload; leak variant retained 600 documents, 200 normalized handles, and about 100 MiB, with RSS increasing by 99.1 MiB.
  • Resource recreation: fixed 0 live resources and observed 200 logical close calls; leak variant retained 200 resources, made 0 logical close calls, and retained 50 MiB.

Exact live-reference counts, close counts, and reachable payload bytes are the primary evidence. RSS is supporting evidence because JVM/CPython startup and native allocator behavior can dominate a bounded run.

How to reproduce

  1. Start a Flink job with one long-lived TaskManager and a Python-backed chat, vector, or resource bridge.
  2. Repeatedly invoke the bridge with non-trivial Python arguments/results, or repeatedly recreate a resource while keeping the TaskManager alive.
  3. Track affected Python objects with weak references and count logical close calls.
  4. Force Java and Python garbage collection after the measured operations.
  5. Compare the existing ownership behavior with an operation scope that recursively closes temporary handles, explicitly transfers long-lived handles, and rolls back failed initialization.

Expected behavior: completed operations and failed resource initialization leave no bridge-owned Python references or partially registered resources behind.

Version and environment

  • Affected code: current main; the ownership pattern also exists in the release-0.3 production line
  • Local A/B: Flink 2.2.0, Pemja 0.5.7, Java 17, CPython 3.11.15, macOS

Are you willing to submit a PR?

  • I'm willing to submit a PR!

Metadata

Metadata

Assignees

Labels

affectVersion/0.2.1affectVersion/0.3.1The bug affects the 0.3.1 version. The features are not needed.bug[Issue Type] Something isn't working as expected.priority/majorDefault priority of the PR or issue.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions