sangkyoonnam opened a new issue, #1187: URL: https://github.com/apache/flink-agents/issues/1187
### Search before asking - [x] I searched in the [issues](https://github.com/apache/flink-agents/issues) and found nothing similar. ### Description In Java `MemoryObjectImpl`, a `set` or `newObject` whose path runs below a field that holds a value is not rejected on purpose. It fails by accident, and depending on the call, after part of the write has already happened. With `a` holding a value: `set("a.b", v)` fails before storing anything; `set("a.b.c", v)` leaves an orphan OBJECT node `a.b` under the value; `newObject("a.b")` stores the node and records a `MemoryUpdate` that makes replay of that action fail on recovery; `newObject("a.b.c")` leaves the orphan `a.b` and records nothing. What the code does, on `main` at `99103da6`. `set` (`MemoryObjectImpl.java:95`) and `newObject` (line 122) call `fillParents` (line 220), which creates missing intermediate OBJECT nodes but never checks an existing ancestor's type. A VALUE item has `subKeys = Collections.emptySet()` (line 252), so the write stops at the first `getSubKeys().add(...)` on that item (lines 106, 148, 237) with `UnsupportedOperationException` and no message. For `newObject("a.b")` that happens after the node is stored (line 139) and `MemoryUpdate("a.b", objectCreation=true)` is appended to the action's update list (line 141). Observed consequences, each from a test on a `KeyedOneInputStreamOperatorTestHarness` over `ActionExecutionOperator` with `HashMapStateBackend` and `EmbeddedRocksDBStateBackend`. The orphan `a.b` is persisted in keyed state and survives snapshot and restore; `isExist("a.b")` is true while `a` is a value with no fields. When an action catches the exception from `newObject("a.b")` and completes, its `ActionState` holds the update. On recovery with `InMemoryActionStateStore`, redelivery of that input replays the update through `MemoryUpdateReplayer.replay` (`MemoryUpdateReplayer.java:51`), which hits the same `UnsupportedOperationException` from `ActionExecutionOperator$Work.execute` (line 1523), and the input produces no output. I ran one redelivery, not a job with a restart strategy; since the replay is deterministic over the same state and the same `ActionState`, I expect each restart to fail the same way. The serial path at line 772 makes the same call; I read it and did not run it. Paths with an empty component take the same unchecked route. Each case below starts from a clean memory. `set("a.", v)` throws `NullPointerException` when `a` is absent, and the same `UnsupportedOperationException` when `a` is a value; `newObject("a.")` does the same after storing `a.` and recording an update. `set(".x", 1)` succeeds, registers both `""` and `x` as root fields, and the next `getFields()` on the root throws `NullPointerException` because nothing is stored at `x`; `newObject(".y")` does the same and records an update. After `newObject("o")`, `get("o").set("", 2)` stores `o.` and registers `o` as a field of `o`, so `get("o").getFields()` throws the same way. For the leading-dot examples, if `x` or `y` already exists, root `getFields()` succeeds and lists `""` as a nested object. `set("a..b", 1)` with `a` absent succeeds and stays readable and listable. From the root, `newObject("")` records an update and registers the root as its own field; `set("", v)` registers tha t field and then throws `IllegalArgumentException("Cannot overwrite object with value: ")`. Nothing in the repository writes such paths; I list them because the fix below has to decide what a valid write path is. Python actions get the same behavior through Pemja: `FlinkMemoryObject` delegates to the Java object and wraps the exception in `MemoryObjectError("Failed to set value at path 'a.b'")` for `set`, or `MemoryObjectError("Failed to create new object at path 'a.b'")` for `new_object`; the cause carries only `java.lang.UnsupportedOperationException`. The `MemoryObject.set` javadoc documents "overwrite a nested object with a primitive value" and "set a MemoryObject directly" as errors, and `MemoryObjectTest.testOverwriteRules` covers the two overwrite conflicts with `IllegalArgumentException`. The ancestor case and empty components are unspecified. I'd treat them the same way: validate the whole path before any mutation in `set` and `newObject`, and throw `IllegalArgumentException("Cannot write field 'a.b.c': 'a' exists but is not an object.")`, naming the absolute path (a nested handle's prefix included), or `"... path has an empty component."`. `overwrite=true` on `newObject` would keep applying to the target field only. An `ActionState` written before the fix can still hold the recorded update; replaying it on fixed code still fails, now with that message. Whether replay should skip such an update is a separate question I haven't tried to answer. ### How to reproduce Each call below runs against a short-term memory where `a` holds the value 1; the comment describes the state right after that call. ```java MemoryObject stm = ctx.getShortTermMemory(); stm.set("a", 1); try { stm.set("a.b", 2); } catch (UnsupportedOperationException e) { // no message; nothing stored } try { stm.set("a.b.c", 2); } catch (UnsupportedOperationException e) { // isExist("a.b") is now true; a is still the value 1 with no fields } try { stm.newObject("a.b"); } catch (UnsupportedOperationException e) { // node a.b stored, MemoryUpdate("a.b", objectCreation=true) recorded; // if the action completes, replaying its ActionState on recovery throws the same exception } ``` The runtime unit test setup (`new MemoryObjectImpl(SHORT_TERM, new CachedMemoryStore(new ForTestMemoryMapState<>()), ROOT_KEY, updates)`) shows the exception, the orphan and the recorded update. The recovery failure shows up with `ActionExecutionOperatorFactory<>(plan, true, new InMemoryActionStateStore(false))`: process `set("a", 1)`, snapshot, process the catching action, close the harness, restore from the snapshot with the same store, redeliver the same input. ### Version and environment main (`99103da6`). JDK 21, Flink 2.3.0. Java runtime; the Python `FlinkMemoryObject` path delegates to it. ### Are you willing to submit a PR? - [x] I'm willing to submit a PR! -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
