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]

Reply via email to