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]

Reply via email to