sangkyoonnam opened a new pull request, #1212: URL: https://github.com/apache/flink-agents/pull/1212
Linked issue: Closes #1192 ### Purpose of change An Event that already carries `upstreamEventId` or `upstreamActionName` (Python: `upstream_event_id`, `upstream_action_name`) now fails at `sendEvent` / `send_event`. Lineage set by user code used to be silently replaced when the Action's outputs were finalized. Reassigning a Python Event's `id` now points to `reconstruct_from`. #### Runtime flow Every Java emission reaches `RunnerContextImpl.sendEvent`, and so does every Python emission that passes `send_event`'s own check, through `sendEventJson`, which deserializes with `Event.fromJson` first. The check runs there, after the mailbox thread check and before attachments are stored or the Event is buffered. Python's `FlinkRunnerContext.send_event` runs the same check before serializing, so the error names the Python fields. `ActionTask.finalizeOutputEvents` still binds lineage to the triggering Event and Action, at completion and, for a Python coroutine Action, at each yield. Outputs restored from action state reach that step without `sendEvent`. The operator sets `AgentRunBeginEvent`'s lineage itself and processes it directly. #### Key decisions - Reject rather than clear: a caller that set lineage expected it to stick, so clearing it would repeat the silent replacement this removes. - Check at the send boundary only. Deserialization and typed reconstruction keep an Event's lineage; only sending such an Event fails. - Keep `ValidationError` and the "Field is frozen" text for `id` reassignment and add the guidance, so handlers matching the class or message still match. ### Behavioral Semantics #### Interaction decisions | Event | Through `sendEvent` / `send_event` | Result (Java and Python) | |---|---|---| | new Event, no lineage | yes | accepted; lineage bound when outputs are finalized (unchanged) | | lineage set by user code | yes | `IllegalArgumentException` / `ValueError` (was: overwritten) | | the current trigger carrying lineage, re-sent or typed-reconstructed | yes | `IllegalArgumentException` / `ValueError` at send (was: rejected at finalization by the self-loop check) | | an earlier Event from another Action, re-sent or reconstructed | yes | `IllegalArgumentException` / `ValueError` (was: overwritten) | | the root `InputEvent`, re-sent or typed-reconstructed | yes | passes this check (it has no lineage); rejected at finalization by the existing self-loop check (unchanged) | | output restored from action state | no | rebound to the replayed trigger (unchanged) | | `AgentRunBeginEvent` | no | lineage set by the operator (unchanged) | #### Behavioral contracts 1. An Event with either lineage field set fails at Java `sendEvent` with an `IllegalArgumentException` that names only the fields set and points to `attributes`. 2. The rejected Event is not buffered. 3. The same holds for lineage that arrives as JSON from Python. 4. Python's `send_event` raises `ValueError` naming only the Python fields set, and forwards nothing to Java. 5. A typed reconstruction (`fromEvent` / `reconstruct_from`) of an Event that carries lineage is rejected. 6. A new Event is accepted and gets lineage when outputs are finalized. 7. An output restored from action state is rebound to the current trigger. 8. Assigning a Python Event's `id` raises `ValidationError` with "Field is frozen" and mentions `reconstruct_from`. #### Failure behavior The error is raised to the Action at the `sendEvent` / `send_event` call, with no retry, and the Event is not buffered. An Action that catches it can continue and emit other Events. If it escapes, the Action fails like any other Action error: its buffered outputs are not persisted, outputs a Python coroutine Action emitted at earlier yields may already be dispatched, and the configured Flink restart strategy applies. ### Tests | Contract | Tests | |---|---| | 1 | `RunnerContextLineageTest.presetUpstreamEventIdIsRejected`, `presetUpstreamActionNameIsRejected` | | 2 | same two tests (`drainEvents` is empty) | | 3 | `RunnerContextLineageTest.lineageArrivingThroughJsonIsRejected` | | 4 | `test_send_event_rejects_preset_lineage[upstream_event_id\|upstream_action_name]` | | 5 | `RunnerContextLineageTest.typedReconstructionOfAnEventFromAnotherActionIsRejected`, `test_send_event_rejects_reconstructed_event_with_lineage` | | 6 | `RunnerContextLineageTest.freshEventIsAccepted`, `test_send_event_forwards_fresh_event_without_lineage`, `ActionTaskTest.resultFinalizesOutputLineage` | | 7 | `ActionTaskTest.resultRebindsLineageOfOutputsRestoredFromActionState` | | 8 | `test_event.py::test_event_id_cannot_be_reassigned` | Coverage by risk: the new rejection (1 to 5) has tests in both languages, with the JSON path (3) in Java only, and the `id` guidance (8) in Python. Each of those tests fails without the change. Contracts 6 and 7 are unchanged behavior and pass with or without it. Not verified: - Recovery through a real action state store with lineage-carrying outputs. - A Java Action in a Flink job; the job runs below use Python Actions. <details> <summary>PyFlink job runs (MiniCluster, Flink 2.3, JDK 17)</summary> A two-Action Python agent on a bounded file source, run against this branch and against main: | Case | main | this branch | |---|---|---| | new Event | succeeds; downstream sees `upstream_action_name='first'` | same | | Python `send_event` with `upstream_action_name` set | succeeds; the value is replaced with `'first'` | job fails with this PR's `ValueError` through pemja | | `sendEventJson` called directly with `upstreamActionName` set | succeeds; the value is replaced | job fails with this PR's `IllegalArgumentException` | | re-send or `reconstruct_from` of the current trigger (carries lineage) | job fails at finalization (self-loop check) | job fails at send with this PR's `ValueError` | | re-send of the root `InputEvent` | job fails at finalization (self-loop check) | same | </details> ### API No signature changes; `RunnerContext.sendEvent` and `send_event` now document the new error. Code that set lineage before sending, or re-sent an earlier Event from another Action, now fails where its lineage used to be overwritten. Re-sending a current trigger that carries lineage, or a reconstruction of it, already failed at finalization and now fails at send with this message. The Python `id` error keeps its type and "Field is frozen" text, but its Pydantic error type changes from `frozen_field` to `frozen_event_id`. ### Documentation - [ ] `doc-needed` - [ ] `doc-not-needed` - [x] `doc-included` The lineage note in the workflow agent guide and the lineage setter Javadoc now describe the rejection. ### Was this patch authored or co-authored using generative AI tooling? - [x] Yes - [ ] No Generated-by: Claude Code 2.1.295 (Claude Opus 5.5) -- 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]
